2010YOUY01 commented on code in PR #26069:
URL: https://github.com/apache/datafusion/pull/26069#discussion_r4191240086


##########
datafusion/physical-optimizer/src/limited_distinct_aggregation.rs:
##########
@@ -44,95 +86,70 @@ impl LimitedDistinctAggregation {
         Self {}
     }
 
-    fn transform_agg(
-        aggr: &AggregateExec,
-        limit: usize,
-    ) -> Option<Transformed<Arc<dyn ExecutionPlan>>> {
-        let new_aggr = aggr.clone().try_optimize_distinct_soft_limit(limit)?;
-        // An already limited aggregate still permits optimizing its partial 
child.
-        Some(new_aggr.update_data(|aggr| Arc::new(aggr) as Arc<dyn 
ExecutionPlan>))
-    }
+    /// Rewrite a limit and its immediately adjacent final/partial aggregate 
pair.
+    fn transform_limit(
+        plan: Arc<dyn ExecutionPlan>,
+    ) -> Result<Transformed<Arc<dyn ExecutionPlan>>> {
+        // Step 1: Identify the plan shape,
+        //
+        // Limit
+        //   Aggregate(final)
+        //     Aggregate(partial)
+
+        // Check the current plan is limit, and extract limit value, input 
plan.
+        let (limit, input) = match (
+            plan.downcast_ref::<LocalLimitExec>(),
+            plan.downcast_ref::<GlobalLimitExec>(),
+        ) {
+            (Some(local), _) => (local.fetch(), local.input()),
+            (_, Some(global)) => match global.fetch() {
+                Some(fetch) => (global.skip() + fetch, global.input()),
+                None => return Ok(Transformed::no(plan)),
+            },
+            _ => return Ok(Transformed::no(plan)),
+        };
 
-    /// transform_limit matches an `AggregateExec` as the child of a 
`LocalLimitExec`
-    /// or `GlobalLimitExec` and pushes the limit into the aggregation as a 
soft limit when
-    /// there is a group by, but no sorting, no aggregate expressions, and no 
filters in the
-    /// aggregation
-    fn transform_limit(plan: Arc<dyn ExecutionPlan>) -> Option<Arc<dyn 
ExecutionPlan>> {
-        let limit: usize;
-        let mut global_fetch: Option<usize> = None;
-        let mut global_skip: usize = 0;
-        let children: Vec<Arc<dyn ExecutionPlan>>;
-        let mut is_global_limit = false;
-        if let Some(local_limit) = plan.downcast_ref::<LocalLimitExec>() {
-            limit = local_limit.fetch();
-            children = local_limit.children().into_iter().cloned().collect();
-        } else {
-            let global_limit = plan.downcast_ref::<GlobalLimitExec>()?;
-            global_fetch = global_limit.fetch();
-            global_fetch?;
-            global_skip = global_limit.skip();
-            // the aggregate must read at least fetch+skip number of rows
-            limit = global_fetch.unwrap() + global_skip;
-            children = global_limit.children().into_iter().cloned().collect();
-            is_global_limit = true
-        }
-        let child = children.iter().exactly_one().ok()?;
-        // ensure there is no output ordering; can this rule be relaxed?
-        if plan.output_ordering().is_some() {
-            return None;
-        }
-        // ensure no ordering is required on the input
-        if plan.required_input_ordering()[0].is_some() {
-            return None;
+        if plan.output_ordering().is_some() || 
plan.required_input_ordering()[0].is_some()
+        {
+            return Ok(Transformed::no(plan));
         }
 
-        // if found_match_aggr is true, match_aggr holds a parent aggregation 
whose group_by
-        // must match that of a child aggregation in order to rewrite the 
child aggregation
-        let mut match_aggr: Arc<dyn ExecutionPlan> = plan;
-        let mut found_match_aggr = false;
-
-        let mut rewrite_applicable = true;
-        let closure = |plan: Arc<dyn ExecutionPlan>| {
-            if !rewrite_applicable {
-                return Ok(Transformed::no(plan));
-            }
-            if let Some(aggr) = plan.downcast_ref::<AggregateExec>() {
-                if found_match_aggr
-                    && let Some(parent_aggr) = 
match_aggr.downcast_ref::<AggregateExec>()
-                    && !parent_aggr.group_expr().eq(aggr.group_expr())

Review Comment:
   bug fix: missing a `as_final` to adapt projection difference



##########
datafusion/physical-optimizer/src/limited_distinct_aggregation.rs:
##########
@@ -44,95 +86,70 @@ impl LimitedDistinctAggregation {
         Self {}
     }
 
-    fn transform_agg(
-        aggr: &AggregateExec,
-        limit: usize,
-    ) -> Option<Transformed<Arc<dyn ExecutionPlan>>> {
-        let new_aggr = aggr.clone().try_optimize_distinct_soft_limit(limit)?;
-        // An already limited aggregate still permits optimizing its partial 
child.
-        Some(new_aggr.update_data(|aggr| Arc::new(aggr) as Arc<dyn 
ExecutionPlan>))
-    }
+    /// Rewrite a limit and its immediately adjacent final/partial aggregate 
pair.
+    fn transform_limit(
+        plan: Arc<dyn ExecutionPlan>,
+    ) -> Result<Transformed<Arc<dyn ExecutionPlan>>> {
+        // Step 1: Identify the plan shape,
+        //
+        // Limit
+        //   Aggregate(final)
+        //     Aggregate(partial)
+
+        // Check the current plan is limit, and extract limit value, input 
plan.
+        let (limit, input) = match (
+            plan.downcast_ref::<LocalLimitExec>(),
+            plan.downcast_ref::<GlobalLimitExec>(),
+        ) {
+            (Some(local), _) => (local.fetch(), local.input()),
+            (_, Some(global)) => match global.fetch() {
+                Some(fetch) => (global.skip() + fetch, global.input()),
+                None => return Ok(Transformed::no(plan)),
+            },
+            _ => return Ok(Transformed::no(plan)),
+        };
 
-    /// transform_limit matches an `AggregateExec` as the child of a 
`LocalLimitExec`
-    /// or `GlobalLimitExec` and pushes the limit into the aggregation as a 
soft limit when
-    /// there is a group by, but no sorting, no aggregate expressions, and no 
filters in the
-    /// aggregation
-    fn transform_limit(plan: Arc<dyn ExecutionPlan>) -> Option<Arc<dyn 
ExecutionPlan>> {
-        let limit: usize;
-        let mut global_fetch: Option<usize> = None;
-        let mut global_skip: usize = 0;
-        let children: Vec<Arc<dyn ExecutionPlan>>;
-        let mut is_global_limit = false;
-        if let Some(local_limit) = plan.downcast_ref::<LocalLimitExec>() {
-            limit = local_limit.fetch();
-            children = local_limit.children().into_iter().cloned().collect();
-        } else {
-            let global_limit = plan.downcast_ref::<GlobalLimitExec>()?;
-            global_fetch = global_limit.fetch();
-            global_fetch?;
-            global_skip = global_limit.skip();
-            // the aggregate must read at least fetch+skip number of rows
-            limit = global_fetch.unwrap() + global_skip;
-            children = global_limit.children().into_iter().cloned().collect();
-            is_global_limit = true
-        }
-        let child = children.iter().exactly_one().ok()?;
-        // ensure there is no output ordering; can this rule be relaxed?
-        if plan.output_ordering().is_some() {
-            return None;
-        }
-        // ensure no ordering is required on the input
-        if plan.required_input_ordering()[0].is_some() {
-            return None;
+        if plan.output_ordering().is_some() || 
plan.required_input_ordering()[0].is_some()
+        {
+            return Ok(Transformed::no(plan));
         }
 
-        // if found_match_aggr is true, match_aggr holds a parent aggregation 
whose group_by
-        // must match that of a child aggregation in order to rewrite the 
child aggregation
-        let mut match_aggr: Arc<dyn ExecutionPlan> = plan;
-        let mut found_match_aggr = false;
-
-        let mut rewrite_applicable = true;
-        let closure = |plan: Arc<dyn ExecutionPlan>| {
-            if !rewrite_applicable {
-                return Ok(Transformed::no(plan));
-            }
-            if let Some(aggr) = plan.downcast_ref::<AggregateExec>() {
-                if found_match_aggr
-                    && let Some(parent_aggr) = 
match_aggr.downcast_ref::<AggregateExec>()
-                    && !parent_aggr.group_expr().eq(aggr.group_expr())
-                {
-                    // a partial and final aggregation with different 
groupings disqualifies
-                    // rewriting the child aggregation
-                    rewrite_applicable = false;
-                    return Ok(Transformed::no(plan));
-                }
-                // either we run into an Aggregate and transform it, or 
disable the rewrite
-                // for subsequent children
-                match Self::transform_agg(aggr, limit) {
-                    None => {}
-                    Some(new_aggr) => {
-                        match_aggr = plan;
-                        found_match_aggr = true;
-                        return Ok(new_aggr);
-                    }
-                }
-            }
-            rewrite_applicable = false;
-            Ok(Transformed::no(plan))
+        let Some(final_agg) = input.downcast_ref::<AggregateExec>() else {
+            return Ok(Transformed::no(plan));
         };
-        let child = child.to_owned().transform_down(closure).ok()?;
-        if !child.transformed {
-            return None;
+        let Some(partial_agg) = 
final_agg.input().downcast_ref::<AggregateExec>() else {
+            return Ok(Transformed::no(plan));
+        };
+
+        // Step 2: Verify partial/final aggregate is compatible, then apply 
optimization
+        match (final_agg.mode(), partial_agg.mode()) {
+            (
+                AggregateMode::Final | AggregateMode::FinalPartitioned,
+                AggregateMode::Partial,
+            ) if final_agg.group_expr() == 
&partial_agg.group_expr().as_final() => {}

Review Comment:
   👉🏼 



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