From 884542d97ba27fa463287e39e978b7b57a1a598b Mon Sep 17 00:00:00 2001 From: Daniel King Date: Thu, 16 Jul 2026 15:03:02 -0400 Subject: [PATCH 01/18] Port PiecewiseSequence run take consumers Signed-off-by: Daniel King --- .../src/arrays/extension/compute/rules.rs | 2 + .../src/arrays/extension/compute/take.rs | 16 ++ vortex-array/src/arrays/list/compute/take.rs | 249 ++++++++++++++++++ vortex-array/src/arrays/listview/rebuild.rs | 228 ++++++++++------ 4 files changed, 410 insertions(+), 85 deletions(-) diff --git a/vortex-array/src/arrays/extension/compute/rules.rs b/vortex-array/src/arrays/extension/compute/rules.rs index 2d02d1ae7a7..0699f2f68f4 100644 --- a/vortex-array/src/arrays/extension/compute/rules.rs +++ b/vortex-array/src/arrays/extension/compute/rules.rs @@ -11,6 +11,7 @@ use crate::arrays::ConstantArray; use crate::arrays::Extension; use crate::arrays::ExtensionArray; use crate::arrays::Filter; +use crate::arrays::dict::TakeReduceAdaptor; use crate::arrays::extension::ExtensionArrayExt; use crate::arrays::filter::FilterReduceAdaptor; use crate::arrays::slice::SliceReduceAdaptor; @@ -50,6 +51,7 @@ pub(crate) const PARENT_RULES: ParentRuleSet = ParentRuleSet::new(&[ ParentRuleSet::lift(&FilterReduceAdaptor(Extension)), ParentRuleSet::lift(&MaskReduceAdaptor(Extension)), ParentRuleSet::lift(&SliceReduceAdaptor(Extension)), + ParentRuleSet::lift(&TakeReduceAdaptor(Extension)), ]); /// Push filter operations into the storage array of an extension array. diff --git a/vortex-array/src/arrays/extension/compute/take.rs b/vortex-array/src/arrays/extension/compute/take.rs index 8fb0da11222..d9ac1b4f600 100644 --- a/vortex-array/src/arrays/extension/compute/take.rs +++ b/vortex-array/src/arrays/extension/compute/take.rs @@ -10,8 +10,24 @@ use crate::array::ArrayView; use crate::arrays::Extension; use crate::arrays::ExtensionArray; use crate::arrays::dict::TakeExecute; +use crate::arrays::dict::TakeReduce; use crate::arrays::extension::ExtensionArrayExt; +impl TakeReduce for Extension { + fn take(array: ArrayView<'_, Extension>, indices: &ArrayRef) -> VortexResult> { + let taken_storage = array.storage_array().take(indices.clone())?; + Ok(Some( + ExtensionArray::new( + array + .ext_dtype() + .with_nullability(taken_storage.dtype().nullability()), + taken_storage, + ) + .into_array(), + )) + } +} + impl TakeExecute for Extension { fn take( array: ArrayView<'_, Extension>, diff --git a/vortex-array/src/arrays/list/compute/take.rs b/vortex-array/src/arrays/list/compute/take.rs index 14e89d4c238..35788796d78 100644 --- a/vortex-array/src/arrays/list/compute/take.rs +++ b/vortex-array/src/arrays/list/compute/take.rs @@ -1,26 +1,36 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors +use itertools::Itertools as _; +use vortex_buffer::BufferMut; use vortex_error::VortexExpect; use vortex_error::VortexResult; +use vortex_error::vortex_ensure; +use vortex_error::vortex_err; use crate::ArrayRef; use crate::IntoArray; use crate::array::ArrayView; use crate::arrays::List; use crate::arrays::ListArray; +use crate::arrays::PiecewiseSequence; +use crate::arrays::PiecewiseSequenceArray; use crate::arrays::Primitive; use crate::arrays::PrimitiveArray; use crate::arrays::dict::TakeExecute; use crate::arrays::list::ListArrayExt; +use crate::arrays::piecewise_sequence::execute_index_arrays; +use crate::arrays::piecewise_sequence::validate_index_ranges; use crate::arrays::primitive::PrimitiveArrayExt; use crate::builders::ArrayBuilder; use crate::builders::PrimitiveBuilder; use crate::dtype::IntegerPType; use crate::dtype::Nullability; +use crate::dtype::UnsignedPType; use crate::executor::ExecutionCtx; use crate::match_each_unsigned_integer_ptype; use crate::match_smallest_offset_type; +use crate::validity::Validity; // TODO(connor)[ListView]: Re-revert to the version where we simply convert to a `ListView` and call // the `ListView::take` compute function once `ListView` is more stable. @@ -37,6 +47,12 @@ impl TakeExecute for List { indices: &ArrayRef, ctx: &mut ExecutionCtx, ) -> VortexResult> { + if let Some(piecewise_indices) = indices.as_opt::() + && let Some(taken) = take_piecewise_sequence(array, piecewise_indices, indices, ctx)? + { + return Ok(Some(taken)); + } + let indices = indices.clone().execute::(ctx)?; let indices = indices.reinterpret_cast(indices.ptype().to_unsigned()); let offsets = array.offsets().clone().execute::(ctx)?; @@ -127,6 +143,191 @@ fn _take( .into_array()) } +fn take_piecewise_sequence( + array: ArrayView<'_, List>, + indices: ArrayView<'_, PiecewiseSequence>, + indices_ref: &ArrayRef, + ctx: &mut ExecutionCtx, +) -> VortexResult> { + let data_validity = array + .list_validity() + .execute_mask(array.as_ref().len(), ctx)?; + if !data_validity.all_true() { + return Ok(None); + } + + let (starts, lengths) = execute_index_arrays(indices, ctx)?; + let offsets = array.offsets().clone().execute::(ctx)?; + let offsets = offsets.reinterpret_cast(offsets.ptype().to_unsigned()); + + match_each_unsigned_integer_ptype!(starts.ptype(), |S| { + match_each_unsigned_integer_ptype!(lengths.ptype(), |L| { + match_each_unsigned_integer_ptype!(offsets.ptype(), |O| { + take_piecewise_sequence_typed::( + array, + starts.as_slice::(), + lengths.as_slice::(), + offsets.as_slice::(), + indices_ref, + ) + }) + }) + }) + .map(Some) +} + +fn take_piecewise_sequence_typed( + array: ArrayView<'_, List>, + starts: &[S], + lengths: &[L], + offsets: &[Offset], + indices_ref: &ArrayRef, +) -> VortexResult +where + S: UnsignedPType, + L: UnsignedPType, + Offset: UnsignedPType, +{ + validate_index_ranges(array.len(), starts, lengths, indices_ref.len())?; + let total_elements = + piecewise_list_elements_len(array.elements().len(), offsets, starts, lengths)?; + + match_smallest_offset_type!(total_elements, |OutputOffset| { + let gathered = gather_piecewise_list::( + array.elements(), + offsets, + starts, + lengths, + indices_ref.len(), + total_elements, + )?; + let validity = array.validity()?.take(indices_ref)?; + + // SAFETY: output offsets are rebuilt from valid monotonic source offsets; output elements + // are exactly the gathered child ranges referenced by those offsets; validity has one bit + // per output row. + Ok( + unsafe { ListArray::new_unchecked(gathered.elements, gathered.offsets, validity) } + .into_array(), + ) + }) +} + +struct GatheredList { + elements: ArrayRef, + offsets: ArrayRef, +} + +fn piecewise_list_elements_len( + elements_len: usize, + offsets: &[Offset], + starts: &[S], + lengths: &[L], +) -> VortexResult +where + S: UnsignedPType, + L: UnsignedPType, + Offset: UnsignedPType, +{ + let mut total = 0usize; + for (&start, &length) in starts.iter().zip_eq(lengths) { + let start: usize = start.as_(); + let length: usize = length.as_(); + let end = start + length; + if length == 0 { + continue; + } + + let element_start: usize = offsets[start].as_(); + let element_end: usize = offsets[end].as_(); + vortex_ensure!( + element_start <= element_end && element_end <= elements_len, + "List offsets range {element_start}..{element_end} exceeds elements length {elements_len}", + ); + total = total + .checked_add(element_end - element_start) + .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; + } + Ok(total) +} + +fn gather_piecewise_list( + elements: &ArrayRef, + offsets: &[Offset], + starts: &[S], + lengths: &[L], + output_len: usize, + total_elements: usize, +) -> VortexResult +where + S: UnsignedPType, + L: UnsignedPType, + Offset: UnsignedPType, + OutputOffset: IntegerPType, +{ + let offsets_capacity = output_len + .checked_add(1) + .ok_or_else(|| vortex_err!("List take offsets length overflow"))?; + let mut new_offsets = BufferMut::::with_capacity(offsets_capacity); + let mut element_starts = BufferMut::::with_capacity(starts.len()); + let mut element_lengths = BufferMut::::with_capacity(lengths.len()); + let mut output_elements = 0usize; + + new_offsets.push(OutputOffset::zero()); + for (&start, &length) in starts.iter().zip_eq(lengths) { + let start: usize = start.as_(); + let length: usize = length.as_(); + let end = start + length; + if length == 0 { + continue; + } + + let element_start: usize = offsets[start].as_(); + let element_end: usize = offsets[end].as_(); + for &offset in &offsets[start + 1..=end] { + let offset: usize = offset.as_(); + let relative = offset + .checked_sub(element_start) + .ok_or_else(|| vortex_err!("List offsets are not monotonic at offset {offset}"))?; + let output_offset = output_elements + .checked_add(relative) + .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; + new_offsets.push(new_offset_value::(output_offset)?); + } + + let element_length = element_end - element_start; + element_starts.push(element_start as u64); + element_lengths.push(element_length as u64); + output_elements = output_elements + .checked_add(element_length) + .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; + } + debug_assert_eq!(output_elements, total_elements); + + let offsets = PrimitiveArray::new(new_offsets.freeze(), Validity::NonNullable).into_array(); + // SAFETY: element ranges are derived from validated source list offsets, and total_elements is + // the sum of the gathered element range lengths. + let element_indices = unsafe { + PiecewiseSequenceArray::new_unchecked( + element_starts.into_array(), + element_lengths.into_array(), + total_elements, + ) + }; + let elements = elements.take(element_indices.into_array())?; + + Ok(GatheredList { elements, offsets }) +} + +fn new_offset_value(value: usize) -> VortexResult { + T::from_usize(value).ok_or_else(|| { + vortex_err!( + "List take offset value {value} does not fit in {}", + T::PTYPE + ) + }) +} + // Kept out-of-line: as a single-callsite generic helper it would otherwise be inlined into every // monomorphization of `_take`, duplicating the entire nullable path across all specializations. #[inline(never)] @@ -217,6 +418,7 @@ mod test { use crate::arrays::BoolArray; use crate::arrays::ListArray; use crate::arrays::ListViewArray; + use crate::arrays::PiecewiseSequenceArray; use crate::arrays::PrimitiveArray; use crate::compute::conformance::take::test_take_conformance; use crate::dtype::DType; @@ -403,6 +605,53 @@ mod test { ); } + #[test] + fn piecewise_sequence_take() { + let mut ctx = array_session().create_execution_ctx(); + let list = ListArray::try_new( + buffer![0i32, 1, 2, 3, 4, 5, 6].into_array(), + buffer![0u32, 2, 5, 5, 7].into_array(), + Validity::NonNullable, + ) + .unwrap() + .into_array(); + let idx = PiecewiseSequenceArray::try_new( + buffer![1u64, 0].into_array(), + buffer![2u64, 1].into_array(), + 3, + ) + .unwrap() + .into_array(); + + let result = list + .take(idx) + .unwrap() + .execute::(&mut ctx) + .unwrap(); + + let element_dtype: Arc = Arc::new(I32.into()); + assert_eq!( + result.execute_scalar(0, &mut ctx).unwrap(), + Scalar::list( + Arc::clone(&element_dtype), + vec![2i32.into(), 3.into(), 4.into()], + Nullability::NonNullable + ) + ); + assert_eq!( + result.execute_scalar(1, &mut ctx).unwrap(), + Scalar::list(Arc::clone(&element_dtype), vec![], Nullability::NonNullable) + ); + assert_eq!( + result.execute_scalar(2, &mut ctx).unwrap(), + Scalar::list( + element_dtype, + vec![0i32.into(), 1.into()], + Nullability::NonNullable + ) + ); + } + #[test] fn test_take_empty_array() { let list = ListArray::try_new( diff --git a/vortex-array/src/arrays/listview/rebuild.rs b/vortex-array/src/arrays/listview/rebuild.rs index 342c9eccf19..1ce8f69d12e 100644 --- a/vortex-array/src/arrays/listview/rebuild.rs +++ b/vortex-array/src/arrays/listview/rebuild.rs @@ -5,16 +5,18 @@ use num_traits::FromPrimitive; use vortex_buffer::BufferMut; use vortex_error::VortexExpect; use vortex_error::VortexResult; +use vortex_error::vortex_err; +use vortex_mask::Mask; -use crate::Canonical; use crate::ExecutionCtx; use crate::IntoArray; +use crate::RecursiveCanonical; use crate::arrays::ConstantArray; use crate::arrays::ListViewArray; +use crate::arrays::PiecewiseSequenceArray; use crate::arrays::PrimitiveArray; use crate::arrays::listview::ListViewArrayExt; use crate::arrays::primitive::PrimitiveArrayExt; -use crate::builders::builder_with_capacity; use crate::builtins::ArrayBuiltins; use crate::dtype::IntegerPType; use crate::dtype::Nullability; @@ -140,19 +142,18 @@ impl ListViewArray { // for sizes as well. match_each_unsigned_integer_ptype!(sizes_ptype.to_unsigned(), |S| { match offsets_ptype.to_unsigned() { - PType::U8 => self.naive_rebuild::(ctx), - PType::U16 => self.naive_rebuild::(ctx), - PType::U32 => self.naive_rebuild::(ctx), - PType::U64 => self.naive_rebuild::(ctx), + PType::U8 => self.rebuild_with_take_or_piecewise::(ctx), + PType::U16 => self.rebuild_with_take_or_piecewise::(ctx), + PType::U32 => self.rebuild_with_take_or_piecewise::(ctx), + PType::U64 => self.rebuild_with_take_or_piecewise::(ctx), _ => unreachable!("invalid offsets PType"), } }) } /// Picks between [`rebuild_with_take`](Self::rebuild_with_take) and - /// [`rebuild_list_by_list`](Self::rebuild_list_by_list) based on element dtype and average - /// list size. - fn naive_rebuild( + /// [`rebuild_with_piecewise`](Self::rebuild_with_piecewise) based on average list size. + fn rebuild_with_take_or_piecewise( &self, ctx: &mut ExecutionCtx, ) -> VortexResult { @@ -167,12 +168,12 @@ impl ListViewArray { if Self::should_use_take(total, self.len()) { self.rebuild_with_take::(ctx) } else { - self.rebuild_list_by_list::(ctx) + self.rebuild_with_piecewise::(ctx) } } /// Returns `true` when we are confident that `rebuild_with_take` will - /// outperform `rebuild_list_by_list`. + /// outperform `rebuild_with_piecewise`. /// /// Take is dramatically faster for small lists (often 10-100×) because it /// avoids per-list builder overhead. LBL is the safer default for larger @@ -243,87 +244,65 @@ impl ListViewArray { .reinterpret_cast(size_ptype) .into_array(); - // SAFETY: same invariants as `rebuild_list_by_list` — offsets are sequential and - // non-overlapping, all (offset, size) pairs reference valid elements, and the validity - // array is preserved from the original. + // SAFETY: offsets are sequential and non-overlapping, all (offset, size) pairs reference + // valid elements, and the validity array is preserved from the original. Ok(unsafe { ListViewArray::new_unchecked(elements, offsets, sizes, self.validity()?) .with_zero_copy_to_list(true) }) } - /// Rebuilds elements list-by-list: canonicalize elements upfront, then for each list `slice` - /// the relevant range and `append_to_builder` into a typed builder. - fn rebuild_list_by_list( + /// Rebuilds elements using one contiguous-run take over the element child. + fn rebuild_with_piecewise( &self, ctx: &mut ExecutionCtx, ) -> VortexResult { - let element_dtype = self - .dtype() - .as_list_element_opt() - .vortex_expect("somehow had a canonical list that was not a list"); - let new_offset_ptype = rebuilt_offset_ptype(self.offsets().dtype().as_ptype()); let size_ptype = self.sizes().dtype().as_ptype(); let offsets_canonical = self.offsets().clone().execute::(ctx)?; let offsets_canonical = offsets_canonical.reinterpret_cast(offsets_canonical.ptype().to_unsigned()); - let offsets_slice = offsets_canonical.as_slice::(); let sizes_canonical = self.sizes().clone().execute::(ctx)?; let sizes_canonical = sizes_canonical.reinterpret_cast(sizes_canonical.ptype().to_unsigned()); - let sizes_slice = sizes_canonical.as_slice::(); - let len = offsets_slice.len(); + let len = offsets_canonical.len(); + let validity = self.validity()?; + let validity_mask = validity.execute_mask(len, ctx)?; - let mut new_offsets = BufferMut::::with_capacity(len); - // TODO(connor)[ListView]: Do we really need to do this? - // The only reason we need to rebuild the sizes here is that the validity may indicate that - // a list is null even though it has a non-zero size. This rebuild will set the size of all - // null lists to 0. - let mut new_sizes = BufferMut::::with_capacity(len); - - // Canonicalize the elements up front as we will be slicing the elements quite a lot. - let elements_canonical = self + let ranges = match_each_unsigned_integer_ptype!(offsets_canonical.ptype(), |O| { + rebuild_ranges::( + offsets_canonical.as_slice::(), + sizes_canonical.as_slice::(), + &validity_mask, + ) + })?; + + let RebuildRanges { + new_offsets, + new_sizes, + starts, + lengths, + elements_len, + } = ranges; + + // SAFETY: range starts and lengths are derived from valid ListView metadata; elements_len + // is the sum of all generated range lengths. + let element_indices = unsafe { + PiecewiseSequenceArray::new_unchecked( + starts.into_array(), + lengths.into_array(), + elements_len, + ) + }; + let elements = self .elements() - .clone() - .execute::(ctx)? + .take(element_indices.into_array())? + .execute::(ctx)? + .0 .into_array(); - // Note that we do not know what the exact capacity should be of the new elements since - // there could be overlaps in the existing `ListViewArray`. - let mut new_elements_builder = - builder_with_capacity(element_dtype.as_ref(), self.elements().len()); - - // Resolve validity to a mask once instead of probing it per row (see `rebuild_with_take`). - let validity = self.validity()?.execute_mask(len, ctx)?; - - let mut n_elements = NewOffset::zero(); - for index in 0..len { - if !validity.value(index) { - // For NULL lists, place them after the previous item's data to maintain the - // no-overlap invariant for zero-copy to `ListArray` arrays. - new_offsets.push(n_elements); - new_sizes.push(S::zero()); - continue; - } - - let offset = offsets_slice[index]; - let size = sizes_slice[index]; - - let start = offset.as_(); - let stop = start + size.as_(); - - new_offsets.push(n_elements); - new_sizes.push(size); - elements_canonical - .slice(start..stop)? - .append_to_builder(new_elements_builder.as_mut(), ctx)?; - - n_elements += num_traits::cast(size).vortex_expect("Cast failed"); - } - // Built unsigned; reinterpret back to the signed-preserving result types. let offsets = PrimitiveArray::new(new_offsets.freeze(), Validity::NonNullable) .reinterpret_cast(new_offset_ptype) @@ -331,24 +310,11 @@ impl ListViewArray { let sizes = PrimitiveArray::new(new_sizes.freeze(), Validity::NonNullable) .reinterpret_cast(size_ptype) .into_array(); - let elements = new_elements_builder.finish(); - debug_assert_eq!( - n_elements.as_(), - elements.len(), - "The accumulated elements somehow had the wrong length" - ); - - // SAFETY: - // - All offsets are sequential and non-overlapping (`n_elements` tracks running total). - // - Each `offset[i] + size[i]` equals `offset[i+1]` for all valid indices (including null - // lists). - // - All elements referenced by (offset, size) pairs exist within the new `elements` array. - // - The validity array is preserved from the original array unchanged - // - The array satisfies the zero-copy-to-list property by having sorted offsets, no gaps, - // and no overlaps. + // SAFETY: offsets are sequential and non-overlapping, all (offset, size) pairs reference + // valid elements, and validity is preserved from the original array. Ok(unsafe { - ListViewArray::new_unchecked(elements, offsets, sizes, self.validity()?) + ListViewArray::new_unchecked(elements, offsets, sizes, validity) .with_zero_copy_to_list(true) }) } @@ -414,6 +380,69 @@ impl ListViewArray { } } +struct RebuildRanges { + new_offsets: BufferMut, + new_sizes: BufferMut, + starts: BufferMut, + lengths: BufferMut, + elements_len: usize, +} + +fn rebuild_ranges( + offsets: &[O], + sizes: &[S], + validity_mask: &Mask, +) -> VortexResult> +where + NewOffset: IntegerPType, + O: IntegerPType, + S: IntegerPType, +{ + let len = offsets.len(); + let mut new_offsets = BufferMut::::with_capacity(len); + let mut new_sizes = BufferMut::::with_capacity(len); + let mut starts = BufferMut::::with_capacity(len); + let mut lengths = BufferMut::::with_capacity(len); + let mut n_elements = NewOffset::zero(); + let mut elements_len = 0usize; + + for (index, is_valid) in validity_mask.iter().enumerate() { + if !is_valid { + new_offsets.push(n_elements); + new_sizes.push(S::zero()); + starts.push(0); + lengths.push(0); + continue; + } + + let size = sizes[index]; + let start: usize = offsets[index].as_(); + let length: usize = size.as_(); + start.checked_add(length).ok_or_else(|| { + vortex_err!( + "ListView rebuild element range overflow for start {start} and length {length}" + ) + })?; + elements_len = elements_len + .checked_add(length) + .ok_or_else(|| vortex_err!("ListView rebuild elements length overflow"))?; + + new_offsets.push(n_elements); + new_sizes.push(size); + starts.push(start as u64); + lengths.push(length as u64); + n_elements += num_traits::cast(size).vortex_expect("Cast failed"); + } + + Ok(RebuildRanges { + new_offsets, + new_sizes, + starts, + lengths, + elements_len, + }) +} + #[cfg(test)] mod tests { use vortex_buffer::BitBuffer; @@ -561,6 +590,35 @@ mod tests { Ok(()) } + #[test] + fn test_rebuild_flatten_large_lists_with_piecewise_indices() -> VortexResult<()> { + let elements = PrimitiveArray::from_iter(0i32..300).into_array(); + let offsets = PrimitiveArray::from_iter(vec![10u32, 0]).into_array(); + let sizes = PrimitiveArray::from_iter(vec![150u32, 130]).into_array(); + let listview = ListViewArray::new(elements, offsets, sizes, Validity::NonNullable); + + let mut ctx = SESSION.create_execution_ctx(); + let flattened = listview.rebuild(ListViewRebuildMode::MakeZeroCopyToList, &mut ctx)?; + + assert_eq!(flattened.elements().len(), 280); + assert_eq!(flattened.offset_at(0), 0); + assert_eq!(flattened.size_at(0), 150); + assert_eq!(flattened.offset_at(1), 150); + assert_eq!(flattened.size_at(1), 130); + + assert_arrays_eq!( + flattened.list_elements_at(0)?, + PrimitiveArray::from_iter(10i32..160), + &mut ctx + ); + assert_arrays_eq!( + flattened.list_elements_at(1)?, + PrimitiveArray::from_iter(0i32..130), + &mut ctx + ); + Ok(()) + } + #[test] fn test_rebuild_with_trailing_nulls_regression() -> VortexResult<()> { // Regression test for issue #5412 From 5027bff29fc89887766fe1b9e2d0c0ba4a8fce97 Mon Sep 17 00:00:00 2001 From: Daniel King Date: Thu, 16 Jul 2026 16:28:40 -0400 Subject: [PATCH 02/18] Use unit multipliers for PiecewiseSequence consumers Signed-off-by: Daniel King --- vortex-array/src/arrays/list/compute/take.rs | 13 ++++++++++--- vortex-array/src/arrays/listview/rebuild.rs | 4 +++- 2 files changed, 13 insertions(+), 4 deletions(-) diff --git a/vortex-array/src/arrays/list/compute/take.rs b/vortex-array/src/arrays/list/compute/take.rs index 35788796d78..9a6c447b730 100644 --- a/vortex-array/src/arrays/list/compute/take.rs +++ b/vortex-array/src/arrays/list/compute/take.rs @@ -11,6 +11,7 @@ use vortex_error::vortex_err; use crate::ArrayRef; use crate::IntoArray; use crate::array::ArrayView; +use crate::arrays::ConstantArray; use crate::arrays::List; use crate::arrays::ListArray; use crate::arrays::PiecewiseSequence; @@ -19,7 +20,7 @@ use crate::arrays::Primitive; use crate::arrays::PrimitiveArray; use crate::arrays::dict::TakeExecute; use crate::arrays::list::ListArrayExt; -use crate::arrays::piecewise_sequence::execute_index_arrays; +use crate::arrays::piecewise_sequence::execute_unit_multiplier_index_arrays; use crate::arrays::piecewise_sequence::validate_index_ranges; use crate::arrays::primitive::PrimitiveArrayExt; use crate::builders::ArrayBuilder; @@ -156,7 +157,9 @@ fn take_piecewise_sequence( return Ok(None); } - let (starts, lengths) = execute_index_arrays(indices, ctx)?; + let Some((starts, lengths)) = execute_unit_multiplier_index_arrays(indices, ctx)? else { + return Ok(None); + }; let offsets = array.offsets().clone().execute::(ctx)?; let offsets = offsets.reinterpret_cast(offsets.ptype().to_unsigned()); @@ -305,12 +308,14 @@ where debug_assert_eq!(output_elements, total_elements); let offsets = PrimitiveArray::new(new_offsets.freeze(), Validity::NonNullable).into_array(); + let multipliers = ConstantArray::new(1u64, element_starts.len()).into_array(); // SAFETY: element ranges are derived from validated source list offsets, and total_elements is - // the sum of the gathered element range lengths. + // the sum of the gathered element range lengths. Multiplier 1 preserves contiguous ranges. let element_indices = unsafe { PiecewiseSequenceArray::new_unchecked( element_starts.into_array(), element_lengths.into_array(), + multipliers, total_elements, ) }; @@ -416,6 +421,7 @@ mod test { use crate::VortexSessionExecute; use crate::array_session; use crate::arrays::BoolArray; + use crate::arrays::ConstantArray; use crate::arrays::ListArray; use crate::arrays::ListViewArray; use crate::arrays::PiecewiseSequenceArray; @@ -618,6 +624,7 @@ mod test { let idx = PiecewiseSequenceArray::try_new( buffer![1u64, 0].into_array(), buffer![2u64, 1].into_array(), + ConstantArray::new(1u64, 2).into_array(), 3, ) .unwrap() diff --git a/vortex-array/src/arrays/listview/rebuild.rs b/vortex-array/src/arrays/listview/rebuild.rs index 1ce8f69d12e..5cf8bcb95b1 100644 --- a/vortex-array/src/arrays/listview/rebuild.rs +++ b/vortex-array/src/arrays/listview/rebuild.rs @@ -288,11 +288,13 @@ impl ListViewArray { } = ranges; // SAFETY: range starts and lengths are derived from valid ListView metadata; elements_len - // is the sum of all generated range lengths. + // is the sum of all generated range lengths. Multiplier 1 preserves contiguous ranges. + let multipliers = ConstantArray::new(1u64, starts.len()).into_array(); let element_indices = unsafe { PiecewiseSequenceArray::new_unchecked( starts.into_array(), lengths.into_array(), + multipliers, elements_len, ) }; From 68a7eb5c507f4bf5f2ccb099daa1511f8b4828c6 Mon Sep 17 00:00:00 2001 From: Daniel King Date: Thu, 16 Jul 2026 19:51:55 -0400 Subject: [PATCH 03/18] Specialize constant PiecewiseSequence runs Signed-off-by: Daniel King --- vortex-array/src/arrays/list/compute/take.rs | 283 +++++++++++++++++- vortex-array/src/arrays/listview/rebuild.rs | 18 +- .../src/arrays/piecewise_sequence/tests.rs | 22 ++ 3 files changed, 301 insertions(+), 22 deletions(-) diff --git a/vortex-array/src/arrays/list/compute/take.rs b/vortex-array/src/arrays/list/compute/take.rs index 9a6c447b730..369a743d4dc 100644 --- a/vortex-array/src/arrays/list/compute/take.rs +++ b/vortex-array/src/arrays/list/compute/take.rs @@ -20,8 +20,10 @@ use crate::arrays::Primitive; use crate::arrays::PrimitiveArray; use crate::arrays::dict::TakeExecute; use crate::arrays::list::ListArrayExt; +use crate::arrays::piecewise_sequence::UnitMultiplierLengths; use crate::arrays::piecewise_sequence::execute_unit_multiplier_index_arrays; use crate::arrays::piecewise_sequence::validate_index_ranges; +use crate::arrays::piecewise_sequence::validate_index_ranges_constant; use crate::arrays::primitive::PrimitiveArrayExt; use crate::builders::ArrayBuilder; use crate::builders::PrimitiveBuilder; @@ -162,21 +164,174 @@ fn take_piecewise_sequence( }; let offsets = array.offsets().clone().execute::(ctx)?; let offsets = offsets.reinterpret_cast(offsets.ptype().to_unsigned()); + let output_len = indices_ref.len(); + + let taken = match &lengths { + UnitMultiplierLengths::Constant(length) => take_piecewise_sequence_constant_dispatch( + array, + &starts, + *length, + &offsets, + indices_ref, + output_len, + )?, + UnitMultiplierLengths::Array(lengths) => take_piecewise_sequence_lengths_dispatch( + array, + &starts, + lengths, + &offsets, + indices_ref, + output_len, + )?, + }; + Ok(Some(taken)) +} +fn take_piecewise_sequence_constant_dispatch( + array: ArrayView<'_, List>, + starts: &PrimitiveArray, + length: usize, + offsets: &PrimitiveArray, + indices_ref: &ArrayRef, + output_len: usize, +) -> VortexResult { match_each_unsigned_integer_ptype!(starts.ptype(), |S| { - match_each_unsigned_integer_ptype!(lengths.ptype(), |L| { - match_each_unsigned_integer_ptype!(offsets.ptype(), |O| { - take_piecewise_sequence_typed::( - array, - starts.as_slice::(), - lengths.as_slice::(), - offsets.as_slice::(), - indices_ref, - ) - }) - }) + take_piecewise_sequence_constant_start_dispatch::( + array, + starts, + length, + offsets, + indices_ref, + output_len, + ) + }) +} + +fn take_piecewise_sequence_constant_start_dispatch( + array: ArrayView<'_, List>, + starts: &PrimitiveArray, + length: usize, + offsets: &PrimitiveArray, + indices_ref: &ArrayRef, + output_len: usize, +) -> VortexResult +where + S: UnsignedPType, +{ + match_each_unsigned_integer_ptype!(offsets.ptype(), |O| { + take_piecewise_sequence_constant_length::( + array, + starts.as_slice::(), + length, + offsets.as_slice::(), + indices_ref, + output_len, + ) + }) +} + +fn take_piecewise_sequence_lengths_dispatch( + array: ArrayView<'_, List>, + starts: &PrimitiveArray, + lengths: &PrimitiveArray, + offsets: &PrimitiveArray, + indices_ref: &ArrayRef, + output_len: usize, +) -> VortexResult { + match_each_unsigned_integer_ptype!(starts.ptype(), |S| { + take_piecewise_sequence_lengths_start_dispatch::( + array, + starts, + lengths, + offsets, + indices_ref, + output_len, + ) + }) +} + +fn take_piecewise_sequence_lengths_start_dispatch( + array: ArrayView<'_, List>, + starts: &PrimitiveArray, + lengths: &PrimitiveArray, + offsets: &PrimitiveArray, + indices_ref: &ArrayRef, + output_len: usize, +) -> VortexResult +where + S: UnsignedPType, +{ + match_each_unsigned_integer_ptype!(lengths.ptype(), |L| { + take_piecewise_sequence_lengths_start_length_dispatch::( + array, + starts, + lengths, + offsets, + indices_ref, + output_len, + ) + }) +} + +fn take_piecewise_sequence_lengths_start_length_dispatch( + array: ArrayView<'_, List>, + starts: &PrimitiveArray, + lengths: &PrimitiveArray, + offsets: &PrimitiveArray, + indices_ref: &ArrayRef, + output_len: usize, +) -> VortexResult +where + S: UnsignedPType, + L: UnsignedPType, +{ + match_each_unsigned_integer_ptype!(offsets.ptype(), |O| { + take_piecewise_sequence_typed::( + array, + starts.as_slice::(), + lengths.as_slice::(), + offsets.as_slice::(), + indices_ref, + output_len, + ) + }) +} + +fn take_piecewise_sequence_constant_length( + array: ArrayView<'_, List>, + starts: &[S], + length: usize, + offsets: &[Offset], + indices_ref: &ArrayRef, + output_len: usize, +) -> VortexResult +where + S: UnsignedPType, + Offset: UnsignedPType, +{ + validate_index_ranges_constant(array.len(), starts, length, output_len)?; + let total_elements = + piecewise_list_elements_len_constant(array.elements().len(), offsets, starts, length)?; + let validity = array.validity()?.take(indices_ref)?; + + match_smallest_offset_type!(total_elements, |OutputOffset| { + let gathered = gather_piecewise_list_constant_length::( + array.elements(), + offsets, + starts, + length, + output_len, + total_elements, + )?; + + // SAFETY: output offsets are rebuilt from valid monotonic source offsets; output elements + // are exactly the gathered child ranges referenced by those offsets; validity has one bit + // per output row. + Ok( + unsafe { ListArray::new_unchecked(gathered.elements, gathered.offsets, validity) } + .into_array(), + ) }) - .map(Some) } fn take_piecewise_sequence_typed( @@ -185,13 +340,14 @@ fn take_piecewise_sequence_typed( lengths: &[L], offsets: &[Offset], indices_ref: &ArrayRef, + output_len: usize, ) -> VortexResult where S: UnsignedPType, L: UnsignedPType, Offset: UnsignedPType, { - validate_index_ranges(array.len(), starts, lengths, indices_ref.len())?; + validate_index_ranges(array.len(), starts, lengths, output_len)?; let total_elements = piecewise_list_elements_len(array.elements().len(), offsets, starts, lengths)?; @@ -201,7 +357,7 @@ where offsets, starts, lengths, - indices_ref.len(), + output_len, total_elements, )?; let validity = array.validity()?.take(indices_ref)?; @@ -221,6 +377,37 @@ struct GatheredList { offsets: ArrayRef, } +fn piecewise_list_elements_len_constant( + elements_len: usize, + offsets: &[Offset], + starts: &[S], + length: usize, +) -> VortexResult +where + S: UnsignedPType, + Offset: UnsignedPType, +{ + let mut total = 0usize; + for &start in starts { + let start: usize = start.as_(); + let end = start + length; + if length == 0 { + continue; + } + + let element_start: usize = offsets[start].as_(); + let element_end: usize = offsets[end].as_(); + vortex_ensure!( + element_start <= element_end && element_end <= elements_len, + "List offsets range {element_start}..{element_end} exceeds elements length {elements_len}", + ); + total = total + .checked_add(element_end - element_start) + .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; + } + Ok(total) +} + fn piecewise_list_elements_len( elements_len: usize, offsets: &[Offset], @@ -254,6 +441,74 @@ where Ok(total) } +fn gather_piecewise_list_constant_length( + elements: &ArrayRef, + offsets: &[Offset], + starts: &[S], + length: usize, + output_len: usize, + total_elements: usize, +) -> VortexResult +where + S: UnsignedPType, + Offset: UnsignedPType, + OutputOffset: IntegerPType, +{ + let offsets_capacity = output_len + .checked_add(1) + .ok_or_else(|| vortex_err!("List take offsets length overflow"))?; + let mut new_offsets = BufferMut::::with_capacity(offsets_capacity); + let mut element_starts = BufferMut::::with_capacity(starts.len()); + let mut element_lengths = BufferMut::::with_capacity(starts.len()); + let mut output_elements = 0usize; + + new_offsets.push(OutputOffset::zero()); + for &start in starts { + let start: usize = start.as_(); + let end = start + length; + if length == 0 { + continue; + } + + let element_start: usize = offsets[start].as_(); + let element_end: usize = offsets[end].as_(); + for &offset in &offsets[start + 1..=end] { + let offset: usize = offset.as_(); + let relative = offset + .checked_sub(element_start) + .ok_or_else(|| vortex_err!("List offsets are not monotonic at offset {offset}"))?; + let output_offset = output_elements + .checked_add(relative) + .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; + new_offsets.push(new_offset_value::(output_offset)?); + } + + let element_length = element_end - element_start; + element_starts.push(element_start as u64); + element_lengths.push(element_length as u64); + output_elements = output_elements + .checked_add(element_length) + .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; + } + debug_assert_eq!(output_elements, total_elements); + + let offsets = PrimitiveArray::new(new_offsets.freeze(), Validity::NonNullable).into_array(); + let multipliers = ConstantArray::new(1u64, element_starts.len()).into_array(); + // SAFETY: element ranges are derived from validated source list offsets, and total_elements is + // the sum of the gathered element range lengths. Multiplier 1 preserves contiguous ranges. + let element_indices = unsafe { + PiecewiseSequenceArray::new_unchecked( + element_starts.into_array(), + element_lengths.into_array(), + multipliers, + total_elements, + ) + }; + let elements = elements.take(element_indices.into_array())?; + + Ok(GatheredList { elements, offsets }) +} + fn gather_piecewise_list( elements: &ArrayRef, offsets: &[Offset], diff --git a/vortex-array/src/arrays/listview/rebuild.rs b/vortex-array/src/arrays/listview/rebuild.rs index 5cf8bcb95b1..15eb9feab38 100644 --- a/vortex-array/src/arrays/listview/rebuild.rs +++ b/vortex-array/src/arrays/listview/rebuild.rs @@ -10,7 +10,6 @@ use vortex_mask::Mask; use crate::ExecutionCtx; use crate::IntoArray; -use crate::RecursiveCanonical; use crate::arrays::ConstantArray; use crate::arrays::ListViewArray; use crate::arrays::PiecewiseSequenceArray; @@ -286,6 +285,14 @@ impl ListViewArray { lengths, elements_len, } = ranges; + let constant_length = lengths + .first() + .copied() + .filter(|first| lengths.iter().all(|length| *length == *first)); + let lengths = match constant_length { + Some(length) => ConstantArray::new(length, starts.len()).into_array(), + None => lengths.into_array(), + }; // SAFETY: range starts and lengths are derived from valid ListView metadata; elements_len // is the sum of all generated range lengths. Multiplier 1 preserves contiguous ranges. @@ -293,17 +300,12 @@ impl ListViewArray { let element_indices = unsafe { PiecewiseSequenceArray::new_unchecked( starts.into_array(), - lengths.into_array(), + lengths, multipliers, elements_len, ) }; - let elements = self - .elements() - .take(element_indices.into_array())? - .execute::(ctx)? - .0 - .into_array(); + let elements = self.elements().take(element_indices.into_array())?; // Built unsigned; reinterpret back to the signed-preserving result types. let offsets = PrimitiveArray::new(new_offsets.freeze(), Validity::NonNullable) diff --git a/vortex-array/src/arrays/piecewise_sequence/tests.rs b/vortex-array/src/arrays/piecewise_sequence/tests.rs index eef5d4af12d..a087404fee2 100644 --- a/vortex-array/src/arrays/piecewise_sequence/tests.rs +++ b/vortex-array/src/arrays/piecewise_sequence/tests.rs @@ -14,6 +14,7 @@ use crate::arrays::BoolArray; use crate::arrays::ConstantArray; use crate::arrays::DecimalArray; use crate::arrays::FixedSizeListArray; +use crate::arrays::ListArray; use crate::arrays::PiecewiseSequenceArray; use crate::arrays::PrimitiveArray; use crate::arrays::VarBinArray; @@ -338,6 +339,27 @@ fn contiguous_take_consumers_support_constant_piecewise_lengths() -> VortexResul Ok(()) } +#[test] +fn list_take_consumes_constant_piecewise_lengths() -> VortexResult<()> { + let mut ctx = array_session().create_execution_ctx(); + let list = ListArray::try_new( + buffer![0i32, 1, 2, 3, 4, 5, 6].into_array(), + buffer![0u32, 2, 5, 5, 7].into_array(), + Validity::NonNullable, + )? + .into_array(); + let list_indices = piecewise_indices_constant_length([1, 0], 2)?; + let expected = ListArray::try_new( + buffer![2i32, 3, 4, 0, 1, 2, 3, 4].into_array(), + buffer![0u32, 3, 3, 5, 8].into_array(), + Validity::NonNullable, + )? + .into_array(); + assert_arrays_eq!(list.take(list_indices)?, expected, &mut ctx); + + Ok(()) +} + #[test] fn fixed_size_list_take_builds_piecewise_element_indices() -> VortexResult<()> { let elements = PrimitiveArray::from_iter(0i32..12).into_array(); From b69e2db44347b536efaf8fcb57d3c7dfbcd4ff7a Mon Sep 17 00:00:00 2001 From: Daniel King Date: Fri, 17 Jul 2026 14:42:25 -0400 Subject: [PATCH 04/18] Inline List PiecewiseSequence range length checks Signed-off-by: Daniel King --- vortex-array/src/arrays/list/compute/take.rs | 51 +++++++++++++------- 1 file changed, 33 insertions(+), 18 deletions(-) diff --git a/vortex-array/src/arrays/list/compute/take.rs b/vortex-array/src/arrays/list/compute/take.rs index 369a743d4dc..9d83cca3a48 100644 --- a/vortex-array/src/arrays/list/compute/take.rs +++ b/vortex-array/src/arrays/list/compute/take.rs @@ -22,8 +22,6 @@ use crate::arrays::dict::TakeExecute; use crate::arrays::list::ListArrayExt; use crate::arrays::piecewise_sequence::UnitMultiplierLengths; use crate::arrays::piecewise_sequence::execute_unit_multiplier_index_arrays; -use crate::arrays::piecewise_sequence::validate_index_ranges; -use crate::arrays::piecewise_sequence::validate_index_ranges_constant; use crate::arrays::primitive::PrimitiveArrayExt; use crate::builders::ArrayBuilder; use crate::builders::PrimitiveBuilder; @@ -309,7 +307,14 @@ where S: UnsignedPType, Offset: UnsignedPType, { - validate_index_ranges_constant(array.len(), starts, length, output_len)?; + let computed_len = starts + .len() + .checked_mul(length) + .ok_or_else(|| vortex_err!("PiecewiseSequenceArray output length overflows usize"))?; + vortex_ensure!( + computed_len == output_len, + "PiecewiseSequenceArray expanded length {computed_len} does not match declared length {output_len}" + ); let total_elements = piecewise_list_elements_len_constant(array.elements().len(), offsets, starts, length)?; let validity = array.validity()?.take(indices_ref)?; @@ -347,7 +352,17 @@ where L: UnsignedPType, Offset: UnsignedPType, { - validate_index_ranges(array.len(), starts, lengths, output_len)?; + let mut computed_len = 0usize; + for &length in lengths { + let length: usize = length.as_(); + computed_len = computed_len + .checked_add(length) + .ok_or_else(|| vortex_err!("PiecewiseSequenceArray output length overflows usize"))?; + } + vortex_ensure!( + computed_len == output_len, + "PiecewiseSequenceArray expanded length {computed_len} does not match declared length {output_len}" + ); let total_elements = piecewise_list_elements_len(array.elements().len(), offsets, starts, lengths)?; @@ -390,13 +405,13 @@ where let mut total = 0usize; for &start in starts { let start: usize = start.as_(); - let end = start + length; if length == 0 { continue; } - let element_start: usize = offsets[start].as_(); - let element_end: usize = offsets[end].as_(); + let offset_range = &offsets[start..][..=length]; + let element_start: usize = offset_range[0].as_(); + let element_end: usize = offset_range[length].as_(); vortex_ensure!( element_start <= element_end && element_end <= elements_len, "List offsets range {element_start}..{element_end} exceeds elements length {elements_len}", @@ -423,13 +438,13 @@ where for (&start, &length) in starts.iter().zip_eq(lengths) { let start: usize = start.as_(); let length: usize = length.as_(); - let end = start + length; if length == 0 { continue; } - let element_start: usize = offsets[start].as_(); - let element_end: usize = offsets[end].as_(); + let offset_range = &offsets[start..][..=length]; + let element_start: usize = offset_range[0].as_(); + let element_end: usize = offset_range[length].as_(); vortex_ensure!( element_start <= element_end && element_end <= elements_len, "List offsets range {element_start}..{element_end} exceeds elements length {elements_len}", @@ -465,14 +480,14 @@ where new_offsets.push(OutputOffset::zero()); for &start in starts { let start: usize = start.as_(); - let end = start + length; if length == 0 { continue; } - let element_start: usize = offsets[start].as_(); - let element_end: usize = offsets[end].as_(); - for &offset in &offsets[start + 1..=end] { + let offset_range = &offsets[start..][..=length]; + let element_start: usize = offset_range[0].as_(); + let element_end: usize = offset_range[length].as_(); + for &offset in &offset_range[1..] { let offset: usize = offset.as_(); let relative = offset .checked_sub(element_start) @@ -535,14 +550,14 @@ where for (&start, &length) in starts.iter().zip_eq(lengths) { let start: usize = start.as_(); let length: usize = length.as_(); - let end = start + length; if length == 0 { continue; } - let element_start: usize = offsets[start].as_(); - let element_end: usize = offsets[end].as_(); - for &offset in &offsets[start + 1..=end] { + let offset_range = &offsets[start..][..=length]; + let element_start: usize = offset_range[0].as_(); + let element_end: usize = offset_range[length].as_(); + for &offset in &offset_range[1..] { let offset: usize = offset.as_(); let relative = offset .checked_sub(element_start) From 90e500790b0e97ac77edae357c0e44a6405a668b Mon Sep 17 00:00:00 2001 From: Daniel King Date: Fri, 17 Jul 2026 16:04:53 -0400 Subject: [PATCH 05/18] Cover null ListView rebuild placeholders Signed-off-by: Daniel King --- vortex-array/src/arrays/listview/rebuild.rs | 38 +++++++++++++++++++++ 1 file changed, 38 insertions(+) diff --git a/vortex-array/src/arrays/listview/rebuild.rs b/vortex-array/src/arrays/listview/rebuild.rs index 15eb9feab38..6e36958522e 100644 --- a/vortex-array/src/arrays/listview/rebuild.rs +++ b/vortex-array/src/arrays/listview/rebuild.rs @@ -544,6 +544,44 @@ mod tests { Ok(()) } + #[test] + fn test_rebuild_flatten_null_row_ignores_invalid_range_payload() -> VortexResult<()> { + let elements = PrimitiveArray::from_iter(vec![1i32, 2, 3]).into_array(); + let offsets = PrimitiveArray::from_iter(vec![0u32, 999, 2]).into_array(); + let sizes = PrimitiveArray::from_iter(vec![2u32, 999, 1]).into_array(); + let validity = Validity::from_iter([true, false, true]); + + // SAFETY: this intentionally models a null row whose physical offset and size payloads are + // invalid. Rebuild must ignore those payloads and only emit safe placeholder ranges for the + // null row. + let listview = unsafe { ListViewArray::new_unchecked(elements, offsets, sizes, validity) }; + + let mut ctx = SESSION.create_execution_ctx(); + let flattened = listview.rebuild(ListViewRebuildMode::MakeZeroCopyToList, &mut ctx)?; + + assert_eq!(flattened.offset_at(0), 0); + assert_eq!(flattened.size_at(0), 2); + assert_eq!(flattened.offset_at(1), 2); + assert_eq!(flattened.size_at(1), 0); + assert_eq!(flattened.offset_at(2), 2); + assert_eq!(flattened.size_at(2), 1); + assert!(flattened.validity()?.execute_is_valid(0, &mut ctx)?); + assert!(!flattened.validity()?.execute_is_valid(1, &mut ctx)?); + assert!(flattened.validity()?.execute_is_valid(2, &mut ctx)?); + + assert_arrays_eq!( + flattened.list_elements_at(0)?, + PrimitiveArray::from_iter([1i32, 2]), + &mut ctx + ); + assert_arrays_eq!( + flattened.list_elements_at(2)?, + PrimitiveArray::from_iter([3i32]), + &mut ctx + ); + Ok(()) + } + #[test] fn test_rebuild_trim_elements_basic() -> VortexResult<()> { // Test trimming both leading and trailing unused elements while preserving gaps in the From e2df924b0e36ff3d4262fc9c2fb08ab36aa6049e Mon Sep 17 00:00:00 2001 From: Daniel King Date: Fri, 17 Jul 2026 16:12:18 -0400 Subject: [PATCH 06/18] Mask null ListView take metadata Signed-off-by: Daniel King --- .../src/arrays/listview/compute/take.rs | 14 +++++++---- .../src/arrays/listview/tests/take.rs | 23 +++++++++++++++++++ 2 files changed, 33 insertions(+), 4 deletions(-) diff --git a/vortex-array/src/arrays/listview/compute/take.rs b/vortex-array/src/arrays/listview/compute/take.rs index 2b6c016d2c3..2389715c2d6 100644 --- a/vortex-array/src/arrays/listview/compute/take.rs +++ b/vortex-array/src/arrays/listview/compute/take.rs @@ -50,14 +50,20 @@ fn apply_take(array: ArrayView<'_, ListView>, indices: &ArrayRef) -> VortexResul // duplicates. let nullable_new_offsets = offsets.take(indices.clone())?; let nullable_new_sizes = sizes.take(indices.clone())?; + let validity_array = new_validity.to_array(indices.len()); - // `take` returns nullable arrays; cast back to non-nullable (filling with zeros to represent - // the null lists — the validity mask tracks nullness separately). + // Null output rows may carry arbitrary physical offset/size payloads from either the indices + // or the source rows. Mask with the final validity before filling so metadata placeholders are + // safe and non-nullable. let new_offsets = match_each_integer_ptype!(nullable_new_offsets.dtype().as_ptype(), |O| { - nullable_new_offsets.fill_null(Scalar::primitive(O::zero(), Nullability::NonNullable))? + nullable_new_offsets + .mask(validity_array.clone())? + .fill_null(Scalar::primitive(O::zero(), Nullability::NonNullable))? }); let new_sizes = match_each_integer_ptype!(nullable_new_sizes.dtype().as_ptype(), |S| { - nullable_new_sizes.fill_null(Scalar::primitive(S::zero(), Nullability::NonNullable))? + nullable_new_sizes + .mask(validity_array.clone())? + .fill_null(Scalar::primitive(S::zero(), Nullability::NonNullable))? }); // SAFETY: Take operation maintains all `ListViewArray` invariants: diff --git a/vortex-array/src/arrays/listview/tests/take.rs b/vortex-array/src/arrays/listview/tests/take.rs index 2af397e262f..8fb7a4af1bb 100644 --- a/vortex-array/src/arrays/listview/tests/take.rs +++ b/vortex-array/src/arrays/listview/tests/take.rs @@ -101,6 +101,29 @@ fn test_take_with_gaps() { ); } +#[test] +fn test_take_null_source_row_zeros_offset_size_payloads() { + let elements = buffer![1i32, 2].into_array(); + let offsets = buffer![0u32, 999].into_array(); + let sizes = buffer![2u32, 999].into_array(); + let validity = Validity::from_iter([true, false]); + let listview = + unsafe { ListViewArray::new_unchecked(elements, offsets, sizes, validity) }.into_array(); + + let result = listview.take(buffer![1u32].into_array()).unwrap(); + let result_list = result + .execute::(&mut SESSION.create_execution_ctx()) + .unwrap(); + + assert_eq!(result_list.offset_at(0), 0); + assert_eq!(result_list.size_at(0), 0); + assert!( + result_list + .is_invalid(0, &mut SESSION.create_execution_ctx()) + .unwrap() + ); +} + #[test] fn test_take_constant_arrays() { // ListView-specific: Test with ConstantArray for offsets/sizes. From 4e4650f10f2c407817fcd123059a277b8e110293 Mon Sep 17 00:00:00 2001 From: Daniel King Date: Fri, 17 Jul 2026 16:19:04 -0400 Subject: [PATCH 07/18] Rename List contiguous slice helper uses Signed-off-by: Daniel King --- vortex-array/src/arrays/list/compute/take.rs | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/vortex-array/src/arrays/list/compute/take.rs b/vortex-array/src/arrays/list/compute/take.rs index 9d83cca3a48..ff9549d6a34 100644 --- a/vortex-array/src/arrays/list/compute/take.rs +++ b/vortex-array/src/arrays/list/compute/take.rs @@ -20,8 +20,8 @@ use crate::arrays::Primitive; use crate::arrays::PrimitiveArray; use crate::arrays::dict::TakeExecute; use crate::arrays::list::ListArrayExt; -use crate::arrays::piecewise_sequence::UnitMultiplierLengths; -use crate::arrays::piecewise_sequence::execute_unit_multiplier_index_arrays; +use crate::arrays::piecewise_sequence::ConstantOrArray; +use crate::arrays::piecewise_sequence::maybe_contiguous_slices; use crate::arrays::primitive::PrimitiveArrayExt; use crate::builders::ArrayBuilder; use crate::builders::PrimitiveBuilder; @@ -157,7 +157,7 @@ fn take_piecewise_sequence( return Ok(None); } - let Some((starts, lengths)) = execute_unit_multiplier_index_arrays(indices, ctx)? else { + let Some((starts, lengths)) = maybe_contiguous_slices(indices, ctx)? else { return Ok(None); }; let offsets = array.offsets().clone().execute::(ctx)?; @@ -165,7 +165,7 @@ fn take_piecewise_sequence( let output_len = indices_ref.len(); let taken = match &lengths { - UnitMultiplierLengths::Constant(length) => take_piecewise_sequence_constant_dispatch( + ConstantOrArray::Constant(length) => take_piecewise_sequence_constant_dispatch( array, &starts, *length, @@ -173,7 +173,7 @@ fn take_piecewise_sequence( indices_ref, output_len, )?, - UnitMultiplierLengths::Array(lengths) => take_piecewise_sequence_lengths_dispatch( + ConstantOrArray::Array(lengths) => take_piecewise_sequence_lengths_dispatch( array, &starts, lengths, From daa506eff97a8c784e7a14cb02c066f474573082 Mon Sep 17 00:00:00 2001 From: Daniel King Date: Fri, 17 Jul 2026 16:37:33 -0400 Subject: [PATCH 08/18] Remove redundant ListView validity clone Signed-off-by: Daniel King --- vortex-array/src/arrays/listview/compute/take.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/vortex-array/src/arrays/listview/compute/take.rs b/vortex-array/src/arrays/listview/compute/take.rs index 2389715c2d6..01ffb80dab7 100644 --- a/vortex-array/src/arrays/listview/compute/take.rs +++ b/vortex-array/src/arrays/listview/compute/take.rs @@ -62,7 +62,7 @@ fn apply_take(array: ArrayView<'_, ListView>, indices: &ArrayRef) -> VortexResul }); let new_sizes = match_each_integer_ptype!(nullable_new_sizes.dtype().as_ptype(), |S| { nullable_new_sizes - .mask(validity_array.clone())? + .mask(validity_array)? .fill_null(Scalar::primitive(S::zero(), Nullability::NonNullable))? }); From e69c305e323dee49aa4d38a94306874b7e000a68 Mon Sep 17 00:00:00 2001 From: Daniel King Date: Fri, 17 Jul 2026 17:06:27 -0400 Subject: [PATCH 09/18] Simplify zero-length list element sizing Signed-off-by: Daniel King --- vortex-array/src/arrays/list/compute/take.rs | 12 ++++-------- 1 file changed, 4 insertions(+), 8 deletions(-) diff --git a/vortex-array/src/arrays/list/compute/take.rs b/vortex-array/src/arrays/list/compute/take.rs index ff9549d6a34..5fab0c0a8f3 100644 --- a/vortex-array/src/arrays/list/compute/take.rs +++ b/vortex-array/src/arrays/list/compute/take.rs @@ -402,13 +402,13 @@ where S: UnsignedPType, Offset: UnsignedPType, { + if length == 0 { + return Ok(0); + } + let mut total = 0usize; for &start in starts { let start: usize = start.as_(); - if length == 0 { - continue; - } - let offset_range = &offsets[start..][..=length]; let element_start: usize = offset_range[0].as_(); let element_end: usize = offset_range[length].as_(); @@ -438,10 +438,6 @@ where for (&start, &length) in starts.iter().zip_eq(lengths) { let start: usize = start.as_(); let length: usize = length.as_(); - if length == 0 { - continue; - } - let offset_range = &offsets[start..][..=length]; let element_start: usize = offset_range[0].as_(); let element_end: usize = offset_range[length].as_(); From f1880b196dcc4d7ec25a6ee4a56bcf3e3af615a4 Mon Sep 17 00:00:00 2001 From: Daniel King Date: Mon, 20 Jul 2026 12:32:39 -0400 Subject: [PATCH 10/18] Use Columnar lengths in List take Signed-off-by: Daniel King --- vortex-array/src/arrays/list/compute/take.rs | 50 +++++++++++++------- 1 file changed, 32 insertions(+), 18 deletions(-) diff --git a/vortex-array/src/arrays/list/compute/take.rs b/vortex-array/src/arrays/list/compute/take.rs index 5fab0c0a8f3..2ee02ae969a 100644 --- a/vortex-array/src/arrays/list/compute/take.rs +++ b/vortex-array/src/arrays/list/compute/take.rs @@ -5,10 +5,13 @@ use itertools::Itertools as _; use vortex_buffer::BufferMut; use vortex_error::VortexExpect; use vortex_error::VortexResult; +use vortex_error::vortex_bail; use vortex_error::vortex_ensure; use vortex_error::vortex_err; use crate::ArrayRef; +use crate::Canonical; +use crate::Columnar; use crate::IntoArray; use crate::array::ArrayView; use crate::arrays::ConstantArray; @@ -20,7 +23,7 @@ use crate::arrays::Primitive; use crate::arrays::PrimitiveArray; use crate::arrays::dict::TakeExecute; use crate::arrays::list::ListArrayExt; -use crate::arrays::piecewise_sequence::ConstantOrArray; +use crate::arrays::piecewise_sequence::constant_unsigned_usize; use crate::arrays::piecewise_sequence::maybe_contiguous_slices; use crate::arrays::primitive::PrimitiveArrayExt; use crate::builders::ArrayBuilder; @@ -164,23 +167,34 @@ fn take_piecewise_sequence( let offsets = offsets.reinterpret_cast(offsets.ptype().to_unsigned()); let output_len = indices_ref.len(); - let taken = match &lengths { - ConstantOrArray::Constant(length) => take_piecewise_sequence_constant_dispatch( - array, - &starts, - *length, - &offsets, - indices_ref, - output_len, - )?, - ConstantOrArray::Array(lengths) => take_piecewise_sequence_lengths_dispatch( - array, - &starts, - lengths, - &offsets, - indices_ref, - output_len, - )?, + let taken = match lengths { + Columnar::Constant(lengths) => { + let length = constant_unsigned_usize(&lengths)?; + take_piecewise_sequence_constant_dispatch( + array, + &starts, + length, + &offsets, + indices_ref, + output_len, + )? + } + Columnar::Canonical(Canonical::Primitive(lengths)) => { + take_piecewise_sequence_lengths_dispatch( + array, + &starts, + &lengths, + &offsets, + indices_ref, + output_len, + )? + } + Columnar::Canonical(lengths) => { + vortex_bail!( + "PiecewiseSequenceArray lengths must be primitive or constant, got {}", + lengths.dtype() + ) + } }; Ok(Some(taken)) } From 09ca1993ef2ecc29739a65159ab8230625b5c021 Mon Sep 17 00:00:00 2001 From: Daniel King Date: Mon, 20 Jul 2026 12:50:06 -0400 Subject: [PATCH 11/18] Rely on PiecewiseSequence validation in List take Signed-off-by: Daniel King --- vortex-array/src/arrays/list/compute/take.rs | 66 +++++--------------- 1 file changed, 15 insertions(+), 51 deletions(-) diff --git a/vortex-array/src/arrays/list/compute/take.rs b/vortex-array/src/arrays/list/compute/take.rs index 2ee02ae969a..e141e8cf786 100644 --- a/vortex-array/src/arrays/list/compute/take.rs +++ b/vortex-array/src/arrays/list/compute/take.rs @@ -5,12 +5,10 @@ use itertools::Itertools as _; use vortex_buffer::BufferMut; use vortex_error::VortexExpect; use vortex_error::VortexResult; -use vortex_error::vortex_bail; use vortex_error::vortex_ensure; use vortex_error::vortex_err; use crate::ArrayRef; -use crate::Canonical; use crate::Columnar; use crate::IntoArray; use crate::array::ArrayView; @@ -169,7 +167,7 @@ fn take_piecewise_sequence( let taken = match lengths { Columnar::Constant(lengths) => { - let length = constant_unsigned_usize(&lengths)?; + let length = constant_unsigned_usize(&lengths); take_piecewise_sequence_constant_dispatch( array, &starts, @@ -179,7 +177,8 @@ fn take_piecewise_sequence( output_len, )? } - Columnar::Canonical(Canonical::Primitive(lengths)) => { + Columnar::Canonical(lengths) => { + let lengths = lengths.into_primitive(); take_piecewise_sequence_lengths_dispatch( array, &starts, @@ -189,12 +188,6 @@ fn take_piecewise_sequence( output_len, )? } - Columnar::Canonical(lengths) => { - vortex_bail!( - "PiecewiseSequenceArray lengths must be primitive or constant, got {}", - lengths.dtype() - ) - } }; Ok(Some(taken)) } @@ -329,8 +322,7 @@ where computed_len == output_len, "PiecewiseSequenceArray expanded length {computed_len} does not match declared length {output_len}" ); - let total_elements = - piecewise_list_elements_len_constant(array.elements().len(), offsets, starts, length)?; + let total_elements = piecewise_list_elements_len_constant(offsets, starts, length)?; let validity = array.validity()?.take(indices_ref)?; match_smallest_offset_type!(total_elements, |OutputOffset| { @@ -377,8 +369,7 @@ where computed_len == output_len, "PiecewiseSequenceArray expanded length {computed_len} does not match declared length {output_len}" ); - let total_elements = - piecewise_list_elements_len(array.elements().len(), offsets, starts, lengths)?; + let total_elements = piecewise_list_elements_len(offsets, starts, lengths)?; match_smallest_offset_type!(total_elements, |OutputOffset| { let gathered = gather_piecewise_list::( @@ -407,7 +398,6 @@ struct GatheredList { } fn piecewise_list_elements_len_constant( - elements_len: usize, offsets: &[Offset], starts: &[S], length: usize, @@ -426,10 +416,6 @@ where let offset_range = &offsets[start..][..=length]; let element_start: usize = offset_range[0].as_(); let element_end: usize = offset_range[length].as_(); - vortex_ensure!( - element_start <= element_end && element_end <= elements_len, - "List offsets range {element_start}..{element_end} exceeds elements length {elements_len}", - ); total = total .checked_add(element_end - element_start) .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; @@ -438,7 +424,6 @@ where } fn piecewise_list_elements_len( - elements_len: usize, offsets: &[Offset], starts: &[S], lengths: &[L], @@ -455,10 +440,6 @@ where let offset_range = &offsets[start..][..=length]; let element_start: usize = offset_range[0].as_(); let element_end: usize = offset_range[length].as_(); - vortex_ensure!( - element_start <= element_end && element_end <= elements_len, - "List offsets range {element_start}..{element_end} exceeds elements length {elements_len}", - ); total = total .checked_add(element_end - element_start) .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; @@ -499,21 +480,15 @@ where let element_end: usize = offset_range[length].as_(); for &offset in &offset_range[1..] { let offset: usize = offset.as_(); - let relative = offset - .checked_sub(element_start) - .ok_or_else(|| vortex_err!("List offsets are not monotonic at offset {offset}"))?; - let output_offset = output_elements - .checked_add(relative) - .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; - new_offsets.push(new_offset_value::(output_offset)?); + let relative = offset - element_start; + let output_offset = output_elements + relative; + new_offsets.push(new_offset_value::(output_offset)); } let element_length = element_end - element_start; element_starts.push(element_start as u64); element_lengths.push(element_length as u64); - output_elements = output_elements - .checked_add(element_length) - .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; + output_elements += element_length; } debug_assert_eq!(output_elements, total_elements); @@ -569,21 +544,15 @@ where let element_end: usize = offset_range[length].as_(); for &offset in &offset_range[1..] { let offset: usize = offset.as_(); - let relative = offset - .checked_sub(element_start) - .ok_or_else(|| vortex_err!("List offsets are not monotonic at offset {offset}"))?; - let output_offset = output_elements - .checked_add(relative) - .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; - new_offsets.push(new_offset_value::(output_offset)?); + let relative = offset - element_start; + let output_offset = output_elements + relative; + new_offsets.push(new_offset_value::(output_offset)); } let element_length = element_end - element_start; element_starts.push(element_start as u64); element_lengths.push(element_length as u64); - output_elements = output_elements - .checked_add(element_length) - .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; + output_elements += element_length; } debug_assert_eq!(output_elements, total_elements); @@ -604,13 +573,8 @@ where Ok(GatheredList { elements, offsets }) } -fn new_offset_value(value: usize) -> VortexResult { - T::from_usize(value).ok_or_else(|| { - vortex_err!( - "List take offset value {value} does not fit in {}", - T::PTYPE - ) - }) +fn new_offset_value(value: usize) -> T { + T::from_usize(value).vortex_expect("output offset fits selected offset type") } // Kept out-of-line: as a single-callsite generic helper it would otherwise be inlined into every From 9ffb47ac792ef10559eeec680b351dc8f2212fb6 Mon Sep 17 00:00:00 2001 From: Daniel King Date: Mon, 20 Jul 2026 15:40:29 -0400 Subject: [PATCH 12/18] Rename List take slice dispatch helpers Signed-off-by: Daniel King --- vortex-array/src/arrays/list/compute/take.rs | 46 +++++++------------- 1 file changed, 16 insertions(+), 30 deletions(-) diff --git a/vortex-array/src/arrays/list/compute/take.rs b/vortex-array/src/arrays/list/compute/take.rs index e141e8cf786..41ab54989d2 100644 --- a/vortex-array/src/arrays/list/compute/take.rs +++ b/vortex-array/src/arrays/list/compute/take.rs @@ -50,7 +50,7 @@ impl TakeExecute for List { ctx: &mut ExecutionCtx, ) -> VortexResult> { if let Some(piecewise_indices) = indices.as_opt::() - && let Some(taken) = take_piecewise_sequence(array, piecewise_indices, indices, ctx)? + && let Some(taken) = take_slices(array, piecewise_indices, indices, ctx)? { return Ok(Some(taken)); } @@ -145,7 +145,7 @@ fn _take( .into_array()) } -fn take_piecewise_sequence( +fn take_slices( array: ArrayView<'_, List>, indices: ArrayView<'_, PiecewiseSequence>, indices_ref: &ArrayRef, @@ -168,7 +168,7 @@ fn take_piecewise_sequence( let taken = match lengths { Columnar::Constant(lengths) => { let length = constant_unsigned_usize(&lengths); - take_piecewise_sequence_constant_dispatch( + take_slices_constant_start_dispatch( array, &starts, length, @@ -179,20 +179,13 @@ fn take_piecewise_sequence( } Columnar::Canonical(lengths) => { let lengths = lengths.into_primitive(); - take_piecewise_sequence_lengths_dispatch( - array, - &starts, - &lengths, - &offsets, - indices_ref, - output_len, - )? + take_slices_start_dispatch(array, &starts, &lengths, &offsets, indices_ref, output_len)? } }; Ok(Some(taken)) } -fn take_piecewise_sequence_constant_dispatch( +fn take_slices_constant_start_dispatch( array: ArrayView<'_, List>, starts: &PrimitiveArray, length: usize, @@ -201,7 +194,7 @@ fn take_piecewise_sequence_constant_dispatch( output_len: usize, ) -> VortexResult { match_each_unsigned_integer_ptype!(starts.ptype(), |S| { - take_piecewise_sequence_constant_start_dispatch::( + take_slices_constant_offset_dispatch::( array, starts, length, @@ -212,7 +205,7 @@ fn take_piecewise_sequence_constant_dispatch( }) } -fn take_piecewise_sequence_constant_start_dispatch( +fn take_slices_constant_offset_dispatch( array: ArrayView<'_, List>, starts: &PrimitiveArray, length: usize, @@ -224,7 +217,7 @@ where S: UnsignedPType, { match_each_unsigned_integer_ptype!(offsets.ptype(), |O| { - take_piecewise_sequence_constant_length::( + take_slices_constant_length::( array, starts.as_slice::(), length, @@ -235,7 +228,7 @@ where }) } -fn take_piecewise_sequence_lengths_dispatch( +fn take_slices_start_dispatch( array: ArrayView<'_, List>, starts: &PrimitiveArray, lengths: &PrimitiveArray, @@ -244,18 +237,11 @@ fn take_piecewise_sequence_lengths_dispatch( output_len: usize, ) -> VortexResult { match_each_unsigned_integer_ptype!(starts.ptype(), |S| { - take_piecewise_sequence_lengths_start_dispatch::( - array, - starts, - lengths, - offsets, - indices_ref, - output_len, - ) + take_slices_length_dispatch::(array, starts, lengths, offsets, indices_ref, output_len) }) } -fn take_piecewise_sequence_lengths_start_dispatch( +fn take_slices_length_dispatch( array: ArrayView<'_, List>, starts: &PrimitiveArray, lengths: &PrimitiveArray, @@ -267,7 +253,7 @@ where S: UnsignedPType, { match_each_unsigned_integer_ptype!(lengths.ptype(), |L| { - take_piecewise_sequence_lengths_start_length_dispatch::( + take_slices_offset_dispatch::( array, starts, lengths, @@ -278,7 +264,7 @@ where }) } -fn take_piecewise_sequence_lengths_start_length_dispatch( +fn take_slices_offset_dispatch( array: ArrayView<'_, List>, starts: &PrimitiveArray, lengths: &PrimitiveArray, @@ -291,7 +277,7 @@ where L: UnsignedPType, { match_each_unsigned_integer_ptype!(offsets.ptype(), |O| { - take_piecewise_sequence_typed::( + take_slices_typed::( array, starts.as_slice::(), lengths.as_slice::(), @@ -302,7 +288,7 @@ where }) } -fn take_piecewise_sequence_constant_length( +fn take_slices_constant_length( array: ArrayView<'_, List>, starts: &[S], length: usize, @@ -345,7 +331,7 @@ where }) } -fn take_piecewise_sequence_typed( +fn take_slices_typed( array: ArrayView<'_, List>, starts: &[S], lengths: &[L], From 9c642f54e4184c0a7f1cbe58f4404c9824533a81 Mon Sep 17 00:00:00 2001 From: Daniel King Date: Mon, 20 Jul 2026 15:57:06 -0400 Subject: [PATCH 13/18] Preserve ListView take metadata Signed-off-by: Daniel King --- .../src/arrays/listview/compute/take.rs | 14 ++++------- .../src/arrays/listview/tests/take.rs | 23 ------------------- 2 files changed, 4 insertions(+), 33 deletions(-) diff --git a/vortex-array/src/arrays/listview/compute/take.rs b/vortex-array/src/arrays/listview/compute/take.rs index 01ffb80dab7..65c2a06818a 100644 --- a/vortex-array/src/arrays/listview/compute/take.rs +++ b/vortex-array/src/arrays/listview/compute/take.rs @@ -50,20 +50,14 @@ fn apply_take(array: ArrayView<'_, ListView>, indices: &ArrayRef) -> VortexResul // duplicates. let nullable_new_offsets = offsets.take(indices.clone())?; let nullable_new_sizes = sizes.take(indices.clone())?; - let validity_array = new_validity.to_array(indices.len()); - // Null output rows may carry arbitrary physical offset/size payloads from either the indices - // or the source rows. Mask with the final validity before filling so metadata placeholders are - // safe and non-nullable. + // `take` returns nullable arrays; cast back to non-nullable (filling with zeros to represent + // the null lists caused by null indices; the validity mask tracks nullness separately). let new_offsets = match_each_integer_ptype!(nullable_new_offsets.dtype().as_ptype(), |O| { - nullable_new_offsets - .mask(validity_array.clone())? - .fill_null(Scalar::primitive(O::zero(), Nullability::NonNullable))? + nullable_new_offsets.fill_null(Scalar::primitive(O::zero(), Nullability::NonNullable))? }); let new_sizes = match_each_integer_ptype!(nullable_new_sizes.dtype().as_ptype(), |S| { - nullable_new_sizes - .mask(validity_array)? - .fill_null(Scalar::primitive(S::zero(), Nullability::NonNullable))? + nullable_new_sizes.fill_null(Scalar::primitive(S::zero(), Nullability::NonNullable))? }); // SAFETY: Take operation maintains all `ListViewArray` invariants: diff --git a/vortex-array/src/arrays/listview/tests/take.rs b/vortex-array/src/arrays/listview/tests/take.rs index 8fb7a4af1bb..2af397e262f 100644 --- a/vortex-array/src/arrays/listview/tests/take.rs +++ b/vortex-array/src/arrays/listview/tests/take.rs @@ -101,29 +101,6 @@ fn test_take_with_gaps() { ); } -#[test] -fn test_take_null_source_row_zeros_offset_size_payloads() { - let elements = buffer![1i32, 2].into_array(); - let offsets = buffer![0u32, 999].into_array(); - let sizes = buffer![2u32, 999].into_array(); - let validity = Validity::from_iter([true, false]); - let listview = - unsafe { ListViewArray::new_unchecked(elements, offsets, sizes, validity) }.into_array(); - - let result = listview.take(buffer![1u32].into_array()).unwrap(); - let result_list = result - .execute::(&mut SESSION.create_execution_ctx()) - .unwrap(); - - assert_eq!(result_list.offset_at(0), 0); - assert_eq!(result_list.size_at(0), 0); - assert!( - result_list - .is_invalid(0, &mut SESSION.create_execution_ctx()) - .unwrap() - ); -} - #[test] fn test_take_constant_arrays() { // ListView-specific: Test with ConstantArray for offsets/sizes. From 7cb426aa5b9758d3053ddf16380cf192b16b30d2 Mon Sep 17 00:00:00 2001 From: Daniel King Date: Mon, 20 Jul 2026 16:07:17 -0400 Subject: [PATCH 14/18] Use valid null ListView ranges in rebuild test Signed-off-by: Daniel King --- vortex-array/src/arrays/listview/rebuild.rs | 24 ++++++++++----------- 1 file changed, 11 insertions(+), 13 deletions(-) diff --git a/vortex-array/src/arrays/listview/rebuild.rs b/vortex-array/src/arrays/listview/rebuild.rs index 6e36958522e..2db3e925f49 100644 --- a/vortex-array/src/arrays/listview/rebuild.rs +++ b/vortex-array/src/arrays/listview/rebuild.rs @@ -545,38 +545,36 @@ mod tests { } #[test] - fn test_rebuild_flatten_null_row_ignores_invalid_range_payload() -> VortexResult<()> { - let elements = PrimitiveArray::from_iter(vec![1i32, 2, 3]).into_array(); - let offsets = PrimitiveArray::from_iter(vec![0u32, 999, 2]).into_array(); - let sizes = PrimitiveArray::from_iter(vec![2u32, 999, 1]).into_array(); + fn test_rebuild_flatten_null_row_uses_valid_empty_range() -> VortexResult<()> { + let elements = PrimitiveArray::from_iter(vec![1i32, 2, 3, 4]).into_array(); + let offsets = PrimitiveArray::from_iter(vec![0u32, 1, 2]).into_array(); + let sizes = PrimitiveArray::from_iter(vec![1u32, 2, 2]).into_array(); let validity = Validity::from_iter([true, false, true]); - // SAFETY: this intentionally models a null row whose physical offset and size payloads are - // invalid. Rebuild must ignore those payloads and only emit safe placeholder ranges for the - // null row. + // SAFETY: all source ranges are valid, including the null row's non-empty range. let listview = unsafe { ListViewArray::new_unchecked(elements, offsets, sizes, validity) }; let mut ctx = SESSION.create_execution_ctx(); let flattened = listview.rebuild(ListViewRebuildMode::MakeZeroCopyToList, &mut ctx)?; assert_eq!(flattened.offset_at(0), 0); - assert_eq!(flattened.size_at(0), 2); - assert_eq!(flattened.offset_at(1), 2); + assert_eq!(flattened.size_at(0), 1); + assert_eq!(flattened.offset_at(1), 1); assert_eq!(flattened.size_at(1), 0); - assert_eq!(flattened.offset_at(2), 2); - assert_eq!(flattened.size_at(2), 1); + assert_eq!(flattened.offset_at(2), 1); + assert_eq!(flattened.size_at(2), 2); assert!(flattened.validity()?.execute_is_valid(0, &mut ctx)?); assert!(!flattened.validity()?.execute_is_valid(1, &mut ctx)?); assert!(flattened.validity()?.execute_is_valid(2, &mut ctx)?); assert_arrays_eq!( flattened.list_elements_at(0)?, - PrimitiveArray::from_iter([1i32, 2]), + PrimitiveArray::from_iter([1i32]), &mut ctx ); assert_arrays_eq!( flattened.list_elements_at(2)?, - PrimitiveArray::from_iter([3i32]), + PrimitiveArray::from_iter([3i32, 4]), &mut ctx ); Ok(()) From 2ef9e0909f8744aff9f1a5fa8fbd4300a2ecda11 Mon Sep 17 00:00:00 2001 From: Daniel King Date: Mon, 20 Jul 2026 16:12:49 -0400 Subject: [PATCH 15/18] Restore ListView take implementation Signed-off-by: Daniel King --- vortex-array/src/arrays/listview/compute/take.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/vortex-array/src/arrays/listview/compute/take.rs b/vortex-array/src/arrays/listview/compute/take.rs index 65c2a06818a..2b6c016d2c3 100644 --- a/vortex-array/src/arrays/listview/compute/take.rs +++ b/vortex-array/src/arrays/listview/compute/take.rs @@ -52,7 +52,7 @@ fn apply_take(array: ArrayView<'_, ListView>, indices: &ArrayRef) -> VortexResul let nullable_new_sizes = sizes.take(indices.clone())?; // `take` returns nullable arrays; cast back to non-nullable (filling with zeros to represent - // the null lists caused by null indices; the validity mask tracks nullness separately). + // the null lists — the validity mask tracks nullness separately). let new_offsets = match_each_integer_ptype!(nullable_new_offsets.dtype().as_ptype(), |O| { nullable_new_offsets.fill_null(Scalar::primitive(O::zero(), Nullability::NonNullable))? }); From 15dfce890dafde41abdcd0186189ec50ab2f5791 Mon Sep 17 00:00:00 2001 From: Daniel King Date: Mon, 20 Jul 2026 16:16:15 -0400 Subject: [PATCH 16/18] Keep ListView rebuild helper name Signed-off-by: Daniel King --- vortex-array/src/arrays/listview/rebuild.rs | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/vortex-array/src/arrays/listview/rebuild.rs b/vortex-array/src/arrays/listview/rebuild.rs index 2db3e925f49..cee2972ec5d 100644 --- a/vortex-array/src/arrays/listview/rebuild.rs +++ b/vortex-array/src/arrays/listview/rebuild.rs @@ -141,10 +141,10 @@ impl ListViewArray { // for sizes as well. match_each_unsigned_integer_ptype!(sizes_ptype.to_unsigned(), |S| { match offsets_ptype.to_unsigned() { - PType::U8 => self.rebuild_with_take_or_piecewise::(ctx), - PType::U16 => self.rebuild_with_take_or_piecewise::(ctx), - PType::U32 => self.rebuild_with_take_or_piecewise::(ctx), - PType::U64 => self.rebuild_with_take_or_piecewise::(ctx), + PType::U8 => self.naive_rebuild::(ctx), + PType::U16 => self.naive_rebuild::(ctx), + PType::U32 => self.naive_rebuild::(ctx), + PType::U64 => self.naive_rebuild::(ctx), _ => unreachable!("invalid offsets PType"), } }) @@ -152,7 +152,7 @@ impl ListViewArray { /// Picks between [`rebuild_with_take`](Self::rebuild_with_take) and /// [`rebuild_with_piecewise`](Self::rebuild_with_piecewise) based on average list size. - fn rebuild_with_take_or_piecewise( + fn naive_rebuild( &self, ctx: &mut ExecutionCtx, ) -> VortexResult { From 9386ccb3048b1c5ec95dda83fb09116f6acc724e Mon Sep 17 00:00:00 2001 From: Daniel King Date: Mon, 20 Jul 2026 16:41:21 -0400 Subject: [PATCH 17/18] Handle nullable lists in PiecewiseSequence take Signed-off-by: Daniel King --- vortex-array/src/arrays/list/compute/take.rs | 437 +++++++++++++++++-- 1 file changed, 410 insertions(+), 27 deletions(-) diff --git a/vortex-array/src/arrays/list/compute/take.rs b/vortex-array/src/arrays/list/compute/take.rs index 41ab54989d2..2dae8b78bc4 100644 --- a/vortex-array/src/arrays/list/compute/take.rs +++ b/vortex-array/src/arrays/list/compute/take.rs @@ -7,6 +7,7 @@ use vortex_error::VortexExpect; use vortex_error::VortexResult; use vortex_error::vortex_ensure; use vortex_error::vortex_err; +use vortex_mask::Mask; use crate::ArrayRef; use crate::Columnar; @@ -151,16 +152,12 @@ fn take_slices( indices_ref: &ArrayRef, ctx: &mut ExecutionCtx, ) -> VortexResult> { - let data_validity = array - .list_validity() - .execute_mask(array.as_ref().len(), ctx)?; - if !data_validity.all_true() { - return Ok(None); - } - let Some((starts, lengths)) = maybe_contiguous_slices(indices, ctx)? else { return Ok(None); }; + let data_validity = array + .list_validity() + .execute_mask(array.as_ref().len(), ctx)?; let offsets = array.offsets().clone().execute::(ctx)?; let offsets = offsets.reinterpret_cast(offsets.ptype().to_unsigned()); let output_len = indices_ref.len(); @@ -175,11 +172,20 @@ fn take_slices( &offsets, indices_ref, output_len, + &data_validity, )? } Columnar::Canonical(lengths) => { let lengths = lengths.into_primitive(); - take_slices_start_dispatch(array, &starts, &lengths, &offsets, indices_ref, output_len)? + take_slices_start_dispatch( + array, + &starts, + &lengths, + &offsets, + indices_ref, + output_len, + &data_validity, + )? } }; Ok(Some(taken)) @@ -192,6 +198,7 @@ fn take_slices_constant_start_dispatch( offsets: &PrimitiveArray, indices_ref: &ArrayRef, output_len: usize, + data_validity: &Mask, ) -> VortexResult { match_each_unsigned_integer_ptype!(starts.ptype(), |S| { take_slices_constant_offset_dispatch::( @@ -201,6 +208,7 @@ fn take_slices_constant_start_dispatch( offsets, indices_ref, output_len, + data_validity, ) }) } @@ -212,6 +220,7 @@ fn take_slices_constant_offset_dispatch( offsets: &PrimitiveArray, indices_ref: &ArrayRef, output_len: usize, + data_validity: &Mask, ) -> VortexResult where S: UnsignedPType, @@ -224,6 +233,7 @@ where offsets.as_slice::(), indices_ref, output_len, + data_validity, ) }) } @@ -235,9 +245,18 @@ fn take_slices_start_dispatch( offsets: &PrimitiveArray, indices_ref: &ArrayRef, output_len: usize, + data_validity: &Mask, ) -> VortexResult { match_each_unsigned_integer_ptype!(starts.ptype(), |S| { - take_slices_length_dispatch::(array, starts, lengths, offsets, indices_ref, output_len) + take_slices_length_dispatch::( + array, + starts, + lengths, + offsets, + indices_ref, + output_len, + data_validity, + ) }) } @@ -248,6 +267,7 @@ fn take_slices_length_dispatch( offsets: &PrimitiveArray, indices_ref: &ArrayRef, output_len: usize, + data_validity: &Mask, ) -> VortexResult where S: UnsignedPType, @@ -260,6 +280,7 @@ where offsets, indices_ref, output_len, + data_validity, ) }) } @@ -271,6 +292,7 @@ fn take_slices_offset_dispatch( offsets: &PrimitiveArray, indices_ref: &ArrayRef, output_len: usize, + data_validity: &Mask, ) -> VortexResult where S: UnsignedPType, @@ -284,6 +306,7 @@ where offsets.as_slice::(), indices_ref, output_len, + data_validity, ) }) } @@ -295,6 +318,7 @@ fn take_slices_constant_length( offsets: &[Offset], indices_ref: &ArrayRef, output_len: usize, + data_validity: &Mask, ) -> VortexResult where S: UnsignedPType, @@ -308,18 +332,35 @@ where computed_len == output_len, "PiecewiseSequenceArray expanded length {computed_len} does not match declared length {output_len}" ); - let total_elements = piecewise_list_elements_len_constant(offsets, starts, length)?; + let all_valid = data_validity.all_true(); + let total_elements = if all_valid { + piecewise_list_elements_len_constant(offsets, starts, length)? + } else { + piecewise_list_elements_len_constant_validity(offsets, starts, length, data_validity)? + }; let validity = array.validity()?.take(indices_ref)?; match_smallest_offset_type!(total_elements, |OutputOffset| { - let gathered = gather_piecewise_list_constant_length::( - array.elements(), - offsets, - starts, - length, - output_len, - total_elements, - )?; + let gathered = if all_valid { + gather_piecewise_list_constant_length::( + array.elements(), + offsets, + starts, + length, + output_len, + total_elements, + )? + } else { + gather_piecewise_list_constant_length_validity::( + array.elements(), + offsets, + starts, + length, + output_len, + total_elements, + data_validity, + )? + }; // SAFETY: output offsets are rebuilt from valid monotonic source offsets; output elements // are exactly the gathered child ranges referenced by those offsets; validity has one bit @@ -338,6 +379,7 @@ fn take_slices_typed( offsets: &[Offset], indices_ref: &ArrayRef, output_len: usize, + data_validity: &Mask, ) -> VortexResult where S: UnsignedPType, @@ -355,17 +397,34 @@ where computed_len == output_len, "PiecewiseSequenceArray expanded length {computed_len} does not match declared length {output_len}" ); - let total_elements = piecewise_list_elements_len(offsets, starts, lengths)?; + let all_valid = data_validity.all_true(); + let total_elements = if all_valid { + piecewise_list_elements_len(offsets, starts, lengths)? + } else { + piecewise_list_elements_len_validity(offsets, starts, lengths, data_validity)? + }; match_smallest_offset_type!(total_elements, |OutputOffset| { - let gathered = gather_piecewise_list::( - array.elements(), - offsets, - starts, - lengths, - output_len, - total_elements, - )?; + let gathered = if all_valid { + gather_piecewise_list::( + array.elements(), + offsets, + starts, + lengths, + output_len, + total_elements, + )? + } else { + gather_piecewise_list_validity::( + array.elements(), + offsets, + starts, + lengths, + output_len, + total_elements, + data_validity, + )? + }; let validity = array.validity()?.take(indices_ref)?; // SAFETY: output offsets are rebuilt from valid monotonic source offsets; output elements @@ -409,6 +468,31 @@ where Ok(total) } +fn piecewise_list_elements_len_constant_validity( + offsets: &[Offset], + starts: &[S], + length: usize, + data_validity: &Mask, +) -> VortexResult +where + S: UnsignedPType, + Offset: UnsignedPType, +{ + if length == 0 { + return Ok(0); + } + + let mut total = 0usize; + for &start in starts { + let start: usize = start.as_(); + let additional = valid_piece_elements_len(offsets, data_validity, start, length)?; + total = total + .checked_add(additional) + .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; + } + Ok(total) +} + fn piecewise_list_elements_len( offsets: &[Offset], starts: &[S], @@ -433,6 +517,53 @@ where Ok(total) } +fn piecewise_list_elements_len_validity( + offsets: &[Offset], + starts: &[S], + lengths: &[L], + data_validity: &Mask, +) -> VortexResult +where + S: UnsignedPType, + L: UnsignedPType, + Offset: UnsignedPType, +{ + let mut total = 0usize; + for (&start, &length) in starts.iter().zip_eq(lengths) { + let start: usize = start.as_(); + let length: usize = length.as_(); + let additional = valid_piece_elements_len(offsets, data_validity, start, length)?; + total = total + .checked_add(additional) + .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; + } + Ok(total) +} + +fn valid_piece_elements_len( + offsets: &[Offset], + data_validity: &Mask, + start: usize, + length: usize, +) -> VortexResult +where + Offset: UnsignedPType, +{ + let offset_range = &offsets[start..][..=length]; + let mut total = 0usize; + for (data_idx, window) in (start..).zip(offset_range.windows(2)) { + if !data_validity.value(data_idx) { + continue; + } + let element_start: usize = window[0].as_(); + let element_end: usize = window[1].as_(); + total = total + .checked_add(element_end - element_start) + .ok_or_else(|| vortex_err!("List take output elements length overflow"))?; + } + Ok(total) +} + fn gather_piecewise_list_constant_length( elements: &ArrayRef, offsets: &[Offset], @@ -495,6 +626,65 @@ where Ok(GatheredList { elements, offsets }) } +fn gather_piecewise_list_constant_length_validity( + elements: &ArrayRef, + offsets: &[Offset], + starts: &[S], + length: usize, + output_len: usize, + total_elements: usize, + data_validity: &Mask, +) -> VortexResult +where + S: UnsignedPType, + Offset: UnsignedPType, + OutputOffset: IntegerPType, +{ + let offsets_capacity = output_len + .checked_add(1) + .ok_or_else(|| vortex_err!("List take offsets length overflow"))?; + let mut new_offsets = BufferMut::::with_capacity(offsets_capacity); + let mut element_starts = BufferMut::::with_capacity(output_len); + let mut element_lengths = BufferMut::::with_capacity(output_len); + let mut output_elements = 0usize; + + new_offsets.push(OutputOffset::zero()); + for &start in starts { + let start: usize = start.as_(); + if length == 0 { + continue; + } + + gather_valid_piece( + offsets, + data_validity, + start, + length, + &mut new_offsets, + &mut element_starts, + &mut element_lengths, + &mut output_elements, + ); + } + debug_assert_eq!(output_elements, total_elements); + + let offsets = PrimitiveArray::new(new_offsets.freeze(), Validity::NonNullable).into_array(); + let multipliers = ConstantArray::new(1u64, element_starts.len()).into_array(); + // SAFETY: element ranges come only from valid source list rows. Source list construction + // validated those row offsets, and null source rows produce no element range. + let element_indices = unsafe { + PiecewiseSequenceArray::new_unchecked( + element_starts.into_array(), + element_lengths.into_array(), + multipliers, + total_elements, + ) + }; + let elements = elements.take(element_indices.into_array())?; + + Ok(GatheredList { elements, offsets }) +} + fn gather_piecewise_list( elements: &ArrayRef, offsets: &[Offset], @@ -559,6 +749,99 @@ where Ok(GatheredList { elements, offsets }) } +fn gather_piecewise_list_validity( + elements: &ArrayRef, + offsets: &[Offset], + starts: &[S], + lengths: &[L], + output_len: usize, + total_elements: usize, + data_validity: &Mask, +) -> VortexResult +where + S: UnsignedPType, + L: UnsignedPType, + Offset: UnsignedPType, + OutputOffset: IntegerPType, +{ + let offsets_capacity = output_len + .checked_add(1) + .ok_or_else(|| vortex_err!("List take offsets length overflow"))?; + let mut new_offsets = BufferMut::::with_capacity(offsets_capacity); + let mut element_starts = BufferMut::::with_capacity(output_len); + let mut element_lengths = BufferMut::::with_capacity(output_len); + let mut output_elements = 0usize; + + new_offsets.push(OutputOffset::zero()); + for (&start, &length) in starts.iter().zip_eq(lengths) { + let start: usize = start.as_(); + let length: usize = length.as_(); + if length == 0 { + continue; + } + + gather_valid_piece( + offsets, + data_validity, + start, + length, + &mut new_offsets, + &mut element_starts, + &mut element_lengths, + &mut output_elements, + ); + } + debug_assert_eq!(output_elements, total_elements); + + let offsets = PrimitiveArray::new(new_offsets.freeze(), Validity::NonNullable).into_array(); + let multipliers = ConstantArray::new(1u64, element_starts.len()).into_array(); + // SAFETY: element ranges come only from valid source list rows. Source list construction + // validated those row offsets, and null source rows produce no element range. + let element_indices = unsafe { + PiecewiseSequenceArray::new_unchecked( + element_starts.into_array(), + element_lengths.into_array(), + multipliers, + total_elements, + ) + }; + let elements = elements.take(element_indices.into_array())?; + + Ok(GatheredList { elements, offsets }) +} + +fn gather_valid_piece( + offsets: &[Offset], + data_validity: &Mask, + start: usize, + length: usize, + new_offsets: &mut BufferMut, + element_starts: &mut BufferMut, + element_lengths: &mut BufferMut, + output_elements: &mut usize, +) where + Offset: UnsignedPType, + OutputOffset: IntegerPType, +{ + let offset_range = &offsets[start..][..=length]; + for (data_idx, window) in (start..).zip(offset_range.windows(2)) { + if !data_validity.value(data_idx) { + new_offsets.push(new_offset_value::(*output_elements)); + continue; + } + + let element_start: usize = window[0].as_(); + let element_end: usize = window[1].as_(); + let element_length = element_end - element_start; + if element_length != 0 { + element_starts.push(element_start as u64); + element_lengths.push(element_length as u64); + *output_elements += element_length; + } + new_offsets.push(new_offset_value::(*output_elements)); + } +} + fn new_offset_value(value: usize) -> T { T::from_usize(value).vortex_expect("output offset fits selected offset type") } @@ -646,6 +929,7 @@ mod test { use rstest::rstest; use vortex_buffer::buffer; + use vortex_error::VortexResult; use crate::IntoArray as _; use crate::VortexSessionExecute; @@ -656,6 +940,7 @@ mod test { use crate::arrays::ListViewArray; use crate::arrays::PiecewiseSequenceArray; use crate::arrays::PrimitiveArray; + use crate::arrays::listview::ListViewArrayExt; use crate::compute::conformance::take::test_take_conformance; use crate::dtype::DType; use crate::dtype::Nullability; @@ -889,6 +1174,104 @@ mod test { ); } + #[test] + fn piecewise_sequence_take_nullable_list_constant_lengths() -> VortexResult<()> { + let mut ctx = array_session().create_execution_ctx(); + let list = ListArray::try_new( + buffer![0i32, 1, 99, 100, 2, 3, 4, 5].into_array(), + buffer![0u32, 2, 4, 7, 8].into_array(), + Validity::Array(BoolArray::from_iter([true, false, true, true]).into_array()), + )? + .into_array(); + let idx = PiecewiseSequenceArray::try_new( + buffer![0u64].into_array(), + ConstantArray::new(4u64, 1).into_array(), + ConstantArray::new(1u64, 1).into_array(), + 4, + )? + .into_array(); + + let result = list.take(idx)?.execute::(&mut ctx)?; + assert_eq!(result.offset_at(0), 0); + assert_eq!(result.size_at(0), 2); + assert_eq!(result.offset_at(1), 2); + assert_eq!(result.size_at(1), 0); + assert_eq!(result.offset_at(2), 2); + assert_eq!(result.size_at(2), 3); + assert_eq!(result.offset_at(3), 5); + assert_eq!(result.size_at(3), 1); + + let element_dtype: Arc = Arc::new(I32.into()); + assert_eq!( + result.execute_scalar(0, &mut ctx)?, + Scalar::list( + Arc::clone(&element_dtype), + vec![0i32.into(), 1.into()], + Nullability::Nullable + ) + ); + assert!(result.is_invalid(1, &mut ctx)?); + assert_eq!( + result.execute_scalar(2, &mut ctx)?, + Scalar::list( + Arc::clone(&element_dtype), + vec![2i32.into(), 3.into(), 4.into()], + Nullability::Nullable + ) + ); + assert_eq!( + result.execute_scalar(3, &mut ctx)?, + Scalar::list(element_dtype, vec![5i32.into()], Nullability::Nullable) + ); + Ok(()) + } + + #[test] + fn piecewise_sequence_take_nullable_list_array_lengths() -> VortexResult<()> { + let mut ctx = array_session().create_execution_ctx(); + let list = ListArray::try_new( + buffer![0i32, 1, 99, 100, 2, 3, 4, 5].into_array(), + buffer![0u32, 2, 4, 7, 8].into_array(), + Validity::Array(BoolArray::from_iter([true, false, true, true]).into_array()), + )? + .into_array(); + let idx = PiecewiseSequenceArray::try_new( + buffer![1u64, 0].into_array(), + buffer![2u64, 1].into_array(), + ConstantArray::new(1u64, 2).into_array(), + 3, + )? + .into_array(); + + let result = list.take(idx)?.execute::(&mut ctx)?; + assert_eq!(result.offset_at(0), 0); + assert_eq!(result.size_at(0), 0); + assert_eq!(result.offset_at(1), 0); + assert_eq!(result.size_at(1), 3); + assert_eq!(result.offset_at(2), 3); + assert_eq!(result.size_at(2), 2); + + let element_dtype: Arc = Arc::new(I32.into()); + assert!(result.is_invalid(0, &mut ctx)?); + assert_eq!( + result.execute_scalar(1, &mut ctx)?, + Scalar::list( + Arc::clone(&element_dtype), + vec![2i32.into(), 3.into(), 4.into()], + Nullability::Nullable + ) + ); + assert_eq!( + result.execute_scalar(2, &mut ctx)?, + Scalar::list( + element_dtype, + vec![0i32.into(), 1.into()], + Nullability::Nullable + ) + ); + Ok(()) + } + #[test] fn test_take_empty_array() { let list = ListArray::try_new( From c8c6f8151878d25bce932ab58d0a627e3ead07ab Mon Sep 17 00:00:00 2001 From: Daniel King Date: Mon, 20 Jul 2026 17:50:58 -0400 Subject: [PATCH 18/18] Fix list PiecewiseSequence clippy lint Signed-off-by: Daniel King --- vortex-array/src/arrays/list/compute/take.rs | 94 ++++++++++---------- 1 file changed, 45 insertions(+), 49 deletions(-) diff --git a/vortex-array/src/arrays/list/compute/take.rs b/vortex-array/src/arrays/list/compute/take.rs index 2dae8b78bc4..8cb5c8707c3 100644 --- a/vortex-array/src/arrays/list/compute/take.rs +++ b/vortex-array/src/arrays/list/compute/take.rs @@ -442,6 +442,13 @@ struct GatheredList { offsets: ArrayRef, } +struct ValidPieceGather { + new_offsets: BufferMut, + element_starts: BufferMut, + element_lengths: BufferMut, + output_elements: usize, +} + fn piecewise_list_elements_len_constant( offsets: &[Offset], starts: &[S], @@ -643,39 +650,33 @@ where let offsets_capacity = output_len .checked_add(1) .ok_or_else(|| vortex_err!("List take offsets length overflow"))?; - let mut new_offsets = BufferMut::::with_capacity(offsets_capacity); - let mut element_starts = BufferMut::::with_capacity(output_len); - let mut element_lengths = BufferMut::::with_capacity(output_len); - let mut output_elements = 0usize; + let mut gather = ValidPieceGather { + new_offsets: BufferMut::::with_capacity(offsets_capacity), + element_starts: BufferMut::::with_capacity(output_len), + element_lengths: BufferMut::::with_capacity(output_len), + output_elements: 0, + }; - new_offsets.push(OutputOffset::zero()); + gather.new_offsets.push(OutputOffset::zero()); for &start in starts { let start: usize = start.as_(); if length == 0 { continue; } - gather_valid_piece( - offsets, - data_validity, - start, - length, - &mut new_offsets, - &mut element_starts, - &mut element_lengths, - &mut output_elements, - ); + gather_valid_piece(offsets, data_validity, start, length, &mut gather); } - debug_assert_eq!(output_elements, total_elements); + debug_assert_eq!(gather.output_elements, total_elements); - let offsets = PrimitiveArray::new(new_offsets.freeze(), Validity::NonNullable).into_array(); - let multipliers = ConstantArray::new(1u64, element_starts.len()).into_array(); + let offsets = + PrimitiveArray::new(gather.new_offsets.freeze(), Validity::NonNullable).into_array(); + let multipliers = ConstantArray::new(1u64, gather.element_starts.len()).into_array(); // SAFETY: element ranges come only from valid source list rows. Source list construction // validated those row offsets, and null source rows produce no element range. let element_indices = unsafe { PiecewiseSequenceArray::new_unchecked( - element_starts.into_array(), - element_lengths.into_array(), + gather.element_starts.into_array(), + gather.element_lengths.into_array(), multipliers, total_elements, ) @@ -767,12 +768,14 @@ where let offsets_capacity = output_len .checked_add(1) .ok_or_else(|| vortex_err!("List take offsets length overflow"))?; - let mut new_offsets = BufferMut::::with_capacity(offsets_capacity); - let mut element_starts = BufferMut::::with_capacity(output_len); - let mut element_lengths = BufferMut::::with_capacity(output_len); - let mut output_elements = 0usize; + let mut gather = ValidPieceGather { + new_offsets: BufferMut::::with_capacity(offsets_capacity), + element_starts: BufferMut::::with_capacity(output_len), + element_lengths: BufferMut::::with_capacity(output_len), + output_elements: 0, + }; - new_offsets.push(OutputOffset::zero()); + gather.new_offsets.push(OutputOffset::zero()); for (&start, &length) in starts.iter().zip_eq(lengths) { let start: usize = start.as_(); let length: usize = length.as_(); @@ -780,27 +783,19 @@ where continue; } - gather_valid_piece( - offsets, - data_validity, - start, - length, - &mut new_offsets, - &mut element_starts, - &mut element_lengths, - &mut output_elements, - ); + gather_valid_piece(offsets, data_validity, start, length, &mut gather); } - debug_assert_eq!(output_elements, total_elements); + debug_assert_eq!(gather.output_elements, total_elements); - let offsets = PrimitiveArray::new(new_offsets.freeze(), Validity::NonNullable).into_array(); - let multipliers = ConstantArray::new(1u64, element_starts.len()).into_array(); + let offsets = + PrimitiveArray::new(gather.new_offsets.freeze(), Validity::NonNullable).into_array(); + let multipliers = ConstantArray::new(1u64, gather.element_starts.len()).into_array(); // SAFETY: element ranges come only from valid source list rows. Source list construction // validated those row offsets, and null source rows produce no element range. let element_indices = unsafe { PiecewiseSequenceArray::new_unchecked( - element_starts.into_array(), - element_lengths.into_array(), + gather.element_starts.into_array(), + gather.element_lengths.into_array(), multipliers, total_elements, ) @@ -815,10 +810,7 @@ fn gather_valid_piece( data_validity: &Mask, start: usize, length: usize, - new_offsets: &mut BufferMut, - element_starts: &mut BufferMut, - element_lengths: &mut BufferMut, - output_elements: &mut usize, + gather: &mut ValidPieceGather, ) where Offset: UnsignedPType, OutputOffset: IntegerPType, @@ -826,7 +818,9 @@ fn gather_valid_piece( let offset_range = &offsets[start..][..=length]; for (data_idx, window) in (start..).zip(offset_range.windows(2)) { if !data_validity.value(data_idx) { - new_offsets.push(new_offset_value::(*output_elements)); + gather + .new_offsets + .push(new_offset_value::(gather.output_elements)); continue; } @@ -834,11 +828,13 @@ fn gather_valid_piece( let element_end: usize = window[1].as_(); let element_length = element_end - element_start; if element_length != 0 { - element_starts.push(element_start as u64); - element_lengths.push(element_length as u64); - *output_elements += element_length; + gather.element_starts.push(element_start as u64); + gather.element_lengths.push(element_length as u64); + gather.output_elements += element_length; } - new_offsets.push(new_offset_value::(*output_elements)); + gather + .new_offsets + .push(new_offset_value::(gather.output_elements)); } }