This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-25605-ac37adf5a9099b821c64ef1b7dbdcdc115e6d720 in repository https://gitbox.apache.org/repos/asf/datafusion.git
commit ca0a6a42d799b856f3894ca14ef4f186be0c5011 Author: Adrian Garcia Badaracco <[email protected]> AuthorDate: Wed Sep 23 21:45:21 2026 +0000 perf: reuse GroupsAccumulatorAdapter scratch buffers across batches (#25605) ## Which issue does this PR close? - Closes https://github.com/apache/datafusion/issues/25604. ## Rationale for this change `GroupsAccumulatorAdapter` runs every aggregate that has no `GroupsAccumulator` of its own. For each input batch, `invoke_per_accumulator` allocated three new buffers, `groups_with_rows`, `offsets` and `batch_indices`, and dropped them at the end of the batch. @alamb suggested to reuse them in https://github.com/apache/datafusion/pull/25123#discussion_r4066870146. ## What changes are included in this PR? 1. **Reuse `groups_with_rows` and `offsets`.** A new `Scratch` struct on the adapter holds the buffers. Each batch clears them and fills them again, so they allocate only when a batch needs more capacity than an earlier batch. `invoke_per_accumulator` takes the buffers out of `self` for the batch and puts them back on every return path, errors included. The body moves unchanged into `invoke_per_accumulator_with_scratch`. 2. **Reuse `batch_indices`.** The adapter moves the buffer into the `UInt32Array` for the `take` kernels without a copy. The kernels only borrow the array, thus the adapter then gets the buffer back without a copy through `into_parts`. 3. **Shrink oversized buffers.** A buffer keeps the capacity of the largest batch that it has seen. After each batch, the adapter shrinks a buffer whose capacity is more than 64 Ki entries and more than 4 times what the batch used. It shrinks to the smaller of the two, so it never keeps more than either limit allows. Thus one unusually large batch does not make the adapter keep that memory for all the batches after it. The adapter charges the change in the capacity of the buffers to `allocation_bytes`, thus `GroupsAccumulator::size` includes them and the memory pool sees them. It releases them when an emit leaves no group, thus an adapter that has emitted all groups still reports 0 bytes. Each adapter belongs to one partition's aggregate stream, and `GroupsAccumulator` methods take `&mut self`, so the buffers are never shared between partitions. I have not measured a speed change. This removes three allocations for each batch, while the same batch calls `Accumulator::update_batch` once for each group that it touches. Thus I expect the change to be small. ## What is the testing strategy for this PR? - `adapter_charges_retained_indices_once_and_releases_them`, `adapter_clears_successful_group_indices_after_later_error` and `adapter_reconciles_allocation_after_later_error_with_grouped_metric` check the exact memory charge, which now includes the scratch buffers. The first one also checks that the adapter keeps some scratch capacity after a batch, that `batch_indices` comes back from the `take` kernels, and that `size()` is 0 after `evaluate(EmitTo::All)`. - `adapter_releases_oversized_scratch` sends a batch with 128 Ki groups and then a batch with one row. It checks that the second batch shrinks each large buffer to exactly 4 times what it used, and that `size()` goes down by the released amount. To make sure that these tests catch an error, I removed each part in turn: - Without the charge for the scratch buffers, the three accounting tests fail. - Without the release on emit, `adapter_charges_retained_indices_once_and_releases_them` fails. - Without the shrink of oversized buffers, without it for `batch_indices` alone, with a release to capacity 0 instead of the shrink, or with a shrink to 64 Ki entries instead of the smaller limit, `adapter_releases_oversized_scratch` fails. The unit tests of `datafusion-functions-aggregate-common`, the `memory_limit` tests in `core_integration` and the sqllogictest suite pass. ## Are there any user-facing changes? No. There are no changes to any public API or to any result. 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 5 <[email protected]> --- .../src/aggregate/groups_accumulator.rs | 184 ++++++++++++++++++++- 1 file changed, 177 insertions(+), 7 deletions(-) diff --git a/datafusion/functions-aggregate-common/src/aggregate/groups_accumulator.rs b/datafusion/functions-aggregate-common/src/aggregate/groups_accumulator.rs index d49cae6a16..dbf2af512b 100644 --- a/datafusion/functions-aggregate-common/src/aggregate/groups_accumulator.rs +++ b/datafusion/functions-aggregate-common/src/aggregate/groups_accumulator.rs @@ -111,6 +111,63 @@ pub struct GroupsAccumulatorAdapter { /// Optional aggregate-owned metric timed once for a grouped update batch. grouped_update_metric: OnceLock<Option<Arc<dyn AggregateMetric>>>, + + /// Buffers reused by every call to `invoke_per_accumulator` + scratch: Scratch, +} + +/// Per-batch buffers of [`GroupsAccumulatorAdapter`], kept between batches so +/// that each batch does not allocate them again. Their capacity is part of +/// [`GroupsAccumulatorAdapter::allocation_bytes`]. +#[derive(Default)] +struct Scratch { + /// Group indexes that have rows in the batch, in order of first + /// appearance in the batch + groups_with_rows: Vec<usize>, + + /// `offsets[i]` is the index into `batch_indices` where the rows for + /// `groups_with_rows[i]` start + offsets: Vec<usize>, + + /// Indices into the batch rows, with the rows of each group contiguous + batch_indices: Vec<u32>, +} + +/// A scratch buffer keeps its capacity after a batch unless that capacity is +/// more than this many entries and more than [`SCRATCH_RETAIN_RATIO`] times +/// what the batch used. Then it shrinks to the smaller of the two. This stops +/// one unusually large batch from holding memory for all the batches after it. +const MAX_RETAINED_SCRATCH_ENTRIES: usize = 64 * 1024; + +/// See [`MAX_RETAINED_SCRATCH_ENTRIES`] +const SCRATCH_RETAIN_RATIO: usize = 4; + +impl Scratch { + fn allocated_size(&self) -> usize { + self.groups_with_rows.allocated_size() + + self.offsets.allocated_size() + + self.batch_indices.allocated_size() + } + + /// Shrink each buffer that is much larger than the batch that just used + /// it. Call this after the batch, before the buffers are cleared for the + /// next batch. + fn release_oversized(&mut self) { + fn release_if_oversized<T>(buffer: &mut Vec<T>) { + let keep = SCRATCH_RETAIN_RATIO * buffer.len(); + if buffer.capacity() > MAX_RETAINED_SCRATCH_ENTRIES + && buffer.capacity() > keep + { + // The batch is done with the contents, so clear them first + // and the shrink has nothing to copy + buffer.clear(); + buffer.shrink_to(MAX_RETAINED_SCRATCH_ENTRIES.min(keep)); + } + } + release_if_oversized(&mut self.groups_with_rows); + release_if_oversized(&mut self.offsets); + release_if_oversized(&mut self.batch_indices); + } } /// Maximum number of prepared group inputs retained while timing an @@ -154,6 +211,7 @@ impl GroupsAccumulatorAdapter { allocation_bytes: 0, metrics: None, grouped_update_metric: OnceLock::new(), + scratch: Scratch::default(), } } @@ -235,10 +293,45 @@ impl GroupsAccumulatorAdapter { assert_eq!(values[0].len(), group_indices.len()); + // Take the scratch buffers out of `self` for the batch, and put them + // back on every return path, including errors + let mut scratch = std::mem::take(&mut self.scratch); + let scratch_size_pre = scratch.allocated_size(); + let result = self.invoke_per_accumulator_with_scratch( + &mut scratch, + values, + group_indices, + opt_filter, + f, + ); + scratch.release_oversized(); + self.adjust_allocation(scratch_size_pre, scratch.allocated_size()); + self.scratch = scratch; + result + } + + /// Body of [`Self::invoke_per_accumulator`], with its scratch buffers + fn invoke_per_accumulator_with_scratch<F>( + &mut self, + scratch: &mut Scratch, + values: &[ArrayRef], + group_indices: &[usize], + opt_filter: Option<&BooleanArray>, + f: F, + ) -> Result<()> + where + F: Fn(&mut dyn Accumulator, &[ArrayRef]) -> Result<()>, + { + let Scratch { + groups_with_rows, + offsets, + batch_indices, + } = scratch; + // groups_with_rows holds a list of group indexes that have any rows // that need to be accumulated, stored in order of first appearance in // this batch - let mut groups_with_rows = vec![]; + groups_with_rows.clear(); // figure out which input rows correspond to which groups. // Note that self.state.indices starts empty for all groups (it is @@ -263,11 +356,18 @@ impl GroupsAccumulatorAdapter { self.add_allocation(indices_allocation_delta); // batch_indices holds indices into values, each group is contiguous - let mut batch_indices = Vec::with_capacity(group_indices.len()); + batch_indices.clear(); + batch_indices + .try_reserve(group_indices.len()) + .map_err(|e| { + arrow_datafusion_err!(arrow::error::ArrowError::MemoryError( + e.to_string() + )) + })?; // offsets[i] is index into batch_indices where the rows for // groups_with_rows[i] start - let mut offsets = Vec::with_capacity(groups_with_rows.len() + 1); + offsets.clear(); offsets.push(0); let mut offset_so_far = 0; @@ -277,13 +377,20 @@ impl GroupsAccumulatorAdapter { offset_so_far += indices.len(); offsets.push(offset_so_far); } - let batch_indices = batch_indices.into(); + // Move the buffer into an array without a copy, for the take kernels + let indices_array = + PrimitiveArray::<UInt32Type>::from(std::mem::take(batch_indices)); // reorder the values and opt_filter by batch_indices so that // all values for each group are contiguous, then invoke the // accumulator once per group with values - let values = take_arrays(values, &batch_indices, None)?; - let opt_filter = get_filter_at_indices(opt_filter, &batch_indices)?; + let values = take_arrays(values, &indices_array, None)?; + let opt_filter = get_filter_at_indices(opt_filter, &indices_array)?; + + // The take kernels only borrow the indices, so the array holds the + // only reference to the buffer and we can get it back without a copy + let (_, indices_buffer, _) = indices_array.into_parts(); + *batch_indices = indices_buffer.into_inner().into_vec().unwrap_or_default(); let grouped_update_metric = self.grouped_metric(); @@ -398,6 +505,15 @@ impl GroupsAccumulatorAdapter { self.allocation_bytes = self.allocation_bytes.saturating_sub(size) } + /// Release the scratch buffers when no group is left, so that an adapter + /// that has emitted every group holds no memory + fn free_scratch_if_empty(&mut self) { + if self.states.is_empty() { + let scratch = std::mem::take(&mut self.scratch); + self.free_allocation(scratch.allocated_size()); + } + } + /// Release the allocation held by a state that is being emitted. fn free_state_allocation(&mut self, state: &AccumulatorState) { self.free_allocation(state.size()); @@ -456,6 +572,7 @@ impl GroupsAccumulator for GroupsAccumulatorAdapter { let result = ScalarValue::iter_to_array(results); self.adjust_allocation(vec_size_pre, self.states.allocated_size()); + self.free_scratch_if_empty(); result } @@ -520,6 +637,7 @@ impl GroupsAccumulator for GroupsAccumulatorAdapter { } } self.adjust_allocation(vec_size_pre, self.states.allocated_size()); + self.free_scratch_if_empty(); Ok(arrays) } @@ -820,7 +938,14 @@ mod tests { let retained_indices = adapter.states[0].indices.allocated_size(); assert!(retained_indices > 0); assert!(adapter.states[0].indices.is_empty()); - assert_eq!(adapter.size(), allocation_before_update + retained_indices); + // The indices buffer comes back from the take kernels without a copy + assert!(adapter.scratch.batch_indices.capacity() >= 4); + let retained_scratch = adapter.scratch.allocated_size(); + assert!(retained_scratch > 0); + assert_eq!( + adapter.size(), + allocation_before_update + retained_indices + retained_scratch + ); let allocation_after_first_update = adapter.size(); adapter.update_batch(&[values], &[0, 0, 0, 0], None, 1)?; @@ -1085,6 +1210,7 @@ mod tests { .iter() .map(AccumulatorState::size) .sum::<usize>() + + accumulator.scratch.allocated_size() ); accumulator.update_batch(&[values], &[0, 1], None, 2)?; @@ -1129,6 +1255,7 @@ mod tests { .iter() .map(AccumulatorState::size) .sum::<usize>() + + accumulator.scratch.allocated_size() ); Ok(()) } @@ -1158,6 +1285,49 @@ mod tests { Ok(()) } + /// One large batch must not make the adapter keep its scratch capacity + /// for the smaller batches after it. + #[test] + fn adapter_releases_oversized_scratch() -> Result<()> { + const NUM_GROUPS: usize = 2 * MAX_RETAINED_SCRATCH_ENTRIES; + let mut adapter = GroupsAccumulatorAdapter::new(|| { + Ok(Box::new(MaxAccumulator::try_new(&DataType::Int64)?) + as Box<dyn Accumulator>) + }); + + // every row is a different group, so the scratch buffers need one + // entry for each row + let group_indices: Vec<usize> = (0..NUM_GROUPS).collect(); + let values: ArrayRef = + Arc::new(Int64Array::from_iter_values(0..NUM_GROUPS as i64)); + adapter.update_batch(&[values], &group_indices, None, NUM_GROUPS)?; + let large_scratch = adapter.scratch.allocated_size(); + assert!(large_scratch >= NUM_GROUPS * size_of::<usize>()); + let size_after_large_batch = adapter.size(); + + // a batch with one row shrinks each large buffer to + // SCRATCH_RETAIN_RATIO times what it used: one group, two offsets and + // one row index + let values: ArrayRef = Arc::new(Int64Array::from(vec![1])); + adapter.update_batch(&[values], &[0], None, NUM_GROUPS)?; + let small_scratch = adapter.scratch.allocated_size(); + assert!(small_scratch < large_scratch); + assert_eq!( + adapter.scratch.groups_with_rows.capacity(), + SCRATCH_RETAIN_RATIO + ); + assert_eq!(adapter.scratch.offsets.capacity(), 2 * SCRATCH_RETAIN_RATIO); + assert_eq!( + adapter.scratch.batch_indices.capacity(), + SCRATCH_RETAIN_RATIO + ); + assert_eq!( + adapter.size(), + size_after_large_batch - large_scratch + small_scratch + ); + Ok(()) + } + #[test] fn adapter_preserving_evaluation_uses_accumulator_contract() -> Result<()> { let mut accumulator = GroupsAccumulatorAdapter::new(|| { --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
