Repository navigation
Implement cache-efficient partial aggregation by Leis et al #20773
Description
Activity
- addedenhancementNew feature or requestNew feature or requestperformanceMake DataFusion fasterMake DataFusion faster
on Mar 7, 2026 Hi @Dandandan, if you're not currently working on this, I'd be happy to take it on.
As I see, these two articles actually talk about how
duckdbimprove performance in high cardinality groups aggregation.But
datafusionandduckdbare too different in aggregation today,
and methods induckdbseems to help few about continuing to improvehigh cardinality groups aggregationof datafusion?At first let's see the difference in
high cardinality groups aggregation.Duckdb 0.7
https://duckdb.org/2022/03/07/aggregate-hashtable#parallel-aggregation
This article describes the total method to improve performance of high cardinality groups aggregation in duckdb 0.7.And I see, there are some mistakes in it, it use a very different way with what mentioned in Leis et al
-
In low cardinality case (the hashtable len in partial aggr <= 10K)
- use single hashtable in
partial aggr - use single thread to perform
final aggr(merge hashtables frompartial aggrs)
- use single hashtable in
-
In high cardinality case (the hashtable len in partial aggr > 10K)
- use partitioned hashtable in
partial aggr - use multiple threads to perform parallel
final aggr(likefinal aggrin datafusion)
- use partitioned hashtable in
Actually, it give up cache-efficient of hashtable in
partial aggr, and use the parallelfinal aggrto improve performance.Duckdb now
https://db.in.tum.de/~leis/papers/morsels.pdf
Duckdb use the similar method in the paper currently.- In
partial aggr, it always keep a very small and fixed hashtable - and when hashtable in
partial aggrbecome too big, it convert the data to partitions, and reset the hashtable. - In
final aggr, like datafusion, always partitioning and merge in parallel
We have tried the similar approach in datafusion before, see #6937 (comment) , but found no obvious improvement.
I guess, it is due to the execution model of datafusion (pull-based and async, hard to be cache-efficient).Datafusion
datafusion use
skip partial aggregationsto improvehigh cardinality groups aggregationlikedoris.- In
partial aggr, we will perform the groups aggregation logic when not exceed the threshold at first - Then we will totally skip the
partial aggrif excceed the threshold.
So, I guess few obvious improvement can got about optimizing
partial aggrindatafusion, because we will totally skip it inhigh cardinality groups aggregation.-
As I see, one of the most obvious bottleneck of aggregation in datafuison is
RepartitionExec.Introduce partitioned hashtable in
partial aggrmaybe can help, I did an experiment about it before:
#12526
The performance improvement got fromremoving RepartitionExec.But taking consider with
skip partial aggregations, I think it not the good and general way to solve the performance problem ofRepartitionExec.I think #15383 may be the better approach about it, and I think it really valuable to push forward.
We have tried the similar approach in datafusion before, see #6937 (comment) , but found no obvious improvement.
Note that this is significantly different:
The algorithm used by DuckDB/in the paper:
- inserts into the hashmap partition by partition / hashmap by hashmap (I think this is the most important point) - this will make the partial aggregation much more cache efficient, even for aggregations that are not that big as the hashmap that is "in progress" will more likely fit in cache / lookups are likely to be in cache
- Avoids the extra partitioning step (both the double hashing as the copying)
The core difference: The spilling in this paper doesn't just spill to the final partition, it partitions to a number of local tables and will insert to them one-by-one (for all indices that match the partition in the current batch).
The problem with spilling to repartition / final aggregation (#6937 (comment)) is that it will reduce keys less, and thus cause much more cross communication / repartition overhead and cost in the upstream aggregation.
And I see, there are some mistakes in it, it use a very different way with what mentioned in Leis et al
The description in DuckDB blog sounds to me very similar/the same as the paper, "spills when ht becomes full" means there is some threshold on the hashtable, after which it will repartition/spill to the local hash tables.
Hi @Dandandan, if you're not currently working on this, I'd be happy to take it on.
Feel feel to have a go at it, though it will not be easy as it touches on multiple parts within the aggregation code and planning (and making sure not to introduce any slow parts).
Reacted by Abhishek@Dandandan , thanks for the heads up! I’ll start digging into the aggregation and planning parts and will be careful to avoid introducing any performance regressions.
take
My concerns are that:
- The paper describe about how to keep the hashtable small and cache-efficient in high cardinality case in
partial aggr - And situation seems in datafuson:
- In
high cardinalitycase,datafusionuseskip aggregations approachto even totally skip thepartial aggr, and skipping possible to be a better approach(dorisuse the similar approach). - In
low cardinalitycase, the hashtable usually small enough without extra processing? And the bottleneck switch to repartitioning which can be solved with a more general and easier way (also help high cardinality).
- In
The description in DuckDB blog sounds to me very similar/the same as the paper, "spills when ht becomes full" means there is some threshold on the hashtable, after which it will repartition/spill to the local hash tables.
Approach in Leis et al may be more similar as current duckdb(the approach is mentioned together with their external aggregtion in https://duckdb.org/2024/03/29/external-aggregation), it is different with duckdb 0.7(method in the blog ):
- Only one small enough fixed hashtable in
partial aggrthread, and only partition the group values - When hashtable full, flush the partitioned group values
- Clear and reuse the single hashtable
- And in
final aggr, whatcurrent duckdbdo is similar asdatafusion, always perform partition-wise parallel merge (duckdb 0.7will perform single thread merge when groups are few)
As I see,
making partial aggr cache-friendlymentioned in paper(and current duckdb) is actuallymaking the hashtable in partial aggr cache-friendly.
And the paper and current duckdb use reset + rebuild method to reach it even inhigh cardinalitycase.Duckdb 0.7 (method in the blog) seems to just use partitioned hashtable to support parallel merge in
final aggr. And actually partitioned hashtable help few withcache-friendlydue to some experiment I did before.The problem with spilling to repartition / final aggregation (#6937 (comment)) is that it will reduce keys less, and thus cause much more cross communication / repartition overhead and cost in the upstream aggregation.
Agree with "flush" is too expansive here, that possible lead to no obvious improvement.
It may be not the good case to measure if the method in Leis et al can also work well indatafusion.Reacted by Daniël Heres and Abhishek- The paper describe about how to keep the hashtable small and cache-efficient in high cardinality case in
Thanks for your thoughts @Rachelint and sharing the differences in the DuckDB blogs / versions.
I agree we should probably take the ideas from the paper and see if (and perhaps mainly how) we can apply it here.
Also there are cases where DF is currently already faster than DuckDB (even without morsels and filter pushdown) - so we're probably currently do better in certain areas.Probably is gonna take quite some experimentation until we can find out what works.
Reacted by kamilleI wonder if we can make things also more cache aware with some sorting / partitioning based on hash as well with more minimal changes with higher cardinality, perhaps getting a similar win in cache efficiency for high cardinality (without the large changes):
- sort/partition hashes so we make the access to the table more linear (and also detect/skip duplicate keys avoiding probes
- sort/partition group indices to make accumulate state access more linear (rather than scattered over entire group state)
Reacted by kamilleI wonder if we can make things also more cache aware with some sorting / partitioning based on hash as well with more minimal changes with higher cardinality, perhaps getting a similar win in cache efficiency for high cardinality (without the large changes):
- sort/partition hashes so we make the access to the table more linear (and also detect/skip duplicate keys avoiding probes
- sort/partition group indices to make accumulate state access more linear (rather than scattered over entire group state)
Yes, after reading duckdb codes, my current thoughts:
- For
partial aggr, skipping is already a good method to handle high cardinality groups inpartial aggr, what we should contintue to do forpartial aggrI think is improving performance ofRepartExec - For
final aggr, we can use the partition-wise method to keep the hashmap always small like what is done induckdb. And luckily, I found the code changes to reach it is not actually large, I am doing experiment about this.
Reacted by Daniël HeresSome idea that came up today during discussing morsel-driven changes (#20481) is that we can as well sort on (file / row group / page) statistics (i.e. cluster min/maxes together for each worker/partition). This should make the the cardinality of the partial aggregate smaller (and thus much faster / less repartitioning going on).
Not an algorithmic change like this paper, but could potentially help a lot of queries both performance and memory-wise.I prototyped the second half of this idea (draft PR #25567) — a Final aggregation that, once its table outgrows the
cache, moves its state and all further input into 64 hash buckets and aggregates them one at a
time — and measured where the time goes. Two findings seem relevant to the design here.1. Re-partitioning a table that was already built is what loses. Emitting the first table,
scattering it and interning its groups again costs about as much as building it did. With a
threshold ofTgroups, every query that ends at 1-4xTgroups per partition pays that without
earning it back, and raisingTonly moves the loss to other queries:threshold slower than main 256k Q13 1.18x, Q14 1.11x 512k Q13 1.60x, Q14 1.26x 1M Q33 1.23x, Q31 1.15x, Q35 1.12x 2M Q16 1.09x while the large aggregations are best with an early switch (Q18 0.39x at 256k, 0.51x at 2M).
DuckDB's ICDE 2024 paper makes the same point against the HyPer / BLU design: "Rather than
partitioning tuples when the hash table is reset ... our implementation directly materializes
tuples into partitions ... we avoid copying tuples more than once", and a full table only has its
pointer array reset "while the tuples stay in place". I would aim for that shape — rows go to their
partition once, from the first batch — rather than "one table, radix-repartition it when it grows".2. The second scatter is 35-55% of the bucketed Final's compute. Per Final input row:
query one table bucketed = routing + aggregation Q31 (2 int keys) 85 ns 44 = 24.5 + 19.5 Q18 (int + string keys) 229 71 = 27 + 44 Q14 (string key) 79-85 91-99 = 42-46 + 49-54 Q12 (string key) 91 97 = 39 + 58 RepartitionExechas already hashed, gathered and copied each of these rows once; the Final does
it again. If the exchange producedpartitions x Ksub-buckets from the same hash, the routing
column would disappear: the mid-size string-key queries would turn from 1.03-1.07x into roughly
0.92x, and the queries that already win would gain about as much again. A standalone harness
(no Arrow; one copy into 64 buckets, hash kept with the row) confirms that partitioning pays for
string keys from 250k groups per thread: 24-byte keys 0.63x / 0.56x / 0.48x / 0.34x of a single
table at 250k / 500k / 1M / 4M groups, 12 threads. Re-hashing in the bucket instead of carrying the
hash costs little for short keys and a lot for long ones (60 bytes at 1M groups: 0.54x -> 0.78x),
which ties in with #11680.Numbers are from an Apple M4 Pro, 12 partitions, ClickBench
hits_partitioned, interleaved runs. I am prototyping the "repartition produces the buckets" variant next and will report what it measures here.My concerns are that:
- The paper describe about how to keep the hashtable small and cache-efficient in high cardinality case in
partial aggr - And situation seems in datafuson: ....
It is my mistake, after further researching, I found cache efficient partial aggr actually very promising (#26116) .
Some observations from my experiments:
- we should just keep the fixed and very small work set in partial aggr stage(like just 8192 groups) , so they can usually be in cache
- reserve all things, work set is fixed and so we can predict the capacity of the hashmap, groups vector and other things
- and finally we get the surprising performance improvement
- The paper describe about how to keep the hashtable small and cache-efficient in high cardinality case in
Is your feature request related to a problem or challenge?
The paper "Morsel-Driven Parallelism"[1] by Leis et al introduces a cache efficient aggregation algorithm.
It is also implemented by DuckDB[2]
We can implement it and see if it improves performance.
[1] https://db.in.tum.de/~leis/papers/morsels.pdf
[2] https://duckdb.org/2022/03/07/aggregate-hashtable#parallel-aggregation
Describe the solution you'd like
The steps are as follows
target_partitionsand usehash % target_partitionsinstead but in future might make it more adaptive)The benefit of it I think is mostly that by local grouping, we process the rows (smaller) hashmaps partition-by-partition, so while doing it they more likely fit in CPU cache.
It also avoids the double hashing in partial aggregation / hash repartition as the latter can be removed.
Care must be taken to make the code efficient, i.e. avoid materializing the batches upfront / accumulate based on indices.
I am not sure why the final aggregation doesn't take the same approach - this could be tested out as well when it is implemented for the partial aggregation.
Describe alternatives you've considered
No response
Additional context
No response