jayzhan211 commented on PR #25567: URL: https://github.com/apache/datafusion/pull/25567#issuecomment-5825599720
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`): - **#25724** — the complete final-stage feature, 6 commits, off by default - **#25725** — the single-stage `SinglePartitioned` extension, stacked on it 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. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
