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]

Reply via email to