Conversation
Blocked storage for group keys and accumulators (count, count distinct, sum, avg, min/max bytes), mixed with flat storage per part.
|
run benchmarks |
|
run benchmarks h2o_medium external_aggr |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (a97e6b6) to dbf62c1 (merge-base) diff Run configurationrun benchmark clickbench_partitionedResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (a97e6b6) to dbf62c1 (merge-base) diff Run configurationrun benchmark external_aggrResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (a97e6b6) to dbf62c1 (merge-base) diff Run configurationrun benchmark tpcdsResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (a97e6b6) to dbf62c1 (merge-base) diff Run configurationrun benchmark h2o_mediumResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (a97e6b6) to dbf62c1 (merge-base) diff Run configurationrun benchmark tpchResults will be posted here when complete File an issue against this benchmark runner |
| /// otherwise the whole aggregation uses [`GroupsAccumulator`]. | ||
| /// | ||
| /// [`GroupsAccumulator`]: crate::groups_accumulator::GroupsAccumulator | ||
| pub trait BlockedGroupsAccumulator: Send + Any { |
There was a problem hiding this comment.
Will add the rest of the functions (evaluate/state preserving) in later pr so it will be easier to review
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (a97e6b6) to dbf62c1 (merge-base) diff Run configurationrun benchmark tpchCPU Details (lscpu)Details
Resource Usagetpch — base (merge-base)
tpch — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (a97e6b6) to dbf62c1 (merge-base) diff Run configurationrun benchmark tpcdsCPU Details (lscpu)Details
Resource Usagetpcds — base (merge-base)
tpcds — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (a97e6b6) to dbf62c1 (merge-base) diff Run configurationrun benchmark external_aggrCPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usageexternal_aggr — base (merge-base)
external_aggr — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (a97e6b6) to dbf62c1 (merge-base) diff Run configurationrun benchmark clickbench_partitionedCPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (a97e6b6) to dbf62c1 (merge-base) diff Run configurationrun benchmark h2o_mediumCPU Details (lscpu)Details
Resource Usageh2o_medium — base (merge-base)
h2o_medium — branch
File an issue against this benchmark runner |
|
run benchmarks env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M" |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (c160407) to dbf62c1 (merge-base) diff Run configurationrun benchmark clickbench_partitioned
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (c160407) to dbf62c1 (merge-base) diff Run configurationrun benchmark tpcds
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (c160407) to dbf62c1 (merge-base) diff Run configurationrun benchmark tpch
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (ce4e7c6) to dbf62c1 (merge-base) diff Run configurationrun benchmark tpcdsCPU Details (lscpu)Details
Resource Usagetpcds — base (merge-base)
tpcds — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (ce4e7c6) to dbf62c1 (merge-base) diff Run configurationrun benchmark external_aggrCPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usageexternal_aggr — base (merge-base)
external_aggr — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (ce4e7c6) to dbf62c1 (merge-base) diff Run configurationrun benchmark clickbench_partitionedCPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (ce4e7c6) to dbf62c1 (merge-base) diff Run configurationrun benchmark clickbench_partitioned
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "16G"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (ce4e7c6) to dbf62c1 (merge-base) diff Run configurationrun benchmark h2o_mediumCPU Details (lscpu)Details
Resource Usageh2o_medium — base (merge-base)
h2o_medium — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (ce4e7c6) to dbf62c1 (merge-base) diff Run configurationrun benchmark h2o_medium
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "16G"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usageh2o_medium — base (merge-base)
h2o_medium — branch
File an issue against this benchmark runner |
|
run benchmarks |
|
run benchmarks h2o_medium external_aggr |
|
run benchmarks tpch10 h2o_medium clickbench_partitioned env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "16G" |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark tpchResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark h2o_mediumResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark external_aggrResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark tpcdsResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark clickbench_partitionedResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark tpchCPU Details (lscpu)Details
Resource Usagetpch — base (merge-base)
tpch — branch
File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark tpch10
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "16G"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark tpcdsCPU Details (lscpu)Details
Resource Usagetpcds — base (merge-base)
tpcds — branch
File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark h2o_medium
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "16G"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark external_aggrCPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usageexternal_aggr — base (merge-base)
external_aggr — branch
File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark clickbench_partitioned
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "16G"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark clickbench_partitionedCPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark h2o_mediumCPU Details (lscpu)Details
Resource Usageh2o_medium — base (merge-base)
h2o_medium — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark tpch10
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "16G"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usagetpch10 — base (merge-base)
tpch10 — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark clickbench_partitioned
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "16G"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing blocked-agg-poc (51fca43) to 6924c6a (merge-base) diff Run configurationrun benchmark h2o_medium
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "16G"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usageh2o_medium — base (merge-base)
h2o_medium — branch
File an issue against this benchmark runner |
changed to block size without 2^18, so it is now covered
first and next block will be used by ordered but also outside of datafusion so they are valuable from the start, and we can avoid adding breaking change later |
jayzhan211
left a comment
There was a problem hiding this comment.
I left some comments and also the test failed
|
|
||
| // Make sure we can hold on the hash tables and all the batches that need to be emitted (except the first one) | ||
| // if we can't hold it than we can't do anything about it. | ||
| self.reservation.try_resize( |
There was a problem hiding this comment.
try_resize(hash_table.memory_size() + pending_memory)? dropped the fallback main had: when holding the materialized state doesn't fit, main reserved only the table and emitted unreserved. With blocked storage (a blocked count is enough, even with flat keys), partial aggregation now fails with ResourcesExhausted. 4 existing tests fail here and pass on the merge-base 6924c6a: both aggregate_grouping_sets_*_with_spill (Failed to allocate additional 872.0 B ... pool_size: 500.0 B), test_partial_hash_stream_emits_whole_batch_when_held_batch_does_not_fit, and test_partial_hash_stream_accounts_held_batch_on_memory_pressure_while_slicing.
Fix (with it, both grouping-sets tests pass locally):
- self.reservation.try_resize(
- hash_table.memory_size()
- // The batches that need to be held while emitting each one
- + pending_memory,
- )?;
+ let hold_pending = match self
+ .reservation
+ .try_resize(hash_table.memory_size() + pending_memory)
+ {
+ Ok(()) => true,
+ // Can't hold the pending blocks: emit them unreserved, like
+ // the flat path emits the whole batch
+ Err(DataFusionError::ResourcesExhausted(_)) => {
+ self.reservation.try_resize(hash_table.memory_size())?;
+ false
+ }
+ Err(e) => return Err(e),
+ };
for (i, batch) in
materialized_group_states.into_iter().enumerate()
{
- if i != 0 {
+ if i != 0 && hold_pending {
self.reservation.try_shrink(batch.memory_size)?;
}The two hash_stream tests also assume a single state batch (first.num_rows() == num_groups, and a shared data_ptr between the first two outputs). Please update them for one batch per block, e.g. assert that the total rows match and that reserved covers the blocks not emitted yet.
There was a problem hiding this comment.
But you can't do that, this is "cheating",
you basically say, ok, the memory did not allow us to reserve memory for the batches we are holding between emits so we just don't count them.
this is just ignoring the problem, not a fix.
Main could do what it did since it only had 1 batch, so it just emitted it without holding on it, but now we have to hold it since we can't emit multiple batches at once
There was a problem hiding this comment.
Agreed that the held blocks must be counted, and my earlier suggestion (emit without reserving) didn't do that. But the error doesn't count them either: at this point they were never reserved. Logged in the failing tests: aggregate_grouping_sets_with_yielding_with_spill has 0 B reserved while the blocks hold 872 B (pool 500 B), and ..._does_not_fit has 3.5 KB reserved for 4 × 12.5 KB blocks, so ? frees nothing and fails a query main finishes. MemoryReservation::grow is infallible, so you can count them honestly, over the limit, until they drain:
- self.reservation.try_resize(
- hash_table.memory_size()
- // The batches that need to be held while emitting each one
- + pending_memory,
- )?;
+ let held = hash_table.memory_size() + pending_memory;
+ match self.reservation.try_resize(held) {
+ Ok(()) => {}
+ // The blocks already exist, so failing would not free
+ // them: count them anyway, over the pool limit, until
+ // they are emitted
+ Err(DataFusionError::ResourcesExhausted(_)) => {
+ self.reservation.grow(held - self.reservation.size())
+ }
+ Err(e) => return Err(e),
+ }With this both aggregate_grouping_sets_*_with_spill pass. The two test_partial_hash_stream_* failures assume one flat state batch (as their own comment anticipates): a Utf8 key + sum in partial_stream_under_memory_limit keeps the slicing coverage, plus an Int32 + count variant asserting no error at 32 KB.
There was a problem hiding this comment.
Only this issue is left!
There was a problem hiding this comment.
grow() should really not be used anywhere, you are suggesting adding a panic here, which is much more severe than returning an error for the query - which is also bad in my opinion, but preferable to exceeding the expected memory, which can lead to SIGKILLs that are far harder to debug.
Erroring out if the memory is umavailable is the way forward in my opinion.
The next step is to ensure this situation never happens if we are able to hold 2 batches in memory
| // Add one to each group's counter for each non null, non filtered value | ||
| // SAFETY: group_index is guaranteed to be in bounds and less than total_num_groups | ||
| unsafe { | ||
| self.counts.update_unchecked( |
There was a problem hiding this comment.
update_batch calls update_unchecked with only a SAFETY comment about the caller's contract. That means a safe public path (create_blocked_groups_accumulator → update_batch) can write out of bounds. Passing BlocksIndex::new(0, 1) with total_num_groups = 1 aborts in debug with unsafe precondition(s) violated: slice::get_unchecked_mut at blocked_vec.rs:282. merge_batch already goes through the checked update_with, and all_in_bounds is one vectorized pass, so the checked path should cost little here too. The flat count on main has the same pattern, so this is fine to handle in a follow-up if you'd rather.
- self.counts.grow_to(total_num_groups, 0);
-
- // Add one to each group's counter for each non null, non filtered value
- // SAFETY: group_index is guaranteed to be in bounds and less than total_num_groups
- unsafe {
- self.counts.update_unchecked(
- group_indices,
- nulls.as_ref(),
- opt_filter,
- |count| *count += 1,
- );
- }
+ // Add one to each group's counter for each non null, non filtered value
+ self.counts.update(
+ total_num_groups,
+ 0,
+ group_indices,
+ nulls.as_ref(),
+ opt_filter,
+ |count| *count += 1,
+ );There was a problem hiding this comment.
I had unchecked here to keep the same optimizations used as before so this can be as a follow up pr but not in this
| self.null_group = None; | ||
| blocks | ||
| } | ||
| BlockedEmitTo::NextBlock => { |
There was a problem hiding this comment.
NextBlock runs map.retain over every entry to renumber it, so draining n groups one block at a time costs O(n²/block_size). With 10M groups and 8192-row blocks that's about 1.2k passes over the full map. Incremental emit is what NextBlock is for, so this will show up once ordered or external callers use it. Nothing calls it yet, so a follow-up is fine.
A suggestion: store absolute block numbers in the map plus a first_block offset. intern subtracts the offset when it returns an index and adds it when it inserts one. NextBlock then removes only the emitted block's keys by lookup, at O(block_size) per block:
BlockedEmitTo::NextBlock => {
let Some(block) = self.values.take_next_block() else {
return Ok(vec![]);
};
let first_block = self.first_block;
let state = &self.random_state;
for &key in &block {
// The null group's slot holds `Default`, so match on the block too
if let Ok(entry) = self.map.find_entry(key.hash(state), |&(g, k)| {
g.block_index() == first_block && k.is_eq(key)
}) {
entry.remove();
}
}
self.first_block += 1;
let len = block.len();
let array = self.build_block(block, 0);
self.shift_null_group(len);
vec![vec![array]]
}All and clear_shrink reset first_block to 0. First(n) would also need to account for the offset.
There was a problem hiding this comment.
the shifting every entry exists in the non blocked implementation (which is why I did that), so to avoid a lot of changes I will keep that for now and this can be in a later PR
There was a problem hiding this comment.
But you are right, just not in this PR
# Conflicts: # datafusion/physical-plan/src/aggregates/single_stream.rs
Difference with my other PR:
That PR goes
BlockedGroupsAccumulatorfirst which means that it only interact with blocked and any non blocked are wrapped in adaptersthe reason for that is:
this PR however does not go that way, it support both blocked and flat without adapter,
the reason for that is to keep good performance for unsupported cases while also use blocked when possible
Important note
Currently it sort every block and spill it which will create more spill files and it reduce performance by increasing batch size merge degree on spill path, the fix is to have sort kernel that get a vector of arrays to sort and return them sorted, I did not add it in this PR to make it easier to review, but will add in later pr,
I already note that this is needed in the related issue comment with my findings
#24704 (comment)
Which issue does this PR close?
Part of:
Rationale for this change
See issue
What changes are included in this PR?
It contain
BlockedGroupsAccumulatortrait, helperBlockedVec, supportcountin blocked so you see the example usage, change the non-ordered aggregate (skipped ordered to make this pr smaller) to work with either blocked or flat, implement BlockedGroupValues for primitive so you will see how it is being usedWhat is the testing strategy for this PR?
added tests + existing
Are there any user-facing changes?
yes, but not breaking ones