jayzhan211 commented on code in PR #25877: URL: https://github.com/apache/datafusion/pull/25877#discussion_r4141161831
########## datafusion/physical-plan/src/aggregates/hash_stream.rs: ########## Review Comment: the fallback `try_resize(hash_table_mem_size)` drops `pending_memory` for blocks that are still held, so memory is under-reported until the next loop iteration ########## datafusion/physical-plan/src/aggregates/aggregate_hash_table/storage.rs: ########## @@ -0,0 +1,575 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Flat or blocked storage for the group keys and accumulator states of an +//! aggregate table. +//! +//! Every part of a table chooses its storage on its own, once, when the table +//! is created: the group keys use [`BlockedGroupValues`] when their type +//! supports it (see [`use_blocked_keys`]), and each aggregate uses a +//! [`BlockedGroupsAccumulator`] when it has one. Everything else keeps today's +//! flat [`GroupValues`] and [`GroupsAccumulator`], unchanged, so any mix of +//! flat and blocked parts works, and a table with no blocked part runs exactly Review Comment: It says a table with no blocked part runs exactly the flat code, but the keys count as a blocked part for every primitive key, so `sum`-only queries change storage too. ########## datafusion/physical-plan/src/aggregates/spill.rs: ########## @@ -248,6 +248,22 @@ impl AggregateSpill { Ok(()) } + /// [`Self::sort_and_spill`] for the state batches of a hash aggregate + /// table: one batch for flat storage, one per block for blocked storage, + /// each written as its own spill file. + pub(super) fn sort_and_spill_batches( + &mut self, + state_batches: Vec<MaterializedBatch>, + ) -> Result<()> { + // TODO: sort across multiple arrays without concat (see #24928 Review Comment: ### Spilling writes one run per block, which regresses existing queries `sort_and_spill_batches` (`spill.rs:254`) calls `sort_and_spill` once per block. Keys become blocked for any single-column primitive `GROUP BY`, so this also affects queries with no blocked accumulator. For `SELECT k, sum(v) ... GROUP BY k` with `--memory-limit 200M` and `target_partitions=2`: | | spill files | elapsed | |---|---|---| | main | 8 | 0.78–0.81s | | this PR | 61–65 | 0.86–0.91s | | this PR + change below | 8 | 0.79–0.80s | (`count` at 120M: 15–16 → 71–72 files, 0.77–0.86s → 0.85–0.92s.) I don't think the TODO needs to wait for #24928. Sorting only the key columns across blocks and gathering with `interleave_record_batch` gives one sorted run, concatenates only the sort keys, and returned the same results as main in my runs: ```rust pub(super) fn sort_and_spill_batches( &mut self, state_batches: Vec<MaterializedBatch>, ) -> Result<()> { let mut batches: Vec<RecordBatch> = state_batches.into_iter().map(|b| b.batch).collect(); if batches.len() <= 1 { return self.sort_and_spill(batches.pop()); } // Sort the key columns of all blocks together, then gather each // output batch from the blocks, so all blocks form one sorted run let sort_columns = self .spill_expr .iter() .map(|expr| { let arrays = batches .iter() .map(|b| expr.expr.evaluate(b)?.into_array(b.num_rows())) .collect::<Result<Vec<_>>>()?; let arrays: Vec<&dyn Array> = arrays.iter().map(|a| a.as_ref()).collect(); Ok(SortColumn { values: concat(&arrays)?, options: Some(expr.options), }) }) .collect::<Result<Vec<_>>>()?; let sorted = lexsort_to_indices(&sort_columns, None)?; drop(sort_columns); // Every block but the last holds the same number of rows let block_rows = batches[0].num_rows(); let batch_refs: Vec<&RecordBatch> = batches.iter().collect(); let max_batch_rows = sorted.len().min(self.batch_size); let sorted_iter = sorted.values().chunks(self.batch_size).map(|chunk| { let indices: Vec<(usize, usize)> = chunk .iter() .map(|&i| (i as usize / block_rows, i as usize % block_rows)) .collect(); Ok(interleave_record_batch(&batch_refs, &indices)?) }); let spill_file = self .spill_manager .spill_record_batch_iter_and_return_max_batch_memory( sorted_iter, self.label, )?; let Some((file, max_record_batch_memory)) = spill_file else { return internal_err!("{}: produced an empty spill", self.label); }; self.spills.push(SortedSpillFile { file, max_record_batch_memory, }); self.min_spill_batch_rows = self.min_spill_batch_rows.min(max_batch_rows); Ok(()) } ``` ########## datafusion/functions-aggregate-common/src/aggregate/groups_accumulator/blocked_vec.rs: ########## @@ -0,0 +1,676 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! [`BlockedVec`]: group state stored in fixed size blocks + +use std::mem::size_of; + +use arrow::array::BooleanArray; +use arrow::buffer::NullBuffer; + +use datafusion_expr_common::blocked_groups_accumulator::BlocksIndex; + +use super::accumulate::accumulate_blocked_indices; + +/// Rows per chunk whose element addresses are resolved before any of them is +/// updated, see [`BlockedVec::update`]. +const RESOLVE_CHUNK: usize = 256; + +/// A growable vector of per group state stored as a list of blocks, used by +/// [`BlockedGroupsAccumulator`] implementations. +/// +/// It stays a single flat `Vec` (growing by doubling) until it reaches +/// `block_size` elements. That `Vec` then becomes block 0 without any copy. +/// Every following block also grows by doubling, up to `block_size`, so a +/// block that only holds a few groups does not allocate a whole block. +/// Accumulators update their state with [`Self::update`] and +/// [`Self::update_with`], which pick the fastest loop for the current layout +/// once per batch. +/// +/// Emitting the first block ([`Self::take_next_block`]) or all blocks +/// ([`Self::take_all`]) moves the block `Vec`s out without copying. +/// +/// `FIXED_BLOCK_SIZE = false` is reserved for the values of nested children +/// (e.g. list elements), where block boundaries are driven by the parent +/// through a future `start_new_block()`. It is not implemented yet, but the +/// layout (one `Vec` per block, the block length being the `Vec` length) +/// already supports it. +/// +/// # Implementation Notes +/// +/// ## Why `T: Copy`? +/// +/// 1. So the [`BlockedVec::allocated_size`] will be accurate since the size of T is known, and not include heap allocations (like `String` or `Vec`) that are not part of the allocated size of the `BlockedVec` +/// 2. So we can provide mutable access to the items (e.g. `IndexMut`) since if `T` is not `Copy` (like when `T` is a `Vec`) we could change the size of it without the [`BlockedVec::allocated_size`] changing, which would be confusing and lead to bugs +/// +/// [`BlockedGroupsAccumulator`]: datafusion_expr_common::blocked_groups_accumulator::BlockedGroupsAccumulator +#[derive(Debug)] +pub struct BlockedVec<T, const FIXED_BLOCK_SIZE: bool = true> { + /// Every block except the last one holds exactly `block_size` elements + blocks: Vec<Vec<T>>, + /// Total number of elements over all blocks + len: usize, + block_size: usize, + /// Running total of the bytes allocated by the blocks' buffers + allocated: usize, +} + +impl<T: Copy, const FIXED_BLOCK_SIZE: bool> BlockedVec<T, FIXED_BLOCK_SIZE> { + /// Creates an empty vector. + /// + /// # Panics + /// If `block_size` is 0. + pub fn new(block_size: usize) -> Self { + if !FIXED_BLOCK_SIZE { + // Reserved for nested child values, whose block sizes are driven + // by the parent via a future `start_new_block()`. + unimplemented!("BlockedVec with dynamic block sizes"); + } + assert!(block_size > 0, "block_size must be positive"); + Self { + blocks: Vec::new(), + len: 0, + block_size, + allocated: 0, + } + } + + /// Maximum number of elements in a block. + pub fn block_size(&self) -> usize { + self.block_size + } + + /// Total number of elements. + #[inline] + pub fn len(&self) -> usize { + self.len + } + + /// Returns `true` if there are no elements. + #[inline] + pub fn is_empty(&self) -> bool { + self.len == 0 + } + + /// Number of blocks, including a partially filled last block. + #[inline] + pub fn num_blocks(&self) -> usize { + self.blocks.len() + } + + /// Grows the vector to `new_len` elements, filling with `value`. + /// Never shrinks: does nothing if `new_len <= self.len()`. + pub fn grow_to(&mut self, new_len: usize, value: T) { + while self.len < new_len { + let last_len = self.last_block_for_append(); + let additional = (new_len - self.len).min(self.block_size - last_len); + let last = self.reserve_last(last_len + additional); + last.resize(last_len + additional, value); + self.len += additional; + } + } + + /// Appends `value` and returns its index. + #[inline] + pub fn push(&mut self, value: T) -> BlocksIndex { + self.len += 1; + let num_blocks = self.blocks.len(); + match self.blocks.last_mut() { + // Block capacities never exceed `block_size`, so this never + // reallocates + Some(last) if last.len() < last.capacity() => { + last.push(value); + BlocksIndex::new(num_blocks - 1, last.len() - 1) + } + _ => self.push_slow(value), + } + } + + #[cold] + fn push_slow(&mut self, value: T) -> BlocksIndex { + let last_len = self.last_block_for_append(); + self.reserve_last(last_len + 1).push(value); + BlocksIndex::new(self.blocks.len() - 1, last_len) + } + + /// Makes sure the last block has room for at least one more element, + /// adding a new block if needed, and returns its length. + fn last_block_for_append(&mut self) -> usize { + match self.blocks.last() { + Some(last) if last.len() < self.block_size => last.len(), + // Every block, including block 0 (the flat phase), grows by + // doubling in `reserve_last` + _ => { + self.blocks.push(Vec::new()); + 0 + } + } + } + + /// Makes sure the last block can hold `required <= block_size` elements, + /// growing it by doubling (capped at `block_size`), and returns it. + fn reserve_last(&mut self, required: usize) -> &mut Vec<T> { + let block_size = self.block_size; + let last = self.blocks.last_mut().expect("at least one block"); + let old_capacity = last.capacity(); + if old_capacity < required { + let new_capacity = (old_capacity * 2).max(required).min(block_size); + last.reserve_exact(new_capacity - last.len()); + self.allocated += (last.capacity() - old_capacity) * size_of::<T>(); + } + last + } + + /// Returns the element at `index`. + /// + /// # Panics + /// If `index` is out of bounds. + #[inline] + pub fn get(&self, index: BlocksIndex) -> T { + self.blocks[index.block_index()][index.index_in_block()] + } + + /// Returns a mutable reference to the element at `index`. + /// + /// # Safety + /// `index` must point to an element of this vector. + #[inline] + pub unsafe fn get_unchecked_mut(&mut self, index: BlocksIndex) -> &mut T { + // SAFETY: guaranteed by the caller + unsafe { + self.blocks + .get_unchecked_mut(index.block_index()) + .get_unchecked_mut(index.index_in_block()) + } + } + + /// Returns all elements as one slice if there is at most one block, so + /// hot loops can run on a flat slice exactly like today's accumulators. + #[inline] + pub fn as_single_block_mut(&mut self) -> Option<&mut [T]> { + match self.blocks.as_mut_slice() { + [] => Some(&mut []), + [block] => Some(block.as_mut_slice()), + _ => None, + } + } + + /// Returns the address of every block, for resolving many indices + /// without going through the outer `Vec` each time. + /// + /// The returned [`BlockPtrs`] is invalidated by any later use of `self`. + fn block_ptrs_mut(&mut self) -> BlockPtrs<T> { + BlockPtrs { + ptrs: self.blocks.iter_mut().map(|b| b.as_mut_ptr()).collect(), + } + } + + /// Grows the vector to `total_num_groups` elements (new ones are + /// `starting_value`), then calls `update_fn` on the element of every + /// row's group, skipping rows that are null in `nulls` or not selected by + /// `opt_filter`. + /// + /// Chooses the loop once per call: with a single block it is a plain flat + /// loop; with several blocks and no nulls or filter, the addresses of each + /// chunk of rows are resolved before any of them is updated, so the cache + /// misses of the updates don't wait on each other. + /// + /// `group_indices` must only point to the first `total_num_groups` + /// groups, which is the contract of + /// [`BlockedGroupsAccumulator::update_batch`] and + /// [`BlockedGroupsAccumulator::merge_batch`]. It is only checked in debug + /// builds, like the flat accumulators' unchecked indexing. + /// + /// [`BlockedGroupsAccumulator::update_batch`]: datafusion_expr_common::blocked_groups_accumulator::BlockedGroupsAccumulator::update_batch + /// [`BlockedGroupsAccumulator::merge_batch`]: datafusion_expr_common::blocked_groups_accumulator::BlockedGroupsAccumulator::merge_batch + pub fn update<F>( Review Comment: Should we consider using `unsafe` and skip bounds check -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
