Skip to content

fix: keep aggregate limits when combining or rebuilding AggregateExec - #25983

Merged
2010YOUY01 merged 2 commits into
apache:mainfrom
jayzhan211:fix/agg-limit-carryover
Oct 5, 2026
Merged

2010YOUY01 merged 2 commits into
apache:mainfrom
jayzhan211:fix/agg-limit-carryover

Conversation

@jayzhan211

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Rationale for this change

  1. Some GROUP BY ... LIMIT queries fail to plan. With one target partition:

    SET datafusion.execution.target_partitions = 1;
    CREATE TABLE t AS SELECT * FROM (VALUES (1, 10), (2, 20), (3, 30), (4, 40), (5, 50), (6, 60)) AS v(a, b);
    SELECT b, a FROM t GROUP BY b, a LIMIT 2;
    SELECT a % 3 FROM t GROUP BY a % 3 LIMIT 2;

    both fail with:

    DataFusion error: CombinePartialFinalAggregate
    caused by
    Internal error: The AggregateKind should stay either (a) General (b) DistinctLimit introduced with the previous optimizer pass `LimitedDistinctAggregation`, it's impossible to have other variant.
    

    When the final aggregate's grouping differs from the partial's (expression keys, or columns not at their output positions), LimitedDistinctAggregation limits only the final aggregate. CombinePartialFinalAggregate treats that pair as impossible and returns an error.

  2. On Parquet tables the TopK aggregate limit is dropped. EXPLAIN SELECT DISTINCT k FROM pt ORDER BY k LIMIT 10 keeps lim=[10] on the aggregates for a memory table but not for a Parquet table, so the query falls back to a full hash aggregation. MIN/MAX TopK aggregates (ORDER BY max(x) DESC NULLS LAST LIMIT n) lose their limit the same way. Post-optimization filter pushdown rebuilds the Parquet scan when it accepts a dynamic filter, and rebuilding the aggregates above it re-applies only the DISTINCT limit, not the TopK ones.

Before #25696 both places carried the limit over with with_limit_options(...limit_options()).

What changes are included in this PR?

  • CombinePartialFinalAggregate: if the final aggregate has a DISTINCT limit, re-apply it to the combined aggregate with try_optimize_distinct_soft_limit, which checks eligibility again. Otherwise keep the combined aggregate as built.
  • AggregateExec::replace_children (ChildrenPropertiesMode::Recompute): also re-apply TopKMinMax and TopKDistinct through try_optimize_topk, using the stored limit, direction and NULL placement.

What is the testing strategy for this PR?

New sqllogictest cases. Each one fails without the fix and passes with it:

  • aggregate.slt: both queries above at target_partitions = 1. An EXPLAIN checks the combined mode=Single ... lim=[2] aggregate, and a row-count check keeps the results deterministic.
  • aggregates_topk.slt: a Parquet-backed table. EXPLAIN checks that lim=[10] stays on both aggregates for SELECT DISTINCT k ... ORDER BY k LIMIT 10, and that lim=[2] stays for a MAX TopK query whose scan gets a semi-join dynamic filter. Result checks are included.

Are there any user-facing changes?

No API changes. The listed queries plan again, and TopK aggregation applies again on Parquet scans that receive dynamic filters.

@github-actions github-actions Bot added optimizer Optimizer rules sqllogictest SQL Logic Tests (.slt) physical-plan Changes to the physical-plan crate labels Oct 3, 2026
@jayzhan211
jayzhan211 marked this pull request as ready for review October 3, 2026 02:01
@jayzhan211
jayzhan211 requested a review from 2010YOUY01 October 3, 2026 02:01
@codecov-commenter

codecov-commenter commented Oct 3, 2026 •

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 82.64%. Comparing base (801b017) to head (290c3fc).
⚠️ Report is 7 commits behind head on main.

Additional details and impacted files
@@            Coverage Diff             @@
##             main   #25983      +/-   ##
==========================================
+ Coverage   82.62%   82.64%   +0.01%     
==========================================
  Files        1147     1147              
  Lines      445320   445520     +200     
  Branches   445320   445520     +200     
==========================================
+ Hits       367964   368184     +220     
+ Misses      55013    54975      -38     
- Partials    22343    22361      +18     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@2010YOUY01 2010YOUY01 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for the fix!

For two issues:

  1. It seem like a deeper bug in another optimizer rule, and I suggested a simpler fix in comment
  2. LGTM

So I suggest to use the simpler fix here for issue 1, and this PR should be good to go, and we can fix the remaining problem in a follow-up PR.

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

group_by: Arc::clone(group_by),
limit: *limit,
},
_ => {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The alternative fix for the internal error is to let this branch return not transformed, though it misses some optimization oppurtunity.

@jayzhan211
jayzhan211 requested a review from 2010YOUY01 October 3, 2026 14:46
@jayzhan211

Copy link
Copy Markdown
Contributor Author

Comments addressed!

@2010YOUY01 2010YOUY01 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

thanks again

@2010YOUY01

Copy link
Copy Markdown
Contributor
  1. It seem like a deeper bug in another optimizer rule, and I suggested a simpler fix in comment

I got a fix for it, I will open a PR after this PR has merged to reuse tests.

I'll merge this PR tomorrow.

@2010YOUY01
2010YOUY01 added this pull request to the merge queue Oct 5, 2026
Merged via the queue into apache:main with commit 982fca6 Oct 5, 2026
42 checks passed
discord9 pushed a commit to discord9/datafusion that referenced this pull request Oct 8, 2026
…limit in more cases (apache#26069)

## Which issue does this PR close?

<!--
We generally require a GitHub issue to be filed for all bug fixes and
enhancements and this helps us generate change logs for our releases.
You can link an issue to this PR using the GitHub syntax. For example
`Closes apache#123` indicates that this PR will close issue apache#123.
-->

This is primarily a simplifying refactor, with an optimizer regression
fix piggy-backed on top (an inefficient plan shape, not a correctness
issue):

-
apache#25983 (comment)

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

<!--
Why are you proposing this change? If this is already explained clearly
in the issue then this section is not needed.
Explaining clearly why changes are proposed helps reviewers understand
your changes and offer better suggestions for fixes.

Please explain the problem you are trying to solve in terms of the
user-visible
behavior, rather than the implementation.

For example, "The code in `foo.rs` doesn't handle nulls" is a symptom of
the
implementation. "COUNT(DISTINCT) returns wrong results when the column
contains
nulls" is the user-visible problem.
-->
## 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
```
AggregateExec(mode=final)
--AggregateExec(mode=partial)

AggregateExec(mode=single)

AggregateExec(mode=final)
--RepartitionExec
----AggregateExec(mode=partial)
```

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:
- initial physical planning
- CombinePartialFinalAggregate
- EnsureDistributions
- and this rule

We can find only case 1 is possible, so the implementation can be
simplified into a naive pattern matching
```
Find the exact below shape, and try push limits

LimitExec
--AggregateExec(mode=final)
----AggregateExec(mode=partial)
```

## What changes are included in this PR?

<!--
There is no need to duplicate the description in the issue here, but it
is sometimes worth providing a summary of the individual changes in this
PR.
-->
- simplify an optimizer rule according to the above rationale
- fix a small bug

## What is the testing strategy for this PR?

<!--
We typically require tests for all PRs in order to:
1. Prevent the code from being accidentally broken by subsequent changes
2. Serve as another way to document the expected behavior of the code

Briefly describe how this PR is tested, and point to the specific tests
you added. For example: 'This new feature is covered by the
`sqllogictest` cases added in `foo.slt`'.

If this PR does not add tests, explain why. For example, if the change
is already covered by existing tests, please mention it.

You should also check the `codecov` bot reply on this PR to confirm the
changed code is exercised.
-->
- for bug fix, changes in `slt`
- To ensure refactor correctness: no e2e test need to be changed

## Are there any user-facing changes?

<!--
If there are user-facing changes then we may require documentation to be
updated before approving the PR.

If there are any breaking changes to public APIs, please add the `api
change` label.
-->
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

optimizer Optimizer rules physical-plan Changes to the physical-plan crate sqllogictest SQL Logic Tests (.slt) v56.0.0

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants