rluvaton commented on code in PR #25877: URL: https://github.com/apache/datafusion/pull/25877#discussion_r4145723601
########## 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: added unsafe methods as well and used them when the original implementation skipped bound checks -- 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]
