2010YOUY01 commented on code in PR #25983:
URL: https://github.com/apache/datafusion/pull/25983#discussion_r4172117375
##########
datafusion/physical-optimizer/src/combine_partial_final_agg.rs:
##########
@@ -103,36 +102,23 @@ impl PhysicalOptimizerRule for
CombinePartialFinalAggregate {
};
// Re-apply distinct limit optimization.
- //
+ //
// The optimizer related to aggregates are (in order):
// - 1. Initial planning: always final/partial two stage
// - 2. `LimitedDistinctAggregation`: push limit into
`AggregateExec`,
// and build `AggregateKind::DistinctLimit`
// - 3. `CombinePartialFinalAggregate`: the current
optimization
//
- // Here it restores previously applied optimization in step 2
- let combined_aggr_kind = match (&agg_exec.kind,
&input_agg_exec.kind) {
- (AggregateKind::General { .. }, AggregateKind::General {
.. }) => {
- combined_agg.kind
- }
- (
- AggregateKind::DistinctLimit { limit, .. },
- AggregateKind::DistinctLimit {
- group_by,
- limit: partial_limit,
- },
- ) if limit == partial_limit =>
AggregateKind::DistinctLimit {
- // The combined aggregate groups raw input, like the
partial.
- group_by: Arc::clone(group_by),
- limit: *limit,
- },
- _ => {
Review Comment:
The alternative fix for the internal error is to let this branch return not
transformed, though it misses some optimization oppurtunity.
##########
datafusion/physical-optimizer/src/combine_partial_final_agg.rs:
##########
@@ -103,36 +102,23 @@ impl PhysicalOptimizerRule for
CombinePartialFinalAggregate {
};
// Re-apply distinct limit optimization.
- //
+ //
// The optimizer related to aggregates are (in order):
// - 1. Initial planning: always final/partial two stage
// - 2. `LimitedDistinctAggregation`: push limit into
`AggregateExec`,
// and build `AggregateKind::DistinctLimit`
// - 3. `CombinePartialFinalAggregate`: the current
optimization
//
- // Here it restores previously applied optimization in step 2
- let combined_aggr_kind = match (&agg_exec.kind,
&input_agg_exec.kind) {
- (AggregateKind::General { .. }, AggregateKind::General {
.. }) => {
- combined_agg.kind
- }
- (
- AggregateKind::DistinctLimit { limit, .. },
- AggregateKind::DistinctLimit {
- group_by,
- limit: partial_limit,
- },
- ) if limit == partial_limit =>
AggregateKind::DistinctLimit {
- // The combined aggregate groups raw input, like the
partial.
- group_by: Arc::clone(group_by),
- limit: *limit,
- },
- _ => {
- return internal_err!(
- "The AggregateKind should stay either (a) General
(b) DistinctLimit introduced with the previous optimizer pass
`LimitedDistinctAggregation`, it's impossible to have other variant"
- );
- }
- };
- combined_agg.kind = combined_aggr_kind;
+ // Here it restores previously applied optimization in step 2.
The
+ // final aggregate decides: step 2 may have limited only the
final
+ // aggregate, e.g. when its grouping differs from the
partial's.
+ if let AggregateKind::DistinctLimit { limit, .. } =
&agg_exec.kind
Review Comment:
This fix looks like patching another bug in previous optimizer rule
`LimitedDistinctAggregation`
Its contract was
```
# Before
Limit(k)
-- AggregateExec(final)
---- AggregateExec(partial)
# After
Limit(k)
-- AggregateExec(final, limit=k)
---- AggregateExec(partial, limit=k)
```
And this fix seem to accept the mixed case, I'm not sure if there is some
workload can get optimized to this shape, and also combining it to a
`AggregateExec(single, limit=k)` is safe
```
# After
Limit(k)
-- AggregateExec(final, limit=k) <-- with limit
---- AggregateExec(partial) <-- no limit
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]