feat: aggregate large hash tables as hash buckets (experimental, off by default) - #25567
jayzhan211 wants to merge 6 commits into
Conversation
…, off by default)
…threshold while its groups do not recur
…ting buckets of unique groups
|
run benchmarks clickbench_partitioned |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff Run configurationrun benchmark clickbench_partitionedResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff Run configurationrun benchmark clickbench_partitionedCPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
|
run benchmark clickbench_partitioned |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff Run configurationrun benchmark clickbench_partitioned
changed:
env:
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff Run configurationrun benchmark clickbench_partitioned
changed:
env:
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"CPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
|
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 |
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #25567 +/- ##
==========================================
+ Coverage 82.42% 82.43% +0.01%
==========================================
Files 1138 1140 +2
Lines 435501 436707 +1206
Branches 435501 436707 +1206
==========================================
+ Hits 358955 360002 +1047
- Misses 54841 54891 +50
- Partials 21705 21814 +109 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
I experiment the similar idea in #23186 before, and see similar improvements, so I think it maybe really promising. However, as the PR description notes, there are still some unresolved issues preventing this idea to be ready. For string group values case
I previously tried moving the entire bucketing step into RepartitionExec, but the results were disappointing: it was noticeably slower than splitting into buckets inside the final aggregation. My current hypothesis is that this may be related to the repartition coalescer being shared across multiple producer tasks, whereas the extra hashing and copying in the final aggregation are local to a single stream. However, I have not confirmed this yet. And finally I have no better idea than just disable the optimization when found string column. Another possible opportunity to improve more performanceIn my experiments, slicing a batch and pushing each slice into a per-bucket coalescer was not particularly fast. In #23186, it even caused a performance regression(so actually I was surprised that this PR still achieves substantial overall improvements while using slice + coalescer). |
|
run benchmark clickbench_partitioned |
|
Hi @Rachelint, your benchmark configuration could not be parsed (#25567 (comment)). Error: Usage: Any benchmark name is accepted: Per-side configuration ( env:
# shared env is inherited by BOTH the build and the run, so build
# flags go here. Builds default to no debuginfo for speed; opt back
# in for hung-job gdb dumps and cap jobs to stay within memory:
CARGO_PROFILE_RELEASE_DEBUG: "1"
CARGO_BUILD_JOBS: "1"
baseline:
ref: v45.0.0
env:
# per-side env only reaches the benchmark run, not the build
DATAFUSION_RUNTIME_MEMORY_LIMIT: 1G
changed:
ref: v46.0.0
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: 2GFile an issue against this benchmark runner |
|
run benchmark clickbench_partitioned |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff Run configurationrun benchmark clickbench_partitioned
changed:
env:
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff Run configurationrun benchmark clickbench_partitioned
changed:
env:
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"CPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
|
run benchmark clickbench_partitioned |
|
Hi @Rachelint, your benchmark configuration could not be parsed (#25567 (comment)). Error: Usage: Any benchmark name is accepted: Per-side configuration ( env:
# shared env is inherited by BOTH the build and the run, so build
# flags go here. Builds default to no debuginfo for speed; opt back
# in for hung-job gdb dumps and cap jobs to stay within memory:
CARGO_PROFILE_RELEASE_DEBUG: "1"
CARGO_BUILD_JOBS: "1"
baseline:
ref: v45.0.0
env:
# per-side env only reaches the benchmark run, not the build
DATAFUSION_RUNTIME_MEMORY_LIMIT: 1G
changed:
ref: v46.0.0
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: 2GFile an issue against this benchmark runner |
|
run benchmark clickbench_partitioned |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff Run configurationrun benchmark clickbench_partitioned
changed:
env:
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff Run configurationrun benchmark clickbench_partitioned
changed:
env:
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"CPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
…e its input holds them
|
Thanks @Rachelint — this is very helpful, and it is useful to know #23186 hit the same wall from the other direction. First, housekeeping: I have split this series into two reviewable PRs and rebased both onto current main (needed anyway after #25538 moved the spill contexts into
This PR is superseded by those; happy to close it once you have had a look. On string group keysTwo commits landed after you looked here, and both target exactly this case:
On #25724's tree, ClickBench per query, off/on back to back, 3 rounds of the fastest of 3 iterations (M4 Pro, 12 partitions, threshold 262144), the string-keyed queries now look like this:
Total 20.65 s → 17.25 s = 0.835x. So Q33 and Q34 are wins here rather than regressions, which makes me think the discriminator is not the column type but the group count: everything at or above ~1.5M groups per partition wins, everything at 0.5-1M loses. If we disabled bucketing whenever a group key is a string we would give up 0.56x, 0.63x and 0.87x on the three largest string aggregations, which are precisely the ones this is meant to help. On
|
| routing, ns/row | finding group ids, ns/row | |
|---|---|---|
| Q14 | 44.1 → 28.7 | 45.4 → 78.7 |
| Q12 | 37.0 → 22.8 | 57.1 → 91.9 |
| Q31 (no strings) | 29.6 → 36.5 | 10.1 → 9.7 |
…but interning got much more expensive, and every query got slower (Q13 1.07 → 1.21, Q16 1.00 → 1.14, Q18 0.53 → 0.63). The conclusion I drew is that the coalescer's garbage-collection copy is not waste: it packs each bucket's strings into one contiguous, bucket-local block, which is exactly what makes that bucket's small table fast afterwards. Sharing the input's buffers instead leaves a bucket's keys scattered across every buffer the input arrived in. Q31 is the control — no strings, so interning is unchanged, and routing gets worse because arrow rescans the whole mask once per bucket.
I mention it because MutableArrayData also Arc-clones the variadic buffers for view types rather than copying the bytes, so a bucket batch built that way would share buffers in the same way. That may be part of what you saw on Q33/Q34 — offered as a hypothesis about the mechanism, not a claim about your implementation, since I have not run it.
The corollary for the arrow ask is that a public InProgressArray would help most if it still compacts per destination; a zero-copy view-sharing scatter appears to lose more downstream than it saves. I would be glad to help push on that if you pick it up.
For what it is worth I also tried moving the bucketing into RepartitionExec, in two variants (per-destination slices of one grouped take, and interleave over 16 held batches), and both were slower than doing it in the final aggregation — which matches your result. With 12 partitions × 64 buckets an 8192-row batch becomes 768 slices of ~10 rows, and the per-slice cost dominates.
What is genuinely still unresolved
Q11/Q12/Q14, at 0.5-1M groups per partition, are 1.02-1.08x. I do not think this is a tuning problem:
- On Q12 the CPU is a wash — the final aggregate costs 76 ms more and the partial aggregate 68 ms less, summed over all partitions, net −0.7 ms. The 7% is wall clock: without buckets the final stage interns every row as it arrives, overlapped with the scan, while with buckets it only routes during the scan and aggregates after the input ends.
- A smarter trigger does not help. A larger threshold, a self-timing trigger, and holding the input back until it proves large were all built and measured; each only moved the loss to other queries. In particular engaging later is worse than never engaging: Q33's final aggregate takes 12.4 s of compute if it never buckets, 3.5 s if it buckets at 65k groups, and 25.8 s if it buckets at 1M, because a table that has already grown has already paid the penalty and then pays the routing on top.
So the decision has to be made before the table grows — a prediction, not a reaction — which is why the option stays off by default in #25724 and why I have not proposed flipping it. If it ever were flipped, two tests that encode today's behaviour would need updating: aggregate_fuzz::streaming_aggregate_test (a flushing partial aggregation legitimately emits a group more than once) and ordered_aggregate_spill.slt (it asserts the final aggregate spills, and with buckets it sometimes no longer needs to).
A benchmark-bot run against #25724 with DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144" would be very welcome — my numbers are from a single laptop, and the cliff this targets depends on cache and core count, so your c4a-highmem-16 with 80 MiB of L3 may well place the crossover somewhere else.
|
Sorry, I’ve been away for the past two days. I’m back home today and will take a detailed look. |
Which issue does this PR close?
Draft, opened to share measurements and get feedback on the direction. The option is off by default, so nothing changes unless it is set.
Rationale for this change
A hash aggregation with millions of groups per partition spends most of its time finding group ids: 69-97% of aggregate compute on ClickBench Q15-18 / Q31-35. The hash table and the group values no longer fit the CPU caches, and once several such tables of hundreds of MB compete (one per partition) the cost per row falls off a cliff (Q18: 42-51 ns per row up to 1M groups per partition, 530 ns at 4-8M).
This PR lets a final aggregation stop growing one table: past a threshold it moves its state and all further input into 64 hash buckets and aggregates the buckets one after another with a small, reused table. Buckets are emitted, released, compacted, spilled and split again independently.
Measured on an Apple M4 Pro, 12 partitions,
hash_aggregate_bucket_threshold = 262144. ClickBenchhits_partitioned, every query run off/on back to back (suite-level runs drift on this machine), 3 rounds of the fastest of 3 iterations:ClickBench total 37.3 s → 28.6 s = 0.766x. Per query:
(Q0 1.17x and Q19 1.21x are 0.6 ms and 21 ms queries that do not aggregate; that is noise.)
TPC-H SF10, same method: 5.52 s → 5.44 s = 0.986x, i.e. neutral. Its aggregations are small or integer-keyed, so most queries do not reach the threshold at all; the ones that move are Q17 0.87x, Q18 0.92x and, in the other direction, Q16 1.12x and Q22 1.10x (50 ms and 60 ms queries).
Peak RSS is typically 0.4-0.9x.
Known cost, which is why this is a draft and off by default:
GROUP BY l_partkey, l_suppkeyat SF10 (the input repeats each group 7x)Where the remaining cost is, per final input row on Q14: one table spends 78 ns finding group ids; bucketed it spends 26 ns, but 41 ns goes into moving the row to its bucket first (hash, gather, copy). So bucketing trades the cache miss for the move, and only wins once the table is big enough that the miss dominates — which is why Q18 (213 ns of group ids with one table) and Q32 (194 ns) win by 2-3x while a query at 0.5M groups per partition does not. A smarter trigger cannot change that arithmetic: a larger threshold, a self-timing trigger and holding input back until the input proves large were all built and measured, and each only moved the loss to different queries.
What changes are included in this PR?
One commit per step:
datafusion.execution.hash_aggregate_bucket_threshold(groups, default 0 = off) and the bucketed final aggregation:aggregates/final_buckets.rs(routing by hash with one gather per batch, per-bucket compaction, per-bucket spill throughSpillManager) and the bucket output inFinalHashAggregateStream. State with nested types is excluded; a soft group limit disables it.aggregates/bucketed_aggregation.rsshares the bucket logic with the single-stage stream, which buckets when it runs on every partition (SinglePartitioned). A loneSinglestream keeps its table: measured 6-25% slower with buckets, because every row is aggregated twice.multi_group_by/bytes_view.rs). Off for the table that compacts a bucket, whose purpose is to make it smaller, and off for a normal aggregation, whose table outlives every batch it has seen. On the multi-column string keys this removes 17-19 ns per row: Q14's group ids 44.7 → 25.7 ns, Q13's inner stage 45.3 → 28.1, Q16 54.8 → 46.9. The same change to the single-column string path (ArrowBytesViewMap) was measured and reverted: it compares the stored bytes of a candidate on every probe, and the copy is what packs those bytes together, so Q33 (URLs, which share a long prefix) lost 0.53x → 0.57x.What is the testing strategy for this PR?
group_values/multi_group_by/bytes_view.rsthat a borrowing builder still reads its values after the arrays they came from are dropped, and thattake_nkeeps every row readable.aggregates/final_buckets.rsandaggregates/hash_stream.rs: bucketed results equal the single table for the final, partial and single-stage streams; buckets split again; buckets spill under a memory limit; repeated groups are compacted; the partial flush stops when groups recur.aggregate_bucketed.slt: 14 queries (FILTER, ROLLUP, orderedarray_agg, DISTINCT, nested aggregation) at 4 and 1 partitions, with the option off and at 100 groups, all sections identical, withbucket_splits/table_flush_countchecked inEXPLAIN ANALYZE.memory_limittests and the aggregate / group by / distinct sqllogictests pass. One tight-memory file (ordered_aggregate_spill.slt, limits of 500K-2M) failed once in a loaded parallel run and passed in 6 isolated runs; memory behaviour under such limits differs with a 100-group threshold.aggregate_fuzz::streaming_aggregate_testdoes not, as expected: it compares partial state rows one by one, and a flushing partial aggregation legitimately emits a group more than once. It would have to merge partial rows first if the default ever became non-zero.cargo fmt, workspace clippy and the extended test suite pass.Are there any user-facing changes?
A new experimental configuration option,
datafusion.execution.hash_aggregate_bucket_threshold, off by default, documented inconfigs.mdtogether with its known costs. New metrics onAggregateExec:bucket_splits,bucket_compactions,table_flush_count.