lance_encoding/encodings/physical/
byte_stream_split.rs1use std::fmt::Debug;
59
60use crate::buffer::LanceBuffer;
61use crate::compression::MiniBlockDecompressor;
62use crate::compression_config::BssMode;
63use crate::data::{BlockInfo, DataBlock, FixedWidthDataBlock};
64use crate::encodings::logical::primitive::miniblock::{
65 MiniBlockChunk, MiniBlockCompressed, MiniBlockCompressionContext, MiniBlockCompressor,
66};
67use crate::format::ProtobufUtils21;
68use crate::format::pb21::CompressiveEncoding;
69use crate::statistics::{GetStat, Stat};
70use arrow_array::{cast::AsArray, types::UInt64Type};
71use lance_core::Result;
72
73#[derive(Debug, Clone)]
79pub struct ByteStreamSplitEncoder {
80 bits_per_value: usize,
81}
82
83impl ByteStreamSplitEncoder {
84 pub fn new(bits_per_value: usize) -> Self {
85 assert!(
86 bits_per_value == 32 || bits_per_value == 64,
87 "ByteStreamSplit only supports 32-bit (f32) or 64-bit (f64) values"
88 );
89 Self { bits_per_value }
90 }
91
92 fn bytes_per_value(&self) -> usize {
93 self.bits_per_value / 8
94 }
95
96 fn max_chunk_size(&self) -> usize {
97 match self.bits_per_value {
102 32 => 1024,
103 64 => 512,
104 _ => unreachable!("ByteStreamSplit only supports 32 or 64 bit values"),
105 }
106 }
107}
108
109impl MiniBlockCompressor for ByteStreamSplitEncoder {
110 fn compress(
111 &self,
112 _context: MiniBlockCompressionContext,
113 page: DataBlock,
114 ) -> Result<(MiniBlockCompressed, CompressiveEncoding)> {
115 match page {
116 DataBlock::FixedWidth(data) => {
117 let num_values = data.num_values;
118 let bytes_per_value = self.bytes_per_value();
119
120 if num_values == 0 {
121 return Ok((
122 MiniBlockCompressed {
123 data: vec![],
124 chunks: vec![],
125 num_values: 0,
126 },
127 ProtobufUtils21::byte_stream_split(ProtobufUtils21::flat(
128 self.bits_per_value as u64,
129 None,
130 )),
131 ));
132 }
133
134 let total_size = num_values as usize * bytes_per_value;
135 let mut global_buffer = vec![0u8; total_size];
136
137 let mut chunks = Vec::new();
138 let data_slice = data.data.as_ref();
139 let mut processed_values = 0usize;
140 let max_chunk_size = self.max_chunk_size();
141
142 while processed_values < num_values as usize {
143 let chunk_size = (num_values as usize - processed_values).min(max_chunk_size);
144 let chunk_offset = processed_values * bytes_per_value;
145
146 for i in 0..chunk_size {
148 let src_offset = (processed_values + i) * bytes_per_value;
149 for j in 0..bytes_per_value {
150 let dst_offset = chunk_offset + j * chunk_size + i;
152 global_buffer[dst_offset] = data_slice[src_offset + j];
153 }
154 }
155
156 let chunk_bytes = chunk_size * bytes_per_value;
157 let log_num_values = if processed_values + chunk_size == num_values as usize {
158 0 } else {
160 chunk_size.ilog2() as u8
161 };
162
163 debug_assert!(chunk_bytes > 0);
164 chunks.push(MiniBlockChunk {
165 buffer_sizes: vec![chunk_bytes as u32],
166 log_num_values,
167 });
168
169 processed_values += chunk_size;
170 }
171
172 let data_buffers = vec![LanceBuffer::from(global_buffer)];
173
174 let encoding = ProtobufUtils21::byte_stream_split(ProtobufUtils21::flat(
176 self.bits_per_value as u64,
177 None,
178 ));
179
180 Ok((
181 MiniBlockCompressed {
182 data: data_buffers,
183 chunks,
184 num_values,
185 },
186 encoding,
187 ))
188 }
189 _ => Err(lance_core::Error::invalid_input_source(
190 "ByteStreamSplit encoding only supports FixedWidth data blocks".into(),
191 )),
192 }
193 }
194}
195
196#[derive(Debug)]
198pub struct ByteStreamSplitDecompressor {
199 bits_per_value: usize,
200}
201
202impl ByteStreamSplitDecompressor {
203 pub fn new(bits_per_value: usize) -> Self {
204 assert!(
205 bits_per_value == 32 || bits_per_value == 64,
206 "ByteStreamSplit only supports 32-bit (f32) or 64-bit (f64) values"
207 );
208 Self { bits_per_value }
209 }
210
211 fn bytes_per_value(&self) -> usize {
212 self.bits_per_value / 8
213 }
214}
215
216impl MiniBlockDecompressor for ByteStreamSplitDecompressor {
217 fn decompress(&self, data: Vec<LanceBuffer>, num_values: u64) -> Result<DataBlock> {
218 if num_values == 0 {
219 return Ok(DataBlock::FixedWidth(FixedWidthDataBlock {
220 data: LanceBuffer::empty(),
221 bits_per_value: self.bits_per_value as u64,
222 num_values: 0,
223 block_info: BlockInfo::new(),
224 }));
225 }
226
227 let bytes_per_value = self.bytes_per_value();
228 let total_bytes = num_values as usize * bytes_per_value;
229
230 if data.len() != 1 {
231 return Err(lance_core::Error::invalid_input_source(
232 format!(
233 "ByteStreamSplit decompression expects 1 buffer, but got {}",
234 data.len()
235 )
236 .into(),
237 ));
238 }
239
240 let input_buffer = &data[0];
241
242 if input_buffer.len() != total_bytes {
243 return Err(lance_core::Error::invalid_input_source(
244 format!(
245 "Expected {} bytes for decompression, but got {}",
246 total_bytes,
247 input_buffer.len()
248 )
249 .into(),
250 ));
251 }
252
253 let mut output = vec![0u8; total_bytes];
254
255 for i in 0..num_values as usize {
257 for j in 0..bytes_per_value {
258 let src_offset = j * num_values as usize + i;
259 output[i * bytes_per_value + j] = input_buffer[src_offset];
260 }
261 }
262
263 Ok(DataBlock::FixedWidth(FixedWidthDataBlock {
264 data: LanceBuffer::from(output),
265 bits_per_value: self.bits_per_value as u64,
266 num_values,
267 block_info: BlockInfo::new(),
268 }))
269 }
270
271 fn decoded_size_bytes(&self, num_values: u64) -> Option<u64> {
272 num_values.checked_mul(self.bytes_per_value() as u64)
273 }
274}
275
276pub fn should_use_bss(data: &FixedWidthDataBlock, mode: BssMode) -> bool {
278 if data.bits_per_value != 32 && data.bits_per_value != 64 {
282 return false;
283 }
284
285 let sensitivity = mode.to_sensitivity();
286
287 if sensitivity <= 0.0 {
289 return false;
290 }
291 if sensitivity >= 1.0 {
292 return true;
293 }
294
295 evaluate_entropy_for_bss(data, sensitivity)
297}
298
299fn evaluate_entropy_for_bss(data: &FixedWidthDataBlock, sensitivity: f32) -> bool {
301 let Some(entropy_stat) = data.get_stat(Stat::BytePositionEntropy) else {
303 return false; };
305
306 let entropies = entropy_stat.as_primitive::<UInt64Type>();
307 if entropies.is_empty() {
308 return false;
309 }
310
311 let sum: u64 = entropies.values().iter().sum();
313 let avg_entropy = sum as f64 / entropies.len() as f64 / 1000.0; let entropy_threshold = sensitivity as f64 * 8.0;
320
321 avg_entropy < entropy_threshold
324}
325
326#[cfg(test)]
327mod tests {
328 use super::*;
329
330 #[test]
331 fn test_round_trip_f32() {
332 let encoder = ByteStreamSplitEncoder::new(32);
333 let decompressor = ByteStreamSplitDecompressor::new(32);
334
335 let values: Vec<f32> = vec![
337 1.0,
338 2.5,
339 -3.7,
340 4.2,
341 0.0,
342 -0.0,
343 f32::INFINITY,
344 f32::NEG_INFINITY,
345 ];
346 let bytes: Vec<u8> = values.iter().flat_map(|v| v.to_le_bytes()).collect();
347
348 let data_block = DataBlock::FixedWidth(FixedWidthDataBlock {
349 data: LanceBuffer::from(bytes),
350 bits_per_value: 32,
351 num_values: values.len() as u64,
352 block_info: BlockInfo::new(),
353 });
354
355 let (compressed, _encoding) = encoder
357 .compress(MiniBlockCompressionContext::new(0, true, true), data_block)
358 .unwrap();
359
360 let decompressed = decompressor
362 .decompress(compressed.data, values.len() as u64)
363 .unwrap();
364 let DataBlock::FixedWidth(decompressed_fixed) = &decompressed else {
365 panic!("Expected FixedWidth DataBlock")
366 };
367
368 let result_bytes = decompressed_fixed.data.as_ref();
370 let result_values: Vec<f32> = result_bytes
371 .chunks_exact(4)
372 .map(|chunk| f32::from_le_bytes(chunk.try_into().unwrap()))
373 .collect();
374
375 assert_eq!(values, result_values);
376 }
377
378 #[test]
379 fn test_round_trip_f64() {
380 let encoder = ByteStreamSplitEncoder::new(64);
381 let decompressor = ByteStreamSplitDecompressor::new(64);
382
383 let values: Vec<f64> = vec![
385 1.0,
386 2.5,
387 -3.7,
388 4.2,
389 0.0,
390 -0.0,
391 f64::INFINITY,
392 f64::NEG_INFINITY,
393 ];
394 let bytes: Vec<u8> = values.iter().flat_map(|v| v.to_le_bytes()).collect();
395
396 let data_block = DataBlock::FixedWidth(FixedWidthDataBlock {
397 data: LanceBuffer::from(bytes),
398 bits_per_value: 64,
399 num_values: values.len() as u64,
400 block_info: BlockInfo::new(),
401 });
402
403 let (compressed, _encoding) = encoder
405 .compress(MiniBlockCompressionContext::new(0, true, true), data_block)
406 .unwrap();
407
408 let decompressed = decompressor
410 .decompress(compressed.data, values.len() as u64)
411 .unwrap();
412 let DataBlock::FixedWidth(decompressed_fixed) = &decompressed else {
413 panic!("Expected FixedWidth DataBlock")
414 };
415
416 let result_bytes = decompressed_fixed.data.as_ref();
418 let result_values: Vec<f64> = result_bytes
419 .chunks_exact(8)
420 .map(|chunk| f64::from_le_bytes(chunk.try_into().unwrap()))
421 .collect();
422
423 assert_eq!(values, result_values);
424 }
425
426 #[test]
427 fn test_empty_data() {
428 let encoder = ByteStreamSplitEncoder::new(32);
429 let decompressor = ByteStreamSplitDecompressor::new(32);
430
431 let data_block = DataBlock::FixedWidth(FixedWidthDataBlock {
432 data: LanceBuffer::empty(),
433 bits_per_value: 32,
434 num_values: 0,
435 block_info: BlockInfo::new(),
436 });
437
438 let (compressed, _encoding) = encoder
440 .compress(MiniBlockCompressionContext::new(0, true, true), data_block)
441 .unwrap();
442
443 let decompressed = decompressor.decompress(compressed.data, 0).unwrap();
445 let DataBlock::FixedWidth(decompressed_fixed) = &decompressed else {
446 panic!("Expected FixedWidth DataBlock")
447 };
448
449 assert_eq!(decompressed_fixed.num_values, 0);
450 assert_eq!(decompressed_fixed.data.len(), 0);
451 }
452}