apache / apache/datafusion

[EPIC] Eliminate Long Polls in HashAggregate via Chunked Storage and Incremental Emission

Aperta
#19,906 5 commenti 0 reazioni 0 assegnatari Vedi su GitHub
enhancement PROPOSAL EPIC
Lingua principale
Rust
Stelle
9.3k
Fork
2.4k
Merge medio
3g 11h
PR unite (30g)
360

Descrizione

### Is your feature request related to a problem or challenge?

When `GroupedHashAggregateStream` finishes processing input and transitions to emitting accumulated groups, it calls `emit(EmitTo::All)` which materializes all groups at once. For high-cardinality aggregations (>500k groups) or complex grouping keys (Strings, Lists), this becomes a CPU-intensive blocking operation that stalls the async runtime for hundreds of milliseconds to seconds.

This "long poll" prevents other tasks on the same thread from running, causing latency spikes and system "hiccups." We've observed `AggregateExec` stalling the runtime for >1s when processing queries with ~10M groups.

**Prior attempts to fix this**
In #18906, we introduced a `DrainingGroups` state to emit groups incrementally. In #19562, we implemented `EmitTo::Next(n)` for true incremental emission at the `GroupColumn` level. This works correctly but revealed a ~15x performance regression on high-cardinality queries.

**Root cause**
The current `GroupColumn` implementations store values in contiguous `Vec`. When emitting first n elements via `take_n()`, all remaining elements must be shifted:
```rust
fn take_n(&mut self, n: usize) -> ArrayRef {
let first_n = self.group_values.drain(0..n).collect(); // O(remaining)!
}
```

Profiling showed costs distributed across `Vec::from_iter` (allocations), `MaybeNullBufferBuilder::take_n` (copying), and `_platform_memmove` (shifting).

**Conclusion**: Incremental emission is the right approach, but requires chunked storage to be efficient.

### Describe the solution you'd like

This epic tracks two complementary changes:

**1. Chunked Storage for GroupColumn**
Replace contiguous `Vec` with a chunked structure so `take_n()` is O(n) instead of O(remaining):
```
Current: Vec [v0, v1, v2, ...vN] → take_n(3) shifts all remaining = O(remaining)
Chunked: [Chunk0] → [Chunk1] → [Chunk2] → take_n(3) advances head = O(n)
```
**2. Incremental Emission in `HashAggregate`**
Update `GroupedHashAggregateStream` to call `emit(EmitTo::Next(batch_size))` iteratively instead of `emit(EmitTo::All)`, yielding between emissions.

**Implementation phases:**

**1. Infrastructure**
- [ ] Add `ChunkedVec` data structure with `push`, `len`, `get`, `take_n`, `iter`, and `size` methods
- [ ] Add `ChunkedNullBufferBuilder` for chunked null bitmap storage with `take_n` support

**2. `GroupColumn` Migrations**
- [ ] Migrate `PrimitiveGroupValueBuilder` to use `ChunkedVec` for group values storage
- [ ] Migrate `BooleanGroupValueBuilder` to use chunked storage
- [ ] Migrate `ByteGroupValueBuilder` to use chunked offsets and buffer
- [ ] Migrate `ByteViewGroupValueBuilder` to use chunked views storage

**3. Emission Path**
- [ ] Optimize `emit()` index adjustment for `EmitTo::First(n)` with chunked storage
- [ ] Update `GroupedHashAggregateStream` to use `emit(EmitTo::Next(batch_size))` iteratively in `ProducingOutput` state

### Describe alternatives you've considered

Increase `target_partitions` which reduces groups per partition and could mitigate the issue. Can work but doesn't solve the fundamental problem and has other trade-offs (e.g. working on smaller partitions can slow down other operators).

### Additional context

Related issues & experiments
- #18907
- #18906
- #19562

Guida per i contributori

Apri la guida per i contributori

Direzione di ricerca

Inizia tracciando GroupedHashAggregateStream, il suo stato ProducingOutput e i builder di GroupColumn descritti nell’issue; esamina prima il lavoro correlato in #18906 e #19562. Il lavoro è completato quando l’infrastruttura indicata per chunked storage e null-buffer è presente, tutti i builder di GroupColumn nominati sono migrati e l’emissione usa ripetutamente EmitTo::Next(batch_size) senza poll bloccanti prolungati.

Scritto dal modello di indicizzazione a partire dal testo della issue.

Valutazione

Stack tecnologico
rust
Ambito
data-engineering, performance
Tipo di issue
Funzionalità
Difficoltà
5/5
Tempo stimato
Più di una settimana
Stato di attività
Attiva
Chiarezza
Abbastanza chiara
Idoneità per principianti
35/100

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.