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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions vortex-btrblocks/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ vortex-edition = { workspace = true }
vortex-error = { workspace = true }
vortex-fastlanes = { workspace = true }
vortex-fsst = { workspace = true }
vortex-mask = { workspace = true }
vortex-onpair = { workspace = true }
vortex-pco = { workspace = true, optional = true }
vortex-runend = { workspace = true }
Expand Down
2 changes: 2 additions & 0 deletions vortex-btrblocks/src/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ use crate::SchemeExt;
use crate::SchemeId;
use crate::allowed_ids::AllowedIds;
use crate::schemes::binary;
use crate::schemes::fixed_size_list;
use crate::schemes::float;
use crate::schemes::integer;
use crate::schemes::string;
Expand Down Expand Up @@ -73,6 +74,7 @@ impl CompressionMode {
float::FloatRLEScheme.id(),
float::NullDominatedSparseScheme.id(),
string::NullDominatedSparseScheme.id(),
fixed_size_list::FixedSizeListSparseScheme.id(),
string::StringDictScheme.id(),
binary::BinaryDictScheme.id(),
]);
Expand Down
163 changes: 163 additions & 0 deletions vortex-btrblocks/src/schemes/fixed_size_list.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,163 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright the Vortex contributors

//! Sparse encoding for fixed-size lists with many null rows.

use vortex_array::ArrayId;
use vortex_array::ArrayRef;
use vortex_array::Canonical;
use vortex_array::ExecutionCtx;
use vortex_array::IntoArray;
use vortex_array::VTable;
use vortex_array::arrays::PrimitiveArray;
use vortex_array::arrays::primitive::PrimitiveArrayExt;
use vortex_array::scalar::Scalar;
use vortex_array::validity::Validity;
use vortex_buffer::Buffer;
use vortex_compressor::scheme::ChildSelection;
use vortex_compressor::scheme::CompressionEstimate;
use vortex_compressor::scheme::DescendantExclusion;
use vortex_compressor::scheme::EstimateVerdict;
use vortex_error::VortexResult;
use vortex_mask::AllOr;
use vortex_mask::Mask;
use vortex_sparse::Sparse;

use crate::ArrayAndStats;
use crate::CascadingCompressor;
use crate::CompressorContext;
use crate::Scheme;
use crate::SchemeExt;
use crate::schemes::integer::IntDictScheme;
use crate::schemes::integer::IntRLEScheme;
use crate::schemes::integer::RunEndScheme;
use crate::schemes::integer::SparseScheme as IntSparseScheme;

/// The smallest estimated ratio of the dense list's size to the sparse one's worth encoding for.
const MIN_ESTIMATED_RATIO: f64 = 1.1;

/// Sparse encoding of a fixed-size list that stores only its valid rows.
///
/// A fixed-size list keeps `list_size` elements for every row, null or not, so compressing its
/// elements has to encode the arbitrary values under null rows too. This scheme instead stores the
/// positions of the valid rows and a list of only those rows, with null as the fill value.
#[derive(Debug, Copy, Clone, PartialEq, Eq)]
pub struct FixedSizeListSparseScheme;

impl Scheme for FixedSizeListSparseScheme {
fn scheme_name(&self) -> &'static str {
"vortex.fixed_size_list.sparse"
}

fn matches(&self, canonical: &Canonical) -> bool {
matches!(canonical, Canonical::FixedSizeList(_)) && canonical.dtype().is_nullable()
}

fn produced_encodings(&self) -> Vec<ArrayId> {
vec![Sparse.id()]
}

/// Children: values=0, indices=1.
fn num_children(&self) -> usize {
2
}

/// Sparse indices (child 1) are monotonically increasing positions with all unique values.
fn descendant_exclusions(&self) -> Vec<DescendantExclusion> {
vec![
DescendantExclusion {
excluded: IntDictScheme.id(),
children: ChildSelection::One(1),
},
DescendantExclusion {
excluded: RunEndScheme.id(),
children: ChildSelection::One(1),
},
DescendantExclusion {
excluded: IntRLEScheme.id(),
children: ChildSelection::One(1),
},
DescendantExclusion {
excluded: IntSparseScheme.id(),
children: ChildSelection::One(1),
},
]
}

/// Compares the list's size with the valid rows plus their bit-packed positions.
fn expected_compression_ratio(
&self,
data: &ArrayAndStats,
_compress_ctx: CompressorContext,
exec_ctx: &mut ExecutionCtx,
) -> CompressionEstimate {
let array = data.array();
let len = array.len();
let Ok(mask) = validity_mask(array, exec_ctx) else {
return CompressionEstimate::Verdict(EstimateVerdict::Skip);
};
let valid = mask.true_count();
if valid == 0 || valid == len {
return CompressionEstimate::Verdict(EstimateVerdict::Skip);
}

let dense_bytes = array.nbytes() as f64;
let row_bytes = dense_bytes / len as f64;
let index_bytes = f64::from(usize::BITS - (len - 1).leading_zeros()) / 8.0;
let sparse_bytes = valid as f64 * (row_bytes + index_bytes);

let ratio = dense_bytes / sparse_bytes;
if ratio < MIN_ESTIMATED_RATIO {
return CompressionEstimate::Verdict(EstimateVerdict::Skip);
}
CompressionEstimate::Verdict(EstimateVerdict::Ratio(ratio))
}

fn compress(
&self,
compressor: &CascadingCompressor,
data: &ArrayAndStats,
compress_ctx: CompressorContext,
exec_ctx: &mut ExecutionCtx,
) -> VortexResult<ArrayRef> {
let array = data.array();
let mask = validity_mask(array, exec_ctx)?;
let AllOr::Some(valid_indices) = mask.indices() else {
// The estimate skips lists without both valid and null rows.
return Ok(array.clone());
};

let values = array.filter(mask.clone())?;
let compressed_values =
compressor.compress_child(&values, &compress_ctx, self.id(), 0, exec_ctx)?;

let indices = PrimitiveArray::new(
valid_indices
.iter()
.map(|&idx| idx as u64)
.collect::<Buffer<u64>>(),
Validity::NonNullable,
)
.narrow(exec_ctx)?;
let compressed_indices = compressor.compress_child(
&indices.into_array(),
&compress_ctx,
self.id(),
1,
exec_ctx,
)?;

Ok(Sparse::try_new(
compressed_indices,
compressed_values,
array.len(),
Scalar::null(array.dtype().clone()),
)?
.into_array())
}
}

/// Returns which rows of `array` are valid.
fn validity_mask(array: &ArrayRef, exec_ctx: &mut ExecutionCtx) -> VortexResult<Mask> {
array.validity()?.execute_mask(array.len(), exec_ctx)
}
1 change: 1 addition & 0 deletions vortex-btrblocks/src/schemes/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ pub mod integer;
pub mod string;

pub mod decimal;
pub mod fixed_size_list;
pub mod temporal;

pub(crate) mod patches;
Expand Down
3 changes: 3 additions & 0 deletions vortex-btrblocks/src/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ use crate::Scheme;
use crate::SchemeExt;
use crate::schemes::binary;
use crate::schemes::decimal;
use crate::schemes::fixed_size_list;
use crate::schemes::float;
use crate::schemes::integer;
use crate::schemes::string;
Expand Down Expand Up @@ -91,6 +92,8 @@ impl Default for CompressionSession {
&decimal::DECIMAL_V1,
// Temporal schemes.
&temporal::TemporalScheme,
// Fixed-size list schemes.
&fixed_size_list::FixedSizeListSparseScheme,
],
}
}
Expand Down
45 changes: 45 additions & 0 deletions vortex-btrblocks/src/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ use vortex_array::IntoArray;
use vortex_array::VortexSessionExecute;
use vortex_array::aggregate_fn::fns::sum::sum;
use vortex_array::arrays::DecimalArray;
use vortex_array::arrays::FixedSizeListArray;
use vortex_array::arrays::ListArray;
use vortex_array::arrays::PrimitiveArray;
use vortex_array::arrays::StructArray;
Expand All @@ -32,6 +33,7 @@ use vortex_array::dtype::DecimalDType;
use vortex_array::dtype::Nullability;
use vortex_array::extension::datetime::TimeUnit;
use vortex_array::validity::Validity;
use vortex_buffer::Buffer;
#[cfg(feature = "zstd")]
use vortex_edition::EDITION_DECLARATIONS;
#[cfg(feature = "zstd")]
Expand All @@ -44,6 +46,7 @@ use vortex_edition::EditionSessionExt;
use vortex_edition::declarations::core::CORE_2026_08_3;
use vortex_error::VortexResult;
use vortex_fastlanes::Delta;
use vortex_sparse::Sparse;

use crate::BtrBlocksCompressor;
use crate::BtrBlocksCompressorBuilder;
Expand Down Expand Up @@ -135,6 +138,48 @@ fn test_delta_unaligned_roundtrip(#[case] input: PrimitiveArray) -> VortexResult
Ok(())
}

/// Fixed-size lists of random bytes, like UUIDs, with the rows `is_valid` picks valid. The null
/// rows hold zeros, as Arrow's fixed-size binary arrays do.
fn random_byte_lists(rows: usize, is_valid: impl Fn(usize) -> bool) -> ArrayRef {
let mut rng = StdRng::seed_from_u64(42);
let elements = (0..rows * 16)
.map(|i| {
if is_valid(i / 16) {
rng.next_u32().to_le_bytes()[0]
} else {
0
}
})
.collect::<Buffer<u8>>();
let validity = Validity::from_iter((0..rows).map(&is_valid));
FixedSizeListArray::new(elements.into_array(), 16, validity, rows).into_array()
}

#[rstest]
#[case::mostly_null(|row| row % 20 == 0, true)]
#[case::half_null(|row| row % 2 == 0, true)]
// Dropping one row in twenty saves less than storing the valid rows' positions costs.
#[case::few_nulls(|row| row % 20 != 0, false)]
#[case::all_valid(|_| true, false)]
fn test_fixed_size_list_drops_null_rows(
#[case] is_valid: fn(usize) -> bool,
#[case] sparse: bool,
) -> VortexResult<()> {
let compressor = BtrBlocksCompressorBuilder::from_session(&SESSION)
.unrestricted()
.build();
let input = random_byte_lists(8_000, is_valid);
let compressed = assert_roundtrip(&compressor, &input)?;
assert_eq!(
compressed.is::<Sparse>(),
sparse,
"{}",
compressed.display_tree()
);

Ok(())
}

#[cfg(feature = "zstd")]
#[rstest]
#[case::array_level(CORE_2026_08_3, vortex_zstd::Zstd.id())]
Expand Down
26 changes: 26 additions & 0 deletions vortex-compressor/src/compressor/cascade.rs
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,32 @@ impl CascadingCompressor {
}
Canonical::Map(map_array) => self.compress_map_array(map_array, compress_ctx, exec_ctx),
Canonical::FixedSizeList(fsl_array) => {
// Canonicalize the elements first, so that schemes see the list's real size and
// their output is accepted only if it beats it.
let elements = fsl_array
.elements()
.clone()
.execute::<CanonicalValidity>(exec_ctx)?
.0
.compact(exec_ctx)?
.into_array();
let fsl_array = FixedSizeListArray::try_new(
elements,
fsl_array.list_size(),
fsl_array.validity()?,
fsl_array.len(),
)?;

// Schemes over the whole list can drop the elements of null rows, which
// compressing the elements alone has to keep.
if let Selection::Compressed(compressed) = self.choose_and_compress(
Canonical::FixedSizeList(fsl_array.clone()),
compress_ctx,
exec_ctx,
)? {
return Ok(compressed);
}

let compressed_elems = self.compress(fsl_array.elements(), exec_ctx)?;

Ok(FixedSizeListArray::try_new(
Expand Down
Loading