Skip to content

Implement cache-efficient partial aggregation by Leis et al #20773

Description

@Dandandan

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

  1. Start with a single hashtable (current approach) without pre-repartitioning
  2. Once the table exceeds a threshold (i.e. roughly CPU cache size), (radix) repartition the maps into a number of thread-local hashmaps (e.g. we can start of using the number of target_partitions and use hash % target_partitions instead but in future might make it more adaptive)
  3. Once they are repartitioned / finalized the output batches can be sent directly to the target partitions (no RepartitionExec needed! anymore)
Image

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

Activity

  1. added theissue type on Mar 7, 2026
  2. Abhisheklearn12 commented on Mar 7, 2026

    @Abhisheklearn12
    Contributor

    Hi @Dandandan, if you're not currently working on this, I'd be happy to take it on.

  3. Rachelint commented on Mar 7, 2026

    @Rachelint
    Contributor

    As I see, these two articles actually talk about how duckdb improve performance in high cardinality groups aggregation.

    But datafusion and duckdb are too different in aggregation today,
    and methods in duckdb seems to help few about continuing to improve high cardinality groups aggregation of 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 from partial aggrs)
    • 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(like final aggr in datafusion)

    Actually, it give up cache-efficient of hashtable in partial aggr, and use the parallel final aggr to 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 aggr become 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 aggregations to improve high cardinality groups aggregation like doris.

    • In partial aggr, we will perform the groups aggregation logic when not exceed the threshold at first
    • Then we will totally skip the partial aggr if excceed the threshold.

    So, I guess few obvious improvement can got about optimizing partial aggr in datafusion, because we will totally skip it in high cardinality groups aggregation.

  4. Rachelint commented on Mar 7, 2026

    @Rachelint
    Contributor

    As I see, one of the most obvious bottleneck of aggregation in datafuison is RepartitionExec.

    Introduce partitioned hashtable in partial aggr maybe can help, I did an experiment about it before:
    #12526
    The performance improvement got from removing RepartitionExec.

    But taking consider with skip partial aggregations, I think it not the good and general way to solve the performance problem of RepartitionExec.

    I think #15383 may be the better approach about it, and I think it really valuable to push forward.

  5. Dandandan commented on Mar 7, 2026

    @Dandandan
    ContributorAuthor

    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)
  6. Dandandan commented on Mar 7, 2026

    @Dandandan
    ContributorAuthor

    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.

  7. Dandandan commented on Mar 7, 2026

    @Dandandan
    ContributorAuthor

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

  8. Abhisheklearn12 commented on Mar 7, 2026

    @Abhisheklearn12
    Contributor

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

  9. Abhisheklearn12 commented on Mar 7, 2026

    @Abhisheklearn12
    Contributor

    take

  10. Rachelint commented on Mar 7, 2026

    @Rachelint
    Contributor

    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 cardinality case, datafusion use skip aggregations approach to even totally skip the partial aggr, and skipping possible to be a better approach(doris use the similar approach).
      • In low cardinality case, 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).

    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 aggr thread, and only partition the group values
    • When hashtable full, flush the partitioned group values
    • Clear and reuse the single hashtable
    • And in final aggr, what current duckdb do is similar as datafusion, always perform partition-wise parallel merge (duckdb 0.7 will perform single thread merge when groups are few)

    As I see, making partial aggr cache-friendly mentioned in paper(and current duckdb) is actually making the hashtable in partial aggr cache-friendly.
    And the paper and current duckdb use reset + rebuild method to reach it even in high cardinality case.

    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-friendly due 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 in datafusion.

  11. Dandandan commented on Mar 7, 2026

    @Dandandan
    ContributorAuthor

    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.

  12. Dandandan commented on Mar 8, 2026

    @Dandandan
    ContributorAuthor

    I 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)
  13. Rachelint commented on Mar 9, 2026

    @Rachelint
    Contributor

    I 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 in partial aggr, what we should contintue to do for partial aggr I think is improving performance of RepartExec
    • For final aggr, we can use the partition-wise method to keep the hashmap always small like what is done in duckdb. And luckily, I found the code changes to reach it is not actually large, I am doing experiment about this.
  14. Dandandan commented on Mar 9, 2026

    @Dandandan
    ContributorAuthor

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

  15. jayzhan211 commented on Sep 21, 2026

    @jayzhan211
    Contributor

    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 of T groups, every query that ends at 1-4x T groups per partition pays that without
    earning it back, and raising T only 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

    RepartitionExec has already hashed, gathered and copied each of these rows once; the Final does
    it again. If the exchange produced partitions x K sub-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.

  16. Rachelint commented on Oct 7, 2026

    @Rachelint
    Contributor

    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
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Labels

enhancementNew feature or requestperformanceMake DataFusion faster

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions