[Epic] High cardinality aggregation performance wishlist
- Dominant language
- Rust
- Stars
- 9.3k
- Forks
- 2.4k
- Avg merge
- 3d 7h
- Merged PRs (30d)
- 344
Description
### Is your feature request related to a problem or challenge?
DataFusion uses a two phase approach to aggregation (see [`Accumulator::state`](https://docs.rs/datafusion/latest/datafusion/physical_plan/trait.Accumulator.html#tymethod.state)) for details:
```text
▲
│ evaluate() is called to
│ produce the final aggregate
│ value per group
│
┌─────────────────────────┐
│GroupBy │
│(AggregateMode::Final) │ state() is called for each
│ │ group and the resulting
└─────────────────────────┘ RecordBatches passed to the
▲
│
┌────────────────┴───────────────┐
│ │
│ │
┌─────────────────────────┐ ┌─────────────────────────┐
│ GroubyBy │ │ GroubyBy │
│(AggregateMode::Partial) │ │(AggregateMode::Partial) │
└─────────────────────────┘ └────────────▲────────────┘
▲ │
│ │ update_batch() is called for
│ │ each input RecordBatch
.─────────. .─────────.
,─' '─. ,─' '─.
; Input : ; Input :
: Partition 0 ; : Partition 1 ;
╲ ╱ ╲ ╱
'─. ,─' '─. ,─'
`───────' `───────'
```
For low cardinality aggregates (where there are a few distinct groups), this works great 👌 👨🍳
However for high cardinality aggregates (where there are many millions of groups), we can do better by optimizing the path. See the background and ASCII art on https://github.com/apache/datafusion/issues/7957 for why the intermediate cardinality increases
This is my wishlist for improving high cardinality aggregates (ideally for the next blog post in a few months #11631 )
Together with the StringView work in https://github.com/apache/datafusion/issues/10918 that @XiangpengHao @a10y and others are working on, I think it would provide some very compelling overall speedups in ClickBench and TPCH queries
Also I hear that @avantgardnerio may be interested in helping here
### Describe the solution you'd like
Here is my wishlist:
- [x] https://github.com/apache/datafusion/issues/6937 (@korowa has a PR up for this one)
- [ ] https://github.com/apache/datafusion/issues/7957 (I have a prototype and some ideas)
- [ ] https://github.com/apache/datafusion/issues/11680
### Describe alternatives you've considered
Do nothing and let DuckDB pass us by ;)
### Additional context
Other potential things to do:
- [ ] https://github.com/apache/datafusion/issues/7000
Contributor guide
Research direction
Start by reading the Accumulator::state documentation and the background in issues 7957, 11680, and 7000. Compare the proposed high-cardinality aggregation work and identify a specific wishlist item before choosing an implementation path. Done means the selected issue has a defined change and evidence that it improves the relevant aggregation performance.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- rust
- Domain
- data-engineering, performance
- Issue type
- Feature
- Difficulty
- 5/5
- Estimated time
- Over a week
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 25/100