Skip to content

feat: aggregate large hash tables as hash buckets (experimental, off by default) - #25567

Draft
jayzhan211 wants to merge 6 commits into
apache:mainfrom
jayzhan211:agg-bucketed-final
Draft

jayzhan211 wants to merge 6 commits into
apache:mainfrom
jayzhan211:agg-bucketed-final

Conversation

@jayzhan211

@jayzhan211 jayzhan211 commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

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. ClickBench hits_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:

faster slower
Q32 0.28x Q12 1.09x
Q34 0.55x Q11 / Q20 1.08x
Q18 0.61x Q14 1.06x
Q33 0.64x Q5 / Q10 / Q38 1.05x
Q4 0.79x Q13 1.02x
Q15 0.82x
Q23 0.84x
Q8 0.88x
Q31 0.90x

(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:

case time peak RSS
GROUP BY l_partkey, l_suppkey at SF10 (the input repeats each group 7x) 0.86x 1.5x
ClickBench Q11 / Q12 / Q14 / Q20 (0.5-1M groups per partition) 1.06-1.09x 0.8-0.9x

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:

  1. 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 through SpillManager) and the bucket output in FinalHashAggregateStream. State with nested types is excluded; a soft group limit disables it.
  2. A partial aggregation table that reaches the threshold is flushed downstream while its groups do not recur (about 1024 sampled group hashes per flush, kept for the last 64 flushes; flushing stops once more than 20% of a sample recurs).
  3. One hash table is reused across the buckets, and buckets of unique groups are not compacted again.
  4. aggregates/bucketed_aggregation.rs shares the bucket logic with the single-stage stream, which buckets when it runs on every partition (SinglePartitioned). A lone Single stream keeps its table: measured 6-25% slower with buckets, because every row is aggregated twice.
  5. A final table whose input holds about one row per group moves into buckets at a quarter of the threshold, so that little work is repeated; buckets are still split again at the full threshold.
  6. The table of a bucket keeps its group keys where its input already holds them. A bucket's rows sit in one contiguous block that is dropped as soon as that bucket has been aggregated, so its table can store a view into that block instead of copying every new key into a buffer of its own (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?

  • Unit tests in group_values/multi_group_by/bytes_view.rs that a borrowing builder still reads its values after the arrays they came from are dropped, and that take_n keeps every row readable.
  • Unit tests in aggregates/final_buckets.rs and aggregates/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, ordered array_agg, DISTINCT, nested aggregation) at 4 and 1 partitions, with the option off and at 100 groups, all sections identical, with bucket_splits / table_flush_count checked in EXPLAIN ANALYZE.
  • With the default temporarily set to 100, the aggregate unit tests, memory_limit tests 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_test does 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 in configs.md together with its known costs. New metrics on AggregateExec: bucket_splits, bucket_compactions, table_flush_count.

@github-actions github-actions Bot added documentation Improvements or additions to documentation sqllogictest SQL Logic Tests (.slt) common Related to common crate physical-plan Changes to the physical-plan crate labels Sep 21, 2026
@Rachelint

Copy link
Copy Markdown
Contributor

run benchmarks clickbench_partitioned

@adriangbot

Copy link
Copy Markdown

🤖 Benchmark running (GKE) | trigger
Instance: c4a-highmem-16 (12 vCPU / 65 GiB) | Linux bench-c5762634475-2546-fdk96 6.12.94+ #1 SMP Tue Aug 4 08:44:15 UTC 2026 aarch64 GNU/Linux

CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  16
On-line CPU(s) list:                     0-15
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     16
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               1 MiB (16 instances)
L1i cache:                               1 MiB (16 instances)
L2 cache:                                32 MiB (16 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-15
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected

Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff

Run configuration
run benchmark clickbench_partitioned

Results will be posted here when complete


File an issue against this benchmark runner

@adriangbot

Copy link
Copy Markdown

🤖 Benchmark completed (GKE) | trigger

Instance: c4a-highmem-16 (12 vCPU / 65 GiB)

Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff

Run configuration
run benchmark clickbench_partitioned
CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  16
On-line CPU(s) list:                     0-15
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     16
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               1 MiB (16 instances)
L1i cache:                               1 MiB (16 instances)
L2 cache:                                32 MiB (16 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-15
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected
Details

Comparing HEAD and agg-bucketed-final
--------------------
Benchmark clickbench_partitioned.json
--------------------
┏━━━━━━━━━━━┳━━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━┓
┃ Query     ┃       HEAD ┃ agg-bucketed-final ┃    Change ┃
┡━━━━━━━━━━━╇━━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━┩
│ QQuery 0  │    1.25 ms │            1.25 ms │ no change │
│ QQuery 1  │   11.80 ms │           11.77 ms │ no change │
│ QQuery 2  │   36.73 ms │           36.65 ms │ no change │
│ QQuery 3  │   31.34 ms │           31.13 ms │ no change │
│ QQuery 4  │  236.65 ms │          232.08 ms │ no change │
│ QQuery 5  │  273.89 ms │          268.90 ms │ no change │
│ QQuery 6  │    1.30 ms │            1.34 ms │ no change │
│ QQuery 7  │   13.10 ms │           13.22 ms │ no change │
│ QQuery 8  │  325.60 ms │          327.84 ms │ no change │
│ QQuery 9  │  461.09 ms │          460.17 ms │ no change │
│ QQuery 10 │   64.86 ms │           64.03 ms │ no change │
│ QQuery 11 │   76.59 ms │           74.03 ms │ no change │
│ QQuery 12 │  262.58 ms │          256.42 ms │ no change │
│ QQuery 13 │  363.25 ms │          361.47 ms │ no change │
│ QQuery 14 │  275.23 ms │          277.07 ms │ no change │
│ QQuery 15 │  279.04 ms │          282.14 ms │ no change │
│ QQuery 16 │  613.87 ms │          617.38 ms │ no change │
│ QQuery 17 │  622.59 ms │          622.64 ms │ no change │
│ QQuery 18 │ 1261.20 ms │         1275.19 ms │ no change │
│ QQuery 19 │   27.19 ms │           27.26 ms │ no change │
│ QQuery 20 │  512.65 ms │          517.19 ms │ no change │
│ QQuery 21 │  508.21 ms │          505.75 ms │ no change │
│ QQuery 22 │  979.33 ms │          973.45 ms │ no change │
│ QQuery 23 │ 3054.95 ms │         3012.20 ms │ no change │
│ QQuery 24 │   40.11 ms │           40.98 ms │ no change │
│ QQuery 25 │  105.05 ms │          104.34 ms │ no change │
│ QQuery 26 │   40.60 ms │           40.77 ms │ no change │
│ QQuery 27 │  503.57 ms │          510.24 ms │ no change │
│ QQuery 28 │ 2933.31 ms │         2850.67 ms │ no change │
│ QQuery 29 │   41.22 ms │           41.57 ms │ no change │
│ QQuery 30 │  305.33 ms │          304.90 ms │ no change │
│ QQuery 31 │  279.29 ms │          274.02 ms │ no change │
│ QQuery 32 │  926.21 ms │          951.96 ms │ no change │
│ QQuery 33 │ 1445.97 ms │         1426.01 ms │ no change │
│ QQuery 34 │ 1480.69 ms │         1465.89 ms │ no change │
│ QQuery 35 │  291.82 ms │          286.93 ms │ no change │
│ QQuery 36 │   66.19 ms │           66.08 ms │ no change │
│ QQuery 37 │   35.23 ms │           34.50 ms │ no change │
│ QQuery 38 │   41.02 ms │           42.52 ms │ no change │
│ QQuery 39 │  146.94 ms │          145.42 ms │ no change │
│ QQuery 40 │   14.22 ms │           14.23 ms │ no change │
│ QQuery 41 │   13.75 ms │           13.39 ms │ no change │
│ QQuery 42 │   13.29 ms │           12.90 ms │ no change │
└───────────┴────────────┴────────────────────┴───────────┘
┏━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━┓
┃ Benchmark Summary                 ┃            ┃
┡━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━┩
│ Total Time (HEAD)                 │ 19018.09ms │
│ Total Time (agg-bucketed-final)   │ 18877.91ms │
│ Average Time (HEAD)               │   442.28ms │
│ Average Time (agg-bucketed-final) │   439.02ms │
│ Queries Faster                    │          0 │
│ Queries Slower                    │          0 │
│ Queries with No Change            │         43 │
│ Queries with Failure              │          0 │
└───────────────────────────────────┴────────────┘

Distribution per query (min / mean ±stddev / max):

Comparing HEAD and agg-bucketed-final
--------------------
Benchmark clickbench_partitioned.json
--------------------
┏━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━━┓
┃ Query     ┃                                  HEAD ┃                    agg-bucketed-final ┃        Change ┃
┡━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━━┩
│ QQuery 0  │          1.25 / 4.03 ±5.49 / 15.01 ms │          1.25 / 3.95 ±5.34 / 14.63 ms │     no change │
│ QQuery 1  │        11.80 / 11.99 ±0.15 / 12.27 ms │        11.77 / 11.91 ±0.13 / 12.13 ms │     no change │
│ QQuery 2  │        36.73 / 37.76 ±0.97 / 39.33 ms │        36.65 / 36.79 ±0.15 / 37.08 ms │     no change │
│ QQuery 3  │        31.34 / 31.97 ±0.69 / 33.32 ms │        31.13 / 31.30 ±0.10 / 31.40 ms │     no change │
│ QQuery 4  │     236.65 / 238.67 ±1.62 / 241.45 ms │     232.08 / 236.28 ±3.53 / 242.57 ms │     no change │
│ QQuery 5  │     273.89 / 277.57 ±5.49 / 288.41 ms │     268.90 / 271.85 ±2.32 / 274.43 ms │     no change │
│ QQuery 6  │           1.30 / 1.45 ±0.23 / 1.91 ms │           1.34 / 1.48 ±0.22 / 1.92 ms │     no change │
│ QQuery 7  │        13.10 / 13.22 ±0.12 / 13.43 ms │        13.22 / 13.43 ±0.35 / 14.13 ms │     no change │
│ QQuery 8  │     325.60 / 330.35 ±2.94 / 334.32 ms │     327.84 / 332.20 ±2.29 / 334.37 ms │     no change │
│ QQuery 9  │     461.09 / 466.84 ±4.69 / 473.46 ms │    460.17 / 474.59 ±10.27 / 491.42 ms │     no change │
│ QQuery 10 │        64.86 / 66.45 ±1.92 / 70.15 ms │        64.03 / 64.70 ±1.06 / 66.82 ms │     no change │
│ QQuery 11 │        76.59 / 77.13 ±0.52 / 77.83 ms │        74.03 / 74.95 ±0.88 / 76.43 ms │     no change │
│ QQuery 12 │     262.58 / 266.31 ±3.27 / 270.62 ms │     256.42 / 263.33 ±5.26 / 271.55 ms │     no change │
│ QQuery 13 │    363.25 / 377.53 ±14.48 / 405.37 ms │    361.47 / 383.16 ±15.26 / 404.47 ms │     no change │
│ QQuery 14 │     275.23 / 279.65 ±2.48 / 282.67 ms │     277.07 / 281.83 ±3.88 / 288.48 ms │     no change │
│ QQuery 15 │     279.04 / 287.67 ±5.64 / 295.47 ms │     282.14 / 288.27 ±5.12 / 294.84 ms │     no change │
│ QQuery 16 │    613.87 / 627.98 ±18.82 / 665.10 ms │    617.38 / 627.54 ±14.45 / 655.98 ms │     no change │
│ QQuery 17 │     622.59 / 626.88 ±3.05 / 631.55 ms │     622.64 / 634.06 ±8.02 / 643.71 ms │     no change │
│ QQuery 18 │  1261.20 / 1274.25 ±6.86 / 1280.83 ms │ 1275.19 / 1298.89 ±18.37 / 1322.00 ms │     no change │
│ QQuery 19 │       27.19 / 34.09 ±12.88 / 59.85 ms │       27.26 / 33.88 ±12.90 / 59.68 ms │     no change │
│ QQuery 20 │    512.65 / 527.73 ±10.62 / 545.93 ms │     517.19 / 518.92 ±2.41 / 523.59 ms │     no change │
│ QQuery 21 │     508.21 / 515.42 ±5.46 / 522.57 ms │     505.75 / 512.83 ±5.17 / 521.20 ms │     no change │
│ QQuery 22 │    979.33 / 987.95 ±9.49 / 1004.67 ms │     973.45 / 980.29 ±5.30 / 987.40 ms │     no change │
│ QQuery 23 │ 3054.95 / 3081.30 ±32.80 / 3144.52 ms │ 3012.20 / 3052.39 ±22.50 / 3070.69 ms │     no change │
│ QQuery 24 │       40.11 / 51.70 ±15.76 / 81.42 ms │       40.98 / 51.29 ±13.96 / 78.41 ms │     no change │
│ QQuery 25 │     105.05 / 108.48 ±4.83 / 117.90 ms │     104.34 / 108.30 ±4.87 / 117.81 ms │     no change │
│ QQuery 26 │        40.60 / 41.75 ±1.36 / 44.30 ms │        40.77 / 41.04 ±0.21 / 41.37 ms │     no change │
│ QQuery 27 │     503.57 / 514.03 ±7.45 / 523.46 ms │     510.24 / 521.69 ±9.98 / 540.01 ms │     no change │
│ QQuery 28 │ 2933.31 / 2959.62 ±24.55 / 2995.24 ms │ 2850.67 / 2881.92 ±21.10 / 2913.44 ms │     no change │
│ QQuery 29 │       41.22 / 47.82 ±11.90 / 71.59 ms │        41.57 / 48.02 ±9.25 / 66.20 ms │     no change │
│ QQuery 30 │     305.33 / 309.13 ±2.19 / 311.99 ms │     304.90 / 310.91 ±4.68 / 318.57 ms │     no change │
│ QQuery 31 │     279.29 / 288.31 ±7.40 / 298.29 ms │    274.02 / 282.91 ±10.75 / 302.99 ms │     no change │
│ QQuery 32 │   926.21 / 966.04 ±35.39 / 1011.38 ms │   951.96 / 985.17 ±21.70 / 1011.43 ms │     no change │
│ QQuery 33 │ 1445.97 / 1468.47 ±21.33 / 1505.08 ms │ 1426.01 / 1458.78 ±27.21 / 1503.08 ms │     no change │
│ QQuery 34 │ 1480.69 / 1492.90 ±10.67 / 1511.74 ms │ 1465.89 / 1479.08 ±17.82 / 1512.21 ms │     no change │
│ QQuery 35 │    291.82 / 303.21 ±13.42 / 328.91 ms │    286.93 / 323.40 ±33.82 / 370.62 ms │  1.07x slower │
│ QQuery 36 │        66.19 / 68.23 ±2.00 / 71.94 ms │       66.08 / 76.07 ±11.12 / 95.98 ms │  1.11x slower │
│ QQuery 37 │        35.23 / 40.06 ±2.61 / 42.55 ms │        34.50 / 35.37 ±0.65 / 36.33 ms │ +1.13x faster │
│ QQuery 38 │        41.02 / 43.92 ±3.75 / 51.13 ms │        42.52 / 46.62 ±6.43 / 59.32 ms │  1.06x slower │
│ QQuery 39 │     146.94 / 153.79 ±4.90 / 161.18 ms │    145.42 / 158.29 ±13.67 / 184.55 ms │     no change │
│ QQuery 40 │        14.22 / 16.15 ±2.37 / 20.75 ms │        14.23 / 14.42 ±0.16 / 14.65 ms │ +1.12x faster │
│ QQuery 41 │        13.75 / 14.26 ±0.33 / 14.65 ms │        13.39 / 15.04 ±2.56 / 20.11 ms │  1.05x slower │
│ QQuery 42 │        13.29 / 13.45 ±0.21 / 13.86 ms │        12.90 / 15.69 ±5.15 / 25.99 ms │  1.17x slower │
└───────────┴───────────────────────────────────────┴───────────────────────────────────────┴───────────────┘
┏━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━┓
┃ Benchmark Summary                 ┃            ┃
┡━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━┩
│ Total Time (HEAD)                 │ 19345.53ms │
│ Total Time (agg-bucketed-final)   │ 19282.82ms │
│ Average Time (HEAD)               │   449.90ms │
│ Average Time (agg-bucketed-final) │   448.44ms │
│ Queries Faster                    │          2 │
│ Queries Slower                    │          5 │
│ Queries with No Change            │         36 │
│ Queries with Failure              │          0 │
└───────────────────────────────────┴────────────┘

Resource Usage

clickbench_partitioned — base (merge-base)

Metric Value
Wall time 100.0s
Peak memory 11.6 GiB
Avg memory 4.5 GiB
CPU user 991.3s
CPU sys 67.6s
Peak spill 0 B

clickbench_partitioned — branch

Metric Value
Wall time 100.0s
Peak memory 10.7 GiB
Avg memory 4.2 GiB
CPU user 986.1s
CPU sys 70.0s
Peak spill 0 B

File an issue against this benchmark runner

@Rachelint

Copy link
Copy Markdown
Contributor

run benchmark clickbench_partitioned
changed:
env:
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"

@adriangbot

Copy link
Copy Markdown

🤖 Benchmark running (GKE) | trigger
Instance: c4a-highmem-16 (12 vCPU / 65 GiB) | Linux bench-c5763080234-2548-s54kw 6.12.94+ #1 SMP Tue Aug 4 08:44:15 UTC 2026 aarch64 GNU/Linux

CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  16
On-line CPU(s) list:                     0-15
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     16
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               1 MiB (16 instances)
L1i cache:                               1 MiB (16 instances)
L2 cache:                                32 MiB (16 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-15
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected

Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff

Run configuration
run 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

@adriangbot

Copy link
Copy Markdown

🤖 Benchmark completed (GKE) | trigger

Instance: c4a-highmem-16 (12 vCPU / 65 GiB)

Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff

Run configuration
run benchmark clickbench_partitioned
changed:
  env:
    DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"
CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  16
On-line CPU(s) list:                     0-15
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     16
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               1 MiB (16 instances)
L1i cache:                               1 MiB (16 instances)
L2 cache:                                32 MiB (16 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-15
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected
Details

Comparing HEAD and agg-bucketed-final
--------------------
Benchmark clickbench_partitioned.json
--------------------
┏━━━━━━━━━━━┳━━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━━┓
┃ Query     ┃       HEAD ┃ agg-bucketed-final ┃        Change ┃
┡━━━━━━━━━━━╇━━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━━┩
│ QQuery 0  │    1.22 ms │            1.22 ms │     no change │
│ QQuery 1  │   11.66 ms │           11.67 ms │     no change │
│ QQuery 2  │   36.58 ms │           36.84 ms │     no change │
│ QQuery 3  │   31.00 ms │           31.08 ms │     no change │
│ QQuery 4  │  232.63 ms │          188.70 ms │ +1.23x faster │
│ QQuery 5  │  274.28 ms │          279.33 ms │     no change │
│ QQuery 6  │    1.27 ms │            1.31 ms │     no change │
│ QQuery 7  │   13.01 ms │           13.02 ms │     no change │
│ QQuery 8  │  331.00 ms │          293.19 ms │ +1.13x faster │
│ QQuery 9  │  470.54 ms │          468.52 ms │     no change │
│ QQuery 10 │   64.14 ms │           67.14 ms │     no change │
│ QQuery 11 │   75.78 ms │           77.79 ms │     no change │
│ QQuery 12 │  261.25 ms │          270.36 ms │     no change │
│ QQuery 13 │  357.77 ms │          372.96 ms │     no change │
│ QQuery 14 │  280.08 ms │          290.89 ms │     no change │
│ QQuery 15 │  280.38 ms │          220.53 ms │ +1.27x faster │
│ QQuery 16 │  611.96 ms │          645.96 ms │  1.06x slower │
│ QQuery 17 │  623.37 ms │          587.96 ms │ +1.06x faster │
│ QQuery 18 │ 1268.95 ms │         1161.90 ms │ +1.09x faster │
│ QQuery 19 │   27.37 ms │           27.61 ms │     no change │
│ QQuery 20 │  516.76 ms │          525.95 ms │     no change │
│ QQuery 21 │  503.75 ms │          516.70 ms │     no change │
│ QQuery 22 │  972.75 ms │          982.58 ms │     no change │
│ QQuery 23 │ 3044.78 ms │         3073.37 ms │     no change │
│ QQuery 24 │   41.09 ms │           40.62 ms │     no change │
│ QQuery 25 │  104.71 ms │          104.73 ms │     no change │
│ QQuery 26 │   41.24 ms │           40.87 ms │     no change │
│ QQuery 27 │  512.60 ms │          509.83 ms │     no change │
│ QQuery 28 │ 2906.57 ms │         2929.78 ms │     no change │
│ QQuery 29 │   41.51 ms │           41.45 ms │     no change │
│ QQuery 30 │  302.58 ms │          285.93 ms │ +1.06x faster │
│ QQuery 31 │  271.58 ms │          258.54 ms │     no change │
│ QQuery 32 │  931.16 ms │          744.51 ms │ +1.25x faster │
│ QQuery 33 │ 1406.73 ms │         1415.45 ms │     no change │
│ QQuery 34 │ 1445.94 ms │         1440.17 ms │     no change │
│ QQuery 35 │  289.77 ms │          229.52 ms │ +1.26x faster │
│ QQuery 36 │   66.77 ms │           66.61 ms │     no change │
│ QQuery 37 │   34.42 ms │           35.25 ms │     no change │
│ QQuery 38 │   41.46 ms │           42.21 ms │     no change │
│ QQuery 39 │  148.28 ms │          147.76 ms │     no change │
│ QQuery 40 │   14.07 ms │           14.73 ms │     no change │
│ QQuery 41 │   13.68 ms │           13.69 ms │     no change │
│ QQuery 42 │   12.80 ms │           13.22 ms │     no change │
└───────────┴────────────┴────────────────────┴───────────────┘
┏━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━┓
┃ Benchmark Summary                 ┃            ┃
┡━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━┩
│ Total Time (HEAD)                 │ 18919.23ms │
│ Total Time (agg-bucketed-final)   │ 18521.47ms │
│ Average Time (HEAD)               │   439.98ms │
│ Average Time (agg-bucketed-final) │   430.73ms │
│ Queries Faster                    │          8 │
│ Queries Slower                    │          1 │
│ Queries with No Change            │         34 │
│ Queries with Failure              │          0 │
└───────────────────────────────────┴────────────┘

Distribution per query (min / mean ±stddev / max):

Comparing HEAD and agg-bucketed-final
--------------------
Benchmark clickbench_partitioned.json
--------------------
┏━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━━┓
┃ Query     ┃                                  HEAD ┃                     agg-bucketed-final ┃        Change ┃
┡━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━━┩
│ QQuery 0  │          1.22 / 3.95 ±5.36 / 14.67 ms │           1.22 / 3.98 ±5.47 / 14.92 ms │     no change │
│ QQuery 1  │        11.66 / 11.86 ±0.13 / 11.98 ms │         11.67 / 11.80 ±0.12 / 11.97 ms │     no change │
│ QQuery 2  │        36.58 / 36.80 ±0.23 / 37.21 ms │         36.84 / 37.11 ±0.25 / 37.53 ms │     no change │
│ QQuery 3  │        31.00 / 31.74 ±0.64 / 32.89 ms │         31.08 / 32.95 ±2.39 / 37.68 ms │     no change │
│ QQuery 4  │     232.63 / 237.41 ±2.93 / 240.84 ms │      188.70 / 190.31 ±1.70 / 193.46 ms │ +1.25x faster │
│ QQuery 5  │     274.28 / 279.42 ±4.45 / 287.14 ms │      279.33 / 283.41 ±3.93 / 288.52 ms │     no change │
│ QQuery 6  │           1.27 / 1.43 ±0.23 / 1.89 ms │            1.31 / 1.47 ±0.24 / 1.93 ms │     no change │
│ QQuery 7  │        13.01 / 13.29 ±0.20 / 13.64 ms │         13.02 / 13.13 ±0.09 / 13.24 ms │     no change │
│ QQuery 8  │     331.00 / 334.00 ±2.52 / 338.65 ms │      293.19 / 296.22 ±3.14 / 301.12 ms │ +1.13x faster │
│ QQuery 9  │     470.54 / 479.02 ±6.97 / 488.02 ms │      468.52 / 473.78 ±3.30 / 478.00 ms │     no change │
│ QQuery 10 │        64.14 / 65.78 ±1.09 / 67.01 ms │         67.14 / 67.31 ±0.22 / 67.72 ms │     no change │
│ QQuery 11 │        75.78 / 76.28 ±0.73 / 77.73 ms │         77.79 / 78.94 ±0.84 / 80.32 ms │     no change │
│ QQuery 12 │     261.25 / 264.09 ±3.65 / 270.87 ms │      270.36 / 278.04 ±8.60 / 294.84 ms │  1.05x slower │
│ QQuery 13 │     357.77 / 364.73 ±5.39 / 371.09 ms │      372.96 / 380.82 ±3.99 / 383.55 ms │     no change │
│ QQuery 14 │     280.08 / 284.85 ±3.83 / 290.69 ms │      290.89 / 294.05 ±2.39 / 296.58 ms │     no change │
│ QQuery 15 │    280.38 / 296.17 ±14.27 / 319.15 ms │      220.53 / 225.94 ±4.69 / 233.14 ms │ +1.31x faster │
│ QQuery 16 │     611.96 / 620.63 ±6.82 / 631.37 ms │      645.96 / 659.22 ±9.21 / 669.69 ms │  1.06x slower │
│ QQuery 17 │     623.37 / 631.83 ±6.67 / 643.65 ms │     587.96 / 605.29 ±16.52 / 628.12 ms │     no change │
│ QQuery 18 │ 1268.95 / 1288.91 ±13.64 / 1308.44 ms │  1161.90 / 1233.77 ±42.80 / 1285.42 ms │     no change │
│ QQuery 19 │       27.37 / 33.33 ±11.50 / 56.33 ms │         27.61 / 27.81 ±0.16 / 28.01 ms │ +1.20x faster │
│ QQuery 20 │    516.76 / 535.42 ±21.38 / 573.19 ms │     525.95 / 549.13 ±34.65 / 617.73 ms │     no change │
│ QQuery 21 │     503.75 / 508.00 ±2.53 / 510.45 ms │     516.70 / 528.99 ±14.29 / 555.67 ms │     no change │
│ QQuery 22 │     972.75 / 981.66 ±6.86 / 992.73 ms │   982.58 / 1000.83 ±14.28 / 1024.97 ms │     no change │
│ QQuery 23 │ 3044.78 / 3080.50 ±26.90 / 3122.96 ms │  3073.37 / 3116.90 ±37.22 / 3163.29 ms │     no change │
│ QQuery 24 │        41.09 / 50.98 ±9.02 / 67.58 ms │         40.62 / 42.80 ±3.40 / 49.55 ms │ +1.19x faster │
│ QQuery 25 │     104.71 / 106.09 ±1.13 / 107.29 ms │      104.73 / 106.39 ±1.15 / 107.88 ms │     no change │
│ QQuery 26 │        41.24 / 42.46 ±1.69 / 45.78 ms │         40.87 / 41.71 ±0.81 / 43.09 ms │     no change │
│ QQuery 27 │     512.60 / 520.48 ±6.61 / 529.01 ms │      509.83 / 518.31 ±4.94 / 524.62 ms │     no change │
│ QQuery 28 │ 2906.57 / 2936.86 ±21.59 / 2959.94 ms │  2929.78 / 2949.44 ±18.94 / 2975.76 ms │     no change │
│ QQuery 29 │        41.51 / 48.12 ±8.50 / 62.68 ms │         41.45 / 44.79 ±5.92 / 56.62 ms │ +1.07x faster │
│ QQuery 30 │     302.58 / 308.29 ±4.12 / 314.10 ms │     285.93 / 302.96 ±26.61 / 355.23 ms │     no change │
│ QQuery 31 │     271.58 / 280.98 ±7.17 / 288.27 ms │     258.54 / 270.49 ±11.70 / 290.46 ms │     no change │
│ QQuery 32 │    931.16 / 943.60 ±20.14 / 983.80 ms │     744.51 / 761.58 ±10.73 / 773.64 ms │ +1.24x faster │
│ QQuery 33 │ 1406.73 / 1454.84 ±32.80 / 1488.31 ms │ 1415.45 / 1548.81 ±113.76 / 1729.83 ms │  1.06x slower │
│ QQuery 34 │ 1445.94 / 1500.67 ±41.37 / 1560.55 ms │  1440.17 / 1525.68 ±68.94 / 1633.22 ms │     no change │
│ QQuery 35 │    289.77 / 300.36 ±14.08 / 327.52 ms │     229.52 / 281.95 ±45.66 / 359.34 ms │ +1.07x faster │
│ QQuery 36 │       66.77 / 77.41 ±10.93 / 97.47 ms │         66.61 / 75.79 ±4.98 / 81.79 ms │     no change │
│ QQuery 37 │        34.42 / 35.49 ±0.83 / 36.97 ms │         35.25 / 39.21 ±7.07 / 53.33 ms │  1.10x slower │
│ QQuery 38 │        41.46 / 47.07 ±5.63 / 56.69 ms │         42.21 / 46.17 ±4.18 / 53.04 ms │     no change │
│ QQuery 39 │    148.28 / 154.96 ±10.70 / 176.30 ms │      147.76 / 157.16 ±5.76 / 164.59 ms │     no change │
│ QQuery 40 │        14.07 / 14.45 ±0.22 / 14.67 ms │         14.73 / 16.61 ±2.39 / 21.07 ms │  1.15x slower │
│ QQuery 41 │        13.68 / 16.69 ±3.55 / 22.80 ms │         13.69 / 14.25 ±0.53 / 15.22 ms │ +1.17x faster │
│ QQuery 42 │        12.80 / 13.07 ±0.21 / 13.42 ms │         13.22 / 14.76 ±2.81 / 20.39 ms │  1.13x slower │
└───────────┴───────────────────────────────────────┴────────────────────────────────────────┴───────────────┘
┏━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━┓
┃ Benchmark Summary                 ┃            ┃
┡━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━┩
│ Total Time (HEAD)                 │ 19313.95ms │
│ Total Time (agg-bucketed-final)   │ 19150.08ms │
│ Average Time (HEAD)               │   449.16ms │
│ Average Time (agg-bucketed-final) │   445.35ms │
│ Queries Faster                    │          9 │
│ Queries Slower                    │          6 │
│ Queries with No Change            │         28 │
│ Queries with Failure              │          0 │
└───────────────────────────────────┴────────────┘

Resource Usage

clickbench_partitioned — base (merge-base)

Metric Value
Wall time 100.0s
Peak memory 10.7 GiB
Avg memory 4.0 GiB
CPU user 989.7s
CPU sys 67.1s
Peak spill 0 B

clickbench_partitioned — branch

Metric Value
Wall time 100.0s
Peak memory 10.0 GiB
Avg memory 3.8 GiB
CPU user 945.9s
CPU sys 76.0s
Peak spill 0 B

File an issue against this benchmark runner

@github-actions

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-common v55.1.0 (current)
       Built [  34.816s] (current)
     Parsing datafusion-common v55.1.0 (current)
      Parsed [   0.061s] (current)
    Building datafusion-common v55.1.0 (baseline)
       Built [  33.193s] (baseline)
     Parsing datafusion-common v55.1.0 (baseline)
      Parsed [   0.060s] (baseline)
    Checking datafusion-common v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.799s] 223 checks: 222 pass, 1 fail, 0 warn, 31 skip

--- failure constructible_struct_adds_field: struct exhaustively constructible through public API adds field ---

Description:
A pub struct that could be exhaustively constructed with a literal using only public API has a new pub field, breaking existing exhaustive literals.
        ref: https://doc.rust-lang.org/reference/expressions/struct-expr.html
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/constructible_struct_adds_field.ron

Failed in:
  field ExecutionOptions.hash_aggregate_bucket_threshold in /home/runner/work/datafusion/datafusion/datafusion/common/src/config.rs:894

     Summary semver requires new major version: 1 major and 0 minor checks failed
    Finished [  70.261s] datafusion-common
    Building datafusion-physical-plan v55.1.0 (current)
       Built [  38.574s] (current)
     Parsing datafusion-physical-plan v55.1.0 (current)
      Parsed [   0.169s] (current)
    Building datafusion-physical-plan v55.1.0 (baseline)
       Built [  38.956s] (baseline)
     Parsing datafusion-physical-plan v55.1.0 (baseline)
      Parsed [   0.166s] (baseline)
    Checking datafusion-physical-plan v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.769s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [  80.015s] datafusion-physical-plan
    Building datafusion-sqllogictest v55.1.0 (current)
       Built [ 100.971s] (current)
     Parsing datafusion-sqllogictest v55.1.0 (current)
      Parsed [   0.024s] (current)
    Building datafusion-sqllogictest v55.1.0 (baseline)
       Built [ 100.939s] (baseline)
     Parsing datafusion-sqllogictest v55.1.0 (baseline)
      Parsed [   0.023s] (baseline)
    Checking datafusion-sqllogictest v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.105s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [ 205.027s] datafusion-sqllogictest

@github-actions github-actions Bot added the auto detected api change Auto detected API change label Sep 21, 2026
@codecov-commenter

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 87.57108% with 153 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.43%. Comparing base (6574a8c) to head (d4c6299).
⚠️ Report is 5 commits behind head on main.

Files with missing lines Patch % Lines
...fusion/physical-plan/src/aggregates/hash_stream.rs 89.36% 10 Missing and 47 partials ⚠️
...sion/physical-plan/src/aggregates/final_buckets.rs 86.69% 8 Missing and 29 partials ⚠️
...ysical-plan/src/aggregates/bucketed_aggregation.rs 84.79% 10 Missing and 23 partials ⚠️
...sion/physical-plan/src/aggregates/single_stream.rs 80.37% 16 Missing and 5 partials ⚠️
...src/aggregates/aggregate_hash_table/final_table.rs 90.00% 3 Missing ⚠️
...plan/src/aggregates/aggregate_hash_table/common.rs 96.72% 0 Missing and 2 partials ⚠️
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.
📢 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.

@Rachelint

Rachelint commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

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

Where the remaining loss comes from (per final input row): Q14 costs 79-85 ns with one table and 91-99 ns bucketed...

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.
That said, if we cannot avoid repartitioning the RecordBatch again inside the final aggregation, there seems to be little or no benefit for string group keys, and it can even cause a regression. This has been troubling me for a while. It is not limited to Q12–Q14 either; Q33 and Q34 also show significant regressions.

And finally I have no better idea than just disable the optimization when found string column.

Another possible opportunity to improve more performance

In 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).
I later switched to MutableArrayData to build the bucket batches column by column, which improved this part of the process. A similar approach might be worth exploring as a follow-up optimization for this PR.
The implementation was rather awkward, however, because MutableArrayData does not support continuously appending after an output array has been materialized. What I really wanted was a column-oriented coalescer—essentially direct access to the InProgressArray abstraction used internally by Arrow’s BatchCoalescer.
I have wanted to fork Arrow and experiment with exposing that interface publicly, but unfortunately I have not had time to do so yet.

@Rachelint

Copy link
Copy Markdown
Contributor

run benchmark clickbench_partitioned
changed:
env:
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"

@adriangbot

Copy link
Copy Markdown

Hi @Rachelint, your benchmark configuration could not be parsed (#25567 (comment)).

Error: invalid configuration: unknown field DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD, expected one of env, baseline, changed at line 3 column 1

Usage:

run benchmark <name>           # run specific benchmark(s)
run benchmarks                 # run default suite
run benchmarks <name1> <name2> # run specific benchmarks

Any benchmark name is accepted: bench.sh suite names (e.g. tpch, clickbench_partitioned, wide_schema) and Criterion bench targets (e.g. sql_planner) are resolved automatically. A name that matches neither fails on the runner.

Per-side configuration (run benchmark tpch followed by):

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: 2G

File an issue against this benchmark runner

@Rachelint

Copy link
Copy Markdown
Contributor

run benchmark clickbench_partitioned
changed:
env:
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"

@adriangbot

Copy link
Copy Markdown

🤖 Benchmark running (GKE) | trigger
Instance: c4a-highmem-16 (12 vCPU / 65 GiB) | Linux bench-c5770623516-2573-qvmnr 6.12.94+ #1 SMP Wed Aug 19 07:47:20 UTC 2026 aarch64 GNU/Linux

CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  16
On-line CPU(s) list:                     0-15
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     16
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               1 MiB (16 instances)
L1i cache:                               1 MiB (16 instances)
L2 cache:                                32 MiB (16 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-15
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected

Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff

Run configuration
run 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

@adriangbot

Copy link
Copy Markdown

🤖 Benchmark completed (GKE) | trigger

Instance: c4a-highmem-16 (12 vCPU / 65 GiB)

Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff

Run configuration
run benchmark clickbench_partitioned
changed:
  env:
    DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"
CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  16
On-line CPU(s) list:                     0-15
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     16
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               1 MiB (16 instances)
L1i cache:                               1 MiB (16 instances)
L2 cache:                                32 MiB (16 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-15
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected
Details

Comparing HEAD and agg-bucketed-final
--------------------
Benchmark clickbench_partitioned.json
--------------------
┏━━━━━━━━━━━┳━━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━━┓
┃ Query     ┃       HEAD ┃ agg-bucketed-final ┃        Change ┃
┡━━━━━━━━━━━╇━━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━━┩
│ QQuery 0  │    1.36 ms │            1.27 ms │ +1.07x faster │
│ QQuery 1  │   12.39 ms │           11.84 ms │     no change │
│ QQuery 2  │   37.96 ms │           37.01 ms │     no change │
│ QQuery 3  │   33.14 ms │           32.22 ms │     no change │
│ QQuery 4  │  256.16 ms │          190.21 ms │ +1.35x faster │
│ QQuery 5  │  287.77 ms │          285.95 ms │     no change │
│ QQuery 6  │    1.33 ms │            1.35 ms │     no change │
│ QQuery 7  │   13.54 ms │           13.59 ms │     no change │
│ QQuery 8  │  365.96 ms │          301.23 ms │ +1.21x faster │
│ QQuery 9  │  527.24 ms │          522.65 ms │     no change │
│ QQuery 10 │   69.62 ms │           69.47 ms │     no change │
│ QQuery 11 │   79.68 ms │           80.31 ms │     no change │
│ QQuery 12 │  293.55 ms │          295.02 ms │     no change │
│ QQuery 13 │  400.05 ms │          400.67 ms │     no change │
│ QQuery 14 │  306.49 ms │          310.58 ms │     no change │
│ QQuery 15 │  334.09 ms │          244.44 ms │ +1.37x faster │
│ QQuery 16 │  685.82 ms │          734.17 ms │  1.07x slower │
│ QQuery 17 │  687.59 ms │          640.91 ms │ +1.07x faster │
│ QQuery 18 │ 1429.75 ms │         1238.31 ms │ +1.15x faster │
│ QQuery 19 │   29.06 ms │           29.18 ms │     no change │
│ QQuery 20 │  529.66 ms │          538.23 ms │     no change │
│ QQuery 21 │  520.43 ms │          530.17 ms │     no change │
│ QQuery 22 │ 1005.49 ms │         1050.47 ms │     no change │
│ QQuery 23 │ 3256.35 ms │         3342.26 ms │     no change │
│ QQuery 24 │   41.86 ms │           41.63 ms │     no change │
│ QQuery 25 │  108.73 ms │          107.56 ms │     no change │
│ QQuery 26 │   43.15 ms │           42.48 ms │     no change │
│ QQuery 27 │  528.14 ms │          531.44 ms │     no change │
│ QQuery 28 │ 2989.46 ms │         2976.58 ms │     no change │
│ QQuery 29 │   42.59 ms │           42.75 ms │     no change │
│ QQuery 30 │  330.69 ms │          296.71 ms │ +1.11x faster │
│ QQuery 31 │  301.33 ms │          257.51 ms │ +1.17x faster │
│ QQuery 32 │ 1054.47 ms │          774.59 ms │ +1.36x faster │
│ QQuery 33 │ 1569.88 ms │         1584.76 ms │     no change │
│ QQuery 34 │ 1589.02 ms │         1637.91 ms │     no change │
│ QQuery 35 │  361.15 ms │          247.61 ms │ +1.46x faster │
│ QQuery 36 │   74.68 ms │           70.52 ms │ +1.06x faster │
│ QQuery 37 │   39.72 ms │           37.01 ms │ +1.07x faster │
│ QQuery 38 │   46.40 ms │           41.09 ms │ +1.13x faster │
│ QQuery 39 │  178.64 ms │          147.09 ms │ +1.21x faster │
│ QQuery 40 │   17.24 ms │           15.33 ms │ +1.12x faster │
│ QQuery 41 │   16.12 ms │           14.49 ms │ +1.11x faster │
│ QQuery 42 │   15.34 ms │           14.18 ms │ +1.08x faster │
└───────────┴────────────┴────────────────────┴───────────────┘
┏━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━┓
┃ Benchmark Summary                 ┃            ┃
┡━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━┩
│ Total Time (HEAD)                 │ 20513.10ms │
│ Total Time (agg-bucketed-final)   │ 19782.78ms │
│ Average Time (HEAD)               │   477.05ms │
│ Average Time (agg-bucketed-final) │   460.06ms │
│ Queries Faster                    │         17 │
│ Queries Slower                    │          1 │
│ Queries with No Change            │         25 │
│ Queries with Failure              │          0 │
└───────────────────────────────────┴────────────┘

Distribution per query (min / mean ±stddev / max):

Comparing HEAD and agg-bucketed-final
--------------------
Benchmark clickbench_partitioned.json
--------------------
┏━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━━┓
┃ Query     ┃                                  HEAD ┃                    agg-bucketed-final ┃        Change ┃
┡━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━━┩
│ QQuery 0  │          1.36 / 4.49 ±5.97 / 16.43 ms │          1.27 / 4.22 ±5.74 / 15.71 ms │ +1.06x faster │
│ QQuery 1  │        12.39 / 12.63 ±0.19 / 12.95 ms │        11.84 / 12.22 ±0.26 / 12.64 ms │     no change │
│ QQuery 2  │        37.96 / 38.11 ±0.13 / 38.29 ms │        37.01 / 37.62 ±0.35 / 38.05 ms │     no change │
│ QQuery 3  │        33.14 / 34.40 ±1.21 / 36.33 ms │        32.22 / 32.79 ±0.38 / 33.29 ms │     no change │
│ QQuery 4  │     256.16 / 262.70 ±4.53 / 269.98 ms │     190.21 / 192.42 ±1.66 / 194.47 ms │ +1.37x faster │
│ QQuery 5  │     287.77 / 297.77 ±6.21 / 304.83 ms │     285.95 / 293.52 ±5.32 / 300.77 ms │     no change │
│ QQuery 6  │           1.33 / 1.49 ±0.24 / 1.96 ms │           1.35 / 1.52 ±0.24 / 1.99 ms │     no change │
│ QQuery 7  │        13.54 / 13.88 ±0.19 / 14.08 ms │        13.59 / 13.66 ±0.04 / 13.71 ms │     no change │
│ QQuery 8  │     365.96 / 384.40 ±9.31 / 391.25 ms │     301.23 / 311.16 ±8.95 / 325.39 ms │ +1.24x faster │
│ QQuery 9  │     527.24 / 538.66 ±9.63 / 553.69 ms │    522.65 / 538.35 ±13.66 / 558.55 ms │     no change │
│ QQuery 10 │        69.62 / 71.90 ±3.15 / 78.02 ms │       69.47 / 75.67 ±11.36 / 98.38 ms │  1.05x slower │
│ QQuery 11 │        79.68 / 81.25 ±1.54 / 83.68 ms │        80.31 / 81.22 ±0.78 / 82.43 ms │     no change │
│ QQuery 12 │     293.55 / 299.33 ±6.05 / 307.83 ms │    295.02 / 306.18 ±12.61 / 330.73 ms │     no change │
│ QQuery 13 │    400.05 / 423.60 ±13.43 / 440.73 ms │    400.67 / 423.49 ±12.74 / 436.14 ms │     no change │
│ QQuery 14 │     306.49 / 318.11 ±9.93 / 336.59 ms │    310.58 / 339.40 ±26.79 / 388.11 ms │  1.07x slower │
│ QQuery 15 │     334.09 / 343.17 ±4.76 / 347.03 ms │     244.44 / 245.60 ±1.15 / 247.64 ms │ +1.40x faster │
│ QQuery 16 │    685.82 / 710.92 ±16.52 / 728.09 ms │    734.17 / 755.31 ±20.46 / 780.88 ms │  1.06x slower │
│ QQuery 17 │    687.59 / 714.22 ±22.00 / 746.63 ms │    640.91 / 660.59 ±17.78 / 688.63 ms │ +1.08x faster │
│ QQuery 18 │ 1429.75 / 1470.29 ±25.55 / 1499.53 ms │ 1238.31 / 1317.36 ±66.23 / 1425.02 ms │ +1.12x faster │
│ QQuery 19 │       29.06 / 38.60 ±11.87 / 59.79 ms │      29.18 / 89.16 ±77.28 / 219.92 ms │  2.31x slower │
│ QQuery 20 │    529.66 / 536.97 ±10.85 / 558.45 ms │     538.23 / 548.94 ±6.87 / 559.26 ms │     no change │
│ QQuery 21 │     520.43 / 526.72 ±3.65 / 530.08 ms │     530.17 / 541.71 ±8.24 / 553.50 ms │     no change │
│ QQuery 22 │ 1005.49 / 1026.53 ±11.95 / 1042.54 ms │ 1050.47 / 1072.57 ±20.38 / 1098.78 ms │     no change │
│ QQuery 23 │ 3256.35 / 3380.79 ±89.45 / 3478.08 ms │ 3342.26 / 3379.49 ±23.16 / 3402.91 ms │     no change │
│ QQuery 24 │        41.86 / 44.19 ±4.43 / 53.06 ms │       41.63 / 49.37 ±12.44 / 74.14 ms │  1.12x slower │
│ QQuery 25 │     108.73 / 117.43 ±7.82 / 130.41 ms │     107.56 / 110.38 ±2.66 / 113.77 ms │ +1.06x faster │
│ QQuery 26 │        43.15 / 43.98 ±0.86 / 45.57 ms │        42.48 / 46.95 ±5.23 / 56.83 ms │  1.07x slower │
│ QQuery 27 │    528.14 / 544.81 ±10.26 / 556.53 ms │    531.44 / 543.09 ±10.40 / 556.71 ms │     no change │
│ QQuery 28 │ 2989.46 / 3055.55 ±41.29 / 3116.91 ms │ 2976.58 / 3033.43 ±45.07 / 3113.84 ms │     no change │
│ QQuery 29 │        42.59 / 49.42 ±6.91 / 59.00 ms │       42.75 / 52.01 ±13.19 / 77.26 ms │  1.05x slower │
│ QQuery 30 │     330.69 / 340.80 ±9.10 / 357.20 ms │    296.71 / 308.41 ±13.05 / 331.62 ms │ +1.11x faster │
│ QQuery 31 │    301.33 / 311.66 ±12.66 / 335.78 ms │     257.51 / 262.06 ±4.88 / 270.12 ms │ +1.19x faster │
│ QQuery 32 │ 1054.47 / 1073.55 ±16.36 / 1099.79 ms │  774.59 / 883.40 ±127.50 / 1131.59 ms │ +1.22x faster │
│ QQuery 33 │ 1569.88 / 1616.50 ±70.68 / 1757.36 ms │ 1584.76 / 1673.30 ±87.21 / 1825.41 ms │     no change │
│ QQuery 34 │ 1589.02 / 1668.18 ±65.97 / 1756.72 ms │ 1637.91 / 1758.02 ±76.04 / 1874.21 ms │  1.05x slower │
│ QQuery 35 │    361.15 / 395.29 ±36.74 / 458.28 ms │   247.61 / 325.66 ±110.93 / 533.18 ms │ +1.21x faster │
│ QQuery 36 │        74.68 / 76.83 ±2.30 / 81.07 ms │        70.52 / 76.66 ±5.05 / 84.49 ms │     no change │
│ QQuery 37 │        39.72 / 45.70 ±7.41 / 58.39 ms │        37.01 / 38.48 ±1.83 / 42.01 ms │ +1.19x faster │
│ QQuery 38 │        46.40 / 51.07 ±3.60 / 54.67 ms │        41.09 / 43.29 ±1.56 / 45.37 ms │ +1.18x faster │
│ QQuery 39 │     178.64 / 188.65 ±6.47 / 198.01 ms │     147.09 / 152.58 ±3.13 / 156.64 ms │ +1.24x faster │
│ QQuery 40 │        17.24 / 20.91 ±4.70 / 29.79 ms │        15.33 / 17.32 ±2.67 / 22.56 ms │ +1.21x faster │
│ QQuery 41 │        16.12 / 19.43 ±6.37 / 32.17 ms │        14.49 / 15.20 ±0.62 / 16.33 ms │ +1.28x faster │
│ QQuery 42 │        15.34 / 16.49 ±1.83 / 20.13 ms │        14.18 / 14.35 ±0.19 / 14.63 ms │ +1.15x faster │
└───────────┴───────────────────────────────────────┴───────────────────────────────────────┴───────────────┘
┏━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━┓
┃ Benchmark Summary                 ┃            ┃
┡━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━┩
│ Total Time (HEAD)                 │ 21211.35ms │
│ Total Time (agg-bucketed-final)   │ 20678.14ms │
│ Average Time (HEAD)               │   493.29ms │
│ Average Time (agg-bucketed-final) │   480.89ms │
│ Queries Faster                    │         17 │
│ Queries Slower                    │          8 │
│ Queries with No Change            │         18 │
│ Queries with Failure              │          0 │
└───────────────────────────────────┴────────────┘

Resource Usage

clickbench_partitioned — base (merge-base)

Metric Value
Wall time 110.0s
Peak memory 10.6 GiB
Avg memory 4.3 GiB
CPU user 1077.7s
CPU sys 84.5s
Peak spill 0 B

clickbench_partitioned — branch

Metric Value
Wall time 105.0s
Peak memory 9.7 GiB
Avg memory 3.9 GiB
CPU user 1005.5s
CPU sys 92.5s
Peak spill 0 B

File an issue against this benchmark runner

@Rachelint

Copy link
Copy Markdown
Contributor

run benchmark clickbench_partitioned
changed:
env:
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"

@adriangbot

Copy link
Copy Markdown

Hi @Rachelint, your benchmark configuration could not be parsed (#25567 (comment)).

Error: invalid configuration: unknown field DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD, expected one of env, baseline, changed at line 3 column 1

Usage:

run benchmark <name>           # run specific benchmark(s)
run benchmarks                 # run default suite
run benchmarks <name1> <name2> # run specific benchmarks

Any benchmark name is accepted: bench.sh suite names (e.g. tpch, clickbench_partitioned, wide_schema) and Criterion bench targets (e.g. sql_planner) are resolved automatically. A name that matches neither fails on the runner.

Per-side configuration (run benchmark tpch followed by):

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: 2G

File an issue against this benchmark runner

@Rachelint

Copy link
Copy Markdown
Contributor

run benchmark clickbench_partitioned
changed:
env:
DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"

@adriangbot

Copy link
Copy Markdown

🤖 Benchmark running (GKE) | trigger
Instance: c4a-highmem-16 (12 vCPU / 65 GiB) | Linux bench-c5770778035-2576-92xmt 6.12.94+ #1 SMP Wed Aug 19 07:47:20 UTC 2026 aarch64 GNU/Linux

CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  16
On-line CPU(s) list:                     0-15
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     16
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               1 MiB (16 instances)
L1i cache:                               1 MiB (16 instances)
L2 cache:                                32 MiB (16 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-15
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected

Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff

Run configuration
run 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

@adriangbot

Copy link
Copy Markdown

🤖 Benchmark completed (GKE) | trigger

Instance: c4a-highmem-16 (12 vCPU / 65 GiB)

Comparing agg-bucketed-final (d4c6299) to 6574a8c (merge-base) diff

Run configuration
run benchmark clickbench_partitioned
changed:
  env:
    DATAFUSION_EXECUTION_HASH_AGGREGATE_BUCKET_THRESHOLD: "262144"
CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  16
On-line CPU(s) list:                     0-15
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     16
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               1 MiB (16 instances)
L1i cache:                               1 MiB (16 instances)
L2 cache:                                32 MiB (16 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-15
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected
Details

Comparing HEAD and agg-bucketed-final
--------------------
Benchmark clickbench_partitioned.json
--------------------
┏━━━━━━━━━━━┳━━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━━┓
┃ Query     ┃       HEAD ┃ agg-bucketed-final ┃        Change ┃
┡━━━━━━━━━━━╇━━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━━┩
│ QQuery 0  │    1.28 ms │            1.22 ms │     no change │
│ QQuery 1  │   12.02 ms │           11.81 ms │     no change │
│ QQuery 2  │   37.20 ms │           36.83 ms │     no change │
│ QQuery 3  │   32.15 ms │           31.40 ms │     no change │
│ QQuery 4  │  249.54 ms │          188.98 ms │ +1.32x faster │
│ QQuery 5  │  283.60 ms │          279.81 ms │     no change │
│ QQuery 6  │    1.29 ms │            1.27 ms │     no change │
│ QQuery 7  │   13.54 ms │           13.22 ms │     no change │
│ QQuery 8  │  345.40 ms │          296.51 ms │ +1.16x faster │
│ QQuery 9  │  497.19 ms │          475.23 ms │     no change │
│ QQuery 10 │   67.36 ms │           67.15 ms │     no change │
│ QQuery 11 │   77.97 ms │           79.07 ms │     no change │
│ QQuery 12 │  286.27 ms │          272.02 ms │     no change │
│ QQuery 13 │  380.65 ms │          378.05 ms │     no change │
│ QQuery 14 │  295.74 ms │          293.74 ms │     no change │
│ QQuery 15 │  300.52 ms │          219.94 ms │ +1.37x faster │
│ QQuery 16 │  644.90 ms │          666.12 ms │     no change │
│ QQuery 17 │  644.54 ms │          578.51 ms │ +1.11x faster │
│ QQuery 18 │ 1318.24 ms │         1197.63 ms │ +1.10x faster │
│ QQuery 19 │   27.97 ms │           27.71 ms │     no change │
│ QQuery 20 │  523.57 ms │          525.05 ms │     no change │
│ QQuery 21 │  516.25 ms │          512.56 ms │     no change │
│ QQuery 22 │  989.71 ms │          998.06 ms │     no change │
│ QQuery 23 │ 3167.07 ms │         3078.10 ms │     no change │
│ QQuery 24 │   42.67 ms │           40.69 ms │     no change │
│ QQuery 25 │  108.06 ms │          105.18 ms │     no change │
│ QQuery 26 │   42.38 ms │           42.56 ms │     no change │
│ QQuery 27 │  521.97 ms │          511.85 ms │     no change │
│ QQuery 28 │ 2988.57 ms │         2907.93 ms │     no change │
│ QQuery 29 │   42.56 ms │           41.64 ms │     no change │
│ QQuery 30 │  319.91 ms │          291.17 ms │ +1.10x faster │
│ QQuery 31 │  295.88 ms │          251.45 ms │ +1.18x faster │
│ QQuery 32 │  994.78 ms │          739.53 ms │ +1.35x faster │
│ QQuery 33 │ 1499.71 ms │         1444.39 ms │     no change │
│ QQuery 34 │ 1519.55 ms │         1474.66 ms │     no change │
│ QQuery 35 │  307.17 ms │          227.91 ms │ +1.35x faster │
│ QQuery 36 │   65.76 ms │           67.49 ms │     no change │
│ QQuery 37 │   35.68 ms │           35.67 ms │     no change │
│ QQuery 38 │   41.29 ms │           45.21 ms │  1.09x slower │
│ QQuery 39 │  147.97 ms │          139.11 ms │ +1.06x faster │
│ QQuery 40 │   14.82 ms │           14.66 ms │     no change │
│ QQuery 41 │   14.06 ms │           14.01 ms │     no change │
│ QQuery 42 │   13.64 ms │           13.48 ms │     no change │
└───────────┴────────────┴────────────────────┴───────────────┘
┏━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━┓
┃ Benchmark Summary                 ┃            ┃
┡━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━┩
│ Total Time (HEAD)                 │ 19730.42ms │
│ Total Time (agg-bucketed-final)   │ 18638.58ms │
│ Average Time (HEAD)               │   458.85ms │
│ Average Time (agg-bucketed-final) │   433.46ms │
│ Queries Faster                    │         10 │
│ Queries Slower                    │          1 │
│ Queries with No Change            │         32 │
│ Queries with Failure              │          0 │
└───────────────────────────────────┴────────────┘

Distribution per query (min / mean ±stddev / max):

Comparing HEAD and agg-bucketed-final
--------------------
Benchmark clickbench_partitioned.json
--------------------
┏━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━━━━┓
┃ Query     ┃                                  HEAD ┃                    agg-bucketed-final ┃        Change ┃
┡━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━━━━┩
│ QQuery 0  │          1.28 / 4.33 ±5.87 / 16.07 ms │          1.22 / 4.10 ±5.61 / 15.32 ms │ +1.06x faster │
│ QQuery 1  │        12.02 / 12.56 ±0.30 / 12.83 ms │        11.81 / 12.08 ±0.19 / 12.40 ms │     no change │
│ QQuery 2  │        37.20 / 37.42 ±0.28 / 37.91 ms │        36.83 / 37.07 ±0.21 / 37.44 ms │     no change │
│ QQuery 3  │        32.15 / 32.98 ±1.12 / 35.13 ms │        31.40 / 31.63 ±0.15 / 31.83 ms │     no change │
│ QQuery 4  │     249.54 / 254.64 ±3.63 / 258.99 ms │     188.98 / 192.00 ±2.94 / 196.37 ms │ +1.33x faster │
│ QQuery 5  │     283.60 / 287.18 ±2.00 / 289.05 ms │     279.81 / 288.79 ±6.93 / 301.01 ms │     no change │
│ QQuery 6  │           1.29 / 1.44 ±0.24 / 1.91 ms │           1.27 / 1.42 ±0.24 / 1.91 ms │     no change │
│ QQuery 7  │        13.54 / 14.08 ±0.31 / 14.45 ms │        13.22 / 13.45 ±0.14 / 13.60 ms │     no change │
│ QQuery 8  │     345.40 / 353.54 ±5.73 / 363.07 ms │     296.51 / 303.74 ±5.20 / 310.74 ms │ +1.16x faster │
│ QQuery 9  │    497.19 / 511.42 ±11.62 / 530.16 ms │     475.23 / 486.98 ±7.64 / 496.38 ms │     no change │
│ QQuery 10 │        67.36 / 69.09 ±1.20 / 70.75 ms │        67.15 / 67.53 ±0.44 / 68.40 ms │     no change │
│ QQuery 11 │        77.97 / 78.86 ±1.18 / 80.95 ms │        79.07 / 82.82 ±6.02 / 94.72 ms │  1.05x slower │
│ QQuery 12 │     286.27 / 289.28 ±3.10 / 295.23 ms │     272.02 / 280.29 ±6.99 / 291.26 ms │     no change │
│ QQuery 13 │    380.65 / 397.07 ±12.57 / 418.61 ms │     378.05 / 385.64 ±4.69 / 390.19 ms │     no change │
│ QQuery 14 │     295.74 / 300.53 ±4.32 / 307.92 ms │     293.74 / 301.52 ±6.52 / 309.41 ms │     no change │
│ QQuery 15 │     300.52 / 307.07 ±5.11 / 314.56 ms │     219.94 / 224.87 ±5.10 / 233.28 ms │ +1.37x faster │
│ QQuery 16 │     644.90 / 651.06 ±4.79 / 656.79 ms │    666.12 / 690.56 ±17.39 / 717.77 ms │  1.06x slower │
│ QQuery 17 │    644.54 / 657.92 ±12.78 / 677.84 ms │    578.51 / 602.95 ±13.12 / 614.57 ms │ +1.09x faster │
│ QQuery 18 │ 1318.24 / 1337.40 ±16.82 / 1364.88 ms │ 1197.63 / 1249.62 ±30.31 / 1283.40 ms │ +1.07x faster │
│ QQuery 19 │        27.97 / 30.47 ±4.03 / 38.51 ms │       27.71 / 41.44 ±27.04 / 95.51 ms │  1.36x slower │
│ QQuery 20 │    523.57 / 544.57 ±28.62 / 597.74 ms │    525.05 / 543.49 ±24.13 / 590.20 ms │     no change │
│ QQuery 21 │     516.25 / 519.35 ±3.50 / 525.95 ms │     512.56 / 524.85 ±7.68 / 535.57 ms │     no change │
│ QQuery 22 │   989.71 / 1001.02 ±7.75 / 1009.40 ms │  998.06 / 1022.72 ±19.04 / 1052.54 ms │     no change │
│ QQuery 23 │ 3167.07 / 3203.15 ±28.28 / 3238.08 ms │ 3078.10 / 3129.53 ±29.48 / 3168.32 ms │     no change │
│ QQuery 24 │        42.67 / 44.90 ±3.56 / 51.99 ms │        40.69 / 42.49 ±1.53 / 45.33 ms │ +1.06x faster │
│ QQuery 25 │     108.06 / 113.54 ±5.63 / 121.27 ms │     105.18 / 107.44 ±4.07 / 115.57 ms │ +1.06x faster │
│ QQuery 26 │        42.38 / 43.59 ±1.13 / 45.72 ms │        42.56 / 48.10 ±7.42 / 61.85 ms │  1.10x slower │
│ QQuery 27 │     521.97 / 527.16 ±7.98 / 542.97 ms │     511.85 / 526.49 ±8.06 / 535.21 ms │     no change │
│ QQuery 28 │ 2988.57 / 3012.42 ±19.03 / 3043.78 ms │ 2907.93 / 2952.74 ±35.00 / 3015.12 ms │     no change │
│ QQuery 29 │       42.56 / 49.66 ±12.74 / 75.07 ms │        41.64 / 42.20 ±0.55 / 43.25 ms │ +1.18x faster │
│ QQuery 30 │     319.91 / 329.24 ±8.55 / 344.89 ms │    291.17 / 313.27 ±32.27 / 375.94 ms │     no change │
│ QQuery 31 │     295.88 / 305.80 ±9.11 / 319.75 ms │     251.45 / 258.27 ±4.24 / 264.30 ms │ +1.18x faster │
│ QQuery 32 │  994.78 / 1019.77 ±17.66 / 1037.84 ms │     739.53 / 748.98 ±5.46 / 756.60 ms │ +1.36x faster │
│ QQuery 33 │ 1499.71 / 1525.98 ±23.35 / 1564.08 ms │ 1444.39 / 1542.37 ±83.97 / 1658.59 ms │     no change │
│ QQuery 34 │ 1519.55 / 1543.88 ±20.58 / 1576.56 ms │ 1474.66 / 1502.14 ±25.52 / 1548.81 ms │     no change │
│ QQuery 35 │    307.17 / 357.33 ±78.56 / 513.83 ms │    227.91 / 259.35 ±28.91 / 302.70 ms │ +1.38x faster │
│ QQuery 36 │        65.76 / 73.68 ±7.39 / 86.29 ms │        67.49 / 73.24 ±6.70 / 84.75 ms │     no change │
│ QQuery 37 │        35.68 / 40.15 ±3.87 / 44.76 ms │        35.67 / 37.89 ±3.13 / 43.96 ms │ +1.06x faster │
│ QQuery 38 │        41.29 / 44.71 ±4.48 / 53.51 ms │        45.21 / 47.79 ±1.99 / 50.46 ms │  1.07x slower │
│ QQuery 39 │     147.97 / 157.92 ±7.21 / 166.97 ms │     139.11 / 150.21 ±8.40 / 158.79 ms │     no change │
│ QQuery 40 │        14.82 / 16.90 ±3.18 / 23.14 ms │        14.66 / 16.65 ±2.96 / 22.51 ms │     no change │
│ QQuery 41 │        14.06 / 14.26 ±0.13 / 14.42 ms │        14.01 / 14.89 ±1.27 / 17.35 ms │     no change │
│ QQuery 42 │        13.64 / 14.96 ±2.18 / 19.30 ms │        13.48 / 14.99 ±2.78 / 20.53 ms │     no change │
└───────────┴───────────────────────────────────────┴───────────────────────────────────────┴───────────────┘
┏━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┳━━━━━━━━━━━━┓
┃ Benchmark Summary                 ┃            ┃
┡━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━╇━━━━━━━━━━━━┩
│ Total Time (HEAD)                 │ 20132.28ms │
│ Total Time (agg-bucketed-final)   │ 19218.58ms │
│ Average Time (HEAD)               │   468.19ms │
│ Average Time (agg-bucketed-final) │   446.94ms │
│ Queries Faster                    │         13 │
│ Queries Slower                    │          5 │
│ Queries with No Change            │         25 │
│ Queries with Failure              │          0 │
└───────────────────────────────────┴────────────┘

Resource Usage

clickbench_partitioned — base (merge-base)

Metric Value
Wall time 105.0s
Peak memory 10.3 GiB
Avg memory 4.2 GiB
CPU user 1029.5s
CPU sys 71.6s
Peak spill 0 B

clickbench_partitioned — branch

Metric Value
Wall time 100.0s
Peak memory 7.9 GiB
Avg memory 3.5 GiB
CPU user 952.6s
CPU sys 78.2s
Peak spill 0 B

File an issue against this benchmark runner

@jayzhan211

Copy link
Copy Markdown
Contributor Author

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 aggregates/spill.rs):

This PR is superseded by those; happy to close it once you have had a look.

On string group keys

Two commits landed after you looked here, and both target exactly this case:

  • the table moves into buckets at a quarter of the threshold when the input holds about one row per group, so far less work is repeated;
  • the table of a bucket keeps its group keys where its input already holds them instead of copying every new key into its own builder. A bucket's rows sit in one contiguous block that is dropped as soon as that bucket has been aggregated, so the table can store a view into that block. This removes 17-19 ns per row on multi-column string keys — Q14's group-id time goes 44.7 → 25.7 ns/row, Q16 54.8 → 46.9.

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:

query group key groups / partition ratio
Q18 UserID, minute, SearchPhrase 4.7 M 0.56x
Q33 URL 1.5 M 0.63x
Q34 1, URL 1.5 M 0.87x
Q11 SearchPhrase ~0.5 M 1.02x
Q12 SearchPhrase 502 k 1.07x
Q14 SearchEngineID, SearchPhrase 539 k 1.08x

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 MutableArrayData and a column-oriented coalescer

I owe you a negative result here, because I tried the closest thing available today and it went the wrong way.

Arrow's BatchCoalescer::push_batch_with_filter has a fused sparse path: at 1/64 selectivity it filters the views of a Utf8View column and reuses the source data buffers, so no string bytes are copied at all, and the take disappears too. I replaced the take + slice + per-bucket push with one selection mask per bucket and that call.

Routing got cheaper exactly as intended:

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.

@Rachelint

Copy link
Copy Markdown
Contributor

Sorry, I’ve been away for the past two days. I’m back home today and will take a detailed look.

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

Labels

auto detected api change Auto detected API change common Related to common crate documentation Improvements or additions to documentation physical-plan Changes to the physical-plan crate sqllogictest SQL Logic Tests (.slt)

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants