Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
66 changes: 60 additions & 6 deletions encodings/fastlanes/src/bitpacking/compute/filter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -157,8 +158,10 @@ fn filter_with_indices<T: NativePType + BitPacking>(
&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::<T>() {
// 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<T>] = &mut unpacked;
let dst: &mut [T] = std::mem::transmute(dst);
Expand All @@ -167,13 +170,11 @@ fn filter_with_indices<T: NativePType + BitPacking>(
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);
}
},
);
Expand All @@ -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;

Expand All @@ -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::<Vec<_>>();
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::<Vec<_>>();
let actual = filter_with_indices::<$T>(&packed, &indices);
let expected = indices
.iter()
.map(|&index| values[index])
.collect::<Vec<_>>();
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::<u32>(&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();
Expand Down
35 changes: 35 additions & 0 deletions encodings/fastlanes/src/bitpacking/compute/mod.rs
Original file line number Diff line number Diff line change
@@ -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;
Expand All @@ -11,6 +17,35 @@ mod slice;
mod stream_predicate;
mod take;

const fn unpack_chunk_threshold<T>() -> usize {
// FastLanes and Vortex benchmarks set conservative crossovers for each physical type.
match size_of::<T>() {
1 => 16,
2 => 32,
4 => 64,
8 => 160,
_ => unreachable!(),
}
}

fn unpack_indices_into<T: NativePType + BitPacking>(
output: &mut BufferMut<T>,
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<F: FnMut(usize, &[usize])>(
mut indices: impl Iterator<Item = usize>,
Expand Down
136 changes: 93 additions & 43 deletions encodings/fastlanes/src/bitpacking/compute/take.rs
Original file line number Diff line number Diff line change
@@ -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;
Expand All @@ -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 <https://github.com/vortex-data/vortex/pull/190#issue-2223752833>
pub(super) const UNPACK_CHUNK_THRESHOLD: usize = 8;
const FULL_ARRAY_DECODE_RATIO: usize = 8;

impl TakeExecute for BitPacked {
fn take(
Expand All @@ -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::<PrimitiveArray>(ctx)?;
return prim.into_array().take(indices.clone()).map(Some);
}
Expand Down Expand Up @@ -103,41 +99,22 @@ fn take_primitive<T: NativePType + BitPacking, I: IntegerPType>(
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::<UNPACK_CHUNK_THRESHOLD>();

// 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<T>] = &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::<T>(packed, bit_width, index)
});
}
if indices_within_chunk.len() > unpack_chunk_threshold::<T>() {
// 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<T>] = &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);
}
});

Expand All @@ -162,7 +139,7 @@ fn take_primitive<T: NativePType + BitPacking, I: IntegerPType>(

#[cfg(test)]
#[expect(clippy::cast_possible_truncation)]
mod test {
mod tests {
use std::sync::LazyLock;

use rand::RngExt;
Expand All @@ -177,19 +154,92 @@ 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<VortexSession> = LazyLock::new(|| {
let session = vortex_array::array_session();
crate::initialize(&session);
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::<Vec<_>>();
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::<Vec<_>>();
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::<Vec<_>>();
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::<Vec<_>>();
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::<PrimitiveArray>(&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();
Expand Down
Loading