zhuqi-lucas opened a new issue, #23937:
URL: https://github.com/apache/datafusion/issues/23937

   ## Is your feature request related to a problem or challenge?
   
   `GroupValues::intern` pays the full cost — hash computation plus a 
hash-table probe — for **every input row**, even when consecutive rows carry 
the same group key. The probe is a random access into a table that, for 
high-cardinality GROUP BY, is far too large to stay cache-resident, so it 
typically costs an L3/DRAM round trip.
   
   Real-world data very often arrives with **runs of identical keys** 
(time-ordered logs, clustered writes). Measured on ClickBench `hits.parquet` 
(sampled row groups, file order — which is exactly what the Partial aggregation 
phase sees, since it runs before repartitioning):
   
   | column | avg run length | last-group hit rate |
   |---|---:|---:|
   | CounterID | 58,693 | ~100% |
   | EventDate | 39,129 | ~100% |
   | OS | 15.6 | 93.6% |
   | SearchPhrase | 11.7 | 91.5% |
   | RegionID | 11.0 | 90.9% |
   | UserID | 10.1 | 90.1% |
   | URL | 1.4 | 30.2% |
   
   Even the high-cardinality `UserID` column shows a 90% hit rate because the 
file is clustered.
   
   ## Describe the solution you'd like
   
   A **consecutive-keys cache** (last-group slot): remember the previous row's 
`(key, group_index)` in the `GroupValues` struct; if the current row's key 
equals it, reuse the group index and skip the hash + probe entirely.
   
   The trade is extremely asymmetric:
   
   - hit: skip a random probe (~15–100 ns when the table doesn't fit in cache)
   - miss: one well-predicted register compare (< 0.5 ns)
   - break-even hit rate: well under 1%
   
   This is the same optimization ClickHouse ships as its "consecutive keys 
optimization" (`LastElementCache` in `ColumnsHashing`), later improved in 
[ClickHouse#57872](https://github.com/ClickHouse/ClickHouse/pull/57872) with an 
adaptive disable when the miss rate is high and a batch path for blocks with a 
single distinct key.
   
   ### Prototype results
   
   A ~30-line prototype in `GroupValuesPrimitive` (compare the value itself — 
for primitives this also skips the hash), with cache invalidation on `emit` / 
`clear_shrink` (both reassign group indices):
   
   | query | GROUP BY | baseline (3×5 iters) | patched | delta |
   |---|---|---|---|---|
   | Q15 | UserID (~17M groups) | 305.6 / 277.5 / 274.4 ms | 266.6 / 258.5 / 
253.2 ms | **−8~9% (stable)** |
   | Q7 | AdvEngineID (~20 groups) | 16–18 ms | 15–15.5 ms | ~−5% |
   | Q27 | CounterID (~6K groups) | ~686 ms | ~687 ms | neutral |
   
   The pattern confirms the mechanism: the win concentrates where the hash 
table is **large** (probes miss cache — Q15); for low-cardinality keys the 
table is already L1-resident, so there is little to save (Q27), even at ~100% 
hit rate.
   
   ### Proposed scope
   
   - [ ] `GroupValuesPrimitive` — prototype done, needs tests for the 
emit/`clear_shrink` invalidation
   - [ ] `GroupValuesBytes` (`Utf8`/`LargeUtf8`/`Binary` via `ArrowBytesMap`) — 
the natural form here is comparing the current row against the *previous input 
row* (adjacent in the values buffer, cache-hot); as a bonus, `ArrowBytesMap` is 
shared with distinct-aggregate accumulators, which would also benefit
   - [ ] `GroupValuesBytesView` (`Utf8View` via `ArrowBytesViewMap`) — same 
shape
   - [ ] `GroupValuesColumn` (multi-column) — run detection needs all columns 
to match (`row[i] == row[i-1]` per column, AND-ed); can be vectorized per 
column. Covers ClickBench Q16/Q17/Q18 (`GROUP BY UserID, SearchPhrase` — 
run-correlated columns)
   - [ ] Adaptive disable on high observed miss rate (protects the 
Final/post-shuffle phase, where hash repartitioning destroys runs), following 
ClickHouse#57872
   
   ## Describe alternatives you've considered
   
   - **Perfect-hash aggregation** (DuckDB) for small dense integer key domains 
— orthogonal; it removes the probe for *low*-cardinality keys, whereas this 
optimization targets *high*-cardinality keys with clustered data.
   - **Streaming/ordered aggregation** (`GroupOrdering::Full/Partial`) — 
requires the plan to prove sortedness; the last-group cache is the 
opportunistic version that exploits incidental clustering with no plan-level 
requirements.
   - Status quo — leaves an ~asymmetric-odds win on the table for log-shaped 
data.
   
   ## Additional context
   
   Run-length measurement and prototype benchmarks above were done on 
`hits.parquet` (99.9M rows). The Partial aggregation phase sees file-ordered 
data, so run structure is intact where it matters most (the phase that 
processes all raw rows); the Final phase sees hash-shuffled data, which is what 
the adaptive disable is for.
   


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