Skip to content

feat(agg): add initial Blocked aggregate api - #25707

Draft
rluvaton wants to merge 9 commits into
apache:mainfrom
rluvaton:add-initial-blocked-agg-poc
Draft

rluvaton wants to merge 9 commits into
apache:mainfrom
rluvaton:add-initial-blocked-agg-poc

Conversation

@rluvaton

@rluvaton rluvaton commented Sep 24, 2026 •

Copy link
Copy Markdown
Member

Created alternative PR in:

Which issue does this PR close?

Part of:

Rationale for this change

See issue

What changes are included in this PR?

It contain BlockedGroupsAccumulator trait, helper BlockedVec, 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

What is the testing strategy for this PR?

Existing tests

Are there any user-facing changes?

yes

@github-actions github-actions Bot added logical-expr Logical plan and expressions physical-expr Changes to the physical-expr crates functions Changes to functions implementation physical-plan Changes to the physical-plan crate labels Sep 24, 2026
evaluate_to_columns wrapped all evaluated blocks as a single entry,
producing batches with one column per block once there was more than
one block.
@github-actions

github-actions Bot commented Sep 24, 2026 •

Copy link
Copy Markdown

Thank you for opening this pull request!

Reviewer note: cargo-semver-checks reported the current version number is not SemVer-compatible with the changes in this pull request (compared against the base branch).

Details
     Cloning apache/main
    Building datafusion-expr v55.1.0 (current)
       Built [  32.566s] (current)
     Parsing datafusion-expr v55.1.0 (current)
      Parsed [   0.089s] (current)
    Building datafusion-expr v55.1.0 (baseline)
       Built [  31.937s] (baseline)
     Parsing datafusion-expr v55.1.0 (baseline)
      Parsed [   0.091s] (baseline)
    Checking datafusion-expr v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   2.245s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [  68.396s] datafusion-expr
    Building datafusion-expr-common v55.1.0 (current)
       Built [  20.458s] (current)
     Parsing datafusion-expr-common v55.1.0 (current)
      Parsed [   0.024s] (current)
    Building datafusion-expr-common v55.1.0 (baseline)
       Built [  20.409s] (baseline)
     Parsing datafusion-expr-common v55.1.0 (baseline)
      Parsed [   0.023s] (baseline)
    Checking datafusion-expr-common v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.335s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [  42.267s] datafusion-expr-common
    Building datafusion-functions-aggregate v55.1.0 (current)
       Built [  32.579s] (current)
     Parsing datafusion-functions-aggregate v55.1.0 (current)
      Parsed [   0.051s] (current)
    Building datafusion-functions-aggregate v55.1.0 (baseline)
       Built [  32.639s] (baseline)
     Parsing datafusion-functions-aggregate v55.1.0 (baseline)
      Parsed [   0.052s] (baseline)
    Checking datafusion-functions-aggregate v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.278s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [  67.219s] datafusion-functions-aggregate
    Building datafusion-functions-aggregate-common v55.1.0 (current)
       Built [  21.852s] (current)
     Parsing datafusion-functions-aggregate-common v55.1.0 (current)
      Parsed [   0.023s] (current)
    Building datafusion-functions-aggregate-common v55.1.0 (baseline)
       Built [  21.723s] (baseline)
     Parsing datafusion-functions-aggregate-common v55.1.0 (baseline)
      Parsed [   0.022s] (baseline)
    Checking datafusion-functions-aggregate-common v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.187s] 223 checks: 222 pass, 1 fail, 0 warn, 31 skip

--- failure function_requires_different_generic_type_params: function now requires a different number of generic type parameters ---

Description:
A function now requires a different number of generic type parameters than it used to. Uses of this function that supplied the previous number of generic types (e.g. via turbofish syntax) will be broken.
        ref: https://doc.rust-lang.org/reference/items/generics.html
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/function_requires_different_generic_type_params.ron

Failed in:
  function accumulate (2 -> 3 generic types) in /home/runner/work/datafusion/datafusion/datafusion/functions-aggregate-common/src/aggregate/groups_accumulator/accumulate.rs:421
  function accumulate_indices (1 -> 2 generic types) in /home/runner/work/datafusion/datafusion/datafusion/functions-aggregate-common/src/aggregate/groups_accumulator/accumulate.rs:591

     Summary semver requires new major version: 1 major and 0 minor checks failed
    Finished [  44.830s] datafusion-functions-aggregate-common
    Building datafusion-physical-expr v55.1.0 (current)
       Built [  29.746s] (current)
     Parsing datafusion-physical-expr v55.1.0 (current)
      Parsed [   0.055s] (current)
    Building datafusion-physical-expr v55.1.0 (baseline)
       Built [  29.827s] (baseline)
     Parsing datafusion-physical-expr v55.1.0 (baseline)
      Parsed [   0.054s] (baseline)
    Checking datafusion-physical-expr v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.518s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [  61.219s] datafusion-physical-expr
    Building datafusion-physical-plan v55.1.0 (current)
       Built [  39.563s] (current)
     Parsing datafusion-physical-plan v55.1.0 (current)
      Parsed [   0.184s] (current)
    Building datafusion-physical-plan v55.1.0 (baseline)
       Built [  40.015s] (baseline)
     Parsing datafusion-physical-plan v55.1.0 (baseline)
      Parsed [   0.187s] (baseline)
    Checking datafusion-physical-plan v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   1.032s] 223 checks: 222 pass, 1 fail, 0 warn, 31 skip

--- failure method_parameter_count_changed: pub method parameter count changed ---

Description:
A publicly-visible method now takes a different number of parameters, not counting the receiver (self) parameter.
        ref: https://doc.rust-lang.org/cargo/reference/semver.html#fn-change-arity
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/method_parameter_count_changed.ron

Failed in:
  datafusion_physical_plan::aggregates::order::GroupOrderingPartial::try_new takes 1 parameters in /home/runner/work/datafusion/datafusion/target/semver-checks/git-apache_main/2be1376b6526d0f4b80015fd1935a7b1bf748fe8/datafusion/physical-plan/src/aggregates/order/partial.rs:120, but now takes 2 parameters in /home/runner/work/datafusion/datafusion/datafusion/physical-plan/src/aggregates/order/partial.rs:123
  datafusion_physical_plan::aggregates::order::GroupOrdering::try_new takes 1 parameters in /home/runner/work/datafusion/datafusion/target/semver-checks/git-apache_main/2be1376b6526d0f4b80015fd1935a7b1bf748fe8/datafusion/physical-plan/src/aggregates/order/mod.rs:44, but now takes 2 parameters in /home/runner/work/datafusion/datafusion/datafusion/physical-plan/src/aggregates/order/mod.rs:45

     Summary semver requires new major version: 1 major and 0 minor checks failed
    Finished [  82.430s] datafusion-physical-plan

@github-actions github-actions Bot added the auto detected api change Auto detected API change label Sep 24, 2026

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Ignore the implementation itself, but just look at the API itself, this can be changed to change into blocks or flat depending on some threshold

@jayzhan211 jayzhan211 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@rluvaton

I'd suggest reversing the order: first get an implementation that measurably fixes the problem in #24704, then shape the API around it. An API designed before we know what the fast, memory-efficient implementation needs tends to get in its way, and it's much harder to change once it's public.

The current draft shows the risk. Every aggregation now goes through the blocked traits, and most of it runs through adapters over flat state, which is slower than main and uses more memory:

  • group_values/blocked.rs:195 and blocked_groups_accumulator.rs:312: emit All through the adapters is O(G²/B). With datafusion-cli, 1 partition, GROUP BY concat('k', v) plus sum(v) at 4M groups takes 0.21 s on main and 30.2 s here. A primitive key with sum is 0.1 s vs 0.4 s.
  • aggregates/spill.rs:220: one spill file per block. aggregate_memory_spill.slt sees spill_count go from 7 to about 800, and the count(DISTINCT) cases at L95/L106 now fail with ResourcesExhausted.

Proposal:

  1. Pick one target from #24704, e.g. peak memory / spill-free limit and emit time for high-cardinality GROUP BY <primitive> with sum/count/min/max.
  2. Implement blocked storage natively for just that case, behind a gate (blocked only when the group values and all accumulators support it, flat path unchanged otherwise), with no adapters.
  3. Show before/after numbers: peak memory, spill count, and ClickBench/TPC-H timings with no regressions elsewhere.
  4. Then extract the API from what that implementation actually needed, and extend it to more accumulators and group-value types.

That keeps each step reviewable and makes sure the API we commit to is one that performs well.

is_reversed: bool,
input_fields: Vec<FieldRef>,
is_nullable: bool,
batch_size: Option<usize>,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The block size shouldn't be set at plan time on AggregateFunctionExpr. The table uses the runtime context.session_config().batch_size(), so a plan executed with a different TaskContext fails at the assert in aggregate_hash_table/common.rs:100 with Internal("... left: 4096, right: 8192: Block size mismatch ..."). Suggest dropping the field and the builder method, and having the caller pass the runtime block size in:

pub fn blocked_groups_accumulator_supported(&self, block_size: usize) -> bool {
pub fn create_blocked_groups_accumulator(
    &self,
    block_size: usize,
) -> Result<Box<dyn BlockedGroupsAccumulator>> {

}

pub trait BlockedGroupsAccumulator: Send + std::any::Any {
fn batch_size(&self) -> usize;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: this is named batch_size(), but everything else (BlockedVec::block_size, BlockedGroupValues::block_size) says block size. Suggest fn block_size(&self) -> usize;, and also renaming BlockedAccumulatorArgs::batch_size.

/// Create an accumulator for `agg_expr` -- a [`BlockedGroupsAccumulator`] if
/// that is supported by the aggregate, or a
/// [`BlockedGroupsAccumulatorAdapter`] if not.
pub(in crate::aggregates) fn create_blocked_group_accumulator(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should the blocked API be opt-in rather than the route for every aggregation? A flat GroupsAccumulator/GroupValues can't serve BlockedEmitTo cheaply: the adapters loop EmitTo::First(block_size), and each call shifts the whole table. Suggest taking the blocked path only when the group values and every accumulator support it natively, and keeping the flat path unchanged otherwise. That also keeps the adapters out of the public shape.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this mean I have to duplicate the hash aggregate streams code and I wanted to avoid this

GroupsAccumulator should not be used and should be removed in the end, we will not release until all GroupValues has been converted

/// Requirements:
/// 1. `n` must be smaller than block_size
/// 2. `n` is not 0
First(usize),

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The First(n) precondition n < block_size leaks to callers: common_ordered.rs:481 already rebuilds EmitTo::First(n) as NextBlocks plus a First remainder. Either allow any n, or define that decomposition once, next to the enum:

impl BlockedEmitTo {
    /// Splits a flat `EmitTo` into block-sized emits.
    pub fn from_emit_to(emit_to: EmitTo, block_size: usize) -> Vec<Self> {
        match emit_to {
            EmitTo::All => vec![Self::All],
            EmitTo::First(n) => {
                let mut emits = vec![Self::NextBlock; n / block_size];
                if n % block_size != 0 {
                    emits.push(Self::First(n % block_size));
                }
                emits
            }
        }
    }
}

/// - `BlockedEmitTo::NextBlock` it should return single item vector with the block or empty vec in case of no blocks
/// - `BlockedEmitTo::First(n)` it should return single item vector with the first n rows in the first block. n must be smaller than block size and length
///
fn evaluate(&mut self, emit_to: BlockedEmitTo) -> Result<Vec<ArrayRef>>;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Return shapes:

  • NextBlock yields 0 or 1 blocks, so returning a Vec hides that; Option<ArrayRef> (or a dedicated evaluate_next_block) would say so.
  • evaluate_preserving (L348) and state_preserving return one flat array rather than blocks, which is inconsistent with evaluate and state.
  • If the plan is an iterator for emit All (the TODO in group_values/blocked.rs), it's cheaper to put it in the trait now than to change it later.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

evaluate_preserving is like that since it provide the indices, I can split the indices by block size and return Vec

Return shapes:

  • NextBlock yields 0 or 1 blocks, so returning a Vec hides that; Option<ArrayRef> (or a dedicated evaluate_next_block) would say so.
    I thought about it originally but in order to not explode the code with functions to implement I thought this will be a better way, what do you think?
  • evaluate_preserving (L348) and state_preserving return one flat array rather than blocks, which is inconsistent with evaluate and state.

evaluate_preserving/state_preserving is like that since it provide the indices, I can split the indices by block size and return Vec if that is what you mean

  • If the plan is an iterator for emit All (the TODO in group_values/blocked.rs), it's cheaper to put it in the trait now than to change it later.

the iterator that emit all is just an idea, this will be useful when materializing will use more memory than needed

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

having iterator create challenges regarding memory size

}
}

pub trait BlockedGroupsAccumulator: Send + std::any::Any {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This trait repeats GroupsAccumulator method for method, and BlockedGroupSelection (L37) repeats GroupSelection; only the index type changes. Could the two share a generic index type, or could the blocked methods be defaulted on GroupsAccumulator? That would avoid two parallel public traits that have to stay in sync.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

BlockedGroupsAccumulator would have to extend GroupsAccumulator and I don't want that since every function call now will be confusion to what is being called and also the plan is to remove GroupsAccumulator from the way I see it and not keep it, so adding dependencies to old implementation would just couple more the implementations

/// margin, `index_in_block < block_size` (the batch size) and the number of blocks is
/// `groups / block_size`. The API keeps taking and returning `usize`.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Hash)]
pub struct BlocksIndex {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

BlocksIndex puts a lot of public surface in expr-common: sub_flat, sub_flat_checked (L199), gte_flat, prev_block, add, increment. Only group-values implementations need these. Suggest keeping the public type to new, block_index, index_in_block and the flat conversions, and moving the rest next to their only user, or marking them #[doc(hidden)].

@rluvaton

Copy link
Copy Markdown
Member Author

@rluvaton

I'd suggest reversing the order: first get an implementation that measurably fixes the problem in #24704, then shape the API around it. An API designed before we know what the fast, memory-efficient implementation needs tends to get in its way, and it's much harder to change once it's public.

The current draft shows the risk. Every aggregation now goes through the blocked traits, and most of it runs through adapters over flat state, which is slower than main and uses more memory:

Yes, but this is to make it easier to review, I can copy the entire aggregation code and replace there with the blocked but it is harder to review

  • group_values/blocked.rs:195 and blocked_groups_accumulator.rs:312: emit All through the adapters is O(G²/B). With datafusion-cli, 1 partition, GROUP BY concat('k', v) plus sum(v) at 4M groups takes 0.21 s on main and 30.2 s here. A primitive key with sum is 0.1 s vs 0.4 s.

I've did not implement in this PR the bytes/bytes view group by, but it is implemented in later pr

  • aggregates/spill.rs:220: one spill file per block. aggregate_memory_spill.slt sees spill_count go from 7 to about 800, and the count(DISTINCT) cases at L95/L106 now fail with ResourcesExhausted.

I'm aware of that and I added a comment and fixed that problem in later PR by adding sort without concat which fixes that problem

Proposal:

  1. Pick one target from [EPIC] Use blocked / chunked memory management in hash aggregation #24704, e.g. peak memory / spill-free limit and emit time for high-cardinality GROUP BY <primitive> with sum/count/min/max.
  1. Implement blocked storage natively for just that case, behind a gate (blocked only when the group values and all accumulators support it, flat path unchanged otherwise), with no adapters.

Yes, but I've done this way to make it easier to review, I can copy the entire aggregation code and replace there with the blocked but it is harder to review since you have no clear way to see what I actually changed

  1. Show before/after numbers: peak memory, spill count, and ClickBench/TPC-H timings with no regressions elsewhere.

This pr is after I've done all of this (but without the flat approach), the pr that implement it entirely is:

  1. Then extract the API from what that implementation actually needed, and extend it to more accumulators and group-value types.

That keeps each step reviewable and makes sure the API we commit to is one that performs well.

@jayzhan211

Copy link
Copy Markdown
Contributor

@rluvaton

Do you think it's possible to break #24928 into several PRs, where each one includes not just the API but also part of the implementation? Ideally each PR would be safe to merge into main on its own with no regression, assuming the follow-ups might not land for several releases.

If that's not workable, I think we'd need a feature branch for this. On a feature branch we could merge the API first and then the rest incrementally.

@rluvaton

Copy link
Copy Markdown
Member Author

@jayzhan211 Ok, so instead of having an adapter and harm performance I created another PR with alternative approach which support both and have less performance hit for unsupported cases

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

Labels

auto detected api change Auto detected API change functions Changes to functions implementation logical-expr Logical plan and expressions physical-expr Changes to the physical-expr crates physical-plan Changes to the physical-plan crate

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants