Skip to content

[EPIC] Use blocked / chunked memory management in hash aggregation #24704

Description

@alamb

Note: this issue consolidates and replaces #7065, which dates from 2023 and predates much of the discussion.

What is going on (symptoms)

High-cardinality GROUP BY queries in DataFusion suffer from several challenges that look unrelated at first, but all trace back to the same root cause. The problems:

  1. Aggregation memory is held until the hash table is fully drained. All group state is emitted as slices of one giant batch, so none of it is released until the last output batch. For a typical two-stage aggregation, this shows up as ~2x peak memory: while the first stage drains its 1 GB of state, the final stage is simultaneously building its own ~1 GB (measured by @2010YOUY01 in PoC: Blocked state management for hash aggregation #22712).

  2. Queries with more than 2 GiB of overall string data in group keys can crash outright. All the bytes of Utf8/Binary group keys are interned into one contiguous buffer addressed by i32 offsets; once more than 2 GiB of key bytes accumulate, the offsets overflow and the query fails with offset overflow, buffer size > 2147483647:

  3. Downstream memory accounting is off / operators spill when they don't need to. Each output batch reports the memory of the entire aggregation output via get_array_memory_size(), which causes unnecessary spilling in operators such as RepartitionExec and TopK:

  4. The async runtime stalls when output begins. Producing output for >500k groups block a tokio worker thread for hundreds of milliseconds to seconds, causing latency spikes for everything else on that thread:

  5. Potential copying performance. As groups accumulate, internal buffers repeatedly double in size and copy all existing data (up to 2 copies per element on average), which is likely expensive and cache/TLB unfriendly and could in theory be avoided.

What is causing the problem

Note you can read more about the current group state in the blog Aggregating Millions of Groups Fast in Apache Arrow DataFusion 28.0.0

GroupedHashAggregateStream stores all per-group state in single contiguous buffers that grow by doubling:

  • Group keys are stored by a GroupValues implementation (e.g. PrimitiveGroupValueBuilder or ByteGroupValueBuilder). The hash table itself only stores group indexes — usize offsets into these buffers.
  • Aggregate state is stored by one GroupsAccumulator per aggregate expression (not per group). Each accumulator manages the state for all groups, typically as a single Vec<T> (plus a null buffer), again indexed by a single group index (a usize).
                                         ┌──────────────┐   ┌──────────────┐   ┌──────────────┐
                                         │┌────────────┐│   │┌────────────┐│   │┌────────────┐│
    ┌─────┐                              ││accumulator ││   ││accumulator ││   ││accumulator ││
    │  5  │                              ││     0      ││   ││     0      ││   ││     0      ││
    ├─────┤                              ││ ┌────────┐ ││   ││ ┌────────┐ ││   ││ ┌────────┐ ││
    │  9  │                              ││ │ state  │ ││   ││ │ state  │ ││   ││ │ state  │ ││
    ├─────┤                              ││ │        │ ││   ││ │        │ ││   ││ │        │ ││
    │     │                              ││ │        │ ││   ││ │        │ ││   ││ │        │ ││
    ├─────┤                              ││ │        │ ││   ││ │        │ ││   ││ │        │ ││
    │  1  │                              ││ │        │ ││   ││ │        │ ││   ││ │        │ ││
    ├─────┤                              ││ │        │ ││   ││ │        │ ││   ││ │        │ ││
    │     │                              ││ │        │ ││   ││ │        │ ││   ││ │        │ ││
    └─────┘                              ││ │        │ ││   ││ │        │ ││   ││ │        │ ││
                                         ││ │        │ ││   ││ │        │ ││   ││ │        │ ││
                                         ││ └────────┘ ││   ││ └────────┘ ││   ││ └────────┘ ││
                                         │└────────────┘│   │└────────────┘│   │└────────────┘│
    Hash Table                           └──────────────┘   └──────────────┘   └──────────────┘


stores "group indexes"                     There is one GroupsAccumulator per aggregate
which are indexes into                     (NOT PER GROUP). Internally, each
the state vectors                          GroupsAccumulator manages the state for
                                           multiple groups

This contiguous layout explains each symptom:

  1. Memory held until drained: output is produced via EmitTo::All: at end of input, all groups are materialized into one giant RecordBatch, which is then handed downstream as batch.slice(..) chunks of batch_size rows. The slices share the giant batch's buffers, so no memory is freed until the last slice is dropped.
  2. Crash on wide string keys: for Utf8/Binary group keys, ByteGroupValueBuilder<i32> stores all key bytes in a single contiguous buffer addressed by i32 offsets, so accumulating more than 2 GiB of key bytes overflows.
  3. Spilling / accounting: every slice of the giant batch reports the full underlying allocation to the memory accounting, not just its own rows.
  4. Runtime stall: materializing all groups in a single EmitTo::All call is one large CPU-bound operation with no await point.
  5. Copying performance: growing the contiguous buffers requires reallocating and copying all existing data.

The high level solution sketch

The fix that everyone seems to agree on is to store group keys and accumulator state in multiple blocks instead of one contiguous Vec. This is the approach used by DuckDB and most other databases with buffer managers, which don't have the luxury of large contiguous arrays. It was originally suggested for DataFusion by @yjshen in 2023.

At a high level, the idea is that:

  • Blocks store some number of rows (most likely target_batch_size)
  • Blocks are not resized; when a block fills up, a new block is allocated (which avoids copying and the growth is predictable and incremental)
  • A group index becomes some form of (block_id, offset) rather than a single group_idx
  • Emission now happens one block at a time (e.g. something like EmitTo::NextBlock), so memory is freed incrementally. The blocks now back individual output RecordBatches (fixing the accounting).
  • Since each block is capped at target_batch_size rows, they are far more likely to stay below 2 GiB of total string values (the i32 offset limit), avoiding the string offset overflow.

The idea is illustrated here:

                                         ┌──────────────┐   ┌──────────────┐
                                         │┌────────────┐│   │┌────────────┐│
    ┌─────────┐                          ││accumulator ││   ││accumulator ││
    │  (0,5)  │                          ││    AGG     ││   ││    SUM     ││
    ├─────────┤                          ││ ┌────────┐ ││   ││ ┌────────┐ ││
    │  (1,3)  │                          ││ │ block  │ ││   ││ │ block  │ ││
    ├─────────┤                          ││ │   0    │ ││   ││ │   0    │ ││
    │         │                          ││ │        │ ││   ││ │        │ ││
    ├─────────┤                          ││ │        │ ││   ││ │        │ ││
    │  (0,1)  │                          ││ │        │ ││   ││ │        │ ││
    ├─────────┤                          ││ └────────┘ ││   ││ └────────┘ ││
    │         │                          ││            ││   ││            ││
    └─────────┘                          ││ ┌────────┐ ││   ││ ┌────────┐ ││
                                         ││ │ block  │ ││   ││ │ block  │ ││
    Hash Table                           ││ │   1    │ ││   ││ │   1    │ ││
                                         ││ │        │ ││   ││ │        │ ││
                                         ││ │        │ ││   ││ │        │ ││
                                         ││ │        │ ││   ││ │        │ ││
                                         ││ └────────┘ ││   ││ └────────┘ ││
                                         │└────────────┘│   │└────────────┘│
                                         └──────────────┘   └──────────────┘


stores "group indexes"                     Each accumulator stores its state in
as (block_id, offset)                      fixed size blocks: a full block is
pairs into the block                       never resized; instead a new block
storage                                    is allocated as needed, and whole
                                           blocks can be emitted / freed
                                           one at a time

Why this is hard to fix

The reason this is so hard to implement is that the entire API is designed around a single usize group index that is assumed to be contiguous and directly addressable:

For example GroupValues::intern:

pub trait GroupValues: Send {
    // Required methods
    fn intern(
        &mut self,
        cols: &[Arc<dyn Array>],
        groups: &mut Vec<usize>,   // <---- groups are identified by contiguous `usize`
    ) -> Result<(), DataFusionError>;

    ...
}

Because this assumption is spread across every GroupValues and GroupsAccumulator implementation — including user-defined aggregates and the FFI bindings — moving to a blocked (block_id, offset) index:

  1. Potentially touches the whole ecosystem at once (aka is a massive change)
  2. Likely involves an extra memory lookup in the hottest critical path: once to find the base pointer for the block_id and once to find the actual value within that block.

We have discussed ways to address both issues:

  1. Incremental rollout: neither option found so far is great — supporting blocked and contiguous layouts in each implementation leads to the dual code paths and generics that made Intermediate result blocked approach to aggregation memory management #15591 so complex, while switching the index semantics outright is a breaking change that is very hard to stage incrementally (see the discussion on Intermediate result blocked approach to aggregation memory management #15591).
  2. Different strategies for small and large aggregates: use direct indexing while the hash table is small (where the extra indirection would hurt most), and switch to two part (block_id, offset) indexes once the table grows past a threshold — at that point accesses are cache misses anyway, so the extra lookup matters less. See Intermediate result blocked approach to aggregation memory management #15591 (comment) and the earlier version of the same idea in PoC: Blocked state management for hash aggregation #22712 (comment).

Past attempts and prototypes

There is a long and distinguished history of trying to address this problem:

PR Year What it showed Outcome
#11758 2024 Generate GroupByHash output in multiple RecordBatches (@JasonLi-cn) Closed unmerged
#11943 2024 First sketch of blocked management (@Rachelint) Closed, superseded by #15591; motivated the aggregation fuzz test framework (#12114)
#15591 2025 Full blocked implementation (@Rachelint): supports_blocked_groups / alter_block_size trait additions, blocked PrimitiveGroupsAccumulator + GroupValuesPrimitive Open. Extensive review concluded the dual-mode (blocked + contiguous) design is too complex, and some aggregates regress ~10%
#20964 2026 BatchedVec<T> bench (@Dandandan): O(1) per-block emission is achievable in small steps Closed (proof of concept)
#22712 2026 PoC on refactored streams (@2010YOUY01): 10–16% faster at medium/high cardinality; memory curve becomes bell-shaped instead of monotonically growing; ClickBench Q5 +61% pending the skip-partial-aggregation fast path Closed (proof of concept; demonstrated the #22710 refactor is necessary first)
#23274 2026 EmitTo::FirstBlock as an API-only first step (@hhhizzz) Closed: the API should follow the blocked physical layout rather than precede it

One seemingly obvious alternative would be to use EmitTo::First(n) to incrementally emit data from the front of the state. However, this does not work either as it is destructive: it requires shifting all remaining elements to the start of the buffers and renumbering every remaining group index. @ahmed-mez tried incremental emission with these mechanics in #19562 and measured it ~15x slower at high cardinality.

As part of #22710, @2010YOUY01 has been refactoring the monolithic GroupedHashAggregateStream into dedicated per-path streams, in large part to make changes like blocked state management feasible to implement and review. I (@alamb) thinks completing this refactor is a prerequisite for beginning the blocked state work in earnest: implementing it against the old multiplexed stream would be far more complex and would conflict with the refactoring itself.

The same blocked approach was also proposed independently by @alchemist51 in #19649, which includes an experiment (based on @Rachelint's #15591) where a high-cardinality query that fails with resource exhaustion in a 16 GB memory pool today completes once blocked state management is enabled.

An in-progress implementation for multi-column group-by (@rluvaton) using only EmitTo::NextBlock is being discussed on #15591. It is not yet a PR, but the work seems to be on the add-blocks-impl branch of their fork.

A complementary approach that reduces partial-stage state without changing these traits is cache-efficient (morsel-driven) partial aggregation, proposed by @Dandandan:

Related issues

Related symptoms that will NOT be addressed by this issue

Note that other operators produce giant contiguous intermediate batches too, and show the same failure modes. This issue only covers aggregation; something similar will be needed for joins:

Activity

  1. comphead commented on Aug 27, 2026

    @comphead
    Contributor

    Thanks for summarizing it

  2. jayzhan211 commented on Aug 27, 2026

    @jayzhan211
    Contributor

    I would like to help on this topic too 🚀

  3. Rachelint commented on Aug 29, 2026

    @Rachelint
    Contributor

    Have an new idea that maybe can support blocked approach in a simpler way
    #15591 (comment)

  4. jayzhan211 commented on Aug 30, 2026

    @jayzhan211
    Contributor

    Have an new idea that maybe can support blocked approach in a simpler way #15591 (comment)

    I like the idea of EmitTo::Last(n), especially for migration toward the end state (blocked storage) — although introducing EmitTo::Last(n) directly would add complexity for downstream users and the existing emission mechanism. So the proposal below keeps its semantics but doesn't touch EmitTo at first.

    Emission: the tail contract. The contract is emit the state of the n highest group indexes; indexes of all remaining groups are unchanged. Nothing shifts and nothing is renumbered — avoiding the ~15x cost #19562 measured for incremental First(n) — and both layouts can serve the same request: today's contiguous storage by truncation (split_off on the tail), blocked storage by popping its last block in O(1) with zero copy. That symmetry is what makes it a migration path. Instead of a new enum variant, it ships as default-implemented trait methods behind a runtime flag (default off):

    // GroupsAccumulator — all default-implemented; no existing impl breaks
    fn supports_tail_emit(&self) -> bool { false }
    fn evaluate_tail(&mut self, n: usize) -> Result<ArrayRef> { /* default: internal_err! */ }
    fn state_tail(&mut self, n: usize) -> Result<Vec<ArrayRef>> { /* default: internal_err! */ }
    
    // GroupValues
    fn supports_tail_emit(&self) -> bool { false }
    fn emit_tail(&mut self, n: usize) -> Result<Vec<ArrayRef>> { /* default: internal_err! */ }

    Zero downstream compile impact (same pattern as supports_convert_to_state), and the stream negotiates: tail drain only if all state owners in the plan support it, otherwise fall back to All — so unconverted UDAF/FFI accumulators keep working unchanged, and a plan mixing blocked built-ins with a merely-truncatable external accumulator still drains incrementally. The count n sits at the call site because the group values and every accumulator must release the same groups each step anyway (or output columns don't align). Scope: end-of-input drain on the unordered path only — First(n) and its renumbering contract stay untouched for ordered streams and the spill/merge read-back. Whether to later fold this into a real EmitTo::Last(n) variant (ideally together with marking EmitTo #[non_exhaustive]) can be a major-release decision after the contract has soaked.

    Storage. As sketched in this issue: fixed-capacity blocks of B = target_batch_size rows (power of two), never resized, allocated on demand, filling strictly front-to-back. Two refinements: block 0 is sized like today's initial capacity, so low-cardinality queries stay effectively flat (aimed at the ~10% regression observed in #15591); and byte/byte-view blocks each own their values buffer and their i32 offsets, so no offset ever spans more than one block — which closes #23694 structurally rather than by widening types.

    Indexing. The group index stays a single dense usize: with power-of-two blocks filling front-to-back, (block, offset) = (idx >> k, idx & (B - 1)) is a bijection. intern, update_batch, total_num_groups, and the hash table are all unchanged — no trait signature changes anywhere, including UDAFs and FFI. The known cost is one extra indirection per state access (the #15591 concern); hot-loop parity at low cardinality is the declared benchmark gate for the storage PRs, so this gets retired with numbers rather than argument.


    Phases toward the end state

    1. Benchmarks first — minimum only. For regressions, nothing new is built: ClickBench + the h2o groupby suite already in bench.sh are the gate, because they are what aggregation PRs are judged against today and h2o spans the dimensions this touches (int vs. string keys, low vs. high cardinality). For the wins, exactly one new artifact: a drain-phase benchmark that pre-builds ~10k and ~10M groups and records five numbers — total drain time, time-to-first-batch, max inter-batch gap (the long-poll proxy, #19906), peak reserved memory, and reserved memory at the 50%-drained mark (≈100% of peak under All, ≈50% under tail — a one-sample proof of incremental release). Gate numbers agreed on this issue before any behavioral PR.

    2. Tail contract. One vertical slice, targeting the refactored streams from #22710 (the declared prerequisite — happy to help review there): the default trait methods + the drain loop (drop the hash table at drain start, truncate-serve, shrink the reservation per step, one await per step) + tail arms for primitive group values, the primitive accumulators, and count — behind a config flag, with a fuzz mode asserting All ≡ tail modulo row order, and extending the aggregation fuzz framework to run under constrained memory pools (the #19649 scenario — a query that exhausts a fixed pool today and completes with incremental release — becomes a regression test). Everything else negotiates back to All. Flip the flag after soak; release notes call out that GROUP BY … LIMIT without ORDER BY may return different (equally valid) groups.

    3. Blocked internals. Blocks<T> + microbenches; blocked PrimitiveGroupsAccumulator and GroupValuesPrimitive (gate: hot-loop parity); blocked ByteGroupValueBuilder with per-block i32 offsets (closes #23694, with a gated >2 GiB regression test — and worth validating against Comet's suite, since #4718 came from there); byte-view; multi-column ideally by helping land @rluvaton's add-blocks-impl in pieces; then the remaining accumulator families. One requirement I missed earlier: the same implementations also serve the ordered streams and the spill read-back, which emit via First(n) — so each blocked conversion must implement First(n) (and All) over blocks too, with the renumbering semantics intact. so the per-implementation sequence is: blocked storage → First(n)/All over blocks → delete contiguous. The cross-block compaction First(n) needs can be written once in Blocks<T> (those paths keep their tables small, so it needn't be fast), leaving each conversion's First arm thin. Default is one atomic conversion PR per implementation; if one grows too big to review, the escape hatch is a temporary sibling type — old contiguous and new blocked coexisting as separate single-mode implementations selected by negotiation, with a tracked deletion — which is not the dual-mode-inside-one-implementation that sank #15591. Dual storage inside an implementation still never exists (the #15591 lesson). Contiguous survives only as the negotiated compatibility tier for external accumulators — a feature kept indefinitely, not a deprecation candidate — and a small permanent coverage gap remains for types that only the Rows-based fallback handles: those stay on All unless someone later blocks that path.

    4. Order-sensitive follow-ups. Ordered and partially-ordered streams keep First(n) and its renumbering contract untouched, as does the spill read-back (which re-aggregates a sorted merge). Mid-stream emission under memory pressure is a separate design — the hash table is live there, so tail eviction would cost a full table sweep — and it is where the interaction with skip-partial aggregation (the pending ClickBench Q5 +61% from #22712) lives. Blocked state also opens a follow-up the epic hints at: block-granular spilling (spill whole blocks under pressure instead of emit-sort-spill of everything), which is a step toward the memory pool proactively requesting release from consumers. If a hard front-first requirement surfaces, the same blocked layout supports a front cursor instead, at the documented cost of forking First's contract and the dense 0..total_num_groups invariant — recorded as the fallback design, not the default.

    5. Convert to EmitTo::Last(n) — the formalization release. At the next major, once the contract has soaked through phases 2–3: fold the trait methods into a real EmitTo::Last(n) variant, mark EmitTo #[non_exhaustive] in the same release so this class of breakage never recurs, deprecate the tail methods for one release cycle with a mechanical migration path, and ship the upgrade guide for external accumulator authors (how to move a tail arm from the methods to the match). The contract is never deprecated — only its provisional surface is folded into the enum. If soak instead shows the methods are fine as a permanent home, this phase collapses to documenting that decision; either way it is decided on schedule, not by drift.


    Open question: does anything require front-first drain? My answer is no for the unordered path — nothing forces emitting from the front there: SQL guarantees no output order without ORDER BY, today's insertion order is incidental, and the paths that genuinely need ordering (ordered/partially-ordered streams, the spill read-back) keep First(n) and its renumbering contract untouched. The end state is indifferent either way — over blocks, with the hash table dropped at drain start, either end pops in O(1) — so a front-first requirement would only threaten the phase-2 migration slice, where the tail is the only cheap direction on contiguous buffers. If anyone knows a constraint I'm missing — spill file ordering, partial-stage eviction, assumptions in add-blocks-impl

  5. jayzhan211 commented on Aug 30, 2026

    @jayzhan211
    Contributor

    #24795 sets the baseline for the follow-up PR, and the numbers show exactly what issue we're facing.

  6. rluvaton commented on Aug 30, 2026

    @rluvaton
    Member

    I don't like EmitTo::Last since sorted input is common, for example when using GroupAccumulators in Window which can be sorted by the partition keys.

  7. rluvaton commented on Aug 30, 2026

    @rluvaton
    Member

    @alamb Thank you for opening this issue!

    Also, couple of things to do right away, is marking the GroupValues and the implementations of it as private, the reason they are public from what I checked is to allow doing micro benchmarks.


    My branch add-blocks-impl currently goes the easy-ish way, what if we could do as many breaking changes as we want, and we only need to support blocked approach, the reason I took that path is that:

    1. I want to see what it would take
    2. Is the blocked approach easily implemented by existing accumulators
    3. To be able to run benchmark on the best case scenario
    4. EmitTo::All/EmitTo::First and emit blocks does not work together very well if the data structure is built in blocks:
      Keeping EmitTo::All means that if we emit the output as a single array it requires concatenation which I don't like
      Keeping EmitTo::First(n) without shifting the data between blocks means:
      1. When n is larger than block size, it will require concatenation which I don't like as it is expensive and open for issues like we have
      2. that blocks are now not fixed size so next emit block will have either partial or would have to concat 2 batches
      3. indexing become harder and force every implementors of acc to deal with that shifting
      Keeping EmitTo::First(n) with shifting data between blocks:
      1. for nested types, like List, the inner values blocks should be managed by the list itself (which is my current impl), so taking the first n values makes shifting the inner values really complex since the block size is dynamic for the inner values and you must move group of values as 1 to another prev block. so this is really really complex
    5. To surface any problems we might have (like the EmitTo::First I explain below) (which is also why I'm doing this manually and not using LLM)

    A problem I have with not having EmitTo::First is that think of this scenario.
    the block size is 1000, we can emit early (because the input is ordered for example)
    and we got all the values in the first 500 groups, now we got 2 more batches for group 501, and when we try to get another batch we don't have enough memory, if we don't have EmitTo::First, we can't emit the first 500 finished groups, we can emit the whole block which in our case is partial and group 501 is not finished, which I don't have solution for this yet.

    Actually, we can keep the EmitTo::First and it would just require shifting the remaining data in blocks rather than 1 big continues allocation, this is more complex to implement (and for nested types is really really complex)

    Side note: for variable length types (like string, binary, list, etc.) the shifting will require more working memory while shifting, but for fixed size types (if done right) this could be done with no extra memory

  8. rluvaton commented on Aug 30, 2026

    @rluvaton
    Member

    We can have EmitTo::All if we have a signature that return Vec in that case
    Emit all can still be valuable since accumulators might include more tracking data (like unique values for distinct) and emitting all will release those and we can move the blocks memory tracking to the caller

    EmitTo::First(n) if we have a signature that return Vec in the case than n is larger than block_size
    but I think we n must not be larger than block size and it should return Block so users would need to call emit next block with n as the remainder (altough allowing n to be larger than block size would allow the users to avoid cleanup after each block (for example if there is some data that is tracked globally and after every emit, that tracked is get cleaned, so to avoid cleaning for each block can emit first n and one cleanup would be performed)

  9. Dandandan commented on Aug 31, 2026

    @Dandandan
    Contributor

    #24815

    Drafted another approach. It seems performance neutral or positive overal (will try to run the benchmarks in the runner).

    Some notable items:

    • It has both a flat and and a blocked mode, to avoid regressing the low cardinality cases
    • Flat and blocked mode with branching hoisted outside of the loop
    • Size is a power of two, so indexing in blocked mode is cheap
  10. Rachelint commented on Aug 31, 2026

    @Rachelint
    Contributor

    @rluvaton I think we should process the sorted case and the normal case(unsorted) in different ways like what is doing currently:

    1. For the sorted case

    Maybe we can continue to use Emit::First(n).
    And if we decide to remove the Emit::First(n) finally (due to always unacceptable performance, and easy to be misused), I think we can make the GroupValues sorted input aware and use Emit::All to replace Emit::First(n).
    The logic can be:

    • Assume the group values keeping a, and it is aware of the input is sorting, and has the api to return current key
    • When following batch a, b, b, c comes, we get current key from GroupValues first
    • Compare the current key a with a,b,b,c, and found the batch will lead to key switching
    • Split the the batch to a and b,b,c,and only put a into GroupValues
    • Call Emit::All and return the freezed a group
    • Input the b,b,c to reset GroupValues, and go to next loop

    2. For the unsorted case

    After many experiment, I found blocked approach is only a better memory management approach to reduce the peak memory usage.
    And I think our main target just be:

    • Support really return aggr state block by block, rather than holding a large guy + slicing
    • Don't lead to unacceptable performance regression

    And I think the totally blocked approach #15591 may be not worthy:

    • Really large code changes
    • A massive impact on the existing architecture
    • And after many many hard works, still lead to unavoidable performance regression

    So after struggling with #15591 , I think something simple like Emit::Last may be a way more worthy trying.

  11. ariel-miculas commented on Aug 31, 2026

    @ariel-miculas
    Contributor

    After many experiment, I found blocked approach is only a better memory management approach to reduce the peak memory usage.

    It's also relevant for excessive spilling in the downstream operator, as described here: #22526 (comment)

    Even more relevant for the partitioning in datafusion-comet because threre's a shuffle after the partial aggregation that will always hit this spilling issue.

  12. 17 remaining items

  13. rluvaton commented on Sep 16, 2026

    @rluvaton
    Member

    @Dandandan tried this in #24815. I haven't reviewed it fully, but the change looks quite small — I think it's worth trying to find a simpler implementation here.

    any implementation of blocked index that require having fixed block size must try to implement the multi group by for list, since the child GroupValues block size is dynamic. so this should be checked if the impl support that

  14. jayzhan211 commented on Sep 16, 2026

    @jayzhan211
    Contributor

    If multi group by for nested types is the concern, how about we split them into separate paths? We could use the blocked approach for non-nested types and work out a specialized design for nested types.

  15. rluvaton commented on Sep 22, 2026

    @rluvaton
    Member

    It is not a concern just saying that having blocks index that is flat (i.e. if the data was single continues memory it would just be an index rather than a BlocksIndex that contain block index and index in block in it) won't work for the nested types.

  16. rluvaton commented on Sep 22, 2026

    @rluvaton
    Member

    @alamb @Dandandan @comphead @jayzhan211 I want to push this forward, the one question remaining is:

    Do we want to change the existing trait for GroupsAccumulator or create a new one for the BlockedGroupsAccumulator with an adapter for GroupsAccumulator so existing implementation would still work and it can be incrementally migrated.

    From my benchmarks doing continues state/evaluate with emit to first which shifts is really expensive (from what I remember can be 4 times slower).

  17. jayzhan211 commented on Sep 22, 2026

    @jayzhan211
    Contributor

    Do we want to change the existing trait for GroupsAccumulator or create a new one for the BlockedGroupsAccumulator with an adapter for GroupsAccumulator so existing implementation would still work and it can be incrementally migrated

    No strong preference on my side — either approach works for me as long as we get the result. I'm also fine if the first version is harder to review; we can look at simplifying it afterwards.

  18. alamb commented on Sep 22, 2026

    @alamb
    ContributorAuthor

    No strong preference on my side — either approach works for me as long as we get the result. I'm also fine if the first version is harder to review; we can look at simplifying it afterwards.

    I agree too. I think for this work I would suggest:

    1. Get a PR up with the base pattern for one or two accumulators, so we can validate the design/pattern carefully
    2. then we can mechanically switch over the other accumulators to the new pattern.
  19. comphead commented on Sep 23, 2026

    @comphead
    Contributor

    cc @sunchao who also experimenting with hybrid hash join.

    Agree with @alamb lets move trait discussion to specific PR, from more wide angle I would consider a way where we can upgrade or rollback without affecting users, exposing them to the painful migration and/or keeping parts of the code used not effectively.

    I'm more leaning towards second trait and once its stable to merge it to GroupsAccumulator but I might be missing some details

  20. rluvaton commented on Sep 24, 2026

    @rluvaton
    Member

    Here is the implementation:

    It contain BlockedGroupsAccumulator trait, support count in blocked so you see the example usage, change the entire aggregate to work with blocked (this is possible due to the added adapters), implement BlockedGroupValues for primitive so you will see how it is being used

  21. rluvaton commented on Sep 29, 2026

    @rluvaton
    Member

    So instead of having an adapter and having large performance hit until all is supported I created another PR with alternative approach which support both and have less performance hit for unsupported cases

    Please review so I can move this forward

  22. Rachelint commented on Oct 7, 2026

    @Rachelint
    Contributor

    I am thinking about how the block approach might fit with the recent aggregation optimizations we’ve been exploring, and wanted to share a couple of observations:

    If we pursue both optimizations, I wonder would the block approach still offer enough additional benefits.

  23. alamb commented on Oct 7, 2026

    @alamb
    ContributorAuthor

    If we pursue both optimizations, I wonder would the block approach still offer enough additional benefits.

    I do think it could lower the peak memory requirements for queries with high cardinality aggregates -- for example plans with subqueries that have large numbers of aggregates that are then fed into their own 🤔

    But the ideas from @jayzhan211 look quite promishing -- maybe the extra benefit wouldn't be worth the effort

    What do you think @rluvaton ?

  24. jayzhan211 commented on Oct 9, 2026

    @jayzhan211
    Contributor

    I think the Utf8 case is still open after #26116 (Partial) and #25724 (Final). Maybe we could look into whether blocked storage, or an alternative approach, would help with the string case.

  25. Rachelint commented on Oct 9, 2026

    @Rachelint
    Contributor

    I do think it could lower the peak memory requirements for queries with high cardinality aggregates...

    Seems in final when bucketing is on, the memory peak will just be 1.015 * state size rather than original 2 * state size; and when bucketing is off, the memory usage is always low? Maybe it is already good enough for reducing memory usage?

    I think the Utf8 case is still open after #26116 (Partial) and #25724 (Final)

    Yes, for utf8 or more complex type, bucketing seems still possible lead to slight performance regression, I think we should continue to optimize.
    But if take its effect of memory usage reducing into consider, the slight regression seems acceptable (compared with blocked approach) ? And it is actually much simpler than blocked approach.

  26. alamb commented on Oct 9, 2026

    @alamb
    ContributorAuthor

    19. I think the Utf8 case is still open after #26116 (Partial) and #25724 (Final). Maybe we could look into whether blocked storage, or an alternative approach, would help with the string case.

    I agree this sounds like a good idea @jayzhan211

    So what are the next steps here? Shall we get #26116 (Partial) and #25724 (Final) ready for review?

    It would probably be good to hear from @rluvaton to know if the combination of those two PRs would solve his usecase

  27. jayzhan211 commented on Oct 10, 2026

    @jayzhan211
    Contributor

    I'll wait a few days for more input. If there's no further discussion, let's move ahead with reviewing #26116 and #25724. @Rachelint, feel free to start getting #26116 ready for review in the meantime if you'd like. I'll also take a look at #25877.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    enhancementNew feature or request

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions