Support zero copy hash repartitioning for Hash Aggregate
- 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?
### Is your feature request related to a problem or challenge?
Currently `RepartitionExec: partitioning=Hash` will be added whenever for aggregates in `FinalPartitioned` and `SinglePartitioned`
The benefit is increased parallelism, but at the cost of copying the entire table (in a not-so efficient way).
We should consider lowering the cost of repartitioning by not having to copy the input.
Dependencies
- [ ] https://github.com/apache/datafusion/issues/15420
### Describe the solution you'd like
Instead of repartitioning the input in `RepartitionExec`, support repartitioning the inputs based on a selection vector.
Instead of `taking` the `RecordBatch`, we can consider doing the following:
* Add a (boolean) selection vector as output column for each output partition. I.e. `true` means the row is selected for the partition.
* The rest of the `RecordBatch` remains unchanged (i.e. no copy).
* CoalesceBatchesExec is no longer needed for the output (reducing another copy)
* In the hash aggregate code handle the selection vector.
### Describe alternatives you've considered
The partitioning could be done inside the hash aggregate (at the cost of more complexity inside it).
### Additional context
_No response_
### Describe the solution you'd like
_No response_
### Describe alternatives you've considered
_No response_
### Additional context
_No response_
Contributor guide
Assessment
This issue has not been assessed yet.