Skip to content

feat: prototype frozen final aggregation buckets - #8

Closed
Rachelint wants to merge 1 commit into
jayzhan211:agg-buckets-final-stagefrom
Rachelint:codex/frozen-final-aggregation-poc
Closed

Rachelint wants to merge 1 commit into
jayzhan211:agg-buckets-final-stagefrom
Rachelint:codex/frozen-final-aggregation-poc

Conversation

@Rachelint

Copy link
Copy Markdown

Which issue does this PR close?

Related to apache#25724. This is a stacked PoC for benchmark validation and does not close an issue.

Rationale for this change

The final hash aggregation in apache#25724 moves its accumulated groups to buckets once the threshold is reached. Retaining those groups could reduce bucket input and the merge after input ends when later rows repeat the same keys. This draft tests that idea with an opt-in implementation and reproducible benchmarks. It is not proposed for merge yet.

What changes are included in this PR?

  • Add the off-by-default hash_aggregate_bucket_adaptive option. After a 65,536-row sample, probing is bypassed when fewer than 25% of rows hit the frozen keys.
  • Retain the original final hash table while routing misses to buckets; merge its state into buckets at EOF so keys remain correct after bypass. Fall back to the original path for unsupported key types or memory pressure.
  • Add counters and timing metrics, correctness tests, six synthetic SQL benchmarks, and results with reproduction commands.

The local optimized synthetic benchmark found this PoC 7%–21% slower than the original bucket path in all six queries. It uses a second encoded-key index because GroupValues has no lookup-only API. This draft exists so the same branch can be measured on the benchmark host with real ClickBench data.

Benchmark host commands

cargo build -p datafusion-benchmarks --release --bin benchmark_runner
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD=0 target/release/benchmark_runner aggregate_frozen --iterations 10 --partitions 12
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD=4000 target/release/benchmark_runner aggregate_frozen --iterations 10 --partitions 12
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD=4000 DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_ADAPTIVE=true target/release/benchmark_runner aggregate_frozen --iterations 10 --partitions 12

Real ClickBench Q11/Q12/Q14 and Q18/Q33/Q34, using the host's hits.parquet and Jay's threshold of 262,144, are still needed for the decision. This local machine did not have the dataset.

What is the testing strategy for this PR?

  • cargo fmt --all: passed.
  • cargo clippy --all-targets --all-features -- -D warnings: passed.
  • Focused final aggregation Rust tests and aggregate_bucketed / information_schema SQL logic tests: passed.
  • git diff --check: passed.
  • The full extended workspace suite was stopped after the initial submodule issue was fixed, as the immediate objective is benchmark validation. ./dev/rust_lint.sh could not start because local /usr/bin/python3 lacks PyYAML and uv is unavailable. These gates remain open before any merge.

Are there any user-facing changes?

Only the off-by-default execution option and additional physical plan metrics. Existing behavior remains the default.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant