Repository navigation
Conversation
f58e128 to
372429f
Compare
d1d1d2f to
0267ad1
Compare
gene-bordegaray
left a comment
There was a problem hiding this comment.
approved with some non blocking suggestoins. Thank you @xavlee
e7f8b42 to
0108b7c
Compare
|
Rereviewed after changes and everything looks good 👍 |
NGA-TRAN
left a comment
There was a problem hiding this comment.
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. |
There was a problem hiding this comment.
Does this mean I have 2 keys (a, b) and data is sorted on (a) only?
| assert_eq!( | ||
| aggregate.group_completion_mode, | ||
| GroupCompletionMode::Partial(vec![0]) | ||
| ); |
| // 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); |
0108b7c to
6e369ce
Compare
2fae45d to
bb3dc16
Compare
bb3dc16 to
64a6e55
Compare
gene-bordegaray
left a comment
There was a problem hiding this comment.
this also looks good, just needs a rebase
3a0ecb6 to
1fd156e
Compare
xudong963
left a comment
There was a problem hiding this comment.
Nit: finish updating the runtime documentation from “ordered” to “group-complete.”
## 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
left a comment
There was a problem hiding this comment.
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
GroupOrderingtoGroupCompletion, 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.
| 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, |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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); |
There was a problem hiding this comment.
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( |
There was a problem hiding this comment.
marker for previous comment: this is a completion-mode -> ordering conversion.
1fd156e to
d112ecc
Compare
|
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 |
|
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:
Let me know if this is what you had in mind One naming question: would you prefer to rename |
|
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
Yes, I agree |
|
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? |
There was a problem hiding this comment.
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
Which issue does this PR relate to?
Rationale for this change
Aggregate execution currently uses
InputOrderModeto 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:
This PR describes the input's clustering guarantee with
GroupClusteringModeand 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?
InputOrderModefield withGroupClusteringMode::{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.GroupOrderingtracker and itsPartialandFullimplementations toGroupClustering,GroupClusteringPartial, andGroupClusteringFull.group_clustering_modein aggregate plans and document the public API migration in the DataFusion 56 upgrade guide.Stack
Are there any user-facing changes?
Yes. This PR changes public Rust APIs:
AggregateExec::input_order_mode()is replaced byAggregateExec::group_clustering_mode(), returning the publicGroupClusteringModeenum.AggregateExec::compute_propertiesaccepts&GroupClusteringModein place of&InputOrderMode.GroupOrdering,GroupOrderingPartial, andGroupOrderingFulltypes inaggregates::orderare renamed toGroupClustering,GroupClusteringPartial, andGroupClusteringFull.GroupClustering::try_newtakes&GroupClusteringMode;new_group_valuestakes&GroupClustering.Aggregate plan output uses
group_clustering_mode=Fullandgroup_clustering_mode=Partial(indices)in place ofordering_mode=Sortedandordering_mode=PartiallySorted(indices). SQL results and aggregate execution-path selection retain their existing behavior.Review this layer
View this PR's diff