diff --git a/encodings/fastlanes/src/delta/array/delta_decompress.rs b/encodings/fastlanes/src/delta/array/delta_decompress.rs index 3ccfede27c1..88881cb35b1 100644 --- a/encodings/fastlanes/src/delta/array/delta_decompress.rs +++ b/encodings/fastlanes/src/delta/array/delta_decompress.rs @@ -15,6 +15,7 @@ use vortex_array::dtype::NativePType; use vortex_array::match_each_unsigned_integer_ptype; use vortex_buffer::Buffer; use vortex_buffer::BufferMut; +use vortex_error::VortexExpect; use vortex_error::VortexResult; use crate::DeltaArray; @@ -54,6 +55,34 @@ pub fn delta_decompress( /// Performs the low-level delta decompression on primitive values. /// /// All chunks must be full 1024-element chunks (deltas length must be a multiple of 1024). +/// Decode one 1024-value chunk into `output`, using `transposed` as scratch. +/// +/// Chunks are independent: chunk `chunk` reads only its own lanes' bases and its own deltas, so +/// kernels can decode just the chunks they touch into cache-resident buffers instead of +/// materializing the whole array. +pub(crate) fn decode_chunk( + bases: &[T], + deltas: &[T], + chunk: usize, + transposed: &mut [T; 1024], + output: &mut [T; 1024], +) where + T: NativePType + Delta + Transpose, +{ + let deltas = deltas[chunk * 1024..(chunk + 1) * 1024] + .as_chunks::<1024>() + .0 + .first() + .vortex_expect("deltas are padded to whole 1024-value chunks"); + let bases = bases[chunk * LANES..(chunk + 1) * LANES] + .as_chunks::() + .0 + .first() + .vortex_expect("bases hold one value per lane per chunk"); + Delta::undelta::(deltas, bases, transposed); + Transpose::untranspose(transposed, output); +} + pub(crate) fn decompress_primitive(bases: &[T], deltas: &[T]) -> Buffer where T: NativePType + Delta + Transpose, diff --git a/encodings/fastlanes/src/delta/compute/compare.rs b/encodings/fastlanes/src/delta/compute/compare.rs new file mode 100644 index 00000000000..a8d79b97076 --- /dev/null +++ b/encodings/fastlanes/src/delta/compute/compare.rs @@ -0,0 +1,219 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use fastlanes::Delta as FastLanesDelta; +use fastlanes::FastLanes; +use fastlanes::Transpose; +use num_traits::AsPrimitive; +use num_traits::PrimInt; +use vortex_array::ArrayRef; +use vortex_array::ArrayView; +use vortex_array::ExecutionCtx; +use vortex_array::IntoArray; +use vortex_array::arrays::BoolArray; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::arrays::primitive::PrimitiveArrayExt; +use vortex_array::dtype::NativePType; +use vortex_array::match_each_integer_ptype; +use vortex_array::match_each_unsigned_integer_ptype; +use vortex_array::scalar_fn::fns::binary::CompareKernel; +use vortex_array::scalar_fn::fns::operators::CompareOperator; +use vortex_buffer::BitBuffer; +use vortex_buffer::BitBufferMut; +use vortex_buffer::BufferMut; +use vortex_error::VortexExpect; +use vortex_error::VortexResult; + +use crate::Delta; +use crate::delta::array::DeltaArrayExt; +use crate::delta::array::DeltaArraySlotsExt; +use crate::delta::array::delta_decompress::decode_chunk; + +const WORDS_PER_CHUNK: usize = 1024 / u64::BITS as usize; + +/// Compares against a constant one chunk at a time, so the decoded values stay in a stack buffer +/// and only the result bits are written out. +impl CompareKernel for Delta { + fn compare( + lhs: ArrayView<'_, Self>, + rhs: &ArrayRef, + operator: CompareOperator, + ctx: &mut ExecutionCtx, + ) -> VortexResult> { + let Some(constant) = rhs.as_constant() else { + return Ok(None); + }; + let Some(constant) = constant.as_primitive_opt() else { + return Ok(None); + }; + let ptype = lhs.dtype().as_ptype(); + if constant.ptype() != ptype { + return Ok(None); + } + + let bits = match_each_integer_ptype!(ptype, |T| { + let rhs: T = constant + .typed_value::() + .vortex_expect("compare adaptor strips null constants"); + match_each_unsigned_integer_ptype!(ptype.to_unsigned(), |U| { + const LANES: usize = U::LANES; + compare_constant::(lhs, rhs.as_(), ptype.is_signed_int(), operator, ctx)? + }) + }); + + let nullability = lhs.dtype().nullability() | rhs.dtype().nullability(); + let validity = lhs.validity()?.union_nullability(nullability); + Ok(Some(BoolArray::new(bits, validity).into_array())) + } +} + +/// Compare every value against `rhs` in the unsigned domain the chunks decode into. +/// +/// Flipping the sign bit maps signed order onto unsigned order, so ordered comparisons of signed +/// values work on their unsigned bit patterns. +fn compare_constant( + array: ArrayView<'_, Delta>, + rhs: U, + signed: bool, + operator: CompareOperator, + ctx: &mut ExecutionCtx, +) -> VortexResult +where + U: NativePType + PrimInt + FastLanesDelta + Transpose, +{ + let bases = array + .bases() + .clone() + .execute::(ctx)? + .reinterpret_cast(U::PTYPE); + let deltas = array + .deltas() + .clone() + .execute::(ctx)? + .reinterpret_cast(U::PTYPE); + let (bases, deltas) = (bases.as_slice::(), deltas.as_slice::()); + + let sign = if signed { + U::one() << (U::PTYPE.bit_width() - 1) + } else { + U::zero() + }; + let rhs = rhs ^ sign; + + let offset = array.offset(); + let len = array.len(); + let num_chunks = (offset + len).div_ceil(1024); + let mut words = BufferMut::::zeroed(num_chunks * WORDS_PER_CHUNK); + let mut transposed = [U::zero(); 1024]; + let mut values = [U::zero(); 1024]; + for (chunk, out) in words + .as_mut_slice() + .as_chunks_mut::() + .0 + .iter_mut() + .enumerate() + { + decode_chunk::(bases, deltas, chunk, &mut transposed, &mut values); + match operator { + CompareOperator::Eq => pack(out, &values, |v| (v ^ sign) == rhs), + CompareOperator::NotEq => pack(out, &values, |v| (v ^ sign) != rhs), + CompareOperator::Lt => pack(out, &values, |v| (v ^ sign) < rhs), + CompareOperator::Lte => pack(out, &values, |v| (v ^ sign) <= rhs), + CompareOperator::Gt => pack(out, &values, |v| (v ^ sign) > rhs), + CompareOperator::Gte => pack(out, &values, |v| (v ^ sign) >= rhs), + } + } + + Ok(BitBufferMut::from_buffer(words.into_byte_buffer(), offset, len).freeze()) +} + +/// Write one bit per value, least significant bit first, 64 values per word. +#[inline] +fn pack( + out: &mut [u64; WORDS_PER_CHUNK], + values: &[U; 1024], + predicate: impl Fn(U) -> bool, +) { + for (word, values) in out.iter_mut().zip(values.as_chunks::<64>().0) { + *word = values.iter().enumerate().fold(0u64, |word, (bit, &value)| { + word | (u64::from(predicate(value)) << bit) + }); + } +} + +#[cfg(test)] +mod tests { + use std::sync::LazyLock; + + use rstest::rstest; + use vortex_array::IntoArray; + use vortex_array::VortexSessionExecute; + use vortex_array::arrays::BoolArray; + use vortex_array::arrays::ConstantArray; + use vortex_array::arrays::PrimitiveArray; + use vortex_array::assert_arrays_eq; + use vortex_array::builtins::ArrayBuiltins; + use vortex_array::scalar_fn::fns::operators::Operator; + use vortex_error::VortexResult; + use vortex_session::VortexSession; + + use crate::Delta; + + static SESSION: LazyLock = LazyLock::new(|| { + let session = vortex_array::array_session(); + crate::initialize(&session); + session + }); + + /// Series-major timestamps: two series that each restart, length not a multiple of 1024. + fn timestamps() -> PrimitiveArray { + PrimitiveArray::from_iter((0..2).flat_map(|_| (0..1500i64).map(|i| i * 30 - 20_000))) + } + + #[rstest] + #[case(Operator::Eq)] + #[case(Operator::NotEq)] + #[case(Operator::Lt)] + #[case(Operator::Lte)] + #[case(Operator::Gt)] + #[case(Operator::Gte)] + fn compare_matches_decoded(#[case] op: Operator) -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + let primitive = timestamps(); + let delta = Delta::try_from_primitive_array(&primitive, &mut ctx)?.into_array(); + // A threshold that is negative, so the signed path matters. + let rhs = ConstantArray::new(-5_000i64, primitive.len()).into_array(); + + let actual = delta + .binary(rhs.clone(), op)? + .execute::(&mut ctx)?; + let expected = primitive + .into_array() + .binary(rhs, op)? + .execute::(&mut ctx)?; + assert_arrays_eq!(actual, expected, &mut ctx); + Ok(()) + } + + #[test] + fn compare_on_slice_and_nulls() -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + let primitive = + PrimitiveArray::from_option_iter((0u32..3000).map(|v| (v % 7 != 0).then_some(v * 3))); + let delta = Delta::try_from_primitive_array(&primitive, &mut ctx)? + .into_array() + .slice(1000..2500)?; + let rhs = ConstantArray::new(5_000u32, delta.len()).into_array(); + + let actual = delta + .binary(rhs.clone(), Operator::Gte)? + .execute::(&mut ctx)?; + let expected = primitive + .into_array() + .slice(1000..2500)? + .binary(rhs, Operator::Gte)? + .execute::(&mut ctx)?; + assert_arrays_eq!(actual, expected, &mut ctx); + Ok(()) + } +} diff --git a/encodings/fastlanes/src/delta/compute/filter.rs b/encodings/fastlanes/src/delta/compute/filter.rs new file mode 100644 index 00000000000..9db888442e7 --- /dev/null +++ b/encodings/fastlanes/src/delta/compute/filter.rs @@ -0,0 +1,279 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use fastlanes::Delta as FastLanesDelta; +use fastlanes::FastLanes; +use fastlanes::Transpose; +use vortex_array::ArrayRef; +use vortex_array::ArrayView; +use vortex_array::ExecutionCtx; +use vortex_array::IntoArray; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::arrays::filter::FilterKernel; +use vortex_array::arrays::primitive::PrimitiveArrayExt; +use vortex_array::dtype::NativePType; +use vortex_array::match_each_unsigned_integer_ptype; +use vortex_buffer::BitBuffer; +use vortex_buffer::Buffer; +use vortex_buffer::BufferMut; +use vortex_error::VortexResult; +use vortex_mask::Mask; +use vortex_mask::MaskValues; + +use crate::Delta; +use crate::delta::array::DeltaArrayExt; +use crate::delta::array::DeltaArraySlotsExt; +use crate::delta::array::delta_decompress::decode_chunk; + +/// When a mask touches at least this fraction of the chunks and selects at least +/// [`FALLBACK_MIN_DENSITY`] of the values, decoding everything and filtering the result with the +/// generic bulk filter is faster than decoding chunk by chunk and gathering. Both thresholds come +/// from benchmarking run, cluster, strided and uniform random masks over 32- and 64-bit values. +const FALLBACK_MIN_TOUCHED: f64 = 0.9; +const FALLBACK_MIN_DENSITY: f64 = 0.15; + +/// Gathers the selected values one chunk at a time, decoding only the chunks that hold a selected +/// value and never materializing the full decoded array. Masks that select a large share of +/// values from nearly every chunk keep the decode-then-filter path, and contiguous masks keep the +/// slice path. +impl FilterKernel for Delta { + fn filter( + array: ArrayView<'_, Self>, + mask: &Mask, + ctx: &mut ExecutionCtx, + ) -> VortexResult> { + let Mask::Values(values) = mask else { + return Ok(None); + }; + let bits = values.bit_buffer(); + // A contiguous mask executes as a slice, which decodes only the selected range. + if is_contiguous(values) { + return Ok(None); + } + let offset = array.offset(); + let total_chunks = (offset + array.len()).div_ceil(1024); + let touched = touched_chunks(bits, offset); + if touched as f64 >= FALLBACK_MIN_TOUCHED * total_chunks as f64 + && values.density() >= FALLBACK_MIN_DENSITY + { + return Ok(None); + } + + let ptype = array.dtype().as_ptype(); + let validity = array.validity()?.filter(mask)?; + let filtered = match_each_unsigned_integer_ptype!(ptype.to_unsigned(), |U| { + const LANES: usize = U::LANES; + let buffer = gather::(array, bits, values.true_count(), ctx)?; + PrimitiveArray::new(buffer, validity) + }); + Ok(Some(filtered.reinterpret_cast(ptype).into_array())) + } +} + +/// Mirrors the filter executor's contiguity probe so declining stays cheap: cached slices or +/// indices answer directly, and otherwise only the candidate run or the tail after it is scanned. +fn is_contiguous(values: &MaskValues) -> bool { + if let Some(slices) = values.cached_slices() { + return slices.len() <= 1; + } + let true_count = values.true_count(); + if let Some(indices) = values.cached_indices() { + return match (indices.first(), indices.last()) { + (Some(first), Some(last)) => last - first + 1 == true_count, + _ => true, + }; + } + let bits = values.bit_buffer(); + let Some(start) = bits.set_indices().next() else { + return true; + }; + let end = start + true_count; + if end > bits.len() { + return false; + } + if end - start <= bits.len() - end { + bits.count_range(start, end) == true_count + } else { + bits.last_set_index() == Some(end - 1) + } +} + +/// Count the chunks that hold at least one selected value. `offset` is the array's position in +/// its first physical chunk. +fn touched_chunks(bits: &BitBuffer, offset: usize) -> usize { + let mut count = 0; + let mut last_counted = None; + for (index, word) in bits.chunks().iter_padded().enumerate() { + if word == 0 { + continue; + } + let base = offset + index * 64; + let first = (base + word.trailing_zeros() as usize) / 1024; + let last = (base + 63 - word.leading_zeros() as usize) / 1024; + for chunk in first..=last { + if last_counted.is_none_or(|counted| chunk > counted) { + count += 1; + last_counted = Some(chunk); + } + } + } + count +} + +/// Walk the mask 64 bits at a time: skip empty words, copy full words in bulk, and pick out the +/// set bits of partial words. Each touched chunk is decoded once. +fn gather( + array: ArrayView<'_, Delta>, + bits: &BitBuffer, + true_count: usize, + ctx: &mut ExecutionCtx, +) -> VortexResult> +where + U: NativePType + FastLanesDelta + Transpose, +{ + let bases = array + .bases() + .clone() + .execute::(ctx)? + .reinterpret_cast(U::PTYPE); + let deltas = array + .deltas() + .clone() + .execute::(ctx)? + .reinterpret_cast(U::PTYPE); + let mut chunks = DecodedChunk::::new(bases.as_slice::(), deltas.as_slice::()); + + let offset = array.offset(); + let mut output = BufferMut::::with_capacity(true_count); + for (index, mut word) in bits.chunks().iter_padded().enumerate() { + if word == 0 { + continue; + } + let base = offset + index * 64; + if word == u64::MAX { + // A full word may straddle a chunk boundary, so copy it in up to two parts. + let (mut position, end) = (base, base + 64); + while position < end { + let chunk = position / 1024; + let chunk_end = end.min((chunk + 1) * 1024); + output.extend_from_slice( + &chunks.get(chunk)[position % 1024..chunk_end - chunk * 1024], + ); + position = chunk_end; + } + } else { + while word != 0 { + let position = base + word.trailing_zeros() as usize; + output.push(chunks.get(position / 1024)[position % 1024]); + word &= word - 1; + } + } + } + Ok(output.freeze()) +} + +/// The most recently decoded chunk, so consecutive reads from one chunk decode it once. +struct DecodedChunk<'a, U, const LANES: usize> { + bases: &'a [U], + deltas: &'a [U], + chunk: Option, + transposed: [U; 1024], + values: [U; 1024], +} + +impl<'a, U, const LANES: usize> DecodedChunk<'a, U, LANES> +where + U: NativePType + FastLanesDelta + Transpose, +{ + fn new(bases: &'a [U], deltas: &'a [U]) -> Self { + Self { + bases, + deltas, + chunk: None, + transposed: [U::default(); 1024], + values: [U::default(); 1024], + } + } + + #[inline] + fn get(&mut self, chunk: usize) -> &[U; 1024] { + if self.chunk != Some(chunk) { + decode_chunk::( + self.bases, + self.deltas, + chunk, + &mut self.transposed, + &mut self.values, + ); + self.chunk = Some(chunk); + } + &self.values + } +} + +#[cfg(test)] +mod tests { + use std::sync::LazyLock; + + use rstest::rstest; + use vortex_array::IntoArray; + use vortex_array::VortexSessionExecute; + use vortex_array::arrays::PrimitiveArray; + use vortex_array::assert_arrays_eq; + use vortex_error::VortexResult; + use vortex_mask::Mask; + use vortex_session::VortexSession; + + use crate::Delta; + + static SESSION: LazyLock = LazyLock::new(|| { + let session = vortex_array::array_session(); + crate::initialize(&session); + session + }); + + #[rstest] + #[case::sparse_scattered(Mask::from_indices(3000, (0..3000).step_by(97)))] + #[case::one_run_slices(Mask::from_slices(3000, vec![(1020, 1100)]))] + #[case::last_chunk(Mask::from_indices(3000, [2047, 2048, 2999]))] + #[case::full_words_across_chunks(Mask::from_slices(3000, vec![(960, 1150), (2000, 2100)]))] + #[case::dense_falls_back(Mask::from_indices(3000, (0..3000).step_by(2)))] + fn filter_matches_decoded(#[case] mask: Mask) -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + let primitive = PrimitiveArray::from_option_iter( + (0..3000i64).map(|v| (v % 11 != 0).then_some(v * 7 - 9_000)), + ); + let delta = Delta::try_from_primitive_array(&primitive, &mut ctx)?.into_array(); + + let actual = delta + .filter(mask.clone())? + .execute::(&mut ctx)?; + let expected = primitive + .into_array() + .filter(mask)? + .execute::(&mut ctx)?; + assert_arrays_eq!(actual, expected, &mut ctx); + Ok(()) + } + + #[test] + fn filter_on_slice() -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + let primitive = PrimitiveArray::from_iter((0..3000u32).map(|v| v * 3)); + let delta = Delta::try_from_primitive_array(&primitive, &mut ctx)? + .into_array() + .slice(700..2900)?; + let mask = Mask::from_indices(delta.len(), (0..delta.len()).step_by(13)); + + let actual = delta + .filter(mask.clone())? + .execute::(&mut ctx)?; + let expected = primitive + .into_array() + .slice(700..2900)? + .filter(mask)? + .execute::(&mut ctx)?; + assert_arrays_eq!(actual, expected, &mut ctx); + Ok(()) + } +} diff --git a/encodings/fastlanes/src/delta/compute/mod.rs b/encodings/fastlanes/src/delta/compute/mod.rs index fa79c62a596..bae66ae1a48 100644 --- a/encodings/fastlanes/src/delta/compute/mod.rs +++ b/encodings/fastlanes/src/delta/compute/mod.rs @@ -2,3 +2,5 @@ // SPDX-FileCopyrightText: Copyright the Vortex contributors mod cast; +mod compare; +mod filter; diff --git a/encodings/fastlanes/src/delta/mod.rs b/encodings/fastlanes/src/delta/mod.rs index 3cf2d4a6815..df64d678796 100644 --- a/encodings/fastlanes/src/delta/mod.rs +++ b/encodings/fastlanes/src/delta/mod.rs @@ -12,4 +12,8 @@ mod compute; mod vtable; pub use vtable::Delta; + +pub(crate) fn initialize(session: &vortex_session::VortexSession) { + vtable::initialize(session); +} pub use vtable::DeltaArray; diff --git a/encodings/fastlanes/src/delta/vtable/kernels.rs b/encodings/fastlanes/src/delta/vtable/kernels.rs new file mode 100644 index 00000000000..6c6345d7eb1 --- /dev/null +++ b/encodings/fastlanes/src/delta/vtable/kernels.rs @@ -0,0 +1,19 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use vortex_array::ArrayVTable; +use vortex_array::arrays::Filter; +use vortex_array::arrays::filter::FilterExecuteAdaptor; +use vortex_array::optimizer::kernels::ArrayKernelsExt; +use vortex_array::scalar_fn::ScalarFnVTable; +use vortex_array::scalar_fn::fns::binary::Binary; +use vortex_array::scalar_fn::fns::binary::CompareExecuteAdaptor; +use vortex_session::VortexSession; + +use crate::Delta; + +pub(crate) fn initialize(session: &VortexSession) { + let kernels = session.kernels(); + kernels.register_execute_parent_kernel(Binary.id(), Delta, CompareExecuteAdaptor(Delta)); + kernels.register_execute_parent_kernel(Filter.id(), Delta, FilterExecuteAdaptor(Delta)); +} diff --git a/encodings/fastlanes/src/delta/vtable/mod.rs b/encodings/fastlanes/src/delta/vtable/mod.rs index 1f2e3823689..7d8b380336c 100644 --- a/encodings/fastlanes/src/delta/vtable/mod.rs +++ b/encodings/fastlanes/src/delta/vtable/mod.rs @@ -38,11 +38,16 @@ use crate::delta::array::delta_decompress::delta_decompress; use crate::delta::array::lane_count; use crate::delta_compress; +mod kernels; mod operations; mod rules; mod slice; mod validity; +pub(crate) fn initialize(session: &VortexSession) { + kernels::initialize(session); +} + /// A [`Delta`]-encoded Vortex array. pub type DeltaArray = Array; diff --git a/encodings/fastlanes/src/lib.rs b/encodings/fastlanes/src/lib.rs index 8f9309970e0..c7810fdd425 100644 --- a/encodings/fastlanes/src/lib.rs +++ b/encodings/fastlanes/src/lib.rs @@ -91,6 +91,7 @@ pub fn initialize(session: &VortexSession) { session.arrays().register(RLE); session.arrays().register(TransposedBool); bitpacking::initialize(session); + delta::initialize(session); r#for::initialize(session); rle::initialize(session);