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
3 changes: 1 addition & 2 deletions rust/lance-encoding/benches/common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@ use lance_encoding::{
},
},
encodings::logical::primitive::{fullzip::PerValueCompressor, miniblock::MiniBlockCompressor},
format::pb21::CompressiveEncoding,
};

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
Expand Down Expand Up @@ -141,7 +140,7 @@ impl CompressionStrategy for BenchCompressionStrategy {
&self,
field: &Field,
data: &DataBlock,
) -> Result<(Box<dyn BlockCompressor>, CompressiveEncoding)> {
) -> Result<Box<dyn BlockCompressor>> {
let params = self.field_params(field);
if self.encoding == BenchEncoding::StructuralU32
&& let Some(compressor) = try_fixed_u8_rle_block(data, &params)?
Expand Down
2 changes: 1 addition & 1 deletion rust/lance-encoding/benches/decoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -651,7 +651,7 @@ where
#[cfg(feature = "bitpacking")]
fn typed_view_unchunk(buffer: LanceBuffer, uncompressed_bits: u64, num_values: u64) -> DataBlock {
InlineBitpacking::new(uncompressed_bits)
.decompress(buffer, num_values)
.decompress(Some(buffer), num_values)
.unwrap()
}

Expand Down
162 changes: 86 additions & 76 deletions rust/lance-encoding/src/compression.rs
Original file line number Diff line number Diff line change
Expand Up @@ -61,10 +61,7 @@ use crate::{
value::{ValueDecompressor, ValueEncoder},
},
},
format::{
ProtobufUtils21,
pb21::{CompressiveEncoding, compressive_encoding::Compression},
},
format::pb21::{CompressiveEncoding, compressive_encoding::Compression},
statistics::{GetStat, Stat},
};

Expand Down Expand Up @@ -98,11 +95,11 @@ const RLE_BLOCK_HEADER_BYTES: u128 = std::mem::size_of::<u64>() as u128;
/// required (e.g. when encoding metadata buffers like a dictionary or for encoding rep/def
/// mini-block chunks)
pub trait BlockCompressor: std::fmt::Debug + Send + Sync {
/// Compress the data into a single buffer
/// Compress the data into zero or one buffers and describe the codec used.
///
/// Also returns a description of the compression that can be used to decompress
/// when reading the data back
fn compress(&self, data: DataBlock) -> Result<LanceBuffer>;
/// `None` represents a metadata-only codec. `Some` represents a physical
/// payload, including a zero-byte payload.
fn compress(&self, data: DataBlock) -> Result<(Option<LanceBuffer>, CompressiveEncoding)>;
}

/// A trait to pick which compression to use for given data
Expand All @@ -118,12 +115,12 @@ pub trait BlockCompressor: std::fmt::Debug + Send + Sync {
/// used for narrow data types (both fixed and variable length) where we can
/// fit many values into an 16KiB block.
pub trait CompressionStrategy: Send + Sync + std::fmt::Debug {
/// Create a block compressor for the given data
/// Create a block compressor for the given data.
fn create_block_compressor(
&self,
field: &Field,
data: &DataBlock,
) -> Result<(Box<dyn BlockCompressor>, CompressiveEncoding)>;
) -> Result<Box<dyn BlockCompressor>>;

/// Create a per-value compressor for the given data
fn create_per_value(
Expand All @@ -140,6 +137,19 @@ pub trait CompressionStrategy: Send + Sync + std::fmt::Debug {
) -> Result<Box<dyn MiniBlockCompressor>>;
}

pub(crate) fn compress_required_block(
strategy: &dyn CompressionStrategy,
field: &Field,
data: DataBlock,
) -> Result<(LanceBuffer, CompressiveEncoding)> {
let compressor = strategy.create_block_compressor(field, &data)?;
let (payload, encoding) = compressor.compress(data)?;
let payload = payload.ok_or_else(|| {
Error::internal("Required block compressor selected a metadata-only codec".to_string())
})?;
Ok((payload, encoding))
}

fn try_bss_for_mini_block(
data: &FixedWidthDataBlock,
params: &CompressionFieldParams,
Expand Down Expand Up @@ -269,7 +279,7 @@ fn try_rle_for_block_with_width(
params: &CompressionFieldParams,
run_length_width: RunLengthWidth,
rle_payload_bytes: u128,
) -> Result<Option<(Box<dyn BlockCompressor>, CompressiveEncoding)>> {
) -> Result<Option<Box<dyn BlockCompressor>>> {
let bits = data.bits_per_value;
if !matches!(bits, 8 | 16 | 32 | 64) {
return Ok(None);
Expand Down Expand Up @@ -305,18 +315,15 @@ fn try_rle_for_block_with_width(
}
}

let compressor = Box::new(RleEncoder::with_run_length_width(run_length_width));
let encoding = ProtobufUtils21::rle(
ProtobufUtils21::flat(bits, None),
ProtobufUtils21::flat(run_length_width.bits_per_value(), None),
);
Ok(Some((compressor, encoding)))
Ok(Some(Box::new(RleEncoder::with_run_length_width(
run_length_width,
))))
}

fn try_fixed_u8_rle_for_block(
data: &FixedWidthDataBlock,
params: &CompressionFieldParams,
) -> Result<Option<(Box<dyn BlockCompressor>, CompressiveEncoding)>> {
) -> Result<Option<Box<dyn BlockCompressor>>> {
if !matches!(data.bits_per_value, 8 | 16 | 32 | 64) {
return Ok(None);
}
Expand All @@ -327,7 +334,7 @@ fn try_fixed_u8_rle_for_block(
fn try_variable_rle_for_block(
data: &FixedWidthDataBlock,
params: &CompressionFieldParams,
) -> Result<Option<(Box<dyn BlockCompressor>, CompressiveEncoding)>> {
) -> Result<Option<Box<dyn BlockCompressor>>> {
if !matches!(data.bits_per_value, 8 | 16 | 32 | 64) {
return Ok(None);
}
Expand Down Expand Up @@ -410,9 +417,7 @@ fn estimate_inline_bitpacking_bytes(data: &FixedWidthDataBlock) -> Option<u64> {
u64::try_from(estimated_bytes).ok()
}

fn try_bitpack_for_block(
data: &FixedWidthDataBlock,
) -> Option<(Box<dyn BlockCompressor>, CompressiveEncoding)> {
fn try_bitpack_for_block(data: &FixedWidthDataBlock) -> Option<Box<dyn BlockCompressor>> {
let bits = data.bits_per_value;
if !matches!(bits, 8 | 16 | 32 | 64) {
return None;
Expand All @@ -430,16 +435,9 @@ fn try_bitpack_for_block(
}

if data.num_values <= 1024 {
let compressor = Box::new(InlineBitpacking::new(bits));
let encoding = ProtobufUtils21::inline_bitpacking(bits, None);
Some((compressor, encoding))
Some(Box::new(InlineBitpacking::new(bits)))
} else {
let compressor = Box::new(OutOfLineBitpacking::new(max_bit_width, bits));
let encoding = ProtobufUtils21::out_of_line_bitpacking(
bits,
ProtobufUtils21::flat(max_bit_width, None),
);
Some((compressor, encoding))
Some(Box::new(OutOfLineBitpacking::new(max_bit_width, bits)))
}
}

Expand Down Expand Up @@ -818,7 +816,7 @@ pub fn try_variable_width_per_value(
pub fn try_fixed_u8_rle_block(
data: &DataBlock,
params: &CompressionFieldParams,
) -> Result<Option<(Box<dyn BlockCompressor>, CompressiveEncoding)>> {
) -> Result<Option<Box<dyn BlockCompressor>>> {
let DataBlock::FixedWidth(data) = data else {
return Ok(None);
};
Expand All @@ -829,17 +827,15 @@ pub fn try_fixed_u8_rle_block(
pub fn try_variable_rle_block(
data: &DataBlock,
params: &CompressionFieldParams,
) -> Result<Option<(Box<dyn BlockCompressor>, CompressiveEncoding)>> {
) -> Result<Option<Box<dyn BlockCompressor>>> {
let DataBlock::FixedWidth(data) = data else {
return Ok(None);
};
try_variable_rle_for_block(data, params)
}

/// Select block bitpacking for applicable fixed-width values.
pub fn try_bitpacking_block(
data: &DataBlock,
) -> Option<(Box<dyn BlockCompressor>, CompressiveEncoding)> {
pub fn try_bitpacking_block(data: &DataBlock) -> Option<Box<dyn BlockCompressor>> {
let DataBlock::FixedWidth(data) = data else {
return None;
};
Expand All @@ -850,35 +846,22 @@ pub fn try_bitpacking_block(
pub fn try_general_block(
data: &DataBlock,
params: &CompressionFieldParams,
) -> Result<Option<(Box<dyn BlockCompressor>, CompressiveEncoding)>> {
let Some((compressor, config)) = try_general_compression(params, data)? else {
) -> Result<Option<Box<dyn BlockCompressor>>> {
let Some((compressor, _config)) = try_general_compression(params, data)? else {
return Ok(None);
};
let inner = match data {
DataBlock::FixedWidth(data) => ProtobufUtils21::flat(data.bits_per_value, None),
DataBlock::VariableWidth(data) => ProtobufUtils21::variable(
ProtobufUtils21::flat(data.bits_per_offset as u64, None),
None,
),
_ => return Ok(None),
};
Ok(Some((compressor, ProtobufUtils21::wrapped(config, inner)?)))
Ok(Some(compressor))
}

/// Store fixed- and variable-width block values without block compression.
pub fn try_raw_block(data: &DataBlock) -> Option<(Box<dyn BlockCompressor>, CompressiveEncoding)> {
pub fn try_raw_block(data: &DataBlock) -> Option<Box<dyn BlockCompressor>> {
match data {
DataBlock::FixedWidth(data) => Some((
Box::new(ValueEncoder::default()) as Box<dyn BlockCompressor>,
ProtobufUtils21::flat(data.bits_per_value, None),
)),
DataBlock::VariableWidth(data) => Some((
Box::new(VariableEncoder::default()) as Box<dyn BlockCompressor>,
ProtobufUtils21::variable(
ProtobufUtils21::flat(data.bits_per_offset as u64, None),
None,
),
)),
DataBlock::FixedWidth(_) => {
Some(Box::new(ValueEncoder::default()) as Box<dyn BlockCompressor>)
}
DataBlock::VariableWidth(_) => {
Some(Box::new(VariableEncoder::default()) as Box<dyn BlockCompressor>)
}
_ => None,
}
}
Expand Down Expand Up @@ -922,7 +905,23 @@ pub trait VariablePerValueDecompressor: std::fmt::Debug + Send + Sync {
}

pub trait BlockDecompressor: std::fmt::Debug + Send + Sync {
fn decompress(&self, data: LanceBuffer, num_values: u64) -> Result<DataBlock>;
fn decompress(&self, data: Option<LanceBuffer>, num_values: u64) -> Result<DataBlock>;

/// Whether this codec consumes one payload buffer.
fn requires_payload(&self) -> bool {
true
}
}

pub(crate) fn require_block_payload(data: Option<LanceBuffer>, codec: &str) -> Result<LanceBuffer> {
data.ok_or_else(|| Error::invalid_input(format!("{codec} requires one payload")))
}

pub(crate) fn require_no_block_payload(data: Option<LanceBuffer>, codec: &str) -> Result<()> {
if data.is_some() {
return Err(Error::invalid_input(format!("{codec} expects no payload")));
}
Ok(())
}

pub trait DecompressionStrategy: std::fmt::Debug + Send + Sync {
Expand Down Expand Up @@ -1389,6 +1388,16 @@ mod tests {
strategy(TestEncoding::StructuralU16, params)
}

fn selected_block_codec(
strategy: &Arc<dyn CompressionStrategy>,
field: &Field,
data: &DataBlock,
) -> (Box<dyn BlockCompressor>, CompressiveEncoding) {
let compressor = strategy.create_block_compressor(field, data).unwrap();
let (_, encoding) = compressor.compress(data.clone()).unwrap();
(compressor, encoding)
}

fn miniblock_context() -> MiniBlockCompressionContext {
MiniBlockCompressionContext::new(0, true, true)
}
Expand Down Expand Up @@ -1643,7 +1652,7 @@ mod tests {
block.compute_stat();
let data = DataBlock::FixedWidth(block);

let (compressor, _encoding) = strategy.create_block_compressor(&field, &data).unwrap();
let compressor = strategy.create_block_compressor(&field, &data).unwrap();
let debug_str = format!("{:?}", compressor);
assert!(
debug_str.contains("OutOfLineBitpacking"),
Expand All @@ -1665,7 +1674,7 @@ mod tests {
block.compute_stat();
let data = DataBlock::FixedWidth(block);

let (compressor, encoding) = strategy.create_block_compressor(&field, &data).unwrap();
let (compressor, encoding) = selected_block_codec(&strategy, &field, &data);

assert!(format!("{compressor:?}").contains("ValueEncoder"));
assert!(matches!(
Expand Down Expand Up @@ -1693,7 +1702,7 @@ mod tests {
block.compute_stat();
let data = DataBlock::FixedWidth(block);

let (compressor, encoding) = strategy.create_block_compressor(&field, &data).unwrap();
let (compressor, encoding) = selected_block_codec(&strategy, &field, &data);
let debug_str = format!("{compressor:?}");
assert!(
debug_str.contains("OutOfLineBitpacking"),
Expand Down Expand Up @@ -2492,18 +2501,17 @@ mod tests {
let expected_num_values = expected_block.num_values;
let num_values = expected_num_values;

let (compressor, encoding) = strategy
let compressor = strategy
.create_block_compressor(&field, &data)
.expect("general compression should be selected");
let (compressed_buffer, encoding) = compressor
.compress(data.clone())
.expect("write path general compression should succeed");
match encoding.compression.as_ref() {
Some(Compression::General(_)) => {}
other => panic!("expected general compression, got {:?}", other),
}

let compressed_buffer = compressor
.compress(data.clone())
.expect("write path general compression should succeed");

let decompressor = DefaultDecompressionStrategy::default()
.create_block_decompressor(&encoding)
.expect("general block decompressor should be created");
Expand Down Expand Up @@ -2538,9 +2546,10 @@ mod tests {
let field = create_test_field("dict_values", DataType::FixedSizeBinary(3));
let data = create_fixed_width_block(24, 1024);

let (_compressor, encoding) = strategy
let compressor = strategy
.create_block_compressor(&field, &data)
.expect("block compressor selection should succeed");
let (_, encoding) = compressor.compress(data).unwrap();

assert!(
!matches!(encoding.compression.as_ref(), Some(Compression::General(_))),
Expand Down Expand Up @@ -2568,9 +2577,10 @@ mod tests {
"test requires block size above automatic general compression threshold"
);

let (_compressor, encoding) = strategy
let compressor = strategy
.create_block_compressor(&field, &data)
.expect("block compressor selection should succeed");
let (_, encoding) = compressor.compress(data).unwrap();

assert!(
!matches!(encoding.compression.as_ref(), Some(Compression::General(_))),
Expand All @@ -2592,10 +2602,9 @@ mod tests {
let data = DataBlock::FixedWidth(block);

let strategy = strategy(TestEncoding::StructuralSparse, CompressionParams::new());
let (compressor, encoding) = strategy.create_block_compressor(&field, &data).unwrap();
let compressor = strategy.create_block_compressor(&field, &data).unwrap();
let (compressed, encoding) = compressor.compress(data).unwrap();
assert_eq!(rle_run_length_bits(&encoding), 32);

let compressed = compressor.compress(data).unwrap();
let decompressor = DefaultDecompressionStrategy::default()
.create_block_decompressor(&encoding)
.unwrap();
Expand Down Expand Up @@ -2626,7 +2635,8 @@ mod tests {
let data = DataBlock::FixedWidth(block);

let strategy = strategy(TestEncoding::StructuralU32, CompressionParams::new());
let (_compressor, encoding) = strategy.create_block_compressor(&field, &data).unwrap();
let compressor = strategy.create_block_compressor(&field, &data).unwrap();
let (_, encoding) = compressor.compress(data).unwrap();
assert_eq!(rle_run_length_bits(&encoding), 8);
}

Expand Down Expand Up @@ -2656,7 +2666,7 @@ mod tests {

let strategy = strategy(TestEncoding::StructuralU32, CompressionParams::new());

let (compressor, _) = strategy
let compressor = strategy
.create_block_compressor(&field, &data_block)
.unwrap();

Expand Down Expand Up @@ -2690,7 +2700,7 @@ mod tests {

let strategy = strategy(TestEncoding::StructuralU16, CompressionParams::new());

let (compressor, _) = strategy
let compressor = strategy
.create_block_compressor(&field, &data_block)
.unwrap();

Expand Down
Loading
Loading