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]

Reply via email to