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]