Skip to content
Merged
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
37 changes: 25 additions & 12 deletions rust/lance-encoding/src/compression.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1367,11 +1367,16 @@ mod tests {
use super::*;
use crate::buffer::LanceBuffer;
use crate::data::{BlockInfo, DataBlock, FixedWidthDataBlock};
use crate::encodings::logical::primitive::miniblock::MiniBlockCompressionContext;
use crate::statistics::ComputeStat;
use crate::testing::extract_array_encoding_chain;
use arrow_schema::{DataType, Field as ArrowField};
use std::collections::HashMap;

fn miniblock_context() -> MiniBlockCompressionContext {
MiniBlockCompressionContext::new(0, true, true)
}

fn create_test_field(name: &str, data_type: DataType) -> Field {
let arrow_field = ArrowField::new(name, data_type, true);
let mut field = Field::try_from(&arrow_field).unwrap();
Expand Down Expand Up @@ -1759,12 +1764,16 @@ mod tests {
let compressor = strategy
.create_miniblock_compressor(&field, &fixed_data)
.unwrap();
let (_block, encoding) = compressor.compress(fixed_data.clone()).unwrap();
let (_block, encoding) = compressor
.compress(miniblock_context(), fixed_data.clone())
.unwrap();
check_uncompressed_encoding(&encoding, false);
let compressor = strategy
.create_miniblock_compressor(&field, &variable_data)
.unwrap();
let (_block, encoding) = compressor.compress(variable_data.clone()).unwrap();
let (_block, encoding) = compressor
.compress(miniblock_context(), variable_data.clone())
.unwrap();
check_uncompressed_encoding(&encoding, true);

// Test pervalue
Expand Down Expand Up @@ -1794,13 +1803,17 @@ mod tests {
let compressor = strategy
.create_miniblock_compressor(&field, &fixed_data)
.unwrap();
let (_block, encoding) = compressor.compress(fixed_data.clone()).unwrap();
let (_block, encoding) = compressor
.compress(miniblock_context(), fixed_data.clone())
.unwrap();
check_uncompressed_encoding(&encoding, false);

let compressor = strategy
.create_miniblock_compressor(&field, &variable_data)
.unwrap();
let (_block, encoding) = compressor.compress(variable_data.clone()).unwrap();
let (_block, encoding) = compressor
.compress(miniblock_context(), variable_data.clone())
.unwrap();
check_uncompressed_encoding(&encoding, true);

// Test pervalue
Expand Down Expand Up @@ -2097,7 +2110,7 @@ mod tests {

let strategy = DefaultCompressionStrategy::new().with_version(LanceFileVersion::V2_3);
let compressor = strategy.create_miniblock_compressor(&field, &data).unwrap();
let (_compressed, encoding) = compressor.compress(data).unwrap();
let (_compressed, encoding) = compressor.compress(miniblock_context(), data).unwrap();
assert_eq!(rle_run_length_bits(&encoding), 16);
}

Expand All @@ -2122,7 +2135,7 @@ mod tests {

let strategy = DefaultCompressionStrategy::new().with_version(version);
let compressor = strategy.create_miniblock_compressor(&field, &data).unwrap();
let (_compressed, encoding) = compressor.compress(data).unwrap();
let (_compressed, encoding) = compressor.compress(miniblock_context(), data).unwrap();
assert_eq!(rle_run_length_bits(&encoding), 8, "version={version}");
}
}
Expand Down Expand Up @@ -2150,7 +2163,7 @@ mod tests {
let debug_str = format!("{compressor:?}");
assert!(debug_str.contains("RleEncoder"));

let (_compressed, encoding) = compressor.compress(data).unwrap();
let (_compressed, encoding) = compressor.compress(miniblock_context(), data).unwrap();
assert_eq!(rle_run_length_bits(&encoding), 16);
}

Expand All @@ -2173,7 +2186,7 @@ mod tests {

let strategy = DefaultCompressionStrategy::new().with_version(LanceFileVersion::V2_3);
let compressor = strategy.create_miniblock_compressor(&field, &data).unwrap();
let (_compressed, encoding) = compressor.compress(data).unwrap();
let (_compressed, encoding) = compressor.compress(miniblock_context(), data).unwrap();
assert_eq!(rle_run_length_bits(&encoding), 16);
}

Expand All @@ -2196,7 +2209,7 @@ mod tests {

let strategy = DefaultCompressionStrategy::new().with_version(LanceFileVersion::V2_3);
let compressor = strategy.create_miniblock_compressor(&field, &data).unwrap();
let (_compressed, encoding) = compressor.compress(data).unwrap();
let (_compressed, encoding) = compressor.compress(miniblock_context(), data).unwrap();
assert_eq!(rle_run_length_bits(&encoding), 8);
}

Expand Down Expand Up @@ -2233,7 +2246,7 @@ mod tests {
let data = DataBlock::FixedWidth(data);

let compressor = strategy.create_miniblock_compressor(&field, &data).unwrap();
let (_compressed, encoding) = compressor.compress(data).unwrap();
let (_compressed, encoding) = compressor.compress(miniblock_context(), data).unwrap();
let rle = expect_rle_encoding(&encoding);

assert!(
Expand Down Expand Up @@ -2281,7 +2294,7 @@ mod tests {
let debug_str = format!("{compressor:?}");
assert!(debug_str.contains("RleEncoder"));

let (_compressed, encoding) = compressor.compress(data).unwrap();
let (_compressed, encoding) = compressor.compress(miniblock_context(), data).unwrap();
let Compression::Rle(rle) = encoding.compression.as_ref().unwrap() else {
panic!("expected RLE encoding");
};
Expand Down Expand Up @@ -2331,7 +2344,7 @@ mod tests {
"expected RLE to beat inline bitpacking after child selection, got: {debug_str}"
);

let (_compressed, encoding) = compressor.compress(data).unwrap();
let (_compressed, encoding) = compressor.compress(miniblock_context(), data).unwrap();
let rle = expect_rle_encoding(&encoding);
assert!(matches!(
rle.values.as_ref().unwrap().compression.as_ref().unwrap(),
Expand Down
8 changes: 6 additions & 2 deletions rust/lance-encoding/src/encodings/logical/primitive.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ use crate::{
encodings::logical::primitive::fullzip::PerValueDataBlock,
};
use crate::{
encodings::logical::primitive::miniblock::MiniBlockCompressed,
encodings::logical::primitive::miniblock::{MiniBlockCompressed, MiniBlockCompressionContext},
statistics::{ComputeStat, GetStat, Stat},
};
use crate::{
Expand Down Expand Up @@ -5282,7 +5282,11 @@ impl PrimitiveStructuralEncoder {
let num_items = data.num_values();

let compressor = compression_strategy.create_miniblock_compressor(field, &data)?;
let (compressed_data, value_encoding) = compressor.compress(data)?;
let common_chunk_buffers =
u64::from(repdef.rep_slicer().is_some()) + u64::from(repdef.def_slicer().is_some());
let compression_context =
MiniBlockCompressionContext::new(common_chunk_buffers, support_large_chunk, true);
let (compressed_data, value_encoding) = compressor.compress(compression_context, data)?;

let max_rep = repdef.def_meaning.iter().filter(|l| l.is_list()).count() as u16;

Expand Down
29 changes: 28 additions & 1 deletion rust/lance-encoding/src/encodings/logical/primitive/miniblock.rs
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,29 @@ pub struct MiniBlockCompressed {
pub num_values: u64,
}

/// Per-page framing details that can affect a mini-block compressor's choice.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct MiniBlockCompressionContext {
common_chunk_buffers: u64,
support_large_chunk: bool,
allow_generic_offsets: bool,
}

impl MiniBlockCompressionContext {
/// Creates the framing context supplied by the owning mini-block page.
pub fn new(
common_chunk_buffers: u64,
support_large_chunk: bool,
allow_generic_offsets: bool,
) -> Self {
Self {
common_chunk_buffers,
support_large_chunk,
allow_generic_offsets,
}
}
}

/// Describes the size of a mini-block chunk of data
///
/// Mini-block chunks are designed to be small (just a few disk sectors)
Expand Down Expand Up @@ -113,7 +136,11 @@ pub trait MiniBlockCompressor: std::fmt::Debug + Send + Sync {
///
/// This method also returns a description of the encoding applied that will be
/// used at decode time to read the data.
fn compress(&self, page: DataBlock) -> Result<(MiniBlockCompressed, CompressiveEncoding)>;
fn compress(
&self,
context: MiniBlockCompressionContext,
page: DataBlock,
) -> Result<(MiniBlockCompressed, CompressiveEncoding)>;
}

#[cfg(test)]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,8 @@ use super::{
SparseValidityMeaning, SparseValiditySet,
};
use crate::encodings::logical::primitive::{
FILL_BYTE, MINIBLOCK_ALIGNMENT, miniblock::MiniBlockCompressed,
FILL_BYTE, MINIBLOCK_ALIGNMENT,
miniblock::{MiniBlockCompressed, MiniBlockCompressionContext},
};

#[derive(Clone, Copy, Default)]
Expand Down Expand Up @@ -563,7 +564,8 @@ pub fn prepare_values(

let num_values = data.num_values();
let compressor = compression_strategy.create_miniblock_compressor(field, &data)?;
let (compressed, value_compression) = compressor.compress(data)?;
let compression_context = MiniBlockCompressionContext::new(0, support_large_chunk, false);
let (compressed, value_compression) = compressor.compress(compression_context, data)?;
let values =
serialize_value_chunks(with_explicit_value_counts(compressed)?, support_large_chunk)?;
Ok(PreparedSparseValues {
Expand Down
9 changes: 7 additions & 2 deletions rust/lance-encoding/src/encodings/physical/binary.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,8 @@ use crate::buffer::LanceBuffer;
use crate::data::{BlockInfo, DataBlock, VariableWidthBlock};
use crate::encodings::logical::primitive::fullzip::{PerValueCompressor, PerValueDataBlock};
use crate::encodings::logical::primitive::miniblock::{
MAX_MINIBLOCK_VALUES, MiniBlockChunk, MiniBlockCompressed, MiniBlockCompressor,
MAX_MINIBLOCK_VALUES, MiniBlockChunk, MiniBlockCompressed, MiniBlockCompressionContext,
MiniBlockCompressor,
};
use crate::format::pb21::CompressiveEncoding;
use crate::format::pb21::compressive_encoding::Compression;
Expand Down Expand Up @@ -245,7 +246,11 @@ impl BinaryMiniBlockEncoder {
}

impl MiniBlockCompressor for BinaryMiniBlockEncoder {
fn compress(&self, data: DataBlock) -> Result<(MiniBlockCompressed, CompressiveEncoding)> {
fn compress(
&self,
_context: MiniBlockCompressionContext,
data: DataBlock,
) -> Result<(MiniBlockCompressed, CompressiveEncoding)> {
match data {
DataBlock::VariableWidth(variable_width) => Ok(self.chunk_data(variable_width)),
_ => Err(Error::invalid_input_source(
Expand Down
8 changes: 6 additions & 2 deletions rust/lance-encoding/src/encodings/physical/bitpacking.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ use crate::compression::{BlockCompressor, BlockDecompressor, MiniBlockDecompress
use crate::data::BlockInfo;
use crate::data::{DataBlock, FixedWidthDataBlock};
use crate::encodings::logical::primitive::miniblock::{
MiniBlockChunk, MiniBlockCompressed, MiniBlockCompressor,
MiniBlockChunk, MiniBlockCompressed, MiniBlockCompressionContext, MiniBlockCompressor,
};
use crate::format::pb21::CompressiveEncoding;
use crate::format::{ProtobufUtils21, pb21};
Expand Down Expand Up @@ -215,7 +215,11 @@ impl InlineBitpacking {
}

impl MiniBlockCompressor for InlineBitpacking {
fn compress(&self, chunk: DataBlock) -> Result<(MiniBlockCompressed, CompressiveEncoding)> {
fn compress(
&self,
_context: MiniBlockCompressionContext,
chunk: DataBlock,
) -> Result<(MiniBlockCompressed, CompressiveEncoding)> {
match chunk {
DataBlock::FixedWidth(fixed_width) => Ok(self.chunk_data(fixed_width)),
_ => Err(Error::invalid_input_source(
Expand Down
20 changes: 15 additions & 5 deletions rust/lance-encoding/src/encodings/physical/byte_stream_split.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,7 @@ use crate::compression::MiniBlockDecompressor;
use crate::compression_config::BssMode;
use crate::data::{BlockInfo, DataBlock, FixedWidthDataBlock};
use crate::encodings::logical::primitive::miniblock::{
MiniBlockChunk, MiniBlockCompressed, MiniBlockCompressor,
MiniBlockChunk, MiniBlockCompressed, MiniBlockCompressionContext, MiniBlockCompressor,
};
use crate::format::ProtobufUtils21;
use crate::format::pb21::CompressiveEncoding;
Expand Down Expand Up @@ -107,7 +107,11 @@ impl ByteStreamSplitEncoder {
}

impl MiniBlockCompressor for ByteStreamSplitEncoder {
fn compress(&self, page: DataBlock) -> Result<(MiniBlockCompressed, CompressiveEncoding)> {
fn compress(
&self,
_context: MiniBlockCompressionContext,
page: DataBlock,
) -> Result<(MiniBlockCompressed, CompressiveEncoding)> {
match page {
DataBlock::FixedWidth(data) => {
let num_values = data.num_values;
Expand Down Expand Up @@ -345,7 +349,9 @@ mod tests {
});

// Compress
let (compressed, _encoding) = encoder.compress(data_block).unwrap();
let (compressed, _encoding) = encoder
.compress(MiniBlockCompressionContext::new(0, true, true), data_block)
.unwrap();

// Decompress
let decompressed = decompressor
Expand Down Expand Up @@ -391,7 +397,9 @@ mod tests {
});

// Compress
let (compressed, _encoding) = encoder.compress(data_block).unwrap();
let (compressed, _encoding) = encoder
.compress(MiniBlockCompressionContext::new(0, true, true), data_block)
.unwrap();

// Decompress
let decompressed = decompressor
Expand Down Expand Up @@ -424,7 +432,9 @@ mod tests {
});

// Compress empty data
let (compressed, _encoding) = encoder.compress(data_block).unwrap();
let (compressed, _encoding) = encoder
.compress(MiniBlockCompressionContext::new(0, true, true), data_block)
.unwrap();

// Decompress empty data
let decompressed = decompressor.decompress(compressed.data, 0).unwrap();
Expand Down
10 changes: 7 additions & 3 deletions rust/lance-encoding/src/encodings/physical/fsst.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ use crate::{
data::{BlockInfo, DataBlock, VariableWidthBlock},
encodings::logical::primitive::{
fullzip::{PerValueCompressor, PerValueDataBlock},
miniblock::{MiniBlockCompressed, MiniBlockCompressor},
miniblock::{MiniBlockCompressed, MiniBlockCompressionContext, MiniBlockCompressor},
},
format::{
ProtobufUtils21,
Expand Down Expand Up @@ -138,7 +138,11 @@ impl FsstMiniBlockEncoder {
}

impl MiniBlockCompressor for FsstMiniBlockEncoder {
fn compress(&self, data: DataBlock) -> Result<(MiniBlockCompressed, CompressiveEncoding)> {
fn compress(
&self,
context: MiniBlockCompressionContext,
data: DataBlock,
) -> Result<(MiniBlockCompressed, CompressiveEncoding)> {
let compressed = FsstCompressed::fsst_compress(data)?;

let data_block = DataBlock::VariableWidth(compressed.data);
Expand All @@ -148,7 +152,7 @@ impl MiniBlockCompressor for FsstMiniBlockEncoder {
as Box<dyn MiniBlockCompressor>;

let (binary_miniblock_compressed, binary_array_encoding) =
binary_compressor.compress(data_block)?;
binary_compressor.compress(context, data_block)?;

Ok((
binary_miniblock_compressed,
Expand Down
Loading
Loading