Skip to content

Commit db6903c

Browse files
committed
add list and row backed blocked group values in multi group by
1 parent 5d6018b commit db6903c

7 files changed

Lines changed: 2310 additions & 95 deletions

File tree

‎datafusion/physical-plan/src/aggregates_blocked/group_values/multi_group_by/boolean.rs‎

Lines changed: 6 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -84,7 +84,7 @@ impl<const IS_FIXED_BLOCK: bool, const NULLABLE: bool> BlockedGroupColumn<IS_FIX
8484
}
8585

8686
fn vectorized_equal_to(
87-
&self,
87+
&mut self,
8888
lhs_rows: &[BlocksIndex],
8989
array: &ArrayRef,
9090
rhs_rows: &[usize],
@@ -219,18 +219,16 @@ impl<const IS_FIXED_BLOCK: bool, const NULLABLE: bool> BlockedGroupColumn<IS_FIX
219219
}
220220
}
221221

222-
fn take_n(&mut self, n: usize,
223-
// adjusted_block_size_iter: Option<Box<dyn ClonableIter<Item=usize>>>,
224-
) -> ArrayRef {
222+
fn take_n(&mut self, n: usize, adjusted_block_size: Option<&[usize]>) -> ArrayRef {
223+
assert_eq!(adjusted_block_size.is_none(), IS_FIXED_BLOCK);
224+
225225
let first_n_nulls = if NULLABLE { self.nulls.take_n(
226226
n,
227-
// adjusted_block_size_iter.clone()
228-
None::<std::iter::Empty<_>>,
227+
adjusted_block_size.map(|s| s.iter().copied()),
229228
) } else { None };
230229
let first_n_values = self.buffer.take_n(
231230
n,
232-
// adjusted_block_size_iter
233-
None::<std::iter::Empty<_>>,
231+
adjusted_block_size.map(|s| s.iter().copied()),
234232
);
235233

236234
Arc::new(BooleanArray::new(first_n_values, first_n_nulls))

‎datafusion/physical-plan/src/aggregates_blocked/group_values/multi_group_by/bytes.rs‎

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -269,7 +269,7 @@ where
269269
}
270270

271271
fn vectorized_equal_to(
272-
&self,
272+
&mut self,
273273
lhs_rows: &[BlocksIndex],
274274
array: &ArrayRef,
275275
rhs_rows: &[usize],
@@ -384,11 +384,13 @@ where
384384
Some(Self::build_array(self.output_type, data))
385385
}
386386

387-
fn take_n(&mut self, n: usize) -> ArrayRef {
387+
fn take_n(&mut self, n: usize, adjusted_block_size: Option<&[usize]>) -> ArrayRef {
388+
assert_eq!(adjusted_block_size.is_none(), IS_FIXED_BLOCK);
389+
388390
debug_assert!(self.len() >= n);
389391
// SAFETY: the offsets were constructed correctly
390392

391-
let data = unsafe { self.data.take_n_unchecked(n, None::<std::iter::Empty<_>>) };
393+
let data = unsafe { self.data.take_n_unchecked(n, adjusted_block_size.map(|s| s.iter().copied())) };
392394

393395
Self::build_array(self.output_type, data)
394396
}

‎datafusion/physical-plan/src/aggregates_blocked/group_values/multi_group_by/bytes_view.rs‎

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -433,7 +433,7 @@ impl<const FIXED_BLOCK_SIZING: bool, B: ByteViewType> BlockedGroupColumn<FIXED_B
433433
}
434434

435435
fn vectorized_equal_to(
436-
&self,
436+
&mut self,
437437
group_indices: &[BlocksIndex],
438438
array: &ArrayRef,
439439
rows: &[usize],
@@ -511,11 +511,13 @@ impl<const FIXED_BLOCK_SIZING: bool, B: ByteViewType> BlockedGroupColumn<FIXED_B
511511
/// remaining bytes are only re-blocked so that every views block still owns whole
512512
/// bytes blocks, which means rewriting the buffer index and offset of the remaining
513513
/// non inlined views
514-
fn take_n(&mut self, n: usize) -> ArrayRef {
514+
fn take_n(&mut self, n: usize, adjusted_block_size: Option<&[usize]>) -> ArrayRef {
515+
assert_eq!(adjusted_block_size.is_none(), FIXED_BLOCK_SIZING);
516+
515517
debug_assert!(self.len() >= n);
516518
let block_size = self.views.block_size();
517-
let taken_views = self.views.take_n(n, None::<std::iter::Empty<_>>);
518-
let nulls = self.nulls.take_n(n, None::<std::iter::Empty<_>>);
519+
let taken_views = self.views.take_n(n, adjusted_block_size.map(|s| s.iter().copied()));
520+
let nulls = self.nulls.take_n(n, adjusted_block_size.map(|s| s.iter().copied()));
519521

520522
// Bytes are appended in value order, so the taken values own a prefix of the
521523
// bytes ending where the last non inlined taken value ends, in bytes block `b` of
@@ -829,22 +831,22 @@ mod tests {
829831
assert_eq!(stored(&builder, 4), owned(&values[4..]));
830832

831833
// taking values releases exactly their bytes and keeps the rest addressable
832-
let taken = builder.take_n(1);
834+
let taken = builder.take_n(1, None);
833835
assert_eq!(strings(&taken), owned(&values[4..5]));
834836
assert_eq!(taken.as_string_view().data_buffers()[0].len(), value_len);
835837
assert_eq!(builder.bytes.len(), 7 * value_len);
836838
// the old bytes block boundary after value 5 is kept, value 8 moved into block 0
837839
assert_eq!(builder.num_bytes_blocks_per_block, Vec::from(vec![3, 2]));
838840
assert_eq!(stored(&builder, 4), owned(&values[5..]));
839-
let taken = builder.take_n(1);
841+
let taken = builder.take_n(1, None);
840842
assert_eq!(strings(&taken), owned(&values[5..6]));
841843
assert_eq!(stored(&builder, 4), owned(&values[6..]));
842844
for (row, value) in values.iter().enumerate().skip(6) {
843845
let index = BlocksIndex::from_index_in_fixed_block_size(row - 6, 4);
844846
assert!(builder.equal_to(index, &input, row), "{value:?}");
845847
}
846848

847-
let taken = builder.take_n(3);
849+
let taken = builder.take_n(3, None);
848850
assert_eq!(strings(&taken), owned(&values[6..9]));
849851
let blocks = Box::new(builder).take_all();
850852
let all: Vec<Option<String>> = blocks.iter().flat_map(strings).collect();
@@ -866,7 +868,7 @@ mod tests {
866868
let values = sample();
867869
let mut builder = fixed_with(4, &values);
868870

869-
let taken = builder.take_n(3);
871+
let taken = builder.take_n(3, None);
870872
assert_eq!(strings(&taken), owned(&values[..3]));
871873
assert_eq!(builder.len(), 6);
872874
assert_eq!(stored(&builder, 4), owned(&values[3..]));
@@ -877,7 +879,7 @@ mod tests {
877879
assert!(!builder.equal_to(BlocksIndex::from_index_in_fixed_block_size(2, 4), &input, 1));
878880
builder.append_val(&input, 0).unwrap();
879881

880-
let taken = builder.take_n(0);
882+
let taken = builder.take_n(0, None);
881883
assert_eq!(taken.len(), 0);
882884

883885
let blocks = Box::new(builder).take_all();

0 commit comments

Comments
 (0)