alamb commented on code in PR #25696:
URL: https://github.com/apache/datafusion/pull/25696#discussion_r4143878776


##########
datafusion/physical-plan/src/aggregates/mod.rs:
##########
@@ -901,6 +927,202 @@ pub struct AggregateExec {
 }
 
 impl AggregateExec {
+    /// Try to use TopK (min/max heap) optimization in AggregateExec.
+    ///
+    /// If applicable, an inner `AggregateKind` will be set, and later 
[`ExecutionPlan::execute`]

Review Comment:
   I think the note about AggregateKind is somewhat too detailed as it refers 
to an implementation detail -- I think just saying that "the optimizer must 
keep the limit operator above because the aggregate may still produce more than 
K rows" is what is important.



##########
datafusion/physical-optimizer/src/topk_aggregation.rs:
##########
@@ -132,25 +73,21 @@ impl TopKAggregation {
                 return Ok(Transformed::no(plan));
             }
             if let Some(aggr) = plan.downcast_ref::<AggregateExec>() {
-                // either we run into an Aggregate and transform it
-                match Self::transform_agg(
-                    aggr,
-                    &cur_col_name,
-                    order_desc,
-                    nulls_first,
-                    limit,
-                ) {
+                match aggr
+                    .clone()
+                    .try_optimize_topk(limit, &sort_col_name, sort_options)

Review Comment:
   it is a very nice idea to move this into the AggregateExec itself



##########
datafusion/physical-plan/src/aggregates/mod.rs:
##########
@@ -952,25 +1175,42 @@ impl AggregateExec {
         let mut new = self.clone();
         match &mut new.kind {
             AggregateKind::General { aggr_expr: old, .. } => *old = aggr_expr,
-            AggregateKind::DistinctLimit { .. } if aggr_expr.is_empty() => {}
-            AggregateKind::DistinctLimit { group_by, .. } => {
-                // An accumulator rewrite cannot inherit DISTINCT's early stop.
+            AggregateKind::DistinctLimit { .. } | AggregateKind::TopKDistinct 
{ .. }
+                if aggr_expr.is_empty() => {}
+            // Ordering optimization can revisit an existing TopK with the
+            // same expressions. Preserve its specialization in that case.
+            AggregateKind::TopKMinMax {
+                aggr_expr: existing,
+                ..
+            } if matches!(aggr_expr.as_ref(), [new] if Arc::ptr_eq(existing, 
new)) => {}
+            AggregateKind::DistinctLimit { group_by, .. }
+            | AggregateKind::TopKMinMax { group_by, .. }
+            | AggregateKind::TopKDistinct { group_by, .. } => {
+                // A replacement expression cannot inherit a specialization
+                // validated for the previous aggregate.
                 new.kind = AggregateKind::General {
                     group_by: Arc::clone(group_by),
-                    // The previous `DistinctLimit` type doesn't include filter
                     filter_expr: vec![None; aggr_expr.len()].into(),
                     aggr_expr,
-                    limit_options: None,
                 };
             }
         }
         new.metrics = ExecutionPlanMetricsSet::new();
         new
     }
 
-    /// Clone this exec, overriding only the limit hint.
+    /// Clone with a validated legacy limit hint, falling back to ordinary
+    /// aggregation when the request is unsupported.
+    #[deprecated(
+        since = "56.0.0",
+        note = "This API is intended for internal use only and was 
inadvertently made public. Do not use this API."

Review Comment:
   is there an alternate API we should direct them to instead? Maybe 
`try_optimize_topk` should be public?
   
   It seems strange to say it wasn't meant to be public...



##########
datafusion/physical-plan/src/aggregates/mod.rs:
##########
@@ -846,30 +847,55 @@ impl LimitOptions {
     }
 }
 
-/// Aggregation state, separating a DISTINCT soft limit from accumulators and 
filters.
+/// Mutually exclusive aggregation implementations and their configuration.
+///
+/// # Public Only for Internal Use:

Review Comment:
   Yes, I agree making a public API on AggregateExec would be better. Maybe it 
is worth filing a ticket to track that idea



##########
datafusion/physical-plan/src/aggregates/mod.rs:
##########
@@ -846,30 +847,55 @@ impl LimitOptions {
     }
 }
 
-/// Aggregation state, separating a DISTINCT soft limit from accumulators and 
filters.
+/// Mutually exclusive aggregation implementations and their configuration.
+///
+/// # Public Only for Internal Use:
+/// `datafusion-physical-optimizer` inspects and combines aggregate kinds.
+/// This enum is not part of the supported public API.
+#[doc(hidden)]
 #[derive(Debug, Clone)]
-enum AggregateKind {
-    /// Ordinary aggregation, including the existing Top-K configuration.
+pub enum AggregateKind {
+    /// Ordinary aggregation, with no limit on the groups retained.
     General {
         group_by: Arc<PhysicalGroupBy>,
         aggr_expr: Arc<[Arc<AggregateFunctionExpr>]>,
         filter_expr: Arc<[Option<Arc<dyn PhysicalExpr>>]>,
-        limit_options: Option<LimitOptions>,
     },
     /// `SELECT DISTINCT k FROM t LIMIT n`: eligible streams may stop after n 
groups.
     /// Other streams consume all input; the parent LIMIT enforces the row 
count.
     DistinctLimit {
         group_by: Arc<PhysicalGroupBy>,
         limit: usize,
     },
+    /// See [`AggregateExec::try_optimize_topk`] for details.
+    TopKMinMax {
+        group_by: Arc<PhysicalGroupBy>,
+        aggr_expr: Arc<AggregateFunctionExpr>,
+        limit: usize,
+        descending: bool,
+        nulls_first: bool,
+    },
+    /// See [`AggregateExec::try_optimize_topk`] for details.
+    TopKDistinct {
+        group_by: Arc<PhysicalGroupBy>,
+        limit: usize,
+        descending: bool,
+    },
 }
 
 /// Hash aggregate execution plan
 #[derive(Debug, Clone)]
 pub struct AggregateExec {
     /// Aggregation mode (full, partial)
     mode: AggregateMode,
-    kind: AggregateKind,
+    /// Aggregation implementation and its configuration.
+    ///
+    /// # Public Only for Internal Use:
+    /// `datafusion-physical-optimizer` updates this when combining aggregates.
+    /// Changes must preserve the expressions, schema, and plan properties.

Review Comment:
   can we hide it behind an accessor? 
   
   Maybe that could be part of a follow on PR to better encapsulate this struct 
(aka remove `pub` from all the fields)



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