diff --git a/Cargo.toml b/Cargo.toml index f26c25fc0f6..a60e17881dc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -158,7 +158,7 @@ datafusion-sqllogictest = { version = "55.0.0" } divan = { package = "codspeed-divan-compat", version = "5.0.0" } enum-iterator = "2.0.0" env_logger = "0.11" -fastlanes = { version = "0.7.0", features = ["runtime"] } +fastlanes = { version = "0.7.1", features = ["runtime"] } fearless_simd = "1.0.0" fearless_simd_macros = "0.1.0" flatbuffers = "25.2.10" diff --git a/encodings/fastlanes/src/bitpacking/compute/filter.rs b/encodings/fastlanes/src/bitpacking/compute/filter.rs index 952e7728e52..12b64493ba0 100644 --- a/encodings/fastlanes/src/bitpacking/compute/filter.rs +++ b/encodings/fastlanes/src/bitpacking/compute/filter.rs @@ -22,7 +22,8 @@ use vortex_mask::Mask; use vortex_mask::MaskValuesRef; use super::chunked_indices; -use super::take::UNPACK_CHUNK_THRESHOLD; +use super::unpack_chunk_threshold; +use super::unpack_indices_into; use crate::BitPacked; use crate::BitPackedArrayExt; use crate::BitPackedData; @@ -157,8 +158,10 @@ fn filter_with_indices( &mut values.as_mut_slice()[values_len..], ); } - } else if indices_within_chunk.len() > UNPACK_CHUNK_THRESHOLD { + } else if indices_within_chunk.len() > unpack_chunk_threshold::() { // Unpack into a temporary chunk and then copy the values. + // SAFETY: The validated bit width fits `T`. The source and destination contain + // one complete FastLanes block. The call initializes every destination value. unsafe { let dst: &mut [MaybeUninit] = &mut unpacked; let dst: &mut [T] = std::mem::transmute(dst); @@ -167,13 +170,11 @@ fn filter_with_indices( values.extend_trusted( indices_within_chunk .iter() + // SAFETY: The preceding unpack initialized the complete temporary block. .map(|&idx| unsafe { unpacked.get_unchecked(idx).assume_init() }), ); } else { - // Otherwise, unpack each element individually. - values.extend_trusted(indices_within_chunk.iter().map(|&idx| unsafe { - BitPacking::unchecked_unpack_single(bit_width, packed, idx) - })); + unpack_indices_into(&mut values, bit_width, packed, indices_within_chunk); } }, ); @@ -193,9 +194,12 @@ mod tests { use vortex_array::validity::Validity; use vortex_buffer::Buffer; use vortex_buffer::buffer; + use vortex_error::VortexResult; use vortex_mask::Mask; use vortex_session::VortexSession; + use super::filter_with_indices; + use super::unpack_chunk_threshold; use crate::BitPackedData; use crate::bitpacking::array::BitPackedArrayExt; @@ -205,6 +209,56 @@ mod tests { session }); + #[test] + fn sparse_extraction_covers_batch_and_full_chunk_paths() -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + + macro_rules! check_type { + ($T:ty, $bit_width:expr) => {{ + let values = (0..2_048) + .map(|index| (index % 127) as $T) + .collect::>(); + let packed = BitPackedData::encode( + &PrimitiveArray::from_iter(values.iter().copied()).into_array(), + $bit_width, + &mut ctx, + )?; + let threshold = unpack_chunk_threshold::<$T>(); + + for selected in [threshold, threshold + 1] { + let indices = (0..selected) + .map(|index| index * 1_024 / selected) + .collect::>(); + let actual = filter_with_indices::<$T>(&packed, &indices); + let expected = indices + .iter() + .map(|&index| values[index]) + .collect::>(); + assert_eq!(actual.as_slice(), expected); + } + }}; + } + + check_type!(u8, 7); + check_type!(u16, 15); + check_type!(u32, 31); + check_type!(u64, 63); + Ok(()) + } + + #[test] + fn sparse_extraction_supports_zero_width() -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + let packed = BitPackedData::encode( + &PrimitiveArray::from_iter([0u32; 2_048]).into_array(), + 0, + &mut ctx, + )?; + let actual = filter_with_indices::(&packed, &[0, 17, 1_023, 1_024, 2_047]); + assert_eq!(actual.as_slice(), &[0; 5]); + Ok(()) + } + #[test] fn take_indices() { let mut ctx = SESSION.create_execution_ctx(); diff --git a/encodings/fastlanes/src/bitpacking/compute/mod.rs b/encodings/fastlanes/src/bitpacking/compute/mod.rs index 38f86f781bb..fba1d8fd403 100644 --- a/encodings/fastlanes/src/bitpacking/compute/mod.rs +++ b/encodings/fastlanes/src/bitpacking/compute/mod.rs @@ -1,6 +1,12 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors +use std::mem::size_of; + +use fastlanes::BitPacking; +use vortex_array::dtype::NativePType; +use vortex_buffer::BufferMut; + mod between; mod cast; mod compare; @@ -11,6 +17,35 @@ mod slice; mod stream_predicate; mod take; +const fn unpack_chunk_threshold() -> usize { + // FastLanes and Vortex benchmarks set conservative crossovers for each physical type. + match size_of::() { + 1 => 16, + 2 => 32, + 4 => 64, + 8 => 160, + _ => unreachable!(), + } +} + +fn unpack_indices_into( + output: &mut BufferMut, + bit_width: usize, + packed: &[T], + indices: &[usize], +) { + let output_len = output.len(); + let destination = &mut output.spare_capacity_mut()[..indices.len()]; + + // SAFETY: `bit_width` comes from validated data and fits `T`. + // `packed` contains one complete block, and each index is block-relative. + // The destination length equals the index length, and `output` reserves enough space. + unsafe { + T::unchecked_unpack_indices(bit_width, packed, indices, destination); + output.set_len(output_len + indices.len()); + } +} + // TODO(connor): This is duplicated in `encodings/fastlanes/src/bitpacking/kernels/mod.rs`. fn chunked_indices( mut indices: impl Iterator, diff --git a/encodings/fastlanes/src/bitpacking/compute/take.rs b/encodings/fastlanes/src/bitpacking/compute/take.rs index 1bb778212c3..bbcf0b9ba56 100644 --- a/encodings/fastlanes/src/bitpacking/compute/take.rs +++ b/encodings/fastlanes/src/bitpacking/compute/take.rs @@ -1,7 +1,6 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors -use std::mem; use std::mem::MaybeUninit; use fastlanes::BitPacking; @@ -23,16 +22,13 @@ use vortex_error::VortexExpect as _; use vortex_error::VortexResult; use super::chunked_indices; +use super::unpack_chunk_threshold; +use super::unpack_indices_into; use crate::BitPacked; use crate::BitPackedArrayExt; use crate::BitWidthsView; -use crate::bitpack_decompress; -// TODO(connor): This is duplicated in `encodings/fastlanes/src/bitpacking/kernels/mod.rs`. -/// assuming the buffer is already allocated (which will happen at most once) then unpacking -/// all 1024 elements takes ~8.8x as long as unpacking a single element on an M2 Macbook Air. -/// see -pub(super) const UNPACK_CHUNK_THRESHOLD: usize = 8; +const FULL_ARRAY_DECODE_RATIO: usize = 8; impl TakeExecute for BitPacked { fn take( @@ -44,7 +40,7 @@ impl TakeExecute for BitPacked { return Ok(None); }; // If the indices are large enough, it's faster to flatten and take the primitive array. - if indices.len() * UNPACK_CHUNK_THRESHOLD > array.len() { + if indices.len() * FULL_ARRAY_DECODE_RATIO > array.len() { let prim = array.array().clone().execute::(ctx)?; return prim.into_array().take(indices.clone()).map(Some); } @@ -103,41 +99,22 @@ fn take_primitive( chunked_indices(indices_iter, offset, |chunk_idx, indices_within_chunk| { let packed = &packed[chunk_idx * chunk_len..][..chunk_len]; - let mut have_unpacked = false; - let (offset_chunks, remainder) = indices_within_chunk.as_chunks::(); - - // this loop only runs if we have at least UNPACK_CHUNK_THRESHOLD offsets - for offset_chunk in offset_chunks { - if !have_unpacked { - unsafe { - let dst: &mut [MaybeUninit] = &mut unpacked; - let dst: &mut [T] = mem::transmute(dst); - BitPacking::unchecked_unpack(bit_width, packed, dst); - } - have_unpacked = true; - } - - for &index in offset_chunk { - output.push(unsafe { unpacked[index].assume_init() }); - } - } - - // if we have a remainder (i.e., < UNPACK_CHUNK_THRESHOLD leftover offsets), we need to handle it - if !remainder.is_empty() { - if have_unpacked { - // we already bulk unpacked this chunk, so we can just push the remaining elements - for &index in remainder { - output.push(unsafe { unpacked[index].assume_init() }); - } - } else { - // we had fewer than UNPACK_CHUNK_THRESHOLD offsets in the first place, - // so we need to unpack each one individually - for &index in remainder { - output.push(unsafe { - bitpack_decompress::unpack_single_primitive::(packed, bit_width, index) - }); - } + if indices_within_chunk.len() > unpack_chunk_threshold::() { + // SAFETY: The validated bit width fits `T`. The source and destination contain one + // complete FastLanes block. The call initializes every destination value. + unsafe { + let dst: &mut [MaybeUninit] = &mut unpacked; + let dst: &mut [T] = std::mem::transmute(dst); + BitPacking::unchecked_unpack(bit_width, packed, dst); } + output.extend_trusted( + indices_within_chunk + .iter() + // SAFETY: The preceding unpack initialized the complete temporary block. + .map(|&index| unsafe { unpacked.get_unchecked(index).assume_init() }), + ); + } else { + unpack_indices_into(&mut output, bit_width, packed, indices_within_chunk); } }); @@ -162,7 +139,7 @@ fn take_primitive( #[cfg(test)] #[expect(clippy::cast_possible_truncation)] -mod test { +mod tests { use std::sync::LazyLock; use rand::RngExt; @@ -177,12 +154,14 @@ mod test { use vortex_array::validity::Validity; use vortex_buffer::Buffer; use vortex_buffer::buffer; + use vortex_error::VortexResult; use vortex_session::VortexSession; use crate::BitPackedArray; use crate::BitPackedData; use crate::bitpacking::array::BitPackedArrayExt; use crate::bitpacking::compute::take::take_primitive; + use crate::bitpacking::compute::take::unpack_chunk_threshold; static SESSION: LazyLock = LazyLock::new(|| { let session = vortex_array::array_session(); @@ -190,6 +169,77 @@ mod test { session }); + #[test] + fn sparse_extraction_covers_batch_and_full_chunk_paths() -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + + macro_rules! check_type { + ($T:ty, $bit_width:expr) => {{ + let values = (0..2_048) + .map(|index| (index % 127) as $T) + .collect::>(); + let packed = BitPackedData::encode( + &PrimitiveArray::from_iter(values.iter().copied()).into_array(), + $bit_width, + &mut ctx, + )?; + let threshold = unpack_chunk_threshold::<$T>(); + + for selected in [threshold, threshold + 1] { + let indices = (0..selected) + .map(|index| (index * 1_024 / selected) as u32) + .collect::>(); + let actual = take_primitive::<$T, u32>( + packed.as_view(), + $bit_width, + &PrimitiveArray::from_iter(indices.iter().copied()), + Validity::NonNullable, + &mut ctx, + )?; + let expected = indices + .iter() + .map(|&index| values[index as usize]) + .collect::>(); + assert_eq!(actual.as_slice::<$T>(), expected); + } + }}; + } + + check_type!(u8, 7); + check_type!(u16, 15); + check_type!(u32, 31); + check_type!(u64, 63); + Ok(()) + } + + #[test] + fn sparse_take_preserves_nullable_order_and_duplicates() -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + let values = (0..8_192).map(|index| index as u32).collect::>(); + let packed = BitPackedData::encode( + &PrimitiveArray::from_iter(values.iter().copied()).into_array(), + 13, + &mut ctx, + )?; + let indices = [ + Some(3_073u32), + Some(2), + None, + Some(2), + Some(1_025), + Some(3_073), + Some(0), + ]; + let actual = packed + .take(PrimitiveArray::from_option_iter(indices).into_array())? + .execute::(&mut ctx)?; + let expected = PrimitiveArray::from_option_iter( + indices.map(|index| index.map(|index| values[index as usize])), + ); + assert_arrays_eq!(actual, expected, &mut ctx); + Ok(()) + } + #[test] fn take_indices() { let mut ctx = SESSION.create_execution_ctx();