rluvaton commented on code in PR #25877: URL: https://github.com/apache/datafusion/pull/25877#discussion_r4134755296
########## datafusion/expr-common/src/blocked_groups_accumulator.rs: ########## @@ -0,0 +1,267 @@ +// 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. + +//! Vectorized [`BlockedGroupsAccumulator`]: [`GroupsAccumulator`] with state +//! stored in blocks. +//! +//! Storing per-group state in fixed-size blocks lets hash aggregation emit and +//! free its state one block at a time instead of materializing every group +//! into one large array. See <https://github.com/apache/datafusion/issues/24704>. +//! +//! [`GroupsAccumulator`]: crate::groups_accumulator::GroupsAccumulator + +use std::any::Any; +use std::sync::Arc; + +use arrow::array::{ArrayRef, BooleanArray}; +use datafusion_common::Result; + +use crate::accumulator::AggregateMetrics; +use crate::groups_accumulator::EmitTo; + +/// Identifies one group in blocked storage: the block that holds the group +/// and the group's position inside that block. +/// +/// This is the only place that knows how a group index is laid out; storage +/// and accumulators only use the accessors below. +#[repr(transparent)] +#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, Default)] +pub struct BlocksIndex(u64); + +impl BlocksIndex { + /// Creates the index of the `index_in_block`-th group of block + /// `block_index`. + #[inline] + pub const fn new(block_index: usize, index_in_block: usize) -> Self { + debug_assert!(block_index <= u32::MAX as usize); + debug_assert!(index_in_block <= u32::MAX as usize); + Self(((block_index as u64) << 32) | index_in_block as u64) + } + + /// Block that holds this group. + #[inline] + pub const fn block_index(self) -> usize { + (self.0 >> 32) as usize + } + + /// Position of this group inside its block. + #[inline] + pub const fn index_in_block(self) -> usize { + self.0 as u32 as usize + } + + /// Creates the index of the group at flat index `flat`, counting from the + /// first group, when every block holds `block_size` groups. + #[inline] + pub const fn from_flat(flat: usize, block_size: usize) -> Self { + Self::new(flat / block_size, flat % block_size) + } + + /// Flat index of this group, counting from the first group, when every + /// block holds `block_size` groups. + #[inline] + pub const fn flat(self, block_size: usize) -> usize { + self.block_index() * block_size + self.index_in_block() + } + + /// Writes [`Self::from_flat`] of every index in `flat` to `out`, + /// replacing its contents. + pub fn from_flat_slice(flat: &[usize], block_size: usize, out: &mut Vec<Self>) { + out.clear(); + if block_size.is_power_of_two() { + let shift = block_size.trailing_zeros(); + let mask = block_size - 1; + out.extend(flat.iter().map(|&f| Self::new(f >> shift, f & mask))); + } else { + out.extend(flat.iter().map(|&f| Self::from_flat(f, block_size))); + } + } + + /// Writes [`Self::flat`] of every index in `indices` to `out`, replacing + /// its contents. + pub fn to_flat_slice(indices: &[Self], block_size: usize, out: &mut Vec<usize>) { + out.clear(); + if block_size.is_power_of_two() { + let shift = block_size.trailing_zeros(); + out.extend( + indices + .iter() + .map(|i| (i.block_index() << shift) | i.index_in_block()), + ); + } else { + out.extend(indices.iter().map(|i| i.flat(block_size))); + } + } +} + +/// Which groups to emit from a [`BlockedGroupsAccumulator`]. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum BlockedEmitTo { + /// Emit every group, as one array per block, and reset the state. + All, + /// Emit the first block (or nothing if there are no groups). + /// + /// The remaining groups move down by one block, like [`EmitTo::First`] + /// with the first block's length. To emit every group, use [`Self::All`]. + NextBlock, + /// Emit the first `n` groups, where `0 < n < block_size`, and shift the + /// remaining groups down by `n`, like [`EmitTo::First`]. + First(usize), +} + +impl BlockedEmitTo { + /// Splits a flat [`EmitTo`] into block-sized emits. + /// + /// `EmitTo::First(n)` becomes `n / block_size` [`Self::NextBlock`]s + /// followed by one [`Self::First`] with the remainder, if any. Each emit + /// shifts the remaining groups, so they must all be applied, in order, + /// before any new group is added. + pub fn from_emit_to(emit_to: EmitTo, block_size: usize) -> Vec<Self> { + match emit_to { + EmitTo::All => vec![Self::All], + EmitTo::First(n) => { + let mut emits = vec![Self::NextBlock; n / block_size]; + if !n.is_multiple_of(block_size) { + emits.push(Self::First(n % block_size)); + } + emits + } + } + } +} + +/// Like [`GroupsAccumulator`], but the per-group state is stored in blocks of +/// [`Self::block_size`] groups, and emitting returns one array per block +/// instead of one large array. +/// +/// Hash aggregation uses a `BlockedGroupsAccumulator` only when every +/// aggregate in the query (and the group keys) support blocked storage; +/// otherwise the whole aggregation uses [`GroupsAccumulator`]. +/// +/// [`GroupsAccumulator`]: crate::groups_accumulator::GroupsAccumulator +pub trait BlockedGroupsAccumulator: Send + Any { Review Comment: Will add the rest of the functions (evaluate/state preserving) in later pr so it will be easier to review -- 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]
