Skip to content

feat: add a validated AggregateExecBuilder for building and rewriting AggregateExec - #77

Closed
adriangb wants to merge 3 commits into
mainfrom
claude/datafusion-25257-builder-api-5pd3bs
Closed

adriangb wants to merge 3 commits into
mainfrom
claude/datafusion-25257-builder-api-5pd3bs

Conversation

@adriangb

@adriangb adriangb commented Sep 16, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Rationale for this change

The methods apache#25257 hides are hidden because they are dangerous. They are just as dangerous for the code inside DataFusion that keeps using them, and hiding them does nothing about that.

Every one of these plans can be built today, and each one only misbehaves once it is executed:

plan what happens when it runs
limit, no MIN/MAX aggregate, no ordering direction internal_err!("Ordering direction required for DISTINCT with limit")
limit with more than one group by expression aggr.group_expr().expr()[0] panics
limit of 0 on a top-k aggregate worst_val().expect("Missing root") panics: a queue of capacity 0 reports itself full with an empty root
limit on an aggregate with a FILTER GroupedTopKAggregateStream never reads filter_expr, so the filter is silently dropped and the query returns wrong results
limit on a COUNT/AVG/… aggregate the TopK stream ignores the aggregate, wrong results
limit with an unsupported key/value type debug_assert! in debug builds only

On top of that, AggregateExec had two nearly identical clone-with-one-change methods with different semantics (with_new_limit_options resets the metrics, with_limit_options keeps them), each hand-copying twelve fields, and a try_new with six positional arguments, two of which are schemas that are easy to transpose (input vs input_schema).

Validating the node when it is built makes these unrepresentable in a built plan, which means the API no longer needs to be dangerous — only internal.

What changes are included in this PR?

AggregateExecBuilder (datafusion/physical-plan/src/aggregates/builder.rs), reached through AggregateExec::builder(mode, input) for a new node and AggregateExec::to_builder() for a rewrite:

let exec = AggregateExec::builder(AggregateMode::Single, input)
    .with_group_by(group_by)
    .with_aggr_exprs(aggr_exprs)
    .with_limit_options(LimitOptions::new(10))
    .build()?;

let with_limit = exec
    .to_builder()
    .with_limit_options(LimitOptions::new_with_order(10, true))
    .build()?;
  • every argument is named, and filter_expr defaults to "no filter per aggregate"
  • build() validates the node once: a limit is checked against the shape of the aggregate (the rows in the table above), and a replacement of the aggregate expressions is checked against both the output schema it inherits and the dynamic filter it inherits
  • a rewrite carries over the derived state of the node it came from — output schema, plan properties, ordering requirements, dynamic filter — instead of asking each caller to copy the fields. Structural setters (with_mode, with_group_by, with_input, with_filter_exprs) drop that state and recompute; the others keep it, so rewrites behave exactly as they did before.

A dynamic filter records which aggregate expressions are MIN/MAX, at which index, and over which column — a MIN pushes down col < bound, a MAX pushes col > bound. Carrying it verbatim across a rewrite that turns a MIN into a MAX would push down the wrong predicate and prune rows the aggregate needs, and the schema check cannot see it because both produce the same output field. with_aggr_exprs therefore re-derives the state and keeps the inherited filter only while it still describes the new expressions, so a rewrite that merely reorders or reverses them keeps the original filter (and with it the link to whichever child accepted it during pushdown).

Migrated to the builder: TopKAggregation, LimitedDistinctAggregation, CombinePartialFinalAggregate, OptimizeAggregateOrder, the protobuf decoder, and every test. The optimizer rules use .build().ok()?, so an invalid limit means "skip this optimization" instead of a plan that fails at runtime. No #[expect(deprecated)] anywhere.

Deprecated and #[doc(hidden)]: with_limit_options, with_new_limit_options, with_new_aggr_exprs.

#[doc(hidden)] but not deprecated: AggregateExec::builder, AggregateExec::to_builder, AggregateExecBuilder, and the limit_options() getter. Building and rewriting an aggregate is how DataFusion's own optimizer rules work, not a public API, so the whole surface is hidden — as apache#25257 does. The getter is not deprecated because reading a limit is safe and has no replacement; deprecating it would only push #[expect(deprecated)] back into CombinePartialFinalAggregate, which lives in another crate and cannot reach the field directly.

Two existing tests were building plans that cannot execute and had to be adjusted: a statistics test that put a limit on a COUNT(a) aggregate, and the protobuf roundtrip test that put an ordered limit on AVG(b). Both now use shapes the optimizer actually produces.

What is the testing strategy for this PR?

Eight new unit tests in aggregates/builder.rs cover the builder itself: defaulted filter expressions, mismatched filter arity, derived state being reused on a rewrite (asserted with Arc::ptr_eq on the plan properties) and metrics being reset, the schema recompute when the mode changes, rejection of incompatible aggregate expressions, each limit validation rule (including the zero limit and the ordered-input case that used to produce an internal error at execution time), and all three dynamic-filter outcomes when the aggregate expressions are replaced — kept, rebuilt, dropped. The drop case uses SUM, whose output field matches MIN's, so it gets past the schema check and actually reaches the dynamic filter. Two doctests on the builder and one on to_builder cover the documented usage.

Both panics in the table above were reproduced before being fixed, not reasoned about: the zero limit panics at heap.rs:133.

Everything else is covered by the existing suites, which all pass:

  • datafusion-physical-plan lib (1874) and doctests
  • datafusion-physical-optimizer (33)
  • datafusion core_integration physical_optimizer (560)
  • datafusion-proto aggregate roundtrips
  • the full 505-file sqllogictest suite
  • cargo clippy --all-targets --all-features --workspace -- -D warnings, cargo fmt --all, and cargo doc -p datafusion-physical-plan under -D warnings

Not run: the sql_planner planning benchmarks — the environment ran out of disk during the release build. The change is plan-time only (execution paths are untouched) and adds one create_schema call per aggregate rewritten by OptimizeAggregateOrder, so it should be in the noise, but it has not been measured.

Are there any user-facing changes?

Yes, and docs/source/library-user-guide/upgrading/56.0.0.md has a section for them.

  • The three methods above are deprecated with a replacement, and this whole API is now #[doc(hidden)].
  • A limit that the aggregate cannot execute is now an error when the plan is built, including when decoding from protobuf, instead of an internal error, a panic, or silently ignored FILTER expressions at execution time. No plan DataFusion's own optimizer produces is affected — the rules already check these conditions before pushing a limit down.
  • cargo-semver-checks classifies this as requiring a major version: #[doc(hidden)] on the pre-existing AggregateExec::limit_options removes it from the public API (major), and the three deprecations are a minor change. That is the intended consequence of hiding an internal API and is expected for the 56.0.0 release; the Check semver job reports it as a note and passes. Nothing else in the four checked crates moved.

Follow-up notes

Validation catches the bad limit shapes, but LimitOptions still lets you construct them. Splitting it into SoftLimit { limit } / TopK { limit, descending } would make two of the failure modes unrepresentable rather than merely rejected. That's a wider rename, so it was left out of this PR — happy to do it here if a reviewer wants to go further.

The deprecated with_new_aggr_exprs keeps the dynamic-filter hazard described above: it copies the state verbatim, and this PR does not change its behaviour. Only the builder path re-derives it.

🤖 Generated with Claude Code

https://claude.ai/code/session_01PwTc51ca2XHDCyVB7MbJoz

`AggregateExec::try_new` takes six positional arguments, and the fields the
physical optimizer changes afterwards were set through methods that each
copied the remaining fields by hand:

- `with_new_limit_options(&self, ..)` clones the node and resets metrics
- `with_limit_options(self, ..)` consumes it and keeps them
- `with_new_aggr_exprs(&self, ..)` clones it for a different field

None of them validate the result, so it is easy to build an aggregate that
only fails once it is executed. Pushing a limit into an aggregate that
cannot execute one is the clearest case: with no MIN/MAX aggregate and no
ordering direction `GroupedTopKAggregateStream` returns an internal error,
with more than one group by expression it indexes out of bounds, and with a
`FILTER` it silently drops the filter and returns wrong results.

Replace them with `AggregateExecBuilder`, reachable as
`AggregateExec::builder(mode, input)` for a new node and
`AggregateExec::to_builder()` for a rewrite:

- every argument is named, so `input`/`input_schema` and
  `aggr_expr`/`filter_expr` cannot be transposed
- `filter_expr` defaults to "no filter per aggregate"
- `build()` validates the node once: limits are checked against the shape of
  the aggregate, and a rewrite of the aggregate expressions is checked
  against the output schema it inherits
- a rewrite carries over the derived state (output schema, plan properties,
  ordering requirements, dynamic filter) of the node it came from instead of
  asking each caller to copy the fields, so it cannot rename output fields

The three optimizer rules that push a limit, `OptimizeAggregateOrder` and
the protobuf decoder now go through the builder; the old methods are
deprecated in favour of it. `AggregateExec::limit_options()` is unchanged:
reading the limit was never the problem.

Invalid limits are now an error at plan time rather than at execution time.
Optimizer rules treat that as "skip this optimization", so plans that used to
be built and then fail are simply not rewritten.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PwTc51ca2XHDCyVB7MbJoz
Building and rewriting an `AggregateExec` is how DataFusion's own physical
optimizer rules work, not something external consumers should reach for.
Hide that whole surface from the rendered docs, as apache#25257
does for the methods it deprecates:

- the deprecated `with_limit_options`, `with_new_limit_options` and
  `with_new_aggr_exprs`
- their replacements, `AggregateExec::builder`, `AggregateExec::to_builder`
  and `AggregateExecBuilder` itself
- the `limit_options` getter

`limit_options` is hidden but deliberately not deprecated: it has no
replacement, reading the limit of an aggregate is safe, and deprecating it
would only force `#[expect(deprecated)]` back into the optimizer rule that
copies a limit between aggregates.

Doc links from the still-public `AggregateExec::try_new` and `LimitOptions`
into the hidden items are now plain code spans, since a link to a hidden
item renders as a dead anchor.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PwTc51ca2XHDCyVB7MbJoz
@github-actions github-actions Bot added documentation Improvements or additions to documentation core optimizer proto physical-plan labels Sep 16, 2026
Comment thread datafusion/physical-plan/src/aggregates/builder.rs
Comment thread datafusion/physical-plan/src/aggregates/mod.rs
@github-actions

github-actions Bot commented Sep 16, 2026 •

Copy link
Copy Markdown

Thank you for opening this pull request!

Reviewer note: cargo-semver-checks reported the current version number is not SemVer-compatible with the changes in this pull request (compared against the base branch).

Details
     Cloning apache/main
    Building datafusion v55.0.0 (current)
       Built [  46.939s] (current)
     Parsing datafusion v55.0.0 (current)
      Parsed [   0.030s] (current)
    Building datafusion v55.0.0 (baseline)
       Built [  45.853s] (baseline)
     Parsing datafusion v55.0.0 (baseline)
      Parsed [   0.030s] (baseline)
    Checking datafusion v55.0.0 -> v55.0.0 (no change; assume patch)
     Checked [   0.711s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [  95.159s] datafusion
    Building datafusion-physical-optimizer v55.0.0 (current)
       Built [  32.251s] (current)
     Parsing datafusion-physical-optimizer v55.0.0 (current)
      Parsed [   0.019s] (current)
    Building datafusion-physical-optimizer v55.0.0 (baseline)
       Built [  32.346s] (baseline)
     Parsing datafusion-physical-optimizer v55.0.0 (baseline)
      Parsed [   0.019s] (baseline)
    Checking datafusion-physical-optimizer v55.0.0 -> v55.0.0 (no change; assume patch)
     Checked [   0.115s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [  65.681s] datafusion-physical-optimizer
    Building datafusion-physical-plan v55.0.0 (current)
       Built [  30.264s] (current)
     Parsing datafusion-physical-plan v55.0.0 (current)
      Parsed [   0.127s] (current)
    Building datafusion-physical-plan v55.0.0 (baseline)
       Built [  30.294s] (baseline)
     Parsing datafusion-physical-plan v55.0.0 (baseline)
      Parsed [   0.125s] (baseline)
    Checking datafusion-physical-plan v55.0.0 -> v55.0.0 (no change; assume patch)
     Checked [   0.789s] 223 checks: 221 pass, 2 fail, 0 warn, 31 skip

--- failure inherent_method_now_doc_hidden: inherent method #[doc(hidden)] added ---

Description:
A method or associated fn is now #[doc(hidden)], removing it from the crate's public API.
        ref: https://doc.rust-lang.org/rustdoc/write-documentation/the-doc-attribute.html#hidden
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/inherent_method_now_doc_hidden.ron

Failed in:
  AggregateExec::limit_options in file /home/runner/work/datafusion/datafusion/target/semver-checks/git-apache_main/758c72efa4d0b008ece9c91b5b6f9a92caacc7c9/datafusion/physical-plan/src/aggregates/mod.rs:1097

--- failure type_method_marked_deprecated: type method #[deprecated] added ---

Description:
A type method is now #[deprecated]. Downstream crates will get a compiler warning when using this method.
        ref: https://doc.rust-lang.org/reference/attributes/diagnostics.html#the-deprecated-attribute
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/type_method_marked_deprecated.ron

Failed in:
  method datafusion_physical_plan::aggregates::AggregateExec::with_new_aggr_exprs in /home/runner/work/datafusion/datafusion/datafusion/physical-plan/src/aggregates/mod.rs:967
  method datafusion_physical_plan::aggregates::AggregateExec::with_new_limit_options in /home/runner/work/datafusion/datafusion/datafusion/physical-plan/src/aggregates/mod.rs:995
  method datafusion_physical_plan::aggregates::AggregateExec::with_limit_options in /home/runner/work/datafusion/datafusion/datafusion/physical-plan/src/aggregates/mod.rs:1182

     Summary semver requires new major version: 1 major and 1 minor checks failed
    Finished [  62.627s] datafusion-physical-plan
    Building datafusion-proto v55.0.0 (current)
       Built [  43.315s] (current)
     Parsing datafusion-proto v55.0.0 (current)
      Parsed [   0.015s] (current)
    Building datafusion-proto v55.0.0 (baseline)
       Built [  42.814s] (baseline)
     Parsing datafusion-proto v55.0.0 (baseline)
      Parsed [   0.016s] (baseline)
    Checking datafusion-proto v55.0.0 -> v55.0.0 (no change; assume patch)
     Checked [   0.116s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [  87.258s] datafusion-proto

…write

Two defects found by review on the builder, both of which it is the
builder's job to catch:

A limit of zero on an aggregate that executes as a top-k builds fine and
then panics. `TopKHeap::new(0, ..)` has capacity 0 and no root, so
`is_full()` is immediately true and the first row read reaches
`worst_val().expect("Missing root")` in `PrimitiveHeap::is_worse`.
`build()` now rejects it. Reproduced before fixing: a `Single` aggregate
over `MIN(a)` with `LimitOptions::new(0)` panics at heap.rs:133.

A dynamic filter records which aggregate expressions are `MIN`/`MAX`, at
which index, and over which column: a `MIN` pushes `col < bound`, a `MAX`
pushes `col > bound`. `to_builder().with_aggr_exprs(..)` kept that state
verbatim, so rewriting `MIN(a)` into `MAX(a)` left the filter tagged
`Min` and pushed down a predicate that prunes the rows the aggregate
needs. The schema check does not catch it — both produce the same output
field. The state is now re-derived from the new expressions and kept only
while it still describes them, so a rewrite that merely reorders or
reverses them keeps the original filter, and with it the link to whichever
child accepted it during pushdown.

`init_dynamic_filter`'s body moves into `AggregateExec::derive_dynamic_filter`
so the builder can re-derive it; the logic is unchanged.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PwTc51ca2XHDCyVB7MbJoz

Copy link
Copy Markdown
Member Author

cargo test datafusion-cli (amd64) is failing, and it is not this PR's

Four tests fail — test_cli, test_aws_options, test_object_store_profiling, test_s3_url_fallback — all in datafusion-cli/tests/cli_integration.rs, and all inside setup_minio_container():

Failed to start MinIO container. Ensure Docker is running and accessible:
failed to pull the image 'minio/minio:RELEASE.2025-02-28T09-55-16Z',
Docker responded with status code 404: pull access denied for minio/minio,
repository does not exist or may require 'docker login': denied

That is the runner being unable to pull the MinIO image from Docker Hub, not a test failure. This PR touches no file under datafusion-cli/.

Confirmed rather than assumed, on two independent PRs:

PR run result
#76 (unrelated, not mine) 19:03 same job red, every other check green
#76 21:35 same job red, every other check green
this PR 04:50 same job red, same four tests, same pull denial

No fix exists to port: the image pull needs Docker Hub credentials (or a mirror) on the runner, which is CI configuration and outside what this PR should change. The tests are not flaky and a re-run will not help while the registry keeps denying the pull, so I have not spent one on it. A fresh run of that job is in flight anyway on the new head (674921f); if it comes back green there, this note is moot.

Everything else on this PR is mine to keep green, and I am still watching it.


Generated by Claude Code

@adriangb adriangb closed this Sep 16, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

core documentation Improvements or additions to documentation optimizer physical-plan proto

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants