diff --git a/vortex-btrblocks/src/builder.rs b/vortex-btrblocks/src/builder.rs index 7bdc6661ff6..d3dea3fdb53 100644 --- a/vortex-btrblocks/src/builder.rs +++ b/vortex-btrblocks/src/builder.rs @@ -5,6 +5,7 @@ use vortex_array::ArrayId; use vortex_decimal_byte_parts::decimal_byte_parts_v2_id; +use vortex_fastlanes::bitpacked_v2_id; use vortex_fastlanes::for_v2_id; use vortex_session::VortexSession; use vortex_utils::aliases::hash_set::HashSet; @@ -94,9 +95,9 @@ impl CompressionMode { fn excluded_encodings(self) -> Vec { match self { Self::All | Self::Default | Self::Compact => Vec::new(), - // Multi-part DecimalByteParts arrays and FoR arrays with per-chunk references have no - // CUDA decode kernel. - Self::Cuda => vec![decimal_byte_parts_v2_id(), for_v2_id()], + // Multi-part DecimalByteParts arrays, FoR arrays with per-chunk references and + // BitPacked arrays with per-block bit widths have no CUDA decode kernel. + Self::Cuda => vec![decimal_byte_parts_v2_id(), for_v2_id(), bitpacked_v2_id()], } } } diff --git a/vortex-btrblocks/src/schemes/integer/bitpacking.rs b/vortex-btrblocks/src/schemes/integer/bitpacking.rs index 4f81462ca3c..35d6bcc3c06 100644 --- a/vortex-btrblocks/src/schemes/integer/bitpacking.rs +++ b/vortex-btrblocks/src/schemes/integer/bitpacking.rs @@ -16,21 +16,69 @@ use vortex_compressor::scheme::CompressionEstimate; use vortex_compressor::scheme::DeferredEstimate; use vortex_compressor::scheme::EstimateVerdict; use vortex_error::VortexResult; +use vortex_error::vortex_bail; use vortex_fastlanes::BitPacked; +use vortex_fastlanes::BitWidths; use vortex_fastlanes::bitpack_compress::bit_width_histogram; +use vortex_fastlanes::bitpack_compress::bitpack_blocked_to_best_bit_widths; use vortex_fastlanes::bitpack_compress::bitpack_encode; use vortex_fastlanes::bitpack_compress::find_best_bit_width; use vortex_fastlanes::bitpacked_v1_id; +use vortex_fastlanes::bitpacked_v2_id; use crate::ArrayAndStats; use crate::CascadingCompressor; use crate::CompressorContext; use crate::Scheme; +use crate::SchemeExt; use crate::compress_patches; +#[derive(Debug, Copy, Clone, PartialEq, Eq)] +enum BitPackingSchemeMode { + V1, + V2, +} + +/// The BitPacking scheme that only produces arrays with one bit width. +pub(crate) static BITPACKING_V1: BitPackingScheme = BitPackingScheme::v1(); + +/// The BitPacking scheme that may give each 1024-element block its own bit width. +pub(crate) static BITPACKING_V2: BitPackingScheme = BitPackingScheme::v2(); + /// BitPacking encoding for non-negative integers. +/// +/// The v1 mode packs every value at one bit width, and serializes as `fastlanes.bitpacked`. The v2 +/// mode always chooses a width for each 1024-element block, and serializes as +/// `fastlanes.bitpacked.v2`. +/// +/// The default uses v1. [`refine`](Scheme::refine) picks v2 when the v2 ID is allowed and v1 +/// otherwise. #[derive(Debug, Copy, Clone, PartialEq, Eq)] -pub struct BitPackingScheme; +pub struct BitPackingScheme { + mode: BitPackingSchemeMode, +} + +impl BitPackingScheme { + /// Creates a BitPacking scheme configured for v1, which uses one bit width. + pub const fn v1() -> Self { + Self { + mode: BitPackingSchemeMode::V1, + } + } + + /// Creates a BitPacking scheme configured for v2, which may use one bit width per block. + pub const fn v2() -> Self { + Self { + mode: BitPackingSchemeMode::V2, + } + } +} + +impl Default for BitPackingScheme { + fn default() -> Self { + Self::v1() + } +} impl Scheme for BitPackingScheme { fn scheme_name(&self) -> &'static str { @@ -42,14 +90,33 @@ impl Scheme for BitPackingScheme { } fn produced_encodings(&self) -> Vec { - // Global-width arrays serialize under the frozen v1 ID. - let mut encodings = vec![bitpacked_v1_id()]; + let mut encodings = match self.mode { + // Global-width arrays serialize under the frozen v1 ID. + BitPackingSchemeMode::V1 => vec![bitpacked_v1_id()], + BitPackingSchemeMode::V2 => vec![bitpacked_v2_id()], + }; if use_experimental_patches() { encodings.push(Patched.id()); } encodings } + fn refine(&self, allowed: &dyn Fn(&ArrayId) -> bool) -> &dyn Scheme { + if allowed(&bitpacked_v2_id()) { + &BITPACKING_V2 + } else { + &BITPACKING_V1 + } + } + + /// Children: block offsets=0 in v2 mode. + fn num_children(&self) -> usize { + match self.mode { + BitPackingSchemeMode::V1 => 0, + BitPackingSchemeMode::V2 => 1, + } + } + fn expected_compression_ratio( &self, data: &ArrayAndStats, @@ -68,11 +135,15 @@ impl Scheme for BitPackingScheme { fn compress( &self, - _compressor: &CascadingCompressor, + compressor: &CascadingCompressor, data: &ArrayAndStats, - _compress_ctx: CompressorContext, + compress_ctx: CompressorContext, exec_ctx: &mut ExecutionCtx, ) -> VortexResult { + if self.mode == BitPackingSchemeMode::V2 { + return self.compress_blocked(compressor, data, compress_ctx, exec_ctx); + } + let primitive_array = data.array_as_primitive(); let histogram = bit_width_histogram(primitive_array, exec_ctx)?; @@ -135,3 +206,106 @@ impl Scheme for BitPackingScheme { Ok(array) } } + +impl BitPackingScheme { + /// Bit-pack each 1024-element block at its own best width. + fn compress_blocked( + &self, + compressor: &CascadingCompressor, + data: &ArrayAndStats, + compress_ctx: CompressorContext, + exec_ctx: &mut ExecutionCtx, + ) -> VortexResult { + let primitive_array = data.array_as_primitive().into_owned(); + let packed = bitpack_blocked_to_best_bit_widths(&primitive_array, exec_ctx)?; + + let packed_stats = packed.statistics().to_owned(); + let ptype = packed.dtype().as_ptype(); + let mut parts = BitPacked::into_parts(packed); + let BitWidths::Blocked(block_offsets) = parts.bit_widths else { + vortex_bail!("Blocked bit-packing must produce block offsets"); + }; + let block_offsets = + compressor.compress_child(&block_offsets, &compress_ctx, self.id(), 0, exec_ctx)?; + + let array = if use_experimental_patches() { + let patches = parts.patches.take(); + // Transpose patches into G-ALP style PatchedArray, wrapping an inner BitPackedArray. + let array = BitPacked::try_new_with_block_offsets( + parts.packed, + ptype, + parts.validity, + None, + block_offsets, + parts.len, + parts.offset, + )? + .into_array(); + + match patches { + None => array, + Some(p) => Patched::from_array_and_patches(array, &p, exec_ctx)? + .with_stats_set(packed_stats) + .into_array(), + } + } else { + // Compress patches and place back into BitPackedArray. + let patches = parts + .patches + .take() + .map(|p| compress_patches(p, exec_ctx)) + .transpose()?; + BitPacked::try_new_with_block_offsets( + parts.packed, + ptype, + parts.validity, + patches, + block_offsets, + parts.len, + parts.offset, + )? + .with_stats_set(packed_stats) + .into_array() + }; + + Ok(array) + } +} + +#[cfg(test)] +mod tests { + use rstest::rstest; + use vortex_array::ArrayId; + use vortex_fastlanes::bitpacked_v1_id; + use vortex_fastlanes::bitpacked_v2_id; + + use super::BITPACKING_V1; + use super::BITPACKING_V2; + use crate::Scheme; + use crate::SchemeExt; + + /// Both variants refine to v2 exactly when the v2 ID is allowed. + #[rstest] + #[case::neither(false, false, false)] + #[case::v1(true, false, false)] + #[case::v2_only(false, true, true)] + #[case::both(true, true, true)] + fn refine_picks_v2_when_v2_id_is_allowed( + #[case] allow_v1: bool, + #[case] allow_v2: bool, + #[case] expect_v2: bool, + #[values(&BITPACKING_V1, &BITPACKING_V2)] scheme: &'static dyn Scheme, + ) { + let allowed = |id: &ArrayId| { + (allow_v1 && *id == bitpacked_v1_id()) || (allow_v2 && *id == bitpacked_v2_id()) + }; + let refined = scheme.refine(&allowed); + assert_eq!(refined.id(), scheme.id()); + let expected: &dyn Scheme = if expect_v2 { + &BITPACKING_V2 + } else { + &BITPACKING_V1 + }; + assert_eq!(refined.produced_encodings(), expected.produced_encodings()); + } +} diff --git a/vortex-btrblocks/src/schemes/integer/for_.rs b/vortex-btrblocks/src/schemes/integer/for_.rs index 43c7ee56d68..34376673c11 100644 --- a/vortex-btrblocks/src/schemes/integer/for_.rs +++ b/vortex-btrblocks/src/schemes/integer/for_.rs @@ -25,7 +25,7 @@ use vortex_fastlanes::FoRArraySlotsExt; use vortex_fastlanes::for_v1_id; use vortex_fastlanes::for_v2_id; -use super::BitPackingScheme; +use super::BITPACKING_V1; use crate::ArrayAndStats; use crate::CascadingCompressor; use crate::CompressorContext; @@ -216,7 +216,7 @@ impl Scheme for FoRScheme { let leaf_ctx = compress_ctx.clone().as_leaf(); let biased_data = ArrayAndStats::new(biased.into_array(), compress_ctx.merged_stats_options()); - let compressed = BitPackingScheme.compress(compressor, &biased_data, leaf_ctx, exec_ctx)?; + let compressed = BITPACKING_V1.compress(compressor, &biased_data, leaf_ctx, exec_ctx)?; // TODO(connor): This should really be `new_unchecked`. let for_compressed = match for_array.constant_reference() { diff --git a/vortex-btrblocks/src/schemes/integer/mod.rs b/vortex-btrblocks/src/schemes/integer/mod.rs index 742b453c8ef..74edd53463f 100644 --- a/vortex-btrblocks/src/schemes/integer/mod.rs +++ b/vortex-btrblocks/src/schemes/integer/mod.rs @@ -15,6 +15,7 @@ mod zigzag; #[cfg(feature = "pco")] mod pco; +pub(crate) use bitpacking::BITPACKING_V1; pub use bitpacking::BitPackingScheme; pub use delta::DeltaScheme; pub(crate) use for_::FOR_V1; diff --git a/vortex-btrblocks/src/session.rs b/vortex-btrblocks/src/session.rs index 4ca2a0e5de3..120345efd62 100644 --- a/vortex-btrblocks/src/session.rs +++ b/vortex-btrblocks/src/session.rs @@ -48,7 +48,7 @@ impl Default for CompressionSession { &integer::FOR_V1, // NOTE: ZigZag should precede BitPacking because we don't want negative numbers. &integer::ZigZagScheme, - &integer::BitPackingScheme, + &integer::BITPACKING_V1, &integer::SparseScheme, &integer::IntDictScheme, &integer::RunEndScheme, diff --git a/vortex-btrblocks/tests/bitpacking_config.rs b/vortex-btrblocks/tests/bitpacking_config.rs new file mode 100644 index 00000000000..37621fed83c --- /dev/null +++ b/vortex-btrblocks/tests/bitpacking_config.rs @@ -0,0 +1,141 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! BitPacking scheme refinement by allowed serialized IDs, and per-block bit widths. + +#![cfg(test)] + +use std::sync::LazyLock; + +use rstest::rstest; +use vortex_array::ArrayContext; +use vortex_array::ArrayId; +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::VortexSessionExecute; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::assert_arrays_eq; +use vortex_array::serde::SerializeOptions; +use vortex_array::serde::SerializedArray; +use vortex_btrblocks::BtrBlocksCompressorBuilder; +use vortex_btrblocks::schemes::integer::BitPackingScheme; +use vortex_buffer::ByteBufferMut; +use vortex_edition::EDITION_DECLARATIONS; +use vortex_edition::EDITION_FAMILIES; +use vortex_edition::EditionSession; +use vortex_edition::EditionSessionExt; +use vortex_edition::declarations::core::CORE_2026_08_3; +use vortex_error::VortexResult; +use vortex_fastlanes::bitpacked_v1_id; +use vortex_fastlanes::bitpacked_v2_id; +use vortex_session::VortexSession; +use vortex_session::registry::ReadContext; + +static BITPACKING_V1: BitPackingScheme = BitPackingScheme::v1(); + +/// Registers the fastlanes encodings, and enables no editions. +static SESSION: LazyLock = LazyLock::new(|| { + let session = vortex_array::array_session(); + vortex_fastlanes::initialize(&session); + session +}); + +/// Like [`SESSION`], with the latest core edition enabled: it allows `fastlanes.bitpacked` but not +/// `fastlanes.bitpacked.v2`. +static CORE_SESSION: LazyLock = LazyLock::new(|| { + let session = vortex_array::array_session().with::(); + for family in EDITION_FAMILIES { + session + .editions() + .declare_family(family) + .expect("first-party edition family"); + } + for declaration in EDITION_DECLARATIONS { + session + .register_edition(declaration) + .expect("first-party edition"); + } + session + .enable_edition(CORE_2026_08_3) + .expect("core edition is registered"); + vortex_fastlanes::initialize(&session); + session +}); + +/// Values whose 1024-value blocks need 1 to 8 bits. +fn drifting() -> ArrayRef { + PrimitiveArray::from_iter((0..8192u32).map(|i| i % (2 << (i / 1024)))).into_array() +} + +/// Values that need 7 bits in every block. +fn uniform() -> ArrayRef { + PrimitiveArray::from_iter((0..8192u32).map(|i| i % 128)).into_array() +} + +/// Compresses `array`, round trips it through serialization, and returns the serialized IDs. +fn compress_roundtrip( + builder: BtrBlocksCompressorBuilder, + array: &ArrayRef, +) -> VortexResult> { + let mut ctx = SESSION.create_execution_ctx(); + let compressed = builder.build().compress(array, &mut ctx)?; + assert_arrays_eq!(array, compressed, &mut ctx); + + let array_ctx = ArrayContext::empty(); + let mut bytes = ByteBufferMut::empty(); + for buffer in compressed.serialize(&array_ctx, &SESSION, &SerializeOptions::default())? { + bytes.extend_from_slice(&buffer); + } + let read = SerializedArray::try_from(bytes.freeze())?.decode( + array.dtype(), + array.len(), + &ReadContext::new(array_ctx.to_ids()), + &SESSION, + )?; + assert_arrays_eq!(array, read, &mut ctx); + Ok(array_ctx.to_ids()) +} + +/// Only BitPacking, with every serialized ID allowed, so it refines to v2. +fn bitpacking_only() -> BtrBlocksCompressorBuilder { + BtrBlocksCompressorBuilder::empty().with_new_scheme(&BITPACKING_V1) +} + +/// v2 produces per-block bit widths even when every block chooses the same width. +#[rstest] +#[case::drifting(drifting())] +#[case::uniform(uniform())] +fn v2_always_serializes_as_v2(#[case] array: ArrayRef) -> VortexResult<()> { + let ids = compress_roundtrip(bitpacking_only(), &array)?; + assert!(ids.contains(&bitpacked_v2_id())); + Ok(()) +} + +#[test] +fn nullable_drifting_roundtrip() -> VortexResult<()> { + let array = PrimitiveArray::from_option_iter( + (0..8192u64).map(|i| (i % 7 != 0).then_some(i % (2 << (i / 1024)))), + ) + .into_array(); + let ids = compress_roundtrip(bitpacking_only(), &array)?; + assert!(ids.contains(&bitpacked_v2_id())); + Ok(()) +} + +#[test] +fn core_edition_keeps_global_width() -> VortexResult<()> { + let ids = compress_roundtrip( + BtrBlocksCompressorBuilder::from_session(&CORE_SESSION), + &drifting(), + )?; + assert!(!ids.contains(&bitpacked_v2_id())); + Ok(()) +} + +#[test] +fn cuda_preset_keeps_global_width() -> VortexResult<()> { + let ids = compress_roundtrip(bitpacking_only().only_cuda_compatible(), &drifting())?; + assert!(ids.contains(&bitpacked_v1_id())); + assert!(!ids.contains(&bitpacked_v2_id())); + Ok(()) +} diff --git a/vortex-btrblocks/tests/decimal_config.rs b/vortex-btrblocks/tests/decimal_config.rs index e91995e0e4f..749aebc561b 100644 --- a/vortex-btrblocks/tests/decimal_config.rs +++ b/vortex-btrblocks/tests/decimal_config.rs @@ -45,6 +45,7 @@ use vortex_session::registry::ReadContext; static DECIMAL_V2: DecimalScheme = DecimalScheme::v2(); static FOR_V1: FoRScheme = FoRScheme::v1(); +static BITPACKING_V1: BitPackingScheme = BitPackingScheme::v1(); /// Registers the decimal and fastlanes encodings, and enables no editions. static SESSION: LazyLock = LazyLock::new(|| { @@ -199,7 +200,7 @@ fn wide_decimal_parts_roundtrip( if compress_children { builder = builder .with_new_scheme(&FOR_V1) - .with_new_scheme(&BitPackingScheme); + .with_new_scheme(&BITPACKING_V1); } let mut ctx = SESSION.create_execution_ctx(); let compressed = builder.build().compress(&array, &mut ctx)?; diff --git a/vortex-btrblocks/tests/for_config.rs b/vortex-btrblocks/tests/for_config.rs index 812e462871c..cf6be6e4826 100644 --- a/vortex-btrblocks/tests/for_config.rs +++ b/vortex-btrblocks/tests/for_config.rs @@ -33,7 +33,7 @@ use vortex_session::VortexSession; use vortex_session::registry::ReadContext; static FOR_V1: FoRScheme = FoRScheme::v1(); -static BITPACKING: BitPackingScheme = BitPackingScheme; +static BITPACKING: BitPackingScheme = BitPackingScheme::v1(); /// Registers the fastlanes encodings, and enables no editions. static SESSION: LazyLock = LazyLock::new(|| { diff --git a/vortex-file/src/tests.rs b/vortex-file/src/tests.rs index 083d6bcdf6e..6cc6d161064 100644 --- a/vortex-file/src/tests.rs +++ b/vortex-file/src/tests.rs @@ -1538,8 +1538,9 @@ async fn test_into_tokio_array_stream() -> VortexResult<()> { async fn test_array_stream_no_double_dict_encode() -> VortexResult<()> { let num_vals = 2048; let mut values = Vec::::with_capacity(num_vals); - values.extend(iter::repeat_n(0, num_vals / 2)); - values.extend(iter::repeat_n(1, num_vals / 2)); + // Far apart, so per-block bit widths can't pack them tighter than a dictionary. + values.extend(iter::repeat_n(1_000_000_000, num_vals / 2)); + values.extend(iter::repeat_n(2_000_000_000, num_vals / 2)); let array = PrimitiveArray::from_iter(values).into_array(); let mut buf = Vec::new();