Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions datafusion/common/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1136,6 +1136,30 @@ config_namespace! {
/// aggregation ratio check and trying to switch to skipping aggregation mode
pub skip_partial_aggregation_probe_rows_threshold: usize, default = 100_000

/// (experimental) Allocated-byte threshold for flushing unordered partial
/// hash aggregation tables. Checked after each input batch; the table can
/// exceed this size by one batch. Flushed groups count toward the existing
/// skip-partial probe. Repeated keys do not disable flushing. Aggregations
/// with a soft group limit or nested aggregate state are excluded.
/// Set to 0 to disable this threshold. If both flush thresholds are
/// enabled, reaching either triggers a flush. Final aggregation is unchanged.
pub partial_aggregation_flush_bytes: usize, default = 0

/// (experimental) Number of distinct group rows above which an unordered
/// partial hash aggregation table is flushed. This counts groups held in
/// the table, not input rows. Uses the same eligibility and skip-partial
/// accounting as partial_aggregation_flush_bytes. Set to 0 to disable this
/// threshold. If both thresholds are enabled, reaching either triggers a flush.
pub partial_aggregation_flush_rows: usize, default = 0

/// (experimental) Flush unordered partial aggregation after a complete input
/// batch when the group count reaches batch_size. Overrides both other flush
/// thresholds and reserves supported group-key and accumulator storage for
/// twice batch_size groups, initially and after each threshold flush.
/// Requires a single grouping set and the same eligibility as the other
/// flush thresholds. String payload buffers use their normal growth policy.
pub partial_aggregation_flush_batch: bool, default = false

/// Should DataFusion use row number estimates at the input to decide
/// whether increasing parallelism is beneficial or not. By default,
/// only exact row numbers (not estimates) are used for this decision.
Expand Down
4 changes: 4 additions & 0 deletions datafusion/expr-common/src/groups_accumulator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -381,6 +381,10 @@ pub trait GroupsAccumulator: Send + std::any::Any {
///
/// May be expensive; check the implementation before calling on hot paths.
fn size(&self) -> usize;

/// Reserves row-addressed state for `capacity` groups where supported.
/// This hint must not change the logical number of groups.
fn reserve_groups(&mut self, _capacity: usize) {}
}

#[cfg(test)]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,11 @@ where

Ok(vec![Arc::new(state_values)])
}
fn reserve_groups(&mut self, capacity: usize) {
self.values
.reserve_exact(capacity.saturating_sub(self.values.len()));
}

fn size(&self) -> usize {
self.values.capacity() * size_of::<T::Native>() + self.null_state.size()
}
Expand Down
7 changes: 7 additions & 0 deletions datafusion/functions-aggregate/src/average.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1247,6 +1247,13 @@ where

Ok(vec![Arc::new(counts) as ArrayRef, Arc::new(sums)])
}
fn reserve_groups(&mut self, capacity: usize) {
self.counts
.reserve_exact(capacity.saturating_sub(self.counts.len()));
self.sums
.reserve_exact(capacity.saturating_sub(self.sums.len()));
}

fn size(&self) -> usize {
// Heap buffers
self.counts.capacity() * size_of::<u64>()
Expand Down
5 changes: 5 additions & 0 deletions datafusion/functions-aggregate/src/count.rs
Original file line number Diff line number Diff line change
Expand Up @@ -823,6 +823,11 @@ impl GroupsAccumulator for CountGroupsAccumulator {

Ok(vec![state_array])
}
fn reserve_groups(&mut self, capacity: usize) {
self.counts
.reserve_exact(capacity.saturating_sub(self.counts.len()));
}

fn size(&self) -> usize {
self.counts.heap_size(&mut DFHeapSizeCtx::default())
}
Expand Down
75 changes: 63 additions & 12 deletions datafusion/physical-expr-common/src/binary_view_map.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,8 +30,12 @@ use datafusion_common::hash_utils::create_hashes;
use datafusion_common::utils::proxy::VecAllocExt;
use datafusion_common::{Result, exec_err};
use std::fmt::Debug;
use std::mem;
use std::sync::Arc;

const BYTE_VIEW_STARTING_BLOCK_SIZE: usize = 8 * 1024;
const BYTE_VIEW_MAX_BLOCK_SIZE: usize = 2 * 1024 * 1024;

/// HashSet optimized for storing string or binary values that can produce that
/// the final set as a `GenericBinaryViewArray` with minimal copies.
#[derive(Debug)]
Expand Down Expand Up @@ -125,9 +129,6 @@ impl ArrowBytesViewSet {
/// This map is used by the special `COUNT DISTINCT` aggregate function to
/// store the distinct values, and by the `GROUP BY` operator to store
/// group values when they are a single string array.
/// Max size of the in-progress buffer before flushing to completed buffers
const BYTE_VIEW_MAX_BLOCK_SIZE: usize = 2 * 1024 * 1024;

pub struct ArrowBytesViewMap<V>
where
V: Debug + PartialEq + Eq + Clone + Copy + Default,
Expand All @@ -147,6 +148,9 @@ where
in_progress: Vec<u8>,
/// Completed buffers containing string data
completed: Vec<Buffer>,
/// Allocation target, doubled for each new payload block up to 2 MiB.
/// Zero keeps lazy payload growth for the small per-group distinct maps.
block_size: usize,

/// random state used to generate hashes
random_state: RandomState,
Expand Down Expand Up @@ -191,18 +195,51 @@ where
views: Vec::new(),
in_progress: Vec::new(),
completed: Vec::new(),
block_size: if map_capacity == 0 {
0
} else {
BYTE_VIEW_STARTING_BLOCK_SIZE
},
random_state: RandomState::default(),
hashes_buffer: vec![],
null: None,
}
}

/// Reserves group slots without preallocating variable-length string payloads.
pub fn reserve_groups(&mut self, capacity: usize) {
if self.block_size == 0 && capacity != 0 {
self.block_size = BYTE_VIEW_STARTING_BLOCK_SIZE;
}
self.map
.reserve(capacity.saturating_sub(self.map.len()), |entry| entry.hash);
self.views
.reserve_exact(capacity.saturating_sub(self.views.len()));
self.hashes_buffer
.reserve_exact(capacity.saturating_sub(self.hashes_buffer.len()));
}

/// Emits keys, retains the hash table and scratch buffer, and reserves fixed
/// view capacity. Payload blocks restart their exponential allocation growth.
pub fn take_with_capacity(&mut self, capacity: usize) -> ArrayRef {
let mut outgoing = Self::new(self.output_type);
mem::swap(self, &mut outgoing);
mem::swap(&mut self.map, &mut outgoing.map);
self.map.clear();
mem::swap(&mut self.hashes_buffer, &mut outgoing.hashes_buffer);
self.hashes_buffer.clear();
self.initial_map_capacity = outgoing.initial_map_capacity;
mem::swap(&mut self.random_state, &mut outgoing.random_state);
self.reserve_groups(capacity);
outgoing.into_state()
}

/// Return the contents of this map and replace it with a new empty map with
/// the same output type
pub fn take(&mut self) -> Self {
let mut new_self =
Self::with_capacity(self.output_type, self.initial_map_capacity);
std::mem::swap(self, &mut new_self);
mem::swap(self, &mut new_self);
new_self
}

Expand Down Expand Up @@ -445,7 +482,7 @@ where
pub fn into_state(mut self) -> ArrayRef {
// Flush any remaining in-progress buffer
if !self.in_progress.is_empty() {
let flushed = std::mem::take(&mut self.in_progress);
let flushed = mem::take(&mut self.in_progress);
self.completed.push(Buffer::from_vec(flushed));
}

Expand Down Expand Up @@ -537,13 +574,23 @@ where
let view = if len <= 12 {
make_view(value, 0, 0)
} else {
// Ensure buffer is big enough
if self.in_progress.len() + len > BYTE_VIEW_MAX_BLOCK_SIZE {
let flushed = std::mem::replace(
&mut self.in_progress,
Vec::with_capacity(BYTE_VIEW_MAX_BLOCK_SIZE),
);
self.completed.push(Buffer::from_vec(flushed));
let block_limit = if self.block_size == 0 {
BYTE_VIEW_MAX_BLOCK_SIZE
} else {
self.in_progress.capacity()
};
if self.in_progress.len() + len > block_limit {
if !self.in_progress.is_empty() {
let flushed = mem::take(&mut self.in_progress);
self.completed.push(Buffer::from_vec(flushed));
}
let capacity = if self.block_size == 0 {
BYTE_VIEW_MAX_BLOCK_SIZE
} else {
self.block_size = (self.block_size * 2).min(BYTE_VIEW_MAX_BLOCK_SIZE);
len.max(self.block_size)
};
self.in_progress = Vec::with_capacity(capacity);
}

let buffer_index = self.completed.len() as u32;
Expand Down Expand Up @@ -947,6 +994,10 @@ mod tests {
map.views.reserve_exact(1);
map.completed.shrink_to_fit();
map.completed.reserve_exact(1);
let completed = &mut map.completed[0];
let mut padded = Vec::with_capacity(completed.len() + 16);
padded.extend_from_slice(completed);
*completed = Buffer::from_vec(padded);

// The map owns these allocations; `values` and its Arrow buffers remain external.
assert!(map.views.capacity() > map.views.len());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -372,6 +372,27 @@ impl<AggrMode> AggregateHashTable<AggrMode> {
/// spilling without finalizing the same group more than once.
pub(in crate::aggregates) fn take_state_batch(
&mut self,
) -> Result<Option<RecordBatch>> {
self.take_state_batch_with_capacity(0)
}

/// Reserves fixed row capacity without changing the number of groups.
pub(in crate::aggregates) fn reserve_groups(&mut self, capacity: usize) {
let state = self.state.building_mut();
state.group_values.reserve_groups(capacity);
state
.batch_group_indices
.reserve_exact(capacity.saturating_sub(state.batch_group_indices.len()));
for acc in &mut state.accumulators {
acc.accumulator.reserve_groups(capacity);
}
}

/// Emits all states and optionally reserves a fixed capacity for the next table.
/// A zero capacity uses the normal memory-releasing path.
pub(in crate::aggregates) fn take_state_batch_with_capacity(
&mut self,
capacity: usize,
) -> Result<Option<RecordBatch>> {
let state_schema = Arc::clone(&self.state_schema);
let accumulator_metrics = Arc::clone(&self.aggregate_accumulator_metrics);
Expand All @@ -382,7 +403,11 @@ impl<AggrMode> AggregateHashTable<AggrMode> {
}

let output = group_by_metrics.time_emitting(|| {
let mut output = state.group_values.emit(EmitTo::All)?;
let mut output = if capacity > 0 {
state.group_values.emit_with_capacity(capacity)?
} else {
state.group_values.emit(EmitTo::All)?
};
for (idx, acc) in state.accumulators.iter_mut().enumerate() {
output.extend(accumulator_metrics.time(
idx,
Expand All @@ -399,9 +424,15 @@ impl<AggrMode> AggregateHashTable<AggrMode> {
// `emit(EmitTo::All)` resets accumulator state. Explicitly shrink the
// key/index buffers too so the memory reservation can be released
// before the batch is sorted for spilling.
state.group_values.clear_shrink(0);
state.batch_group_indices.clear();
state.batch_group_indices.shrink_to_fit();
if capacity == 0 {
state.group_values.clear_shrink(0);
state.batch_group_indices.shrink_to_fit();
} else {
for acc in &mut state.accumulators {
acc.accumulator.reserve_groups(capacity);
}
}

Ok(Some(batch))
}
Expand Down
12 changes: 12 additions & 0 deletions datafusion/physical-plan/src/aggregates/group_values/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,18 @@ pub trait GroupValues: Send {
/// Emits the group values
fn emit(&mut self, emit_to: EmitTo) -> Result<Vec<ArrayRef>>;

/// Reserves room for at least `capacity` groups in supported storage layouts.
/// Does not change the number of groups or the string payload growth policy.
fn reserve_groups(&mut self, _capacity: usize) {}

/// Emits all keys and reserves a fixed capacity for the next table.
fn emit_with_capacity(&mut self, capacity: usize) -> Result<Vec<ArrayRef>> {
let output = self.emit(EmitTo::All)?;
self.clear_shrink(capacity);
self.reserve_groups(capacity);
Ok(output)
}

/// Materializes selected group values without changing the stored values or
/// their group indices.
///
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,9 +29,10 @@ use datafusion_common::Result;
use datafusion_common::utils::split_vec_min_alloc;
use datafusion_expr::GroupSelection;
use std::marker::PhantomData;
use std::mem::{replace, size_of};
use std::mem::{self, replace, size_of};
use std::sync::Arc;

const BYTE_VIEW_STARTING_BLOCK_SIZE: usize = 8 * 1024;
const BYTE_VIEW_MAX_BLOCK_SIZE: usize = 2 * 1024 * 1024;

/// An implementation of [`GroupColumn`] for binary view and utf8 view types.
Expand All @@ -55,19 +56,16 @@ pub struct ByteViewGroupValueBuilder<B: ByteViewType> {
/// The progressing block
///
/// New values will be inserted into it until its capacity
/// is not enough(detail can see `max_block_size`).
/// cannot hold the next value.
in_progress: Vec<u8>,

/// The completed blocks
completed: Vec<Buffer>,

/// The max size of `in_progress`
///
/// `in_progress` will be flushed into `completed`, and create new `in_progress`
/// when found its remaining capacity(`max_block_size` - `len(in_progress)`),
/// is no enough to store the appended value.
///
/// Currently it is fixed at 2MB.
/// Allocation target, doubled for each new payload block.
block_size: usize,

/// Maximum allocation target. A single larger value gets its own block.
max_block_size: usize,

/// Nulls
Expand All @@ -89,6 +87,7 @@ impl<B: ByteViewType> ByteViewGroupValueBuilder<B> {
views: Vec::new(),
in_progress: Vec::new(),
completed: Vec::new(),
block_size: BYTE_VIEW_STARTING_BLOCK_SIZE,
max_block_size: BYTE_VIEW_MAX_BLOCK_SIZE,
nulls: NullBufferBuilder::empty(),
_phantom: PhantomData {},
Expand All @@ -98,6 +97,7 @@ impl<B: ByteViewType> ByteViewGroupValueBuilder<B> {
/// Set the max block size
fn with_max_block_size(mut self, max_block_size: usize) -> Self {
self.max_block_size = max_block_size;
self.block_size = self.block_size.min(max_block_size);
self
}

Expand Down Expand Up @@ -236,16 +236,12 @@ impl<B: ByteViewType> ByteViewGroupValueBuilder<B> {

fn ensure_in_progress_big_enough(&mut self, value_len: usize) {
debug_assert!(value_len > 12);
let require_cap = self.in_progress.len() + value_len;

// If current block isn't big enough, flush it and create a new in progress block
if require_cap > self.max_block_size {
let flushed_block = replace(
&mut self.in_progress,
Vec::with_capacity(self.max_block_size),
);
let buffer = Buffer::from_vec(flushed_block);
self.completed.push(buffer);
if value_len > self.in_progress.capacity() - self.in_progress.len() {
if !self.in_progress.is_empty() {
self.flush_in_progress();
}
self.block_size = (self.block_size * 2).min(self.max_block_size);
self.in_progress = Vec::with_capacity(value_len.max(self.block_size));
}
}

Expand Down Expand Up @@ -564,10 +560,7 @@ impl<B: ByteViewType> ByteViewGroupValueBuilder<B> {
}

fn flush_in_progress(&mut self) {
let flushed_block = replace(
&mut self.in_progress,
Vec::with_capacity(self.max_block_size),
);
let flushed_block = mem::take(&mut self.in_progress);
let buffer = Buffer::from_vec(flushed_block);
self.completed.push(buffer);
}
Expand Down Expand Up @@ -626,6 +619,14 @@ impl<B: ByteViewType> GroupColumn for ByteViewGroupValueBuilder<B> {
self.vectorized_append_inner(array, rows)
}

fn reserve_groups(&mut self, capacity: usize) {
self.views
.reserve_exact(capacity.saturating_sub(self.views.len()));
if self.nulls.is_empty() {
self.nulls = NullBufferBuilder::new(capacity);
}
}

fn len(&self) -> usize {
self.views.len()
}
Expand Down
Loading
Loading