Skip to content

refactor: separate aggregate group completion from input ordering - #24697

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

xavlee wants to merge 1 commit 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

AggregateExec currently uses InputOrderMode for two related decisions: (1) describing how the grouping expressions are ordered and (2) selecting the runtime mechanism that recognizes completed groups.

Group completion requires knowing that a group cannot reappear after its current run ends. Sorted input supplies that guarantee, and grouped input can supply the same guarantee without imposing an order on the group values:

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

This PR represents that execution capability independently so later planning work can derive it from grouped equivalence properties.

What changes are included in this PR?

  • Introduce the crate-private GroupCompletionMode::{None, Partial, Full} variants.
  • Derive GroupCompletionMode from InputOrderMode during AggregateExec construction.
  • Pass GroupCompletionMode through aggregate tables, streams, and spill replay paths.
  • Use GroupCompletionMode when selecting hash and ordered aggregate implementations.
  • Keep InputOrderMode as the source of required-input and output-ordering metadata.

The initial conversion is:

InputOrderMode::Linear                    -> GroupCompletionMode::None
InputOrderMode::PartiallySorted(indices) -> GroupCompletionMode::Partial(indices)
InputOrderMode::Sorted                    -> GroupCompletionMode::Full

Stack

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

Are these changes tested?

Aggregate planning tests cover the None, Partial, and Full conversions and retain the behavior characterized by #24737.

Are there any user-facing changes?

No. GroupCompletionMode is crate-private, and the public GroupOrdering constructor remains unchanged.

Review this layer

View only this PR layer

@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
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.

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 [  54.413s] (current)
     Parsing datafusion v55.1.0 (current)
      Parsed [   0.032s] (current)
    Building datafusion v55.1.0 (baseline)
       Built [  53.827s] (baseline)
     Parsing datafusion v55.1.0 (baseline)
      Parsed [   0.031s] (baseline)
    Checking datafusion v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.661s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [ 111.050s] datafusion
    Building datafusion-physical-plan v55.1.0 (current)
       Built [  35.908s] (current)
     Parsing datafusion-physical-plan v55.1.0 (current)
      Parsed [   0.177s] (current)
    Building datafusion-physical-plan v55.1.0 (baseline)
       Built [  36.020s] (baseline)
     Parsing datafusion-physical-plan v55.1.0 (baseline)
      Parsed [   0.176s] (baseline)
    Checking datafusion-physical-plan v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.626s] 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/705d2e66e6058a9682221d3e6bffabc69ba10efe/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/705d2e66e6058a9682221d3e6bffabc69ba10efe/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/705d2e66e6058a9682221d3e6bffabc69ba10efe/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/705d2e66e6058a9682221d3e6bffabc69ba10efe/datafusion/physical-plan/src/aggregates/order/partial.rs:66

     Summary semver requires new major version: 3 major and 0 minor checks failed
    Finished [  74.084s] datafusion-physical-plan
    Building datafusion-sqllogictest v55.1.0 (current)
       Built [  92.461s] (current)
     Parsing datafusion-sqllogictest v55.1.0 (current)
      Parsed [   0.015s] (current)
    Building datafusion-sqllogictest v55.1.0 (baseline)
       Built [  92.353s] (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.109s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [ 187.461s] datafusion-sqllogictest

@github-actions github-actions Bot added the auto detected api change Auto detected API change label Oct 4, 2026
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.

6 participants