Skip to content

refactor: use group clustering terminology in aggregate execution - #24697

Open
xavlee wants to merge 4 commits into
apache:mainfrom
xavlee:refactor/aggregate-group-completion-mode
Open

xavlee wants to merge 4 commits into
apache:mainfrom
xavlee:refactor/aggregate-group-completion-mode

Conversation

@xavlee

@xavlee xavlee commented Aug 26, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR relate to?

Rationale for this change

Aggregate execution currently uses InputOrderMode to select execution paths and recognize completed groups.

A group can be completed incrementally when its rows form one contiguous run: once the run ends, that group cannot appear again in the same input partition. Sorted input supplies this guarantee, and unsorted input can supply it too:

AAABBBCCCC -> sorted; groups can be completed incrementally
CCCAAABBB  -> unsorted; groups can still be completed incrementally

This PR describes the input's clustering guarantee with GroupClusteringMode and uses clustering terminology throughout the aggregate execution machinery. Aggregate construction derives the mode directly from the existing ordering analysis. The follow-up in #24497 will also derive it from grouped equivalence properties.

What changes are included in this PR?

  • Replace the aggregate's InputOrderMode field with GroupClusteringMode::{None, Partial, Full}:
    • None: Groups have no apparent clustering guarantee.
    • Partial(indices): rows are contiguous on a subset of grouping expressions; a change in that tuple completes all groups in the preceding run.
    • Full: rows are contiguous on the complete grouping tuple; a change in that tuple completes the preceding group.
  • Rename the stateful GroupOrdering tracker and its Partial and Full implementations to GroupClustering, GroupClusteringPartial, and GroupClusteringFull.
  • Rename the ordered aggregate streams, shared table, and adjacent-key group-values implementation to use clustered terminology,
  • Preserve actual sort expressions and options through equivalence properties and input-ordering requirements.
  • Display group_clustering_mode in aggregate plans and document the public API migration in the DataFusion 56 upgrade guide.

Stack

  1. #24737 — test: cover unsorted contiguous groups in one partition (merged)
  2. #24697 — refactor: use group clustering terminology in aggregate execution ← this PR
  3. #24698 — feat: add Grouped equivalence properties
  4. #24497 — feat: stream aggregates over grouped input

Are there any user-facing changes?

Yes. This PR changes public Rust APIs:

  • AggregateExec::input_order_mode() is replaced by AggregateExec::group_clustering_mode(), returning the public GroupClusteringMode enum.
  • AggregateExec::compute_properties accepts &GroupClusteringMode in place of &InputOrderMode.
  • The public GroupOrdering, GroupOrderingPartial, and GroupOrderingFull types in aggregates::order are renamed to GroupClustering, GroupClusteringPartial, and GroupClusteringFull.
  • GroupClustering::try_new takes &GroupClusteringMode; new_group_values takes &GroupClustering.

Aggregate plan output uses group_clustering_mode=Full and group_clustering_mode=Partial(indices) in place of ordering_mode=Sorted and ordering_mode=PartiallySorted(indices). SQL results and aggregate execution-path selection retain their existing behavior.

Review this layer

View this PR's diff

@github-actions github-actions Bot added the physical-plan Changes to the physical-plan crate label Aug 26, 2026
@codecov-commenter

codecov-commenter commented Aug 26, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 86.00583% with 48 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.66%. Comparing base (c3ef346) to head (d112ecc).

Files with missing lines Patch % Lines
datafusion/physical-plan/src/aggregates/mod.rs 83.13% 3 Missing and 11 partials ⚠️
...cal-plan/src/aggregates/clustered_single_stream.rs 79.66% 12 Missing ⚠️
...tafusion/physical-plan/src/aggregates/order/mod.rs 85.71% 2 Missing and 4 partials ⚠️
...hysical-plan/src/aggregates/grouped_hash_stream.rs 73.33% 1 Missing and 3 partials ⚠️
...sion/physical-plan/src/aggregates/order/partial.rs 92.50% 0 Missing and 3 partials ⚠️
...ggregates/aggregate_hash_table/common_clustered.rs 86.66% 0 Missing and 2 partials ⚠️
...plan/src/aggregates/aggregate_hash_table/common.rs 0.00% 0 Missing and 1 partial ⚠️
...c/aggregates/aggregate_hash_table/partial_table.rs 0.00% 0 Missing and 1 partial ⚠️
...ical-plan/src/aggregates/clustered_final_stream.rs 95.23% 1 Missing ⚠️
...n/physical-plan/src/aggregates/group_values/mod.rs 87.50% 0 Missing and 1 partial ⚠️
... and 3 more
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #24697      +/-   ##
==========================================
- Coverage   82.66%   82.66%   -0.01%     
==========================================
  Files        1147     1147              
  Lines      446357   446276      -81     
  Branches   446357   446276      -81     
==========================================
- Hits       368971   368903      -68     
- Misses      54997    54998       +1     
+ Partials    22389    22375      -14     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@xavlee
xavlee force-pushed the refactor/aggregate-group-completion-mode branch 4 times, most recently from f58e128 to 372429f Compare August 27, 2026 20:54
@xavlee
xavlee force-pushed the refactor/aggregate-group-completion-mode branch 2 times, most recently from d1d1d2f to 0267ad1 Compare August 31, 2026 17:50

@gene-bordegaray gene-bordegaray left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

approved with some non blocking suggestoins. Thank you @xavlee

Comment thread datafusion/physical-plan/src/aggregates/order/mod.rs Outdated
Comment thread datafusion/physical-plan/src/aggregates/order/mod.rs Outdated
Comment thread datafusion/physical-plan/src/aggregates/ordered_final_stream.rs Outdated
Comment thread datafusion/physical-plan/src/aggregates/ordered_partial_stream.rs Outdated
Comment thread datafusion/physical-plan/src/aggregates/mod.rs Outdated
Comment thread datafusion/physical-plan/src/aggregates/mod.rs
@xavlee
xavlee force-pushed the refactor/aggregate-group-completion-mode branch 4 times, most recently from e7f8b42 to 0108b7c Compare September 2, 2026 15:34
@gene-bordegaray

Copy link
Copy Markdown
Contributor

Rereviewed after changes and everything looks good 👍

@NGA-TRAN NGA-TRAN left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Looks great. It is a pleasure to review. Nice explanation.

I wonder if we should add an attribute in the explain to show this property? Maybe in a follow-up PR? If it is not that invasive or very easy to review, maybe adding that in this PR?

///
/// For example, with `GROUP BY (a, b)`, `Partial(vec![0])` means all rows
/// for each value of `a` are contiguous, while an `(a, b)` tuple may recur
/// within that range.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does this mean I have 2 keys (a, b) and data is sorted on (a) only?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Exactly

assert_eq!(
aggregate.group_completion_mode,
GroupCompletionMode::Partial(vec![0])
);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nice

// This captures the behavior before #24438. When the source can declare
// `(key, time_bin)` group-contiguous, the corresponding case can use
// `EmissionType::Incremental`.
assert_eq!(aggregate.cache().emission_type, EmissionType::Final);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

👍

@xavlee
xavlee force-pushed the refactor/aggregate-group-completion-mode branch from 0108b7c to 6e369ce Compare September 4, 2026 19:21
@xavlee
xavlee force-pushed the refactor/aggregate-group-completion-mode branch 2 times, most recently from 2fae45d to bb3dc16 Compare September 7, 2026 20:23
@xavlee
xavlee force-pushed the refactor/aggregate-group-completion-mode branch from bb3dc16 to 64a6e55 Compare September 16, 2026 14:24

@gene-bordegaray gene-bordegaray left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this also looks good, just needs a rebase

@xavlee
xavlee force-pushed the refactor/aggregate-group-completion-mode branch 2 times, most recently from 3a0ecb6 to 1fd156e Compare September 23, 2026 02:04
@xavlee
xavlee marked this pull request as ready for review September 23, 2026 19:30

@xudong963 xudong963 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: finish updating the runtime documentation from “ordered” to “group-complete.”

zhuqi-lucas pushed a commit to zhuqi-lucas/arrow-datafusion that referenced this pull request Sep 28, 2026
## Which issue does this PR relate to?

- Part of apache#24438.

## Rationale for this change

A source can concatenate several sorted logical runs into one DataFusion
output partition. The resulting stream may be globally unsorted while
every distinct `(key, time_bin)` tuple still occupies one contiguous
range.

This PR records how aggregate planning handles that layout before
grouped input properties are available. It provides the behavioral
baseline for the remaining PRs in the stack.

## What changes are included in this PR?

- Add a single-partition `TestMemoryExec` fixture containing two sorted
logical runs whose `(key, time_bin)` order resets at the record-batch
boundary.
- Aggregate by the complete `(key, time_bin)` tuple and verify the
result.
- Assert that planning selects `InputOrderMode::Linear`,
`EmissionType::Final`, and `SingleHashAggregateStream`.

## Stack

1. [apache#24737 — test: cover unsorted contiguous groups in one
partition](apache#24737) ← **this
PR**
2. [apache#24697 — refactor: separate aggregate group completion from input
ordering](apache#24697)
3. [apache#24698 — feat: add grouped equivalence
properties](apache#24698)
4. [apache#24497 — feat: stream aggregates over grouped
input](apache#24497)

## Are these changes tested?

The new aggregate test executes the single-partition input and
snapshot-checks all four grouped sums.

## Are there any user-facing changes?

No.

## Review this layer

[View only this PR
layer](apache@36969e7)

@2010YOUY01 2010YOUY01 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you, this is a neat idea!

My main concern is that we now use separate flags for group clustering and group ordering. This could become hard to maintain if we extend the optimization further.

Ideal solution

Unify input group contiguity and ordering into a single representation. Group clustering could be the more general case, with ordering represented as an optional stronger property.

Practical alternative

Since the existing aggregation optimization only exploits the clustering property, we could probably generalize the current design by replacing the ordering-specific flags with clustered-group terminology.

  • Rename GroupOrdering to GroupCompletion, since that better reflects what it represents today.
  • ...replace other order flag to group clustering flags, perhaps also do some renaming like ClusteredPartialAggregateStream.

If we later introduce optimizations that specifically depend on ordering, we can extend the model then.

Comment on lines +893 to +899
input_order_mode: InputOrderMode,
/// Describes when the executor can determine that groups are complete.
///
/// Input ordering describes a subset of the cases in which groups can be
/// safely emitted before the input ends. Full group completion requires only
/// that rows for each complete grouping tuple are contiguous.
group_completion_mode: GroupCompletionMode,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is it possible to unify these two modes into a single struct? For example, could we represent both the clustering property and the ordering property using only InputOrderingMode?

Roughly, I feel clustering is the more general property, while sort order is a stronger guarantee, so the implementation might model ordering as an optional additional property.

This is not an issue for now, because the existing ordering optimization in aggregation only relies on clustering; no optimization currently relies on the stronger ordering guarantee.

However, if we want to extend this in the future to support both:

  • optimizations that rely only on clustering, and
  • additional optimizations that exploit sort order,

then representing these properties as a combination of flags could become error-prone and difficult to extend.

There are already conversions between GroupCompletionMode and InputOrderMode in this PR, and I find them quite hard to interpret.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I combined these two modes into GroupCompletionMode since that's the more general property the runtime needs to determine how to emit groups.

input_order_mode = InputOrderMode::Linear;
}

let group_completion_mode = GroupCompletionMode::from(&input_order_mode);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

marker for previous comment: this is a order_mode -> completion_mode conversion

}

/// Create a `GroupOrdering` for the specified group-completion mode.
pub(crate) fn try_new_for_group_completion(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

marker for previous comment: this is a completion-mode -> ordering conversion.

@xavlee
xavlee force-pushed the refactor/aggregate-group-completion-mode branch from 1fd156e to d112ecc Compare October 4, 2026 22:53
@github-actions github-actions Bot added documentation Improvements or additions to documentation core Core DataFusion crate sqllogictest SQL Logic Tests (.slt) labels Oct 4, 2026
@github-actions

github-actions Bot commented Oct 4, 2026 •

Copy link
Copy Markdown

Thank you for opening this pull request!

Reviewer note: cargo-semver-checks reported the current version number is not SemVer-compatible with the changes in this pull request (compared against the base branch).

Details
     Cloning apache/main
    Building datafusion v55.1.0 (current)
       Built [  56.207s] (current)
     Parsing datafusion v55.1.0 (current)
      Parsed [   0.032s] (current)
    Building datafusion v55.1.0 (baseline)
       Built [  55.705s] (baseline)
     Parsing datafusion v55.1.0 (baseline)
      Parsed [   0.032s] (baseline)
    Checking datafusion v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.666s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [ 114.240s] datafusion
    Building datafusion-physical-plan v55.1.0 (current)
       Built [  36.976s] (current)
     Parsing datafusion-physical-plan v55.1.0 (current)
      Parsed [   0.176s] (current)
    Building datafusion-physical-plan v55.1.0 (baseline)
       Built [  36.656s] (baseline)
     Parsing datafusion-physical-plan v55.1.0 (baseline)
      Parsed [   0.177s] (baseline)
    Checking datafusion-physical-plan v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.622s] 223 checks: 220 pass, 3 fail, 0 warn, 31 skip

--- failure enum_missing: pub enum removed or renamed ---

Description:
A publicly-visible enum cannot be imported by its prior path. A `pub use` may have been removed, or the enum itself may have been renamed or removed entirely.
        ref: https://doc.rust-lang.org/cargo/reference/semver.html#item-remove
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/enum_missing.ron

Failed in:
  enum datafusion_physical_plan::aggregates::order::GroupOrdering, previously in file /home/runner/work/datafusion/datafusion/target/semver-checks/git-apache_main/18b6ab8911f21b21caa67d7570f5d4ac633a60d2/datafusion/physical-plan/src/aggregates/order/mod.rs:33

--- failure inherent_method_missing: pub method removed or renamed ---

Description:
A publicly-visible method or associated fn is no longer available under its prior name. It may have been renamed or removed entirely.
        ref: https://doc.rust-lang.org/cargo/reference/semver.html#item-remove
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/inherent_method_missing.ron

Failed in:
  AggregateExec::input_order_mode, previously in file /home/runner/work/datafusion/datafusion/target/semver-checks/git-apache_main/18b6ab8911f21b21caa67d7570f5d4ac633a60d2/datafusion/physical-plan/src/aggregates/mod.rs:1872

--- failure struct_missing: pub struct removed or renamed ---

Description:
A publicly-visible struct cannot be imported by its prior path. A `pub use` may have been removed, or the struct itself may have been renamed or removed entirely.
        ref: https://doc.rust-lang.org/cargo/reference/semver.html#item-remove
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/struct_missing.ron

Failed in:
  struct datafusion_physical_plan::aggregates::order::GroupOrderingFull, previously in file /home/runner/work/datafusion/datafusion/target/semver-checks/git-apache_main/18b6ab8911f21b21caa67d7570f5d4ac633a60d2/datafusion/physical-plan/src/aggregates/order/full.rs:58
  struct datafusion_physical_plan::aggregates::order::GroupOrderingPartial, previously in file /home/runner/work/datafusion/datafusion/target/semver-checks/git-apache_main/18b6ab8911f21b21caa67d7570f5d4ac633a60d2/datafusion/physical-plan/src/aggregates/order/partial.rs:66

     Summary semver requires new major version: 3 major and 0 minor checks failed
    Finished [  75.742s] datafusion-physical-plan
    Building datafusion-sqllogictest v55.1.0 (current)
       Built [  94.524s] (current)
     Parsing datafusion-sqllogictest v55.1.0 (current)
      Parsed [   0.016s] (current)
    Building datafusion-sqllogictest v55.1.0 (baseline)
       Built [  94.741s] (baseline)
     Parsing datafusion-sqllogictest v55.1.0 (baseline)
      Parsed [   0.016s] (baseline)
    Checking datafusion-sqllogictest v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.106s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [ 191.712s] datafusion-sqllogictest

@github-actions github-actions Bot added the auto detected api change Auto detected API change label Oct 4, 2026
@xavlee

xavlee commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor Author

Hi @2010YOUY01, thanks for the review 🙂. If no other optimization currently depends on the existing sorting guarantee, then the practical alternative you mentioned seems tenable -- we can extend the model should a use case for ordering arises.

I've updated this PR in the three ways you mentioned:

  • renamed ordering to completion_mode
  • folded order/completion mode into the single GroupCompletionMode instead of maintaining both
  • Renamed ordered_*_stream to clustered_*_stream

Let me know if this is what you had in mind

One naming question: would you prefer to rename GroupCompletionMode to GroupClusteringMode to align with the clustering terminology you suggested?

@2010YOUY01

2010YOUY01 commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor

Thank you, the structure inside AggregateExec LGTM. Though I might not be able to review this entire PR now.

EDIT: Note that a large mechanical diff isn't an issue for me when reviewing—we can easily verify it with AI.

What's less obvious to me is how to represent the clustered BBAAACCCC input characteristic in an ideal way. That's why I'd prefer to leave it to someone with a deeper understanding.

One naming question: would you prefer to rename GroupCompletionMode to GroupClusteringMode to align with the clustering terminology you suggested?

Yes, I agree GroupClusteringMode looks better, and I think it's also important to have a unified term for this family of optimization in the codebase.

@NGA-TRAN

NGA-TRAN commented Oct 7, 2026

Copy link
Copy Markdown
Contributor

Thanks so much, @2010YOUY01, for the suggestions. Since you may not have time to review the full PR right now, would it be okay for @xudong963, @alamb, or @jayzhan211 to review the entire PR based on your feedback and recommendations?

@xavlee xavlee changed the title refactor: separate aggregate group completion from input ordering refactor: use group clustering terminology in aggregate execution Oct 7, 2026
@jayzhan211
jayzhan211 self-requested a review October 8, 2026 12:46

@jayzhan211 jayzhan211 left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @xavlee , the idea makes sense to me

What's less obvious to me is how to represent the clustered BBAAACCCC input characteristic in an ideal way. That's why I'd prefer to leave it to someone with a deeper understanding.

I think we can revisit this for the follow-up PRs

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

auto detected api change Auto detected API change core Core DataFusion crate documentation Improvements or additions to documentation physical-plan Changes to the physical-plan crate sqllogictest SQL Logic Tests (.slt)

Projects

None yet

Development

Successfully merging this pull request may close these issues.

7 participants