Repository navigation
refactor: Simplify LimitedDistinctAggregation optimizer rule, push limit in more cases - #26069
Conversation
| 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()) |
There was a problem hiding this comment.
bug fix: missing a as_final to adapt projection difference
| ( | ||
| AggregateMode::Final | AggregateMode::FinalPartitioned, | ||
| AggregateMode::Partial, | ||
| ) if final_agg.group_expr() == &partial_agg.group_expr().as_final() => {} |
There was a problem hiding this comment.
Maybe this check could be a method on AggregateExec -- mostly so that it could be given a name and better documented
Something like:
if final_agg.matches_partial(partial_agg) {
...
}(not sure if that is a good name for it)
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #26069 +/- ##
==========================================
+ Coverage 82.66% 82.72% +0.06%
==========================================
Files 1147 1147
Lines 446357 449144 +2787
Branches 446357 449144 +2787
==========================================
+ Hits 368971 371555 +2584
+ Misses 54997 54940 -57
- Partials 22389 22649 +260 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
alamb
left a comment
There was a problem hiding this comment.
Thank you @2010YOUY01 -- I think this is a nice cleanup and improves the output plans and the code readability
| 02)--AggregateExec: mode=Final, gby=[b@0 as b, a@1 as a], aggr=[], lim=[2] | ||
| 03)----AggregateExec: mode=Partial, gby=[b@1 as b, a@0 as a], aggr=[] | ||
| 04)------DataSourceExec: partitions=1, partition_sizes=[1] | ||
| 02)--AggregateExec: mode=Single, gby=[b@1 as b, a@0 as a], aggr=[], lim=[2] |
There was a problem hiding this comment.
this does seem a better plan (use only the final single partition execution plan)
| 03)----AggregateExec: mode=Final, gby=[c2@0 as c2, c3@1 as c3, __grouping_id@2 as __grouping_id], aggr=[], lim=[3] | ||
| 04)------CoalescePartitionsExec | ||
| 05)--------AggregateExec: mode=Partial, gby=[(NULL as c2, NULL as c3), (c2@0 as c2, NULL as c3), (c2@0 as c2, c3@1 as c3)], aggr=[] | ||
| 05)--------AggregateExec: mode=Partial, gby=[(NULL as c2, NULL as c3), (c2@0 as c2, NULL as c3), (c2@0 as c2, c3@1 as c3)], aggr=[], lim=[3] |
There was a problem hiding this comment.
This looks correct to me -- to push the limit down into each partial aggregate / partition
| /// Scan | ||
| /// ``` | ||
| /// | ||
| /// # Invariants before and after the rewrite |
There was a problem hiding this comment.
I am not sure I would use the term "invariant" here though it is technically accurate
This comment mostly explains the effect of this optimizer rule (what it does).
I normally think of Invariant as a property that will not change as any operation is applied to it. I don't think there is any reason to require that some future optimizer rule preserves the same property
For example, what if a future rule (or a user defined rule) has some additional special operator that requires the full intermediate results (aka undoes the soft limit on the Partial AggregateExec)? I don't see any reason to try and prevent that 🤔
There was a problem hiding this comment.
Indeed, 'invariant' is not precise here.
Here is the updated version in c8986fd
It should describe what it is today, but I think we could make this kind of assumption/promise description part of a template, require it as a comment section for all optimizer rules, and try to keep it maintained.
These assumptions are already there, but many of them are undocumented. I’m not sure exactly how we should do this yet, so I’ll keep thinking about it in the background.
/// # What this rule assumes
///
/// This rule assumes the logical aggregate only have one shape showed below, this
/// is what the current physical planning produces.
///
/// ```txt
/// Limit
/// AggregateExec(mode=Final)
/// AggregateExec(mode=Partial)
/// ```
///
/// If future changes or extensions produce a different shape, this rule skips the
/// rewrite rather than reporting an error, potentially missing an optimization
/// opportunity.
///
/// # What this rule promises
///
/// Immediately after an eligible rewrite, both stages have a soft-limit hint.
/// Otherwise, this rule leaves both stages unchanged.
///
/// If a later rule removes the limit, it won't affect correctness, but it may
/// miss an optimization opportunity.| fn transform_limit( | ||
| plan: Arc<dyn ExecutionPlan>, | ||
| ) -> Result<Transformed<Arc<dyn ExecutionPlan>>> { | ||
| // Step 1: Identify the plan shape, |
There was a problem hiding this comment.
thank you -- these comments really make the code much easier to follow
| ( | ||
| AggregateMode::Final | AggregateMode::FinalPartitioned, | ||
| AggregateMode::Partial, | ||
| ) if final_agg.group_expr() == &partial_agg.group_expr().as_final() => {} |
There was a problem hiding this comment.
Maybe this check could be a method on AggregateExec -- mostly so that it could be given a name and better documented
Something like:
if final_agg.matches_partial(partial_agg) {
...
}(not sure if that is a good name for it)
LimitedDistinctAggregation optimizer ruleLimitedDistinctAggregation optimizer rule, push limit in more cases
|
run benchmark sql_planner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing cleanup-limit-aggr (784d746) to 982fca6 (merge-base) diff Run configurationrun benchmark sql_plannerResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing cleanup-limit-aggr (784d746) to 982fca6 (merge-base) diff Run configurationrun benchmark sql_plannerCPU Details (lscpu)Details
Resource Usagesql_planner — base (merge-base)
sql_planner — branch
File an issue against this benchmark runner |
|
run benchmarks sql_planner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing cleanup-limit-aggr (784d746) to 982fca6 (merge-base) diff Run configurationrun benchmark sql_plannerResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing cleanup-limit-aggr (784d746) to 982fca6 (merge-base) diff Run configurationrun benchmark sql_plannerCPU Details (lscpu)Details
Resource Usagesql_planner — base (merge-base)
sql_planner — branch
File an issue against this benchmark runner |
| /// ``` | ||
| /// | ||
| /// # Invariants before and after the rewrite | ||
| /// # What this rule assumes |
|
Thank you @2010YOUY01 |
Which issue does this PR close?
This is primarily a simplifying refactor, with an optimizer regression fix piggy-backed on top (an inefficient plan shape, not a correctness issue):
The optimizer fix itself is only a one-line change. If this simplification is not desirable, I'll open a replacement PR containing just the direct fix.
Rationale for this change
Part 1: Optimizer fix
See sqllogictest diff for the reproducer
I'll mark the fix in comments
Part 2: Optimizer simplification
During different stages in physical optimization, a logical aggregation can be either
So the existing implementation is using a nested dfs (the closure has another dfs inside, to match non-consecutive aggregates) try to match all of the 3 cases.
However, after looking at related optimizer rules:
We can find only case 1 is possible, so the implementation can be simplified into a naive pattern matching
What changes are included in this PR?
What is the testing strategy for this PR?
sltAre there any user-facing changes?