1use std::{
5 any::Any,
6 collections::{HashMap, VecDeque},
7 env,
8 fmt::Debug,
9 iter,
10 ops::Range,
11 sync::Arc,
12 vec,
13};
14
15use crate::{
16 constants::{
17 STRUCTURAL_ENCODING_FULLZIP, STRUCTURAL_ENCODING_META_KEY, STRUCTURAL_ENCODING_MINIBLOCK,
18 },
19 data::DictionaryDataBlock,
20 encodings::logical::primitive::blob::{BlobDescriptionPageScheduler, BlobPageScheduler},
21 format::{
22 ProtobufUtils21,
23 pb21::{self, CompressiveEncoding, PageLayout, compressive_encoding::Compression},
24 },
25};
26use arrow_array::{Array, ArrayRef, PrimitiveArray, cast::AsArray, make_array, types::UInt64Type};
27use arrow_buffer::{BooleanBuffer, BooleanBufferBuilder, NullBuffer, ScalarBuffer};
28use arrow_schema::{DataType, Field as ArrowField};
29use bytes::Bytes;
30use futures::{FutureExt, TryStreamExt, future::BoxFuture, stream::FuturesOrdered};
31use itertools::Itertools;
32use lance_arrow::DataTypeExt;
33use lance_arrow::deepcopy::deep_copy_nulls;
34use lance_core::{
35 cache::{CacheKey, Context, DeepSizeOf},
36 error::{Error, LanceOptionExt},
37 utils::bit::pad_bytes,
38};
39use log::trace;
40
41use crate::{
42 compression::{
43 BlockDecompressor, CompressionStrategy, DecompressionStrategy, MiniBlockDecompressor,
44 },
45 data::{AllNullDataBlock, DataBlock, VariableWidthBlock},
46 utils::bytepack::BytepackedIntegerEncoder,
47};
48use crate::{
49 compression::{FixedPerValueDecompressor, VariablePerValueDecompressor},
50 encodings::logical::primitive::fullzip::PerValueDataBlock,
51};
52use crate::{
53 encodings::logical::primitive::miniblock::MiniBlockChunk, utils::bytepack::ByteUnpacker,
54};
55use crate::{
56 encodings::logical::primitive::miniblock::MiniBlockCompressed,
57 statistics::{ComputeStat, GetStat, Stat},
58};
59use crate::{
60 repdef::{
61 CompositeRepDefUnraveler, ControlWordIterator, ControlWordParser, DefinitionInterpretation,
62 RepDefSlicer, build_control_word_iterator,
63 },
64 utils::accumulation::AccumulationQueue,
65};
66use lance_core::{Result, datatypes::Field, utils::tokio::spawn_cpu};
67
68use crate::constants::{
69 COMPRESSION_LEVEL_META_KEY, COMPRESSION_META_KEY, DICT_DIVISOR_META_KEY,
70 DICT_SIZE_RATIO_META_KEY, DICT_VALUES_COMPRESSION_ENV_VAR,
71 DICT_VALUES_COMPRESSION_LEVEL_ENV_VAR, DICT_VALUES_COMPRESSION_LEVEL_META_KEY,
72 DICT_VALUES_COMPRESSION_META_KEY,
73};
74use crate::version::LanceFileVersion;
75use crate::{
76 EncodingsIo,
77 buffer::LanceBuffer,
78 data::{BlockInfo, DataBlockBuilder, FixedWidthDataBlock},
79 decoder::{
80 ColumnInfo, DecodePageTask, DecodedArray, DecodedPage, FilterExpression, LoadedPageShard,
81 MessageType, PageEncoding, PageInfo, ScheduledScanLine, SchedulerContext,
82 StructuralDecodeArrayTask, StructuralFieldDecoder, StructuralFieldScheduler,
83 StructuralPageDecoder, StructuralSchedulingJob, UnloadedPageShard,
84 },
85 encoder::{
86 EncodeTask, EncodedColumn, EncodedPage, EncodingOptions, FieldEncoder, OutOfLineBuffers,
87 },
88 repdef::{LevelBuffer, RepDefBuilder, RepDefUnraveler},
89};
90
91pub mod blob;
92pub mod constant;
93pub mod dict;
94pub mod fullzip;
95pub mod miniblock;
96
97const FILL_BYTE: u8 = 0xFE;
98const DEFAULT_DICT_DIVISOR: u64 = 2;
99const DEFAULT_DICT_MAX_CARDINALITY: u64 = 100_000;
100const DEFAULT_DICT_SIZE_RATIO: f64 = 0.8;
101const DEFAULT_DICT_VALUES_COMPRESSION: &str = "lz4";
102
103struct PageLoadTask {
104 decoder_fut: BoxFuture<'static, Result<Box<dyn StructuralPageDecoder>>>,
105 num_rows: u64,
106}
107
108trait StructuralPageScheduler: std::fmt::Debug + Send {
111 fn initialize<'a>(
113 &'a mut self,
114 io: &Arc<dyn EncodingsIo>,
115 ) -> BoxFuture<'a, Result<Arc<dyn CachedPageData>>>;
116 fn load(&mut self, data: &Arc<dyn CachedPageData>);
118 fn schedule_ranges(
127 &self,
128 ranges: &[Range<u64>],
129 io: &Arc<dyn EncodingsIo>,
130 ) -> Result<Vec<PageLoadTask>>;
131}
132
133#[derive(Debug)]
135struct ChunkMeta {
136 num_values: u64,
137 chunk_size_bytes: u64,
138 offset_bytes: u64,
139}
140
141#[derive(Debug, Clone)]
143struct DecodedMiniBlockChunk {
144 rep: Option<ScalarBuffer<u16>>,
145 def: Option<ScalarBuffer<u16>>,
146 values: DataBlock,
147}
148
149#[derive(Debug)]
157struct DecodeMiniBlockTask {
158 rep_decompressor: Option<Arc<dyn BlockDecompressor>>,
159 def_decompressor: Option<Arc<dyn BlockDecompressor>>,
160 value_decompressor: Arc<dyn MiniBlockDecompressor>,
161 dictionary_data: Option<Arc<DataBlock>>,
162 def_meaning: Arc<[DefinitionInterpretation]>,
163 num_buffers: u64,
164 max_visible_level: u16,
165 instructions: Vec<(ChunkDrainInstructions, LoadedChunk)>,
166 has_large_chunk: bool,
167}
168
169impl DecodeMiniBlockTask {
170 fn decode_levels(
171 rep_decompressor: &dyn BlockDecompressor,
172 levels: LanceBuffer,
173 num_levels: u16,
174 ) -> Result<ScalarBuffer<u16>> {
175 let rep = rep_decompressor.decompress(levels, num_levels as u64)?;
176 let rep = rep.as_fixed_width().unwrap();
177 debug_assert_eq!(rep.num_values, num_levels as u64);
178 debug_assert_eq!(rep.bits_per_value, 16);
179 Ok(rep.data.borrow_to_typed_slice::<u16>())
180 }
181
182 fn extend_levels(
189 range: Range<u64>,
190 levels: &mut Option<LevelBuffer>,
191 level_buf: &Option<impl AsRef<[u16]>>,
192 dest_offset: usize,
193 ) {
194 if let Some(level_buf) = level_buf {
195 if levels.is_none() {
196 let mut new_levels_vec =
199 LevelBuffer::with_capacity(dest_offset + (range.end - range.start) as usize);
200 new_levels_vec.extend(iter::repeat_n(0, dest_offset));
201 *levels = Some(new_levels_vec);
202 }
203 levels.as_mut().unwrap().extend(
204 level_buf.as_ref()[range.start as usize..range.end as usize]
205 .iter()
206 .copied(),
207 );
208 } else if let Some(levels) = levels {
209 let num_values = (range.end - range.start) as usize;
210 levels.extend(iter::repeat_n(0, num_values));
213 }
214 }
215
216 fn map_range(
253 range: Range<u64>,
254 rep: Option<&impl AsRef<[u16]>>,
255 def: Option<&impl AsRef<[u16]>>,
256 max_rep: u16,
257 max_visible_def: u16,
258 total_items: u64,
261 preamble_action: PreambleAction,
262 ) -> (Range<u64>, Range<u64>) {
263 if let Some(rep) = rep {
264 let mut rep = rep.as_ref();
265 let mut items_in_preamble = 0_u64;
268 let first_row_start = match preamble_action {
269 PreambleAction::Skip | PreambleAction::Take => {
270 let first_row_start = if let Some(def) = def.as_ref() {
271 let mut first_row_start = None;
272 for (idx, (rep, def)) in rep.iter().zip(def.as_ref()).enumerate() {
273 if *rep == max_rep {
274 first_row_start = Some(idx as u64);
275 break;
276 }
277 if *def <= max_visible_def {
278 items_in_preamble += 1;
279 }
280 }
281 first_row_start
282 } else {
283 let first_row_start =
284 rep.iter().position(|&r| r == max_rep).map(|r| r as u64);
285 items_in_preamble = first_row_start.unwrap_or(rep.len() as u64);
286 first_row_start
287 };
288 if first_row_start.is_none() {
291 assert!(preamble_action == PreambleAction::Take);
292 return (0..total_items, 0..rep.len() as u64);
293 }
294 let first_row_start = first_row_start.unwrap();
295 rep = &rep[first_row_start as usize..];
296 first_row_start
297 }
298 PreambleAction::Absent => {
299 debug_assert!(rep[0] == max_rep);
300 0
301 }
302 };
303
304 if range.start == range.end {
306 debug_assert!(preamble_action == PreambleAction::Take);
307 debug_assert!(items_in_preamble <= total_items);
308 return (0..items_in_preamble, 0..first_row_start);
309 }
310 assert!(range.start < range.end);
311
312 let mut rows_seen = 0;
313 let mut new_start = 0;
314 let mut new_levels_start = 0;
315
316 if let Some(def) = def {
317 let def = &def.as_ref()[first_row_start as usize..];
318
319 let mut lead_invis_seen = 0;
321
322 if range.start > 0 {
323 if def[0] > max_visible_def {
324 lead_invis_seen += 1;
325 }
326 for (idx, (rep, def)) in rep.iter().zip(def).skip(1).enumerate() {
327 if *rep == max_rep {
328 rows_seen += 1;
329 if rows_seen == range.start {
330 new_start = idx as u64 + 1 - lead_invis_seen;
331 new_levels_start = idx as u64 + 1;
332 break;
333 }
334 }
335 if *def > max_visible_def {
336 lead_invis_seen += 1;
337 }
338 }
339 }
340
341 rows_seen += 1;
342
343 let mut new_end = u64::MAX;
344 let mut new_levels_end = rep.len() as u64;
345 let new_start_is_visible = def[new_levels_start as usize] <= max_visible_def;
346 let mut tail_invis_seen = if new_start_is_visible { 0 } else { 1 };
347 for (idx, (rep, def)) in rep[(new_levels_start + 1) as usize..]
348 .iter()
349 .zip(&def[(new_levels_start + 1) as usize..])
350 .enumerate()
351 {
352 if *rep == max_rep {
353 rows_seen += 1;
354 if rows_seen == range.end + 1 {
355 new_end = idx as u64 + new_start + 1 - tail_invis_seen;
356 new_levels_end = idx as u64 + new_levels_start + 1;
357 break;
358 }
359 }
360 if *def > max_visible_def {
361 tail_invis_seen += 1;
362 }
363 }
364
365 if new_end == u64::MAX {
366 new_levels_end = rep.len() as u64;
367 let total_invis_seen = lead_invis_seen + tail_invis_seen;
368 new_end = rep.len() as u64 - total_invis_seen;
369 }
370
371 assert_ne!(new_end, u64::MAX);
372
373 if preamble_action == PreambleAction::Skip {
375 new_start += items_in_preamble;
376 new_end += items_in_preamble;
377 new_levels_start += first_row_start;
378 new_levels_end += first_row_start;
379 } else if preamble_action == PreambleAction::Take {
380 debug_assert_eq!(new_start, 0);
381 debug_assert_eq!(new_levels_start, 0);
382 new_end += items_in_preamble;
383 new_levels_end += first_row_start;
384 }
385
386 debug_assert!(new_end <= total_items);
387 (new_start..new_end, new_levels_start..new_levels_end)
388 } else {
389 if range.start > 0 {
395 for (idx, rep) in rep.iter().skip(1).enumerate() {
396 if *rep == max_rep {
397 rows_seen += 1;
398 if rows_seen == range.start {
399 new_start = idx as u64 + 1;
400 break;
401 }
402 }
403 }
404 }
405 let mut new_end = rep.len() as u64;
406 if range.end < total_items {
408 for (idx, rep) in rep[(new_start + 1) as usize..].iter().enumerate() {
409 if *rep == max_rep {
410 rows_seen += 1;
411 if rows_seen == range.end {
412 new_end = idx as u64 + new_start + 1;
413 break;
414 }
415 }
416 }
417 }
418
419 if preamble_action == PreambleAction::Skip {
421 new_start += first_row_start;
422 new_end += first_row_start;
423 } else if preamble_action == PreambleAction::Take {
424 debug_assert_eq!(new_start, 0);
425 new_end += first_row_start;
426 }
427
428 debug_assert!(new_end <= total_items);
429 (new_start..new_end, new_start..new_end)
430 }
431 } else {
432 (range.clone(), range)
435 }
436 }
437
438 fn read_buffer_sizes<const LARGE: bool>(
440 buf: &[u8],
441 offset: &mut usize,
442 num_buffers: u64,
443 ) -> Vec<u32> {
444 let read_size = if LARGE { 4 } else { 2 };
445 (0..num_buffers)
446 .map(|_| {
447 let bytes = &buf[*offset..*offset + read_size];
448 let size = if LARGE {
449 u32::from_le_bytes([bytes[0], bytes[1], bytes[2], bytes[3]])
450 } else {
451 u16::from_le_bytes([bytes[0], bytes[1]]) as u32
453 };
454 *offset += read_size;
455 size
456 })
457 .collect()
458 }
459
460 fn decode_miniblock_chunk(
462 &self,
463 buf: &LanceBuffer,
464 items_in_chunk: u64,
465 ) -> Result<DecodedMiniBlockChunk> {
466 let mut offset = 0;
467 let num_levels = u16::from_le_bytes([buf[offset], buf[offset + 1]]);
468 offset += 2;
469
470 let rep_size = if self.rep_decompressor.is_some() {
471 let rep_size = u16::from_le_bytes([buf[offset], buf[offset + 1]]);
472 offset += 2;
473 Some(rep_size)
474 } else {
475 None
476 };
477 let def_size = if self.def_decompressor.is_some() {
478 let def_size = u16::from_le_bytes([buf[offset], buf[offset + 1]]);
479 offset += 2;
480 Some(def_size)
481 } else {
482 None
483 };
484
485 let buffer_sizes = if self.has_large_chunk {
486 Self::read_buffer_sizes::<true>(buf, &mut offset, self.num_buffers)
487 } else {
488 Self::read_buffer_sizes::<false>(buf, &mut offset, self.num_buffers)
489 };
490
491 offset += pad_bytes::<MINIBLOCK_ALIGNMENT>(offset);
492
493 let rep = rep_size.map(|rep_size| {
494 let rep = buf.slice_with_length(offset, rep_size as usize);
495 offset += rep_size as usize;
496 offset += pad_bytes::<MINIBLOCK_ALIGNMENT>(offset);
497 rep
498 });
499
500 let def = def_size.map(|def_size| {
501 let def = buf.slice_with_length(offset, def_size as usize);
502 offset += def_size as usize;
503 offset += pad_bytes::<MINIBLOCK_ALIGNMENT>(offset);
504 def
505 });
506
507 let buffers = buffer_sizes
508 .into_iter()
509 .map(|buf_size| {
510 let buf = buf.slice_with_length(offset, buf_size as usize);
511 offset += buf_size as usize;
512 offset += pad_bytes::<MINIBLOCK_ALIGNMENT>(offset);
513 buf
514 })
515 .collect::<Vec<_>>();
516
517 let values = self
518 .value_decompressor
519 .decompress(buffers, items_in_chunk)?;
520
521 let rep = rep
522 .map(|rep| {
523 Self::decode_levels(
524 self.rep_decompressor.as_ref().unwrap().as_ref(),
525 rep,
526 num_levels,
527 )
528 })
529 .transpose()?;
530 let def = def
531 .map(|def| {
532 Self::decode_levels(
533 self.def_decompressor.as_ref().unwrap().as_ref(),
534 def,
535 num_levels,
536 )
537 })
538 .transpose()?;
539
540 Ok(DecodedMiniBlockChunk { rep, def, values })
541 }
542}
543
544impl DecodePageTask for DecodeMiniBlockTask {
545 fn decode(self: Box<Self>) -> Result<DecodedPage> {
546 let mut repbuf: Option<LevelBuffer> = None;
548 let mut defbuf: Option<LevelBuffer> = None;
549
550 let max_rep = self.def_meaning.iter().filter(|l| l.is_list()).count() as u16;
551
552 let estimated_size_bytes = self
554 .instructions
555 .iter()
556 .map(|(_, chunk)| chunk.data.len())
557 .sum::<usize>()
558 * 2;
559 let mut data_builder =
560 DataBlockBuilder::with_capacity_estimate(estimated_size_bytes as u64);
561
562 let mut level_offset = 0;
564
565 let needs_caching: Vec<bool> = self
567 .instructions
568 .windows(2)
569 .map(|w| w[0].1.chunk_idx == w[1].1.chunk_idx)
570 .chain(std::iter::once(false)) .collect();
572
573 let mut chunk_cache: Option<(usize, DecodedMiniBlockChunk)> = None;
575
576 for (idx, (instructions, chunk)) in self.instructions.iter().enumerate() {
578 let should_cache_this_chunk = needs_caching[idx];
579
580 let decoded_chunk = match &chunk_cache {
581 Some((cached_chunk_idx, cached_chunk)) if *cached_chunk_idx == chunk.chunk_idx => {
582 cached_chunk.clone()
584 }
585 _ => {
586 let decoded = self.decode_miniblock_chunk(&chunk.data, chunk.items_in_chunk)?;
588
589 if should_cache_this_chunk {
591 chunk_cache = Some((chunk.chunk_idx, decoded.clone()));
592 }
593 decoded
594 }
595 };
596
597 let DecodedMiniBlockChunk { rep, def, values } = decoded_chunk;
598
599 let row_range_start =
601 instructions.rows_to_skip + instructions.chunk_instructions.rows_to_skip;
602 let row_range_end = row_range_start + instructions.rows_to_take;
603
604 let (item_range, level_range) = Self::map_range(
606 row_range_start..row_range_end,
607 rep.as_ref(),
608 def.as_ref(),
609 max_rep,
610 self.max_visible_level,
611 chunk.items_in_chunk,
612 instructions.preamble_action,
613 );
614 if item_range.end - item_range.start > chunk.items_in_chunk {
615 return Err(lance_core::Error::internal(format!(
616 "Item range {:?} is greater than chunk items in chunk {:?}",
617 item_range, chunk.items_in_chunk
618 )));
619 }
620
621 Self::extend_levels(level_range.clone(), &mut repbuf, &rep, level_offset);
623 Self::extend_levels(level_range.clone(), &mut defbuf, &def, level_offset);
624 level_offset += (level_range.end - level_range.start) as usize;
625 data_builder.append(&values, item_range);
626 }
627
628 let mut data = data_builder.finish();
629
630 let unraveler =
631 RepDefUnraveler::new(repbuf, defbuf, self.def_meaning.clone(), data.num_values());
632
633 if let Some(dictionary) = &self.dictionary_data {
634 let DataBlock::FixedWidth(indices) = data else {
636 return Err(lance_core::Error::internal(format!(
637 "Expected FixedWidth DataBlock for dictionary indices, got {:?}",
638 data
639 )));
640 };
641 data = DataBlock::Dictionary(DictionaryDataBlock::from_parts(
642 indices,
643 dictionary.as_ref().clone(),
644 ));
645 }
646
647 Ok(DecodedPage {
648 data,
649 repdef: unraveler,
650 })
651 }
652}
653
654#[derive(Debug)]
657struct LoadedChunk {
658 data: LanceBuffer,
659 items_in_chunk: u64,
660 byte_range: Range<u64>,
661 chunk_idx: usize,
662}
663
664impl Clone for LoadedChunk {
665 fn clone(&self) -> Self {
666 Self {
667 data: self.data.clone(),
669 items_in_chunk: self.items_in_chunk,
670 byte_range: self.byte_range.clone(),
671 chunk_idx: self.chunk_idx,
672 }
673 }
674}
675
676#[derive(Debug)]
679struct MiniBlockDecoder {
680 rep_decompressor: Option<Arc<dyn BlockDecompressor>>,
681 def_decompressor: Option<Arc<dyn BlockDecompressor>>,
682 value_decompressor: Arc<dyn MiniBlockDecompressor>,
683 def_meaning: Arc<[DefinitionInterpretation]>,
684 loaded_chunks: VecDeque<LoadedChunk>,
685 instructions: VecDeque<ChunkInstructions>,
686 offset_in_current_chunk: u64,
687 num_rows: u64,
688 num_buffers: u64,
689 dictionary: Option<Arc<DataBlock>>,
690 has_large_chunk: bool,
691}
692
693impl StructuralPageDecoder for MiniBlockDecoder {
696 fn drain(&mut self, num_rows: u64) -> Result<Box<dyn DecodePageTask>> {
697 let mut items_desired = num_rows;
698 let mut need_preamble = false;
699 let mut skip_in_chunk = self.offset_in_current_chunk;
700 let mut drain_instructions = Vec::new();
701 while items_desired > 0 || need_preamble {
702 let (instructions, consumed) = self
703 .instructions
704 .front()
705 .unwrap()
706 .drain_from_instruction(&mut items_desired, &mut need_preamble, &mut skip_in_chunk);
707
708 while self.loaded_chunks.front().unwrap().chunk_idx
709 != instructions.chunk_instructions.chunk_idx
710 {
711 self.loaded_chunks.pop_front();
712 }
713 drain_instructions.push((instructions, self.loaded_chunks.front().unwrap().clone()));
714 if consumed {
715 self.instructions.pop_front();
716 }
717 }
718 self.offset_in_current_chunk = skip_in_chunk;
721
722 let max_visible_level = self
723 .def_meaning
724 .iter()
725 .take_while(|l| !l.is_list())
726 .map(|l| l.num_def_levels())
727 .sum::<u16>();
728
729 Ok(Box::new(DecodeMiniBlockTask {
730 instructions: drain_instructions,
731 def_decompressor: self.def_decompressor.clone(),
732 rep_decompressor: self.rep_decompressor.clone(),
733 value_decompressor: self.value_decompressor.clone(),
734 dictionary_data: self.dictionary.clone(),
735 def_meaning: self.def_meaning.clone(),
736 num_buffers: self.num_buffers,
737 max_visible_level,
738 has_large_chunk: self.has_large_chunk,
739 }))
740 }
741
742 fn num_rows(&self) -> u64 {
743 self.num_rows
744 }
745}
746
747#[derive(Debug)]
748struct CachedComplexAllNullState {
749 rep: Option<ScalarBuffer<u16>>,
750 def: Option<ScalarBuffer<u16>>,
751}
752
753impl DeepSizeOf for CachedComplexAllNullState {
754 fn deep_size_of_children(&self, _ctx: &mut Context) -> usize {
755 self.rep.as_ref().map(|buf| buf.len() * 2).unwrap_or(0)
756 + self.def.as_ref().map(|buf| buf.len() * 2).unwrap_or(0)
757 }
758}
759
760impl CachedPageData for CachedComplexAllNullState {
761 fn as_arc_any(self: Arc<Self>) -> Arc<dyn Any + Send + Sync + 'static> {
762 self
763 }
764}
765
766#[derive(Debug)]
775pub struct ComplexAllNullScheduler {
776 buffer_offsets_and_sizes: Arc<[(u64, u64)]>,
778 def_meaning: Arc<[DefinitionInterpretation]>,
779 repdef: Option<Arc<CachedComplexAllNullState>>,
780 max_visible_level: u16,
781 rep_decompressor: Option<Arc<dyn BlockDecompressor>>,
782 def_decompressor: Option<Arc<dyn BlockDecompressor>>,
783 num_rep_values: u64,
784 num_def_values: u64,
785}
786
787impl ComplexAllNullScheduler {
788 pub fn new(
789 buffer_offsets_and_sizes: Arc<[(u64, u64)]>,
790 def_meaning: Arc<[DefinitionInterpretation]>,
791 rep_decompressor: Option<Arc<dyn BlockDecompressor>>,
792 def_decompressor: Option<Arc<dyn BlockDecompressor>>,
793 num_rep_values: u64,
794 num_def_values: u64,
795 ) -> Self {
796 let max_visible_level = def_meaning
797 .iter()
798 .take_while(|l| !l.is_list())
799 .map(|l| l.num_def_levels())
800 .sum::<u16>();
801 Self {
802 buffer_offsets_and_sizes,
803 def_meaning,
804 repdef: None,
805 max_visible_level,
806 rep_decompressor,
807 def_decompressor,
808 num_rep_values,
809 num_def_values,
810 }
811 }
812}
813
814impl StructuralPageScheduler for ComplexAllNullScheduler {
815 fn initialize<'a>(
816 &'a mut self,
817 io: &Arc<dyn EncodingsIo>,
818 ) -> BoxFuture<'a, Result<Arc<dyn CachedPageData>>> {
819 let (rep_pos, rep_size) = self.buffer_offsets_and_sizes[0];
821 let (def_pos, def_size) = self.buffer_offsets_and_sizes[1];
822 let has_rep = rep_size > 0;
823 let has_def = def_size > 0;
824
825 let mut reads = Vec::with_capacity(2);
826 if has_rep {
827 reads.push(rep_pos..rep_pos + rep_size);
828 }
829 if has_def {
830 reads.push(def_pos..def_pos + def_size);
831 }
832
833 let data = io.submit_request(reads, 0);
834 let rep_decompressor = self.rep_decompressor.clone();
835 let def_decompressor = self.def_decompressor.clone();
836 let num_rep_values = self.num_rep_values;
837 let num_def_values = self.num_def_values;
838
839 async move {
840 let data = data.await?;
841 let mut data_iter = data.into_iter();
842
843 let decompress_levels = |compressed_bytes: Bytes,
844 decompressor: &Arc<dyn BlockDecompressor>,
845 num_values: u64,
846 level_type: &str|
847 -> Result<ScalarBuffer<u16>> {
848 let compressed_buffer = LanceBuffer::from_bytes(compressed_bytes, 1);
849 let decompressed = decompressor.decompress(compressed_buffer, num_values)?;
850 match decompressed {
851 DataBlock::FixedWidth(block) => {
852 if block.num_values != num_values {
853 return Err(Error::invalid_input_source(format!(
854 "Unexpected {} level count after decompression: expected {}, got {}",
855 level_type, num_values, block.num_values
856 )
857 .into()));
858 }
859 if block.bits_per_value != 16 {
860 return Err(Error::invalid_input_source(format!(
861 "Unexpected {} level bit width after decompression: expected 16, got {}",
862 level_type, block.bits_per_value
863 )
864 .into()));
865 }
866 Ok(block.data.borrow_to_typed_slice::<u16>())
867 }
868 _ => Err(Error::invalid_input_source(format!(
869 "Expected fixed-width data block for {} levels",
870 level_type
871 )
872 .into())),
873 }
874 };
875
876 let rep = if has_rep {
877 let rep = data_iter.next().unwrap();
878 if let Some(rep_decompressor) = rep_decompressor.as_ref() {
879 Some(decompress_levels(
880 rep,
881 rep_decompressor,
882 num_rep_values,
883 "repetition",
884 )?)
885 } else {
886 let rep = LanceBuffer::from_bytes(rep, 2);
887 let rep = rep.borrow_to_typed_slice::<u16>();
888 Some(rep)
889 }
890 } else {
891 None
892 };
893
894 let def = if has_def {
895 let def = data_iter.next().unwrap();
896 if let Some(def_decompressor) = def_decompressor.as_ref() {
897 Some(decompress_levels(
898 def,
899 def_decompressor,
900 num_def_values,
901 "definition",
902 )?)
903 } else {
904 let def = LanceBuffer::from_bytes(def, 2);
905 let def = def.borrow_to_typed_slice::<u16>();
906 Some(def)
907 }
908 } else {
909 None
910 };
911
912 let repdef = Arc::new(CachedComplexAllNullState { rep, def });
913
914 self.repdef = Some(repdef.clone());
915
916 Ok(repdef as Arc<dyn CachedPageData>)
917 }
918 .boxed()
919 }
920
921 fn load(&mut self, data: &Arc<dyn CachedPageData>) {
922 self.repdef = Some(
923 data.clone()
924 .as_arc_any()
925 .downcast::<CachedComplexAllNullState>()
926 .unwrap(),
927 );
928 }
929
930 fn schedule_ranges(
931 &self,
932 ranges: &[Range<u64>],
933 _io: &Arc<dyn EncodingsIo>,
934 ) -> Result<Vec<PageLoadTask>> {
935 let ranges = VecDeque::from_iter(ranges.iter().cloned());
936 let num_rows = ranges.iter().map(|r| r.end - r.start).sum::<u64>();
937 let decoder = Box::new(ComplexAllNullPageDecoder {
938 ranges,
939 rep: self.repdef.as_ref().unwrap().rep.clone(),
940 def: self.repdef.as_ref().unwrap().def.clone(),
941 num_rows,
942 def_meaning: self.def_meaning.clone(),
943 max_visible_level: self.max_visible_level,
944 }) as Box<dyn StructuralPageDecoder>;
945 let page_load_task = PageLoadTask {
946 decoder_fut: std::future::ready(Ok(decoder)).boxed(),
947 num_rows,
948 };
949 Ok(vec![page_load_task])
950 }
951}
952
953#[derive(Debug)]
954pub struct ComplexAllNullPageDecoder {
955 ranges: VecDeque<Range<u64>>,
956 rep: Option<ScalarBuffer<u16>>,
957 def: Option<ScalarBuffer<u16>>,
958 num_rows: u64,
959 def_meaning: Arc<[DefinitionInterpretation]>,
960 max_visible_level: u16,
961}
962
963impl ComplexAllNullPageDecoder {
964 fn drain_ranges(&mut self, num_rows: u64) -> Vec<Range<u64>> {
965 let mut rows_desired = num_rows;
966 let mut ranges = Vec::with_capacity(self.ranges.len());
967 while rows_desired > 0 {
968 let front = self.ranges.front_mut().unwrap();
969 let avail = front.end - front.start;
970 if avail > rows_desired {
971 ranges.push(front.start..front.start + rows_desired);
972 front.start += rows_desired;
973 rows_desired = 0;
974 } else {
975 ranges.push(self.ranges.pop_front().unwrap());
976 rows_desired -= avail;
977 }
978 }
979 ranges
980 }
981}
982
983impl StructuralPageDecoder for ComplexAllNullPageDecoder {
984 fn drain(&mut self, num_rows: u64) -> Result<Box<dyn DecodePageTask>> {
985 let drained_ranges = self.drain_ranges(num_rows);
986 Ok(Box::new(DecodeComplexAllNullTask {
987 ranges: drained_ranges,
988 rep: self.rep.clone(),
989 def: self.def.clone(),
990 def_meaning: self.def_meaning.clone(),
991 max_visible_level: self.max_visible_level,
992 }))
993 }
994
995 fn num_rows(&self) -> u64 {
996 self.num_rows
997 }
998}
999
1000#[derive(Debug)]
1003pub struct DecodeComplexAllNullTask {
1004 ranges: Vec<Range<u64>>,
1005 rep: Option<ScalarBuffer<u16>>,
1006 def: Option<ScalarBuffer<u16>>,
1007 def_meaning: Arc<[DefinitionInterpretation]>,
1008 max_visible_level: u16,
1009}
1010
1011impl DecodeComplexAllNullTask {
1012 fn decode_level(
1013 &self,
1014 levels: &Option<ScalarBuffer<u16>>,
1015 num_values: u64,
1016 ) -> Option<Vec<u16>> {
1017 levels.as_ref().map(|levels| {
1018 let mut referenced_levels = Vec::with_capacity(num_values as usize);
1019 for range in &self.ranges {
1020 referenced_levels.extend(
1021 levels[range.start as usize..range.end as usize]
1022 .iter()
1023 .copied(),
1024 );
1025 }
1026 referenced_levels
1027 })
1028 }
1029}
1030
1031impl DecodePageTask for DecodeComplexAllNullTask {
1032 fn decode(self: Box<Self>) -> Result<DecodedPage> {
1033 let num_values = self.ranges.iter().map(|r| r.end - r.start).sum::<u64>();
1034 let rep = self.decode_level(&self.rep, num_values);
1035 let def = self.decode_level(&self.def, num_values);
1036
1037 let num_values = if let Some(def) = &def {
1041 def.iter().filter(|&d| *d <= self.max_visible_level).count() as u64
1042 } else {
1043 num_values
1044 };
1045
1046 let data = DataBlock::AllNull(AllNullDataBlock { num_values });
1047 let unraveler = RepDefUnraveler::new(rep, def, self.def_meaning, num_values);
1048 Ok(DecodedPage {
1049 data,
1050 repdef: unraveler,
1051 })
1052 }
1053}
1054
1055#[derive(Debug, Default)]
1060pub struct SimpleAllNullScheduler {}
1061
1062impl StructuralPageScheduler for SimpleAllNullScheduler {
1063 fn initialize<'a>(
1064 &'a mut self,
1065 _io: &Arc<dyn EncodingsIo>,
1066 ) -> BoxFuture<'a, Result<Arc<dyn CachedPageData>>> {
1067 std::future::ready(Ok(Arc::new(NoCachedPageData) as Arc<dyn CachedPageData>)).boxed()
1068 }
1069
1070 fn load(&mut self, _cache: &Arc<dyn CachedPageData>) {}
1071
1072 fn schedule_ranges(
1073 &self,
1074 ranges: &[Range<u64>],
1075 _io: &Arc<dyn EncodingsIo>,
1076 ) -> Result<Vec<PageLoadTask>> {
1077 let num_rows = ranges.iter().map(|r| r.end - r.start).sum::<u64>();
1078 let decoder =
1079 Box::new(SimpleAllNullPageDecoder { num_rows }) as Box<dyn StructuralPageDecoder>;
1080 let page_load_task = PageLoadTask {
1081 decoder_fut: std::future::ready(Ok(decoder)).boxed(),
1082 num_rows,
1083 };
1084 Ok(vec![page_load_task])
1085 }
1086}
1087
1088#[derive(Debug)]
1091struct SimpleAllNullDecodePageTask {
1092 num_values: u64,
1093}
1094impl DecodePageTask for SimpleAllNullDecodePageTask {
1095 fn decode(self: Box<Self>) -> Result<DecodedPage> {
1096 let unraveler = RepDefUnraveler::new(
1097 None,
1098 Some(vec![1; self.num_values as usize]),
1099 Arc::new([DefinitionInterpretation::NullableItem]),
1100 self.num_values,
1101 );
1102 Ok(DecodedPage {
1103 data: DataBlock::AllNull(AllNullDataBlock {
1104 num_values: self.num_values,
1105 }),
1106 repdef: unraveler,
1107 })
1108 }
1109}
1110
1111#[derive(Debug)]
1112pub struct SimpleAllNullPageDecoder {
1113 num_rows: u64,
1114}
1115
1116impl StructuralPageDecoder for SimpleAllNullPageDecoder {
1117 fn drain(&mut self, num_rows: u64) -> Result<Box<dyn DecodePageTask>> {
1118 Ok(Box::new(SimpleAllNullDecodePageTask {
1119 num_values: num_rows,
1120 }))
1121 }
1122
1123 fn num_rows(&self) -> u64 {
1124 self.num_rows
1125 }
1126}
1127
1128#[derive(Debug, Clone)]
1129struct MiniBlockSchedulerDictionary {
1130 dictionary_decompressor: Arc<dyn BlockDecompressor>,
1132 dictionary_buf_position_and_size: (u64, u64),
1133 dictionary_data_alignment: u64,
1134 num_dictionary_items: u64,
1135}
1136
1137#[derive(Debug)]
1139struct MiniBlockRepIndexBlock {
1140 first_row: u64,
1144 starts_including_trailer: u64,
1147 has_preamble: bool,
1149 has_trailer: bool,
1151}
1152
1153impl DeepSizeOf for MiniBlockRepIndexBlock {
1154 fn deep_size_of_children(&self, _context: &mut Context) -> usize {
1155 0
1156 }
1157}
1158
1159#[derive(Debug)]
1164struct MiniBlockRepIndex {
1165 blocks: Vec<MiniBlockRepIndexBlock>,
1166}
1167
1168impl DeepSizeOf for MiniBlockRepIndex {
1169 fn deep_size_of_children(&self, context: &mut Context) -> usize {
1170 self.blocks.deep_size_of_children(context)
1171 }
1172}
1173
1174impl MiniBlockRepIndex {
1175 pub fn default_from_chunks(chunks: &[ChunkMeta]) -> Self {
1180 let mut blocks = Vec::with_capacity(chunks.len());
1181 let mut offset: u64 = 0;
1182
1183 for c in chunks {
1184 blocks.push(MiniBlockRepIndexBlock {
1185 first_row: offset,
1186 starts_including_trailer: c.num_values,
1187 has_preamble: false,
1188 has_trailer: false,
1189 });
1190
1191 offset += c.num_values;
1192 }
1193
1194 Self { blocks }
1195 }
1196
1197 pub fn decode_from_bytes(rep_bytes: &[u8], stride: usize) -> Self {
1203 let buffer = crate::buffer::LanceBuffer::from(rep_bytes.to_vec());
1205 let u64_slice = buffer.borrow_to_typed_slice::<u64>();
1206 let n = u64_slice.len() / stride;
1207
1208 let mut blocks = Vec::with_capacity(n);
1209 let mut chunk_has_preamble = false;
1210 let mut offset: u64 = 0;
1211
1212 for i in 0..n {
1214 let base_idx = i * stride;
1215 let ends = u64_slice[base_idx];
1216 let partial = u64_slice[base_idx + 1];
1217
1218 let has_trailer = partial > 0;
1219 let starts_including_trailer =
1221 ends + (has_trailer as u64) - (chunk_has_preamble as u64);
1222
1223 blocks.push(MiniBlockRepIndexBlock {
1224 first_row: offset,
1225 starts_including_trailer,
1226 has_preamble: chunk_has_preamble,
1227 has_trailer,
1228 });
1229
1230 chunk_has_preamble = has_trailer;
1231 offset += starts_including_trailer;
1232 }
1233
1234 Self { blocks }
1235 }
1236}
1237
1238#[derive(Debug)]
1240struct MiniBlockCacheableState {
1241 chunk_meta: Vec<ChunkMeta>,
1243 rep_index: MiniBlockRepIndex,
1245 dictionary: Option<Arc<DataBlock>>,
1247}
1248
1249impl DeepSizeOf for MiniBlockCacheableState {
1250 fn deep_size_of_children(&self, context: &mut Context) -> usize {
1251 self.rep_index.deep_size_of_children(context)
1252 + self
1253 .dictionary
1254 .as_ref()
1255 .map(|dict| dict.data_size() as usize)
1256 .unwrap_or(0)
1257 }
1258}
1259
1260impl CachedPageData for MiniBlockCacheableState {
1261 fn as_arc_any(self: Arc<Self>) -> Arc<dyn Any + Send + Sync + 'static> {
1262 self
1263 }
1264}
1265
1266#[derive(Debug)]
1293pub struct MiniBlockScheduler {
1294 buffer_offsets_and_sizes: Vec<(u64, u64)>,
1296 priority: u64,
1297 items_in_page: u64,
1298 repetition_index_depth: u16,
1299 num_buffers: u64,
1300 rep_decompressor: Option<Arc<dyn BlockDecompressor>>,
1301 def_decompressor: Option<Arc<dyn BlockDecompressor>>,
1302 value_decompressor: Arc<dyn MiniBlockDecompressor>,
1303 def_meaning: Arc<[DefinitionInterpretation]>,
1304 dictionary: Option<MiniBlockSchedulerDictionary>,
1305 page_meta: Option<Arc<MiniBlockCacheableState>>,
1307 has_large_chunk: bool,
1308}
1309
1310impl MiniBlockScheduler {
1311 fn try_new(
1312 buffer_offsets_and_sizes: &[(u64, u64)],
1313 priority: u64,
1314 items_in_page: u64,
1315 layout: &pb21::MiniBlockLayout,
1316 decompressors: &dyn DecompressionStrategy,
1317 ) -> Result<Self> {
1318 let rep_decompressor = layout
1319 .rep_compression
1320 .as_ref()
1321 .map(|rep_compression| {
1322 decompressors
1323 .create_block_decompressor(rep_compression)
1324 .map(Arc::from)
1325 })
1326 .transpose()?;
1327 let def_decompressor = layout
1328 .def_compression
1329 .as_ref()
1330 .map(|def_compression| {
1331 decompressors
1332 .create_block_decompressor(def_compression)
1333 .map(Arc::from)
1334 })
1335 .transpose()?;
1336 let def_meaning = layout
1337 .layers
1338 .iter()
1339 .map(|l| ProtobufUtils21::repdef_layer_to_def_interp(*l))
1340 .collect::<Vec<_>>();
1341 let value_decompressor = decompressors.create_miniblock_decompressor(
1342 layout.value_compression.as_ref().unwrap(),
1343 decompressors,
1344 )?;
1345
1346 let dictionary = if let Some(dictionary_encoding) = layout.dictionary.as_ref() {
1347 let num_dictionary_items = layout.num_dictionary_items;
1348 let dictionary_decompressor = decompressors
1349 .create_block_decompressor(dictionary_encoding)?
1350 .into();
1351 let dictionary_data_alignment = match dictionary_encoding.compression.as_ref().unwrap()
1352 {
1353 Compression::Variable(_) => 4,
1354 Compression::Flat(_) => 16,
1355 Compression::General(_) => 1,
1356 Compression::InlineBitpacking(_) | Compression::OutOfLineBitpacking(_) => {
1357 crate::encoder::MIN_PAGE_BUFFER_ALIGNMENT
1358 }
1359 _ => {
1360 return Err(Error::invalid_input_source(
1361 format!(
1362 "Unsupported mini-block dictionary encoding: {:?}",
1363 dictionary_encoding.compression.as_ref().unwrap()
1364 )
1365 .into(),
1366 ));
1367 }
1368 };
1369 Some(MiniBlockSchedulerDictionary {
1370 dictionary_decompressor,
1371 dictionary_buf_position_and_size: buffer_offsets_and_sizes[2],
1372 dictionary_data_alignment,
1373 num_dictionary_items,
1374 })
1375 } else {
1376 None
1377 };
1378
1379 Ok(Self {
1380 buffer_offsets_and_sizes: buffer_offsets_and_sizes.to_vec(),
1381 rep_decompressor,
1382 def_decompressor,
1383 value_decompressor: value_decompressor.into(),
1384 repetition_index_depth: layout.repetition_index_depth as u16,
1385 num_buffers: layout.num_buffers,
1386 priority,
1387 items_in_page,
1388 dictionary,
1389 def_meaning: def_meaning.into(),
1390 page_meta: None,
1391 has_large_chunk: layout.has_large_chunk,
1392 })
1393 }
1394
1395 fn lookup_chunks(&self, chunk_indices: &[usize]) -> Vec<LoadedChunk> {
1396 let page_meta = self.page_meta.as_ref().unwrap();
1397 chunk_indices
1398 .iter()
1399 .map(|&chunk_idx| {
1400 let chunk_meta = &page_meta.chunk_meta[chunk_idx];
1401 let bytes_start = chunk_meta.offset_bytes;
1402 let bytes_end = bytes_start + chunk_meta.chunk_size_bytes;
1403 LoadedChunk {
1404 byte_range: bytes_start..bytes_end,
1405 items_in_chunk: chunk_meta.num_values,
1406 chunk_idx,
1407 data: LanceBuffer::empty(),
1408 }
1409 })
1410 .collect()
1411 }
1412}
1413
1414#[derive(Debug, PartialEq, Eq, Clone, Copy)]
1415enum PreambleAction {
1416 Take,
1417 Skip,
1418 Absent,
1419}
1420
1421#[derive(Clone, Debug, PartialEq, Eq)]
1456struct ChunkInstructions {
1457 chunk_idx: usize,
1459 preamble: PreambleAction,
1465 rows_to_skip: u64,
1469 rows_to_take: u64,
1472 take_trailer: bool,
1479}
1480
1481#[derive(Debug, PartialEq, Eq)]
1499struct ChunkDrainInstructions {
1500 chunk_instructions: ChunkInstructions,
1501 rows_to_skip: u64,
1502 rows_to_take: u64,
1503 preamble_action: PreambleAction,
1504}
1505
1506impl ChunkInstructions {
1507 fn schedule_instructions(
1513 rep_index: &MiniBlockRepIndex,
1514 user_ranges: &[Range<u64>],
1515 ) -> Vec<Self> {
1516 let mut chunk_instructions = Vec::with_capacity(user_ranges.len());
1520
1521 for user_range in user_ranges {
1522 let mut rows_needed = user_range.end - user_range.start;
1523 let mut need_preamble = false;
1524
1525 let mut block_index = match rep_index
1528 .blocks
1529 .binary_search_by_key(&user_range.start, |block| block.first_row)
1530 {
1531 Ok(idx) => {
1532 let mut idx = idx;
1535 while idx > 0 && rep_index.blocks[idx - 1].first_row == user_range.start {
1536 idx -= 1;
1537 }
1538 idx
1539 }
1540 Err(idx) => idx - 1,
1542 };
1543
1544 let mut to_skip = user_range.start - rep_index.blocks[block_index].first_row;
1545
1546 while rows_needed > 0 || need_preamble {
1547 if block_index >= rep_index.blocks.len() {
1549 log::warn!(
1550 "schedule_instructions inconsistency: block_index >= rep_index.blocks.len(), exiting early"
1551 );
1552 break;
1553 }
1554
1555 let chunk = &rep_index.blocks[block_index];
1556 let rows_avail = chunk.starts_including_trailer.saturating_sub(to_skip);
1557
1558 if rows_avail == 0 && to_skip == 0 {
1562 if chunk.has_preamble && need_preamble {
1564 chunk_instructions.push(Self {
1565 chunk_idx: block_index,
1566 preamble: PreambleAction::Take,
1567 rows_to_skip: 0,
1568 rows_to_take: 0,
1569 take_trailer: chunk.has_trailer,
1573 });
1574 if chunk.starts_including_trailer > 0
1578 || block_index == rep_index.blocks.len() - 1
1579 {
1580 need_preamble = false;
1581 }
1582 }
1583 block_index += 1;
1585 continue;
1586 }
1587
1588 if rows_avail == 0 && to_skip > 0 {
1592 to_skip -= chunk.starts_including_trailer;
1595 block_index += 1;
1596 continue;
1597 }
1598
1599 let rows_to_take = rows_avail.min(rows_needed);
1600 rows_needed -= rows_to_take;
1601
1602 let mut take_trailer = false;
1603 let preamble = if chunk.has_preamble {
1604 if need_preamble {
1605 PreambleAction::Take
1606 } else {
1607 PreambleAction::Skip
1608 }
1609 } else {
1610 PreambleAction::Absent
1611 };
1612
1613 if rows_to_take == rows_avail && chunk.has_trailer {
1615 take_trailer = true;
1616 need_preamble = true;
1617 } else {
1618 need_preamble = false;
1619 };
1620
1621 chunk_instructions.push(Self {
1622 preamble,
1623 chunk_idx: block_index,
1624 rows_to_skip: to_skip,
1625 rows_to_take,
1626 take_trailer,
1627 });
1628
1629 to_skip = 0;
1630 block_index += 1;
1631 }
1632 }
1633
1634 if user_ranges.len() > 1 {
1638 let mut merged_instructions = Vec::with_capacity(chunk_instructions.len());
1640 let mut instructions_iter = chunk_instructions.into_iter();
1641 merged_instructions.push(instructions_iter.next().unwrap());
1642 for instruction in instructions_iter {
1643 let last = merged_instructions.last_mut().unwrap();
1644 if last.chunk_idx == instruction.chunk_idx
1645 && last.rows_to_take + last.rows_to_skip == instruction.rows_to_skip
1646 {
1647 last.rows_to_take += instruction.rows_to_take;
1648 last.take_trailer |= instruction.take_trailer;
1649 } else {
1650 merged_instructions.push(instruction);
1651 }
1652 }
1653 merged_instructions
1654 } else {
1655 chunk_instructions
1656 }
1657 }
1658
1659 fn drain_from_instruction(
1660 &self,
1661 rows_desired: &mut u64,
1662 need_preamble: &mut bool,
1663 skip_in_chunk: &mut u64,
1664 ) -> (ChunkDrainInstructions, bool) {
1665 debug_assert!(!*need_preamble || *skip_in_chunk == 0);
1667 let rows_avail = self.rows_to_take - *skip_in_chunk;
1668 let has_preamble = self.preamble != PreambleAction::Absent;
1669 let preamble_action = match (*need_preamble, has_preamble) {
1670 (true, true) => PreambleAction::Take,
1671 (true, false) => panic!("Need preamble but there isn't one"),
1672 (false, true) => PreambleAction::Skip,
1673 (false, false) => PreambleAction::Absent,
1674 };
1675
1676 let rows_taking = if *rows_desired >= rows_avail {
1679 *need_preamble = self.take_trailer;
1687 rows_avail
1688 } else {
1689 *need_preamble = false;
1692 *rows_desired
1693 };
1694 let rows_skipped = *skip_in_chunk;
1695
1696 let consumed_chunk = if *rows_desired >= rows_avail {
1698 *rows_desired -= rows_avail;
1699 *skip_in_chunk = 0;
1700 true
1701 } else {
1702 *skip_in_chunk += *rows_desired;
1703 *rows_desired = 0;
1704 false
1705 };
1706
1707 (
1708 ChunkDrainInstructions {
1709 chunk_instructions: self.clone(),
1710 rows_to_skip: rows_skipped,
1711 rows_to_take: rows_taking,
1712 preamble_action,
1713 },
1714 consumed_chunk,
1715 )
1716 }
1717}
1718
1719enum Words {
1720 U16(ScalarBuffer<u16>),
1721 U32(ScalarBuffer<u32>),
1722}
1723
1724struct WordsIter<'a> {
1725 iter: Box<dyn Iterator<Item = u32> + 'a>,
1726}
1727
1728impl Words {
1729 pub fn len(&self) -> usize {
1730 match self {
1731 Self::U16(b) => b.len(),
1732 Self::U32(b) => b.len(),
1733 }
1734 }
1735
1736 pub fn iter(&self) -> WordsIter<'_> {
1737 match self {
1738 Self::U16(buf) => WordsIter {
1739 iter: Box::new(buf.iter().map(|&x| x as u32)),
1740 },
1741 Self::U32(buf) => WordsIter {
1742 iter: Box::new(buf.iter().copied()),
1743 },
1744 }
1745 }
1746
1747 pub fn from_bytes(bytes: Bytes, has_large_chunk: bool) -> Result<Self> {
1748 let bytes_per_value = if has_large_chunk { 4 } else { 2 };
1749 assert_eq!(bytes.len() % bytes_per_value, 0);
1750 let buffer = LanceBuffer::from_bytes(bytes, bytes_per_value as u64);
1751 if has_large_chunk {
1752 Ok(Self::U32(buffer.borrow_to_typed_slice::<u32>()))
1753 } else {
1754 Ok(Self::U16(buffer.borrow_to_typed_slice::<u16>()))
1755 }
1756 }
1757}
1758
1759impl<'a> Iterator for WordsIter<'a> {
1760 type Item = u32;
1761
1762 fn next(&mut self) -> Option<Self::Item> {
1763 self.iter.next()
1764 }
1765}
1766
1767impl StructuralPageScheduler for MiniBlockScheduler {
1768 fn initialize<'a>(
1769 &'a mut self,
1770 io: &Arc<dyn EncodingsIo>,
1771 ) -> BoxFuture<'a, Result<Arc<dyn CachedPageData>>> {
1772 let (meta_buf_position, meta_buf_size) = self.buffer_offsets_and_sizes[0];
1776 let value_buf_position = self.buffer_offsets_and_sizes[1].0;
1777 let mut bufs_needed = 1;
1778 if self.dictionary.is_some() {
1779 bufs_needed += 1;
1780 }
1781 if self.repetition_index_depth > 0 {
1782 bufs_needed += 1;
1783 }
1784 let mut required_ranges = Vec::with_capacity(bufs_needed);
1785 required_ranges.push(meta_buf_position..meta_buf_position + meta_buf_size);
1786 if let Some(ref dictionary) = self.dictionary {
1787 required_ranges.push(
1788 dictionary.dictionary_buf_position_and_size.0
1789 ..dictionary.dictionary_buf_position_and_size.0
1790 + dictionary.dictionary_buf_position_and_size.1,
1791 );
1792 }
1793 if self.repetition_index_depth > 0 {
1794 let (rep_index_pos, rep_index_size) = self.buffer_offsets_and_sizes.last().unwrap();
1795 required_ranges.push(*rep_index_pos..*rep_index_pos + *rep_index_size);
1796 }
1797 let io_req = io.submit_request(required_ranges, 0);
1798
1799 async move {
1800 let mut buffers = io_req.await?.into_iter().fuse();
1801 let meta_bytes = buffers.next().unwrap();
1802 let dictionary_bytes = self.dictionary.as_ref().and_then(|_| buffers.next());
1803 let rep_index_bytes = buffers.next();
1804
1805 let words = Words::from_bytes(meta_bytes, self.has_large_chunk)?;
1807 let mut chunk_meta = Vec::with_capacity(words.len());
1808
1809 let mut rows_counter = 0;
1810 let mut offset_bytes = value_buf_position;
1811 for (word_idx, word) in words.iter().enumerate() {
1812 let log_num_values = word & 0x0F;
1813 let divided_bytes = word >> 4;
1814 let num_bytes = (divided_bytes as usize + 1) * MINIBLOCK_ALIGNMENT;
1815 debug_assert!(num_bytes > 0);
1816 let num_values = if word_idx < words.len() - 1 {
1817 debug_assert!(log_num_values > 0);
1818 1 << log_num_values
1819 } else {
1820 debug_assert!(
1821 log_num_values == 0
1822 || (1 << log_num_values) == (self.items_in_page - rows_counter)
1823 );
1824 self.items_in_page - rows_counter
1825 };
1826 rows_counter += num_values;
1827
1828 chunk_meta.push(ChunkMeta {
1829 num_values,
1830 chunk_size_bytes: num_bytes as u64,
1831 offset_bytes,
1832 });
1833 offset_bytes += num_bytes as u64;
1834 }
1835
1836 let rep_index = if let Some(rep_index_data) = rep_index_bytes {
1838 assert!(rep_index_data.len() % 8 == 0);
1839 let stride = self.repetition_index_depth as usize + 1;
1840 MiniBlockRepIndex::decode_from_bytes(&rep_index_data, stride)
1841 } else {
1842 MiniBlockRepIndex::default_from_chunks(&chunk_meta)
1843 };
1844
1845 let mut page_meta = MiniBlockCacheableState {
1846 chunk_meta,
1847 rep_index,
1848 dictionary: None,
1849 };
1850
1851 if let Some(ref mut dictionary) = self.dictionary {
1853 let dictionary_data = dictionary_bytes.unwrap();
1854 page_meta.dictionary =
1855 Some(Arc::new(dictionary.dictionary_decompressor.decompress(
1856 LanceBuffer::from_bytes(
1857 dictionary_data,
1858 dictionary.dictionary_data_alignment,
1859 ),
1860 dictionary.num_dictionary_items,
1861 )?));
1862 };
1863 let page_meta = Arc::new(page_meta);
1864 self.page_meta = Some(page_meta.clone());
1865 Ok(page_meta as Arc<dyn CachedPageData>)
1866 }
1867 .boxed()
1868 }
1869
1870 fn load(&mut self, data: &Arc<dyn CachedPageData>) {
1871 self.page_meta = Some(
1872 data.clone()
1873 .as_arc_any()
1874 .downcast::<MiniBlockCacheableState>()
1875 .unwrap(),
1876 );
1877 }
1878
1879 fn schedule_ranges(
1880 &self,
1881 ranges: &[Range<u64>],
1882 io: &Arc<dyn EncodingsIo>,
1883 ) -> Result<Vec<PageLoadTask>> {
1884 let num_rows = ranges.iter().map(|r| r.end - r.start).sum();
1885
1886 let page_meta = self.page_meta.as_ref().unwrap();
1887
1888 let chunk_instructions =
1889 ChunkInstructions::schedule_instructions(&page_meta.rep_index, ranges);
1890
1891 debug_assert_eq!(
1892 num_rows,
1893 chunk_instructions
1894 .iter()
1895 .map(|ci| ci.rows_to_take)
1896 .sum::<u64>()
1897 );
1898
1899 let chunks_needed = chunk_instructions
1900 .iter()
1901 .map(|ci| ci.chunk_idx)
1902 .unique()
1903 .collect::<Vec<_>>();
1904
1905 let mut loaded_chunks = self.lookup_chunks(&chunks_needed);
1906 let chunk_ranges = loaded_chunks
1907 .iter()
1908 .map(|c| c.byte_range.clone())
1909 .collect::<Vec<_>>();
1910 let loaded_chunk_data = io.submit_request(chunk_ranges, self.priority);
1911
1912 let rep_decompressor = self.rep_decompressor.clone();
1913 let def_decompressor = self.def_decompressor.clone();
1914 let value_decompressor = self.value_decompressor.clone();
1915 let num_buffers = self.num_buffers;
1916 let has_large_chunk = self.has_large_chunk;
1917 let dictionary = page_meta
1918 .dictionary
1919 .as_ref()
1920 .map(|dictionary| dictionary.clone());
1921 let def_meaning = self.def_meaning.clone();
1922
1923 let res = async move {
1924 let loaded_chunk_data = loaded_chunk_data.await?;
1925 for (loaded_chunk, chunk_data) in loaded_chunks.iter_mut().zip(loaded_chunk_data) {
1926 loaded_chunk.data = LanceBuffer::from_bytes(chunk_data, 1);
1927 }
1928
1929 Ok(Box::new(MiniBlockDecoder {
1930 rep_decompressor,
1931 def_decompressor,
1932 value_decompressor,
1933 def_meaning,
1934 loaded_chunks: VecDeque::from_iter(loaded_chunks),
1935 instructions: VecDeque::from(chunk_instructions),
1936 offset_in_current_chunk: 0,
1937 dictionary,
1938 num_rows,
1939 num_buffers,
1940 has_large_chunk,
1941 }) as Box<dyn StructuralPageDecoder>)
1942 }
1943 .boxed();
1944 let page_load_task = PageLoadTask {
1945 decoder_fut: res,
1946 num_rows,
1947 };
1948 Ok(vec![page_load_task])
1949 }
1950}
1951
1952#[derive(Debug, Clone, Copy)]
1953struct FullZipRepIndexDetails {
1954 buf_position: u64,
1955 bytes_per_value: u64, }
1957
1958#[derive(Debug)]
1959enum PerValueDecompressor {
1960 Fixed(Arc<dyn FixedPerValueDecompressor>),
1961 Variable(Arc<dyn VariablePerValueDecompressor>),
1962}
1963
1964#[derive(Debug)]
1965struct FullZipDecodeDetails {
1966 value_decompressor: PerValueDecompressor,
1967 def_meaning: Arc<[DefinitionInterpretation]>,
1968 ctrl_word_parser: ControlWordParser,
1969 max_rep: u16,
1970 max_visible_def: u16,
1971}
1972
1973#[derive(Debug, Clone)]
1985enum FullZipReadSource {
1986 Remote(Arc<dyn EncodingsIo>),
1988 PrefetchedPage { base_offset: u64, data: LanceBuffer },
1990}
1991
1992impl FullZipReadSource {
1993 fn fetch(
1997 &self,
1998 ranges: &[Range<u64>],
1999 priority: u64,
2000 ) -> BoxFuture<'static, Result<VecDeque<LanceBuffer>>> {
2001 match self {
2002 Self::Remote(io) => {
2003 let io = io.clone();
2004 let ranges = ranges.to_vec();
2005 async move {
2006 let data = io.submit_request(ranges, priority).await?;
2007 Ok(data
2008 .into_iter()
2009 .map(|bytes| LanceBuffer::from_bytes(bytes, 1))
2010 .collect::<VecDeque<_>>())
2011 }
2012 .boxed()
2013 }
2014 Self::PrefetchedPage { base_offset, data } => {
2015 let base_offset = *base_offset;
2016 let data = data.clone();
2017 let page_end = base_offset + data.len() as u64;
2018 std::future::ready(
2019 ranges
2020 .iter()
2021 .map(|range| {
2022 if range.start > range.end
2023 || range.start < base_offset
2024 || range.end > page_end
2025 {
2026 return Err(Error::internal(format!(
2027 "Requested range {:?} is outside page range {}..{}",
2028 range, base_offset, page_end
2029 )));
2030 }
2031 let start = (range.start - base_offset) as usize;
2032 let len = (range.end - range.start) as usize;
2033 Ok(data.slice_with_length(start, len))
2034 })
2035 .collect::<Result<VecDeque<_>>>(),
2036 )
2037 .boxed()
2038 }
2039 }
2040 }
2041}
2042
2043#[derive(Debug)]
2051pub struct FullZipScheduler {
2052 data_buf_position: u64,
2053 data_buf_size: u64,
2054 rep_index: Option<FullZipRepIndexDetails>,
2055 priority: u64,
2056 rows_in_page: u64,
2057 bits_per_offset: u8,
2058 details: Arc<FullZipDecodeDetails>,
2059 cached_state: Option<Arc<FullZipCacheableState>>,
2061 enable_cache: bool,
2063}
2064
2065impl FullZipScheduler {
2066 fn try_new(
2067 buffer_offsets_and_sizes: &[(u64, u64)],
2068 priority: u64,
2069 rows_in_page: u64,
2070 layout: &pb21::FullZipLayout,
2071 decompressors: &dyn DecompressionStrategy,
2072 ) -> Result<Self> {
2073 let (data_buf_position, data_buf_size) = buffer_offsets_and_sizes[0];
2074 let rep_index = buffer_offsets_and_sizes.get(1).map(|(pos, len)| {
2075 let num_reps = rows_in_page + 1;
2076 let bytes_per_rep = len / num_reps;
2077 debug_assert_eq!(len % num_reps, 0);
2078 debug_assert!(
2079 bytes_per_rep == 1
2080 || bytes_per_rep == 2
2081 || bytes_per_rep == 4
2082 || bytes_per_rep == 8
2083 );
2084 FullZipRepIndexDetails {
2085 buf_position: *pos,
2086 bytes_per_value: bytes_per_rep,
2087 }
2088 });
2089
2090 let value_decompressor = match layout.details {
2091 Some(pb21::full_zip_layout::Details::BitsPerValue(_)) => {
2092 let decompressor = decompressors.create_fixed_per_value_decompressor(
2093 layout.value_compression.as_ref().unwrap(),
2094 )?;
2095 PerValueDecompressor::Fixed(decompressor.into())
2096 }
2097 Some(pb21::full_zip_layout::Details::BitsPerOffset(_)) => {
2098 let decompressor = decompressors.create_variable_per_value_decompressor(
2099 layout.value_compression.as_ref().unwrap(),
2100 )?;
2101 PerValueDecompressor::Variable(decompressor.into())
2102 }
2103 None => {
2104 panic!("Full-zip layout must have a `details` field");
2105 }
2106 };
2107 let ctrl_word_parser = ControlWordParser::new(
2108 layout.bits_rep.try_into().unwrap(),
2109 layout.bits_def.try_into().unwrap(),
2110 );
2111 let def_meaning = layout
2112 .layers
2113 .iter()
2114 .map(|l| ProtobufUtils21::repdef_layer_to_def_interp(*l))
2115 .collect::<Vec<_>>();
2116
2117 let max_rep = def_meaning.iter().filter(|d| d.is_list()).count() as u16;
2118 let max_visible_def = def_meaning
2119 .iter()
2120 .filter(|d| !d.is_list())
2121 .map(|d| d.num_def_levels())
2122 .sum();
2123
2124 let bits_per_offset = match layout.details {
2125 Some(pb21::full_zip_layout::Details::BitsPerValue(_)) => 32,
2126 Some(pb21::full_zip_layout::Details::BitsPerOffset(bits_per_offset)) => {
2127 bits_per_offset as u8
2128 }
2129 None => panic!("Full-zip layout must have a `details` field"),
2130 };
2131
2132 let details = Arc::new(FullZipDecodeDetails {
2133 value_decompressor,
2134 def_meaning: def_meaning.into(),
2135 ctrl_word_parser,
2136 max_rep,
2137 max_visible_def,
2138 });
2139 Ok(Self {
2140 data_buf_position,
2141 data_buf_size,
2142 rep_index,
2143 details,
2144 priority,
2145 rows_in_page,
2146 bits_per_offset,
2147 cached_state: None,
2148 enable_cache: false,
2149 })
2150 }
2151
2152 fn covers_entire_page(ranges: &[Range<u64>], rows_in_page: u64) -> bool {
2153 if ranges.is_empty() {
2154 return false;
2155 }
2156 let mut expected_start = 0;
2157 for range in ranges {
2158 if range.start != expected_start || range.end > rows_in_page || range.end < range.start
2159 {
2160 return false;
2161 }
2162 expected_start = range.end;
2163 }
2164 expected_start == rows_in_page
2165 }
2166
2167 fn create_page_load_task(
2168 io_future: BoxFuture<'static, Result<Vec<Bytes>>>,
2169 num_rows: u64,
2170 details: Arc<FullZipDecodeDetails>,
2171 bits_per_offset: u8,
2172 ) -> PageLoadTask {
2173 let load_task = async move {
2174 let buffers = io_future.await?;
2175 let data = buffers
2176 .into_iter()
2177 .map(|bytes| LanceBuffer::from_bytes(bytes, 1))
2178 .collect::<VecDeque<_>>();
2179 Self::create_decoder(details, data, num_rows, bits_per_offset)
2180 }
2181 .boxed();
2182 PageLoadTask {
2183 decoder_fut: load_task,
2184 num_rows,
2185 }
2186 }
2187
2188 fn create_decoder(
2190 details: Arc<FullZipDecodeDetails>,
2191 data: VecDeque<LanceBuffer>,
2192 num_rows: u64,
2193 bits_per_offset: u8,
2194 ) -> Result<Box<dyn StructuralPageDecoder>> {
2195 match &details.value_decompressor {
2196 PerValueDecompressor::Fixed(decompressor) => {
2197 let bits_per_value = decompressor.bits_per_value();
2198 if bits_per_value % 8 != 0 {
2199 return Err(lance_core::Error::not_supported_source("Bit-packed full-zip encoding (non-byte-aligned values) is not yet implemented".into()));
2200 }
2201 let bytes_per_value = bits_per_value / 8;
2202 let total_bytes_per_value =
2203 bytes_per_value as usize + details.ctrl_word_parser.bytes_per_word();
2204 if total_bytes_per_value == 0 {
2205 return Err(lance_core::Error::internal(
2206 "Invalid encoding: per-row byte width must be greater than 0",
2207 ));
2208 }
2209 Ok(Box::new(FixedFullZipDecoder {
2210 details,
2211 data,
2212 num_rows,
2213 offset_in_current: 0,
2214 bytes_per_value: bytes_per_value as usize,
2215 total_bytes_per_value,
2216 }) as Box<dyn StructuralPageDecoder>)
2217 }
2218 PerValueDecompressor::Variable(_decompressor) => {
2219 Ok(Box::new(VariableFullZipDecoder::new(
2220 details,
2221 data,
2222 num_rows,
2223 bits_per_offset,
2224 bits_per_offset,
2225 )?))
2226 }
2227 }
2228 }
2229
2230 fn extract_byte_ranges_from_pairs(
2233 buffer: LanceBuffer,
2234 bytes_per_value: u64,
2235 data_buf_position: u64,
2236 ) -> Vec<Range<u64>> {
2237 ByteUnpacker::new(buffer, bytes_per_value as usize)
2238 .chunks(2)
2239 .into_iter()
2240 .map(|mut c| {
2241 let start = c.next().unwrap() + data_buf_position;
2242 let end = c.next().unwrap() + data_buf_position;
2243 start..end
2244 })
2245 .collect::<Vec<_>>()
2246 }
2247
2248 fn extract_byte_ranges_from_cached(
2251 buffer: &LanceBuffer,
2252 ranges: &[Range<u64>],
2253 bytes_per_value: u64,
2254 data_buf_position: u64,
2255 ) -> Vec<Range<u64>> {
2256 ranges
2257 .iter()
2258 .map(|r| {
2259 let start_offset = (r.start * bytes_per_value) as usize;
2260 let end_offset = (r.end * bytes_per_value) as usize;
2261
2262 let start_slice = &buffer[start_offset..start_offset + bytes_per_value as usize];
2263 let start_val =
2264 ByteUnpacker::new(start_slice.iter().copied(), bytes_per_value as usize)
2265 .next()
2266 .unwrap();
2267
2268 let end_slice = &buffer[end_offset..end_offset + bytes_per_value as usize];
2269 let end_val =
2270 ByteUnpacker::new(end_slice.iter().copied(), bytes_per_value as usize)
2271 .next()
2272 .unwrap();
2273
2274 (data_buf_position + start_val)..(data_buf_position + end_val)
2275 })
2276 .collect()
2277 }
2278
2279 fn compute_rep_index_ranges(
2281 ranges: &[Range<u64>],
2282 rep_index: &FullZipRepIndexDetails,
2283 ) -> Vec<Range<u64>> {
2284 ranges
2285 .iter()
2286 .flat_map(|r| {
2287 let first_val_start =
2288 rep_index.buf_position + (r.start * rep_index.bytes_per_value);
2289 let first_val_end = first_val_start + rep_index.bytes_per_value;
2290 let last_val_start = rep_index.buf_position + (r.end * rep_index.bytes_per_value);
2291 let last_val_end = last_val_start + rep_index.bytes_per_value;
2292 [first_val_start..first_val_end, last_val_start..last_val_end]
2293 })
2294 .collect()
2295 }
2296
2297 fn schedule_ranges_rep(
2299 &self,
2300 ranges: &[Range<u64>],
2301 io: &Arc<dyn EncodingsIo>,
2302 rep_index: FullZipRepIndexDetails,
2303 ) -> Result<Vec<PageLoadTask>> {
2304 let num_rows = ranges.iter().map(|r| r.end - r.start).sum();
2305 let data_buf_position = self.data_buf_position;
2306 let priority = self.priority;
2307 let details = self.details.clone();
2308 let bits_per_offset = self.bits_per_offset;
2309
2310 if Self::covers_entire_page(ranges, self.rows_in_page) {
2311 let full_range = self.data_buf_position..(self.data_buf_position + self.data_buf_size);
2312 let page_data = io.submit_single(full_range.clone(), priority);
2313 let load_task = async move {
2314 let page_data = page_data.await?;
2315 let source = FullZipReadSource::PrefetchedPage {
2316 base_offset: full_range.start,
2317 data: LanceBuffer::from_bytes(page_data, 1),
2318 };
2319 let read_ranges = vec![full_range];
2320 let data = source.fetch(&read_ranges, priority).await?;
2321 Self::create_decoder(details, data, num_rows, bits_per_offset)
2322 }
2323 .boxed();
2324 let page_load_task = PageLoadTask {
2325 decoder_fut: load_task,
2326 num_rows,
2327 };
2328 return Ok(vec![page_load_task]);
2329 }
2330
2331 if let Some(cached_state) = &self.cached_state {
2332 let byte_ranges = Self::extract_byte_ranges_from_cached(
2333 &cached_state.rep_index_buffer,
2334 ranges,
2335 rep_index.bytes_per_value,
2336 data_buf_position,
2337 );
2338 let io_future = io.submit_request(byte_ranges, priority);
2339 let page_load_task =
2340 Self::create_page_load_task(io_future, num_rows, details, bits_per_offset);
2341 return Ok(vec![page_load_task]);
2342 }
2343
2344 let rep_ranges = Self::compute_rep_index_ranges(ranges, &rep_index);
2345 let rep_data = io.submit_request(rep_ranges, priority);
2346 let io_clone = io.clone();
2347 let load_task = async move {
2348 let rep_data = rep_data.await?;
2349 let rep_buffer = LanceBuffer::concat(
2350 &rep_data
2351 .into_iter()
2352 .map(|d| LanceBuffer::from_bytes(d, 1))
2353 .collect::<Vec<_>>(),
2354 );
2355 let byte_ranges = Self::extract_byte_ranges_from_pairs(
2356 rep_buffer,
2357 rep_index.bytes_per_value,
2358 data_buf_position,
2359 );
2360 let source = FullZipReadSource::Remote(io_clone);
2361 let data = source.fetch(&byte_ranges, priority).await?;
2362 Self::create_decoder(details, data, num_rows, bits_per_offset)
2363 }
2364 .boxed();
2365 let page_load_task = PageLoadTask {
2366 decoder_fut: load_task,
2367 num_rows,
2368 };
2369 Ok(vec![page_load_task])
2370 }
2371
2372 fn schedule_ranges_simple(
2376 &self,
2377 ranges: &[Range<u64>],
2378 io: &Arc<dyn EncodingsIo>,
2379 ) -> Result<Vec<PageLoadTask>> {
2380 let num_rows = ranges.iter().map(|r| r.end - r.start).sum();
2382
2383 let PerValueDecompressor::Fixed(decompressor) = &self.details.value_decompressor else {
2384 unreachable!()
2385 };
2386
2387 let bits_per_value = decompressor.bits_per_value();
2389 assert_eq!(bits_per_value % 8, 0);
2390 let bytes_per_value = bits_per_value / 8;
2391 let bytes_per_cw = self.details.ctrl_word_parser.bytes_per_word();
2392 let total_bytes_per_value = bytes_per_value + bytes_per_cw as u64;
2393 let byte_ranges = ranges
2394 .iter()
2395 .map(|r| {
2396 debug_assert!(r.end <= self.rows_in_page);
2397 let start = self.data_buf_position + r.start * total_bytes_per_value;
2398 let end = self.data_buf_position + r.end * total_bytes_per_value;
2399 start..end
2400 })
2401 .collect::<Vec<_>>();
2402
2403 let io_future = io.submit_request(byte_ranges, self.priority);
2404 let page_load_task = Self::create_page_load_task(
2405 io_future,
2406 num_rows,
2407 self.details.clone(),
2408 self.bits_per_offset,
2409 );
2410 Ok(vec![page_load_task])
2411 }
2412}
2413
2414#[derive(Debug)]
2416struct FullZipCacheableState {
2417 rep_index_buffer: LanceBuffer,
2419}
2420
2421impl DeepSizeOf for FullZipCacheableState {
2422 fn deep_size_of_children(&self, _context: &mut Context) -> usize {
2423 self.rep_index_buffer.len()
2424 }
2425}
2426
2427impl CachedPageData for FullZipCacheableState {
2428 fn as_arc_any(self: Arc<Self>) -> Arc<dyn Any + Send + Sync + 'static> {
2429 self
2430 }
2431}
2432
2433impl StructuralPageScheduler for FullZipScheduler {
2434 fn initialize<'a>(
2435 &'a mut self,
2436 io: &Arc<dyn EncodingsIo>,
2437 ) -> BoxFuture<'a, Result<Arc<dyn CachedPageData>>> {
2438 if self.enable_cache
2439 && let Some(rep_index) = self.rep_index
2440 {
2441 let total_size = (self.rows_in_page + 1) * rep_index.bytes_per_value;
2442 let rep_index_range = rep_index.buf_position..(rep_index.buf_position + total_size);
2443 let io_clone = io.clone();
2444 return async move {
2445 let rep_index_data = io_clone.submit_request(vec![rep_index_range], 0).await?;
2446 let state = Arc::new(FullZipCacheableState {
2447 rep_index_buffer: LanceBuffer::from_bytes(rep_index_data[0].clone(), 1),
2448 });
2449 self.cached_state = Some(state.clone());
2450 Ok(state as Arc<dyn CachedPageData>)
2451 }
2452 .boxed();
2453 }
2454 std::future::ready(Ok(Arc::new(NoCachedPageData) as Arc<dyn CachedPageData>)).boxed()
2455 }
2456
2457 fn load(&mut self, cache: &Arc<dyn CachedPageData>) {
2461 if let Ok(cached_state) = cache
2463 .clone()
2464 .as_arc_any()
2465 .downcast::<FullZipCacheableState>()
2466 {
2467 self.cached_state = Some(cached_state);
2469 }
2470 }
2471
2472 fn schedule_ranges(
2473 &self,
2474 ranges: &[Range<u64>],
2475 io: &Arc<dyn EncodingsIo>,
2476 ) -> Result<Vec<PageLoadTask>> {
2477 if let Some(rep_index) = self.rep_index {
2478 self.schedule_ranges_rep(ranges, io, rep_index)
2479 } else {
2480 self.schedule_ranges_simple(ranges, io)
2481 }
2482 }
2483}
2484
2485#[derive(Debug)]
2493struct FixedFullZipDecoder {
2494 details: Arc<FullZipDecodeDetails>,
2495 data: VecDeque<LanceBuffer>,
2496 offset_in_current: usize,
2497 bytes_per_value: usize,
2498 total_bytes_per_value: usize,
2499 num_rows: u64,
2500}
2501
2502impl FixedFullZipDecoder {
2503 fn slice_next_task(&mut self, num_rows: u64) -> FullZipDecodeTaskItem {
2504 debug_assert!(num_rows > 0);
2505 let cur_buf = self.data.front_mut().unwrap();
2506 let start = self.offset_in_current;
2507 if self.details.ctrl_word_parser.has_rep() {
2508 let mut rows_started = 0;
2511 let mut num_items = 0;
2514 while self.offset_in_current < cur_buf.len() {
2515 let control = self.details.ctrl_word_parser.parse_desc(
2516 &cur_buf[self.offset_in_current..],
2517 self.details.max_rep,
2518 self.details.max_visible_def,
2519 );
2520 if control.is_new_row {
2521 if rows_started == num_rows {
2522 break;
2523 }
2524 rows_started += 1;
2525 }
2526 num_items += 1;
2527 if control.is_visible {
2528 self.offset_in_current += self.total_bytes_per_value;
2529 } else {
2530 self.offset_in_current += self.details.ctrl_word_parser.bytes_per_word();
2531 }
2532 }
2533
2534 let task_slice = cur_buf.slice_with_length(start, self.offset_in_current - start);
2535 if self.offset_in_current == cur_buf.len() {
2536 self.data.pop_front();
2537 self.offset_in_current = 0;
2538 }
2539
2540 FullZipDecodeTaskItem {
2541 data: PerValueDataBlock::Fixed(FixedWidthDataBlock {
2542 data: task_slice,
2543 bits_per_value: self.bytes_per_value as u64 * 8,
2544 num_values: num_items,
2545 block_info: BlockInfo::new(),
2546 }),
2547 rows_in_buf: rows_started,
2548 }
2549 } else {
2550 let cur_buf = self.data.front_mut().unwrap();
2553 let bytes_avail = cur_buf.len() - self.offset_in_current;
2554 let offset_in_cur = self.offset_in_current;
2555
2556 let bytes_needed = num_rows as usize * self.total_bytes_per_value;
2557 let mut rows_taken = num_rows;
2558 let task_slice = if bytes_needed >= bytes_avail {
2559 self.offset_in_current = 0;
2560 rows_taken = bytes_avail as u64 / self.total_bytes_per_value as u64;
2561 self.data
2562 .pop_front()
2563 .unwrap()
2564 .slice_with_length(offset_in_cur, bytes_avail)
2565 } else {
2566 self.offset_in_current += bytes_needed;
2567 cur_buf.slice_with_length(offset_in_cur, bytes_needed)
2568 };
2569 FullZipDecodeTaskItem {
2570 data: PerValueDataBlock::Fixed(FixedWidthDataBlock {
2571 data: task_slice,
2572 bits_per_value: self.bytes_per_value as u64 * 8,
2573 num_values: rows_taken,
2574 block_info: BlockInfo::new(),
2575 }),
2576 rows_in_buf: rows_taken,
2577 }
2578 }
2579 }
2580}
2581
2582impl StructuralPageDecoder for FixedFullZipDecoder {
2583 fn drain(&mut self, num_rows: u64) -> Result<Box<dyn DecodePageTask>> {
2584 let mut task_data = Vec::with_capacity(self.data.len());
2585 let mut remaining = num_rows;
2586 while remaining > 0 {
2587 let task_item = self.slice_next_task(remaining);
2588 remaining -= task_item.rows_in_buf;
2589 task_data.push(task_item);
2590 }
2591 Ok(Box::new(FixedFullZipDecodeTask {
2592 details: self.details.clone(),
2593 data: task_data,
2594 bytes_per_value: self.bytes_per_value,
2595 num_rows: num_rows as usize,
2596 }))
2597 }
2598
2599 fn num_rows(&self) -> u64 {
2600 self.num_rows
2601 }
2602}
2603
2604#[derive(Debug)]
2609struct VariableFullZipDecoder {
2610 details: Arc<FullZipDecodeDetails>,
2611 decompressor: Arc<dyn VariablePerValueDecompressor>,
2612 data: LanceBuffer,
2613 offsets: LanceBuffer,
2614 rep: ScalarBuffer<u16>,
2615 def: ScalarBuffer<u16>,
2616 repdef_starts: Vec<usize>,
2617 data_starts: Vec<usize>,
2618 offset_starts: Vec<usize>,
2619 visible_item_counts: Vec<u64>,
2620 bits_per_offset: u8,
2621 current_idx: usize,
2622 num_rows: u64,
2623}
2624
2625fn corrupt_file_named(name: &str, message: impl Into<String>) -> Error {
2626 Error::corrupt_file(name.into(), message)
2627}
2628
2629impl VariableFullZipDecoder {
2630 fn new(
2631 details: Arc<FullZipDecodeDetails>,
2632 data: VecDeque<LanceBuffer>,
2633 num_rows: u64,
2634 in_bits_per_length: u8,
2635 out_bits_per_offset: u8,
2636 ) -> Result<Self> {
2637 let decompressor = match details.value_decompressor {
2638 PerValueDecompressor::Variable(ref d) => d.clone(),
2639 _ => unreachable!(),
2640 };
2641
2642 assert_eq!(in_bits_per_length % 8, 0);
2643 assert!(out_bits_per_offset == 32 || out_bits_per_offset == 64);
2644
2645 let mut decoder = Self {
2646 details,
2647 decompressor,
2648 data: LanceBuffer::empty(),
2649 offsets: LanceBuffer::empty(),
2650 rep: LanceBuffer::empty().borrow_to_typed_slice(),
2651 def: LanceBuffer::empty().borrow_to_typed_slice(),
2652 bits_per_offset: out_bits_per_offset,
2653 repdef_starts: Vec::with_capacity(num_rows as usize + 1),
2654 data_starts: Vec::with_capacity(num_rows as usize + 1),
2655 offset_starts: Vec::with_capacity(num_rows as usize + 1),
2656 visible_item_counts: Vec::with_capacity(num_rows as usize + 1),
2657 current_idx: 0,
2658 num_rows,
2659 };
2660
2661 decoder.unzip(data, in_bits_per_length, out_bits_per_offset, num_rows)?;
2682
2683 Ok(decoder)
2684 }
2685
2686 fn slice_batch_data_and_rebase_offsets_typed<T>(
2687 data: &LanceBuffer,
2688 offsets: &LanceBuffer,
2689 ) -> Result<(LanceBuffer, LanceBuffer)>
2690 where
2691 T: arrow_buffer::ArrowNativeType
2692 + Copy
2693 + PartialOrd
2694 + std::ops::Sub<Output = T>
2695 + std::fmt::Display
2696 + TryInto<usize>,
2697 {
2698 let offsets_slice = offsets.borrow_to_typed_slice::<T>();
2699 let offsets_slice = offsets_slice.as_ref();
2700 if offsets_slice.is_empty() {
2701 return Err(Error::internal(
2702 "Variable offsets cannot be empty".to_string(),
2703 ));
2704 }
2705
2706 let base = offsets_slice[0];
2707 let end = *offsets_slice.last().unwrap();
2708 if end < base {
2709 return Err(Error::internal(format!(
2710 "Invalid variable offsets: end ({end}) is less than base ({base})"
2711 )));
2712 }
2713
2714 let data_start = base.try_into().map_err(|_| {
2715 Error::internal(format!("Variable offset ({base}) does not fit into usize"))
2716 })?;
2717 let data_end = end.try_into().map_err(|_| {
2718 Error::internal(format!("Variable offset ({end}) does not fit into usize"))
2719 })?;
2720 if data_end > data.len() {
2721 return Err(Error::internal(format!(
2722 "Invalid variable offsets: end ({data_end}) exceeds data len ({})",
2723 data.len()
2724 )));
2725 }
2726
2727 let mut rebased_offsets = Vec::with_capacity(offsets_slice.len());
2728 for &offset in offsets_slice {
2729 if offset < base {
2730 return Err(Error::internal(format!(
2731 "Invalid variable offsets: offset ({offset}) is less than base ({base})"
2732 )));
2733 }
2734 rebased_offsets.push(offset - base);
2735 }
2736
2737 let sliced_data = data.slice_with_length(data_start, data_end - data_start);
2738 let sliced_data = LanceBuffer::copy_slice(&sliced_data);
2740 let rebased_offsets = LanceBuffer::reinterpret_vec(rebased_offsets);
2741 Ok((sliced_data, rebased_offsets))
2742 }
2743
2744 fn slice_batch_data_and_rebase_offsets(
2745 data: &LanceBuffer,
2746 offsets: &LanceBuffer,
2747 bits_per_offset: u8,
2748 ) -> Result<(LanceBuffer, LanceBuffer)> {
2749 match bits_per_offset {
2750 32 => Self::slice_batch_data_and_rebase_offsets_typed::<u32>(data, offsets),
2751 64 => Self::slice_batch_data_and_rebase_offsets_typed::<u64>(data, offsets),
2752 _ => Err(Error::internal(format!(
2753 "Unsupported bits_per_offset={bits_per_offset}"
2754 ))),
2755 }
2756 }
2757
2758 fn parse_length(data: &[u8], bits_per_offset: u8) -> Result<u64> {
2765 let width = bits_per_offset as usize / 8;
2766 if data.len() < width {
2767 return Err(corrupt_file_named(
2768 "variable_full_zip",
2769 format!(
2770 "truncated length prefix: {} byte(s) remain in the page buffer but a \
2771 {}-bit length prefix requires {}",
2772 data.len(),
2773 bits_per_offset,
2774 width
2775 ),
2776 ));
2777 }
2778 Ok(match bits_per_offset {
2779 8 => data[0] as u64,
2780 16 => u16::from_le_bytes(data[..2].try_into().unwrap()) as u64,
2781 32 => u32::from_le_bytes(data[..4].try_into().unwrap()) as u64,
2782 64 => u64::from_le_bytes(data[..8].try_into().unwrap()),
2783 _ => unreachable!(),
2784 })
2785 }
2786
2787 fn unzip(
2788 &mut self,
2789 data: VecDeque<LanceBuffer>,
2790 in_bits_per_length: u8,
2791 out_bits_per_offset: u8,
2792 num_rows: u64,
2793 ) -> Result<()> {
2794 let mut rep = Vec::with_capacity(num_rows as usize);
2796 let mut def = Vec::with_capacity(num_rows as usize);
2797 let bytes_cw = self.details.ctrl_word_parser.bytes_per_word() * num_rows as usize;
2798
2799 let bytes_per_offset = out_bits_per_offset as usize / 8;
2802 let bytes_offsets = bytes_per_offset * (num_rows as usize + 1);
2803 let mut offsets_data = Vec::with_capacity(bytes_offsets);
2804
2805 let bytes_per_length = in_bits_per_length as usize / 8;
2806 let bytes_lengths = bytes_per_length * num_rows as usize;
2807
2808 let bytes_data = data.iter().map(|d| d.len()).sum::<usize>();
2809 let mut unzipped_data =
2812 Vec::with_capacity((bytes_data - bytes_cw).saturating_sub(bytes_lengths));
2813
2814 let mut current_offset = 0_u64;
2815 let mut visible_item_count = 0_u64;
2816 for databuf in data.into_iter() {
2817 let mut databuf = databuf.as_ref();
2818 while !databuf.is_empty() {
2819 let data_start = unzipped_data.len();
2820 let offset_start = offsets_data.len();
2821 let repdef_start = rep.len().max(def.len());
2824 let ctrl_desc = self.details.ctrl_word_parser.parse_desc(
2826 databuf,
2827 self.details.max_rep,
2828 self.details.max_visible_def,
2829 );
2830 self.details
2831 .ctrl_word_parser
2832 .parse(databuf, &mut rep, &mut def);
2833 databuf = &databuf[self.details.ctrl_word_parser.bytes_per_word()..];
2834
2835 if ctrl_desc.is_new_row {
2836 self.repdef_starts.push(repdef_start);
2837 self.data_starts.push(data_start);
2838 self.offset_starts.push(offset_start);
2839 self.visible_item_counts.push(visible_item_count);
2840 }
2841 if ctrl_desc.is_visible {
2842 visible_item_count += 1;
2843 if ctrl_desc.is_valid_item {
2844 let length = Self::parse_length(databuf, in_bits_per_length)?;
2845 match out_bits_per_offset {
2846 32 => offsets_data
2847 .extend_from_slice(&(current_offset as u32).to_le_bytes()),
2848 64 => offsets_data.extend_from_slice(¤t_offset.to_le_bytes()),
2849 _ => unreachable!(),
2850 };
2851 databuf = &databuf[bytes_per_offset..];
2852 unzipped_data.extend_from_slice(&databuf[..length as usize]);
2853 databuf = &databuf[length as usize..];
2854 current_offset += length;
2855 } else {
2856 match out_bits_per_offset {
2858 32 => offsets_data
2859 .extend_from_slice(&(current_offset as u32).to_le_bytes()),
2860 64 => offsets_data.extend_from_slice(¤t_offset.to_le_bytes()),
2861 _ => unreachable!(),
2862 }
2863 }
2864 }
2865 }
2866 }
2867 self.repdef_starts.push(rep.len().max(def.len()));
2868 self.data_starts.push(unzipped_data.len());
2869 self.offset_starts.push(offsets_data.len());
2870 self.visible_item_counts.push(visible_item_count);
2871 match out_bits_per_offset {
2872 32 => offsets_data.extend_from_slice(&(current_offset as u32).to_le_bytes()),
2873 64 => offsets_data.extend_from_slice(¤t_offset.to_le_bytes()),
2874 _ => unreachable!(),
2875 };
2876 self.rep = ScalarBuffer::from(rep);
2877 self.def = ScalarBuffer::from(def);
2878 self.data = LanceBuffer::from(unzipped_data);
2879 self.offsets = LanceBuffer::from(offsets_data);
2880 Ok(())
2881 }
2882}
2883
2884impl StructuralPageDecoder for VariableFullZipDecoder {
2885 fn drain(&mut self, num_rows: u64) -> Result<Box<dyn DecodePageTask>> {
2886 let start = self.current_idx;
2887 let end = start + num_rows as usize;
2888
2889 let offset_start = self.offset_starts[start];
2890 let offset_end = self.offset_starts[end] + (self.bits_per_offset as usize / 8);
2891 let offsets = self
2892 .offsets
2893 .slice_with_length(offset_start, offset_end - offset_start);
2894 let (data, offsets) =
2896 Self::slice_batch_data_and_rebase_offsets(&self.data, &offsets, self.bits_per_offset)?;
2897
2898 let repdef_start = self.repdef_starts[start];
2899 let repdef_end = self.repdef_starts[end];
2900 let rep = if self.rep.is_empty() {
2901 self.rep.clone()
2902 } else {
2903 self.rep.slice(repdef_start, repdef_end - repdef_start)
2904 };
2905 let def = if self.def.is_empty() {
2906 self.def.clone()
2907 } else {
2908 self.def.slice(repdef_start, repdef_end - repdef_start)
2909 };
2910
2911 let visible_item_counts_start = self.visible_item_counts[start];
2912 let visible_item_counts_end = self.visible_item_counts[end];
2913 let num_visible_items = visible_item_counts_end - visible_item_counts_start;
2914
2915 self.current_idx += num_rows as usize;
2916
2917 Ok(Box::new(VariableFullZipDecodeTask {
2918 details: self.details.clone(),
2919 decompressor: self.decompressor.clone(),
2920 data,
2921 offsets,
2922 bits_per_offset: self.bits_per_offset,
2923 num_visible_items,
2924 rep,
2925 def,
2926 }))
2927 }
2928
2929 fn num_rows(&self) -> u64 {
2930 self.num_rows
2931 }
2932}
2933
2934#[derive(Debug)]
2935struct VariableFullZipDecodeTask {
2936 details: Arc<FullZipDecodeDetails>,
2937 decompressor: Arc<dyn VariablePerValueDecompressor>,
2938 data: LanceBuffer,
2939 offsets: LanceBuffer,
2940 bits_per_offset: u8,
2941 num_visible_items: u64,
2942 rep: ScalarBuffer<u16>,
2943 def: ScalarBuffer<u16>,
2944}
2945
2946impl DecodePageTask for VariableFullZipDecodeTask {
2947 fn decode(self: Box<Self>) -> Result<DecodedPage> {
2948 let block = VariableWidthBlock {
2949 data: self.data,
2950 offsets: self.offsets,
2951 bits_per_offset: self.bits_per_offset,
2952 num_values: self.num_visible_items,
2953 block_info: BlockInfo::new(),
2954 };
2955 let decomopressed = self.decompressor.decompress(block)?;
2956 let rep = if self.rep.is_empty() {
2957 None
2958 } else {
2959 Some(self.rep.to_vec())
2960 };
2961 let def = if self.def.is_empty() {
2962 None
2963 } else {
2964 Some(self.def.to_vec())
2965 };
2966 let unraveler = RepDefUnraveler::new(
2967 rep,
2968 def,
2969 self.details.def_meaning.clone(),
2970 self.num_visible_items,
2971 );
2972 Ok(DecodedPage {
2973 data: decomopressed,
2974 repdef: unraveler,
2975 })
2976 }
2977}
2978
2979#[derive(Debug)]
2980struct FullZipDecodeTaskItem {
2981 data: PerValueDataBlock,
2982 rows_in_buf: u64,
2983}
2984
2985#[derive(Debug)]
2988struct FixedFullZipDecodeTask {
2989 details: Arc<FullZipDecodeDetails>,
2990 data: Vec<FullZipDecodeTaskItem>,
2991 num_rows: usize,
2992 bytes_per_value: usize,
2993}
2994
2995impl DecodePageTask for FixedFullZipDecodeTask {
2996 fn decode(self: Box<Self>) -> Result<DecodedPage> {
2997 let estimated_size_bytes = self
2999 .data
3000 .iter()
3001 .map(|task_item| task_item.data.data_size() as usize)
3002 .sum::<usize>()
3003 * 2;
3004 let mut data_builder =
3005 DataBlockBuilder::with_capacity_estimate(estimated_size_bytes as u64);
3006
3007 if self.details.ctrl_word_parser.bytes_per_word() == 0 {
3008 for task_item in self.data.into_iter() {
3012 let PerValueDataBlock::Fixed(fixed_data) = task_item.data else {
3013 unreachable!()
3014 };
3015 let PerValueDecompressor::Fixed(decompressor) = &self.details.value_decompressor
3016 else {
3017 unreachable!()
3018 };
3019 debug_assert_eq!(fixed_data.num_values, task_item.rows_in_buf);
3020 let decompressed = decompressor.decompress(fixed_data, task_item.rows_in_buf)?;
3021 data_builder.append(&decompressed, 0..task_item.rows_in_buf);
3022 }
3023
3024 let unraveler = RepDefUnraveler::new(
3025 None,
3026 None,
3027 self.details.def_meaning.clone(),
3028 self.num_rows as u64,
3029 );
3030
3031 Ok(DecodedPage {
3032 data: data_builder.finish(),
3033 repdef: unraveler,
3034 })
3035 } else {
3036 let mut rep = Vec::with_capacity(self.num_rows);
3038 let mut def = Vec::with_capacity(self.num_rows);
3039
3040 for task_item in self.data.into_iter() {
3041 let PerValueDataBlock::Fixed(fixed_data) = task_item.data else {
3042 unreachable!()
3043 };
3044 let mut buf_slice = fixed_data.data.as_ref();
3045 let num_values = fixed_data.num_values as usize;
3046 let mut values = Vec::with_capacity(
3049 fixed_data.data.len()
3050 - (self.details.ctrl_word_parser.bytes_per_word() * num_values),
3051 );
3052 let mut visible_items = 0;
3053 for _ in 0..num_values {
3054 self.details
3056 .ctrl_word_parser
3057 .parse(buf_slice, &mut rep, &mut def);
3058 buf_slice = &buf_slice[self.details.ctrl_word_parser.bytes_per_word()..];
3059
3060 let is_visible = def
3061 .last()
3062 .map(|d| *d <= self.details.max_visible_def)
3063 .unwrap_or(true);
3064 if is_visible {
3065 values.extend_from_slice(buf_slice[..self.bytes_per_value].as_ref());
3067 buf_slice = &buf_slice[self.bytes_per_value..];
3068 visible_items += 1;
3069 }
3070 }
3071
3072 let values_buf = LanceBuffer::from(values);
3074 let fixed_data = FixedWidthDataBlock {
3075 bits_per_value: self.bytes_per_value as u64 * 8,
3076 block_info: BlockInfo::new(),
3077 data: values_buf,
3078 num_values: visible_items,
3079 };
3080 let PerValueDecompressor::Fixed(decompressor) = &self.details.value_decompressor
3081 else {
3082 unreachable!()
3083 };
3084 let decompressed = decompressor.decompress(fixed_data, visible_items)?;
3085 data_builder.append(&decompressed, 0..visible_items);
3086 }
3087
3088 let repetition = if rep.is_empty() { None } else { Some(rep) };
3089 let definition = if def.is_empty() { None } else { Some(def) };
3090
3091 let unraveler = RepDefUnraveler::new(
3092 repetition,
3093 definition,
3094 self.details.def_meaning.clone(),
3095 self.num_rows as u64,
3096 );
3097 let data = data_builder.finish();
3098
3099 Ok(DecodedPage {
3100 data,
3101 repdef: unraveler,
3102 })
3103 }
3104 }
3105}
3106
3107#[derive(Debug)]
3108struct StructuralPrimitiveFieldSchedulingJob<'a> {
3109 scheduler: &'a StructuralPrimitiveFieldScheduler,
3110 ranges: Vec<Range<u64>>,
3111 page_idx: usize,
3112 range_idx: usize,
3113 global_row_offset: u64,
3114}
3115
3116impl<'a> StructuralPrimitiveFieldSchedulingJob<'a> {
3117 pub fn new(scheduler: &'a StructuralPrimitiveFieldScheduler, ranges: Vec<Range<u64>>) -> Self {
3118 Self {
3119 scheduler,
3120 ranges,
3121 page_idx: 0,
3122 range_idx: 0,
3123 global_row_offset: 0,
3124 }
3125 }
3126}
3127
3128impl StructuralSchedulingJob for StructuralPrimitiveFieldSchedulingJob<'_> {
3129 fn schedule_next(&mut self, context: &mut SchedulerContext) -> Result<Vec<ScheduledScanLine>> {
3130 if self.range_idx >= self.ranges.len() {
3131 return Ok(Vec::new());
3132 }
3133 let mut range = self.ranges[self.range_idx].clone();
3135 let priority = range.start;
3136
3137 let mut cur_page = &self.scheduler.page_schedulers[self.page_idx];
3138 trace!(
3139 "Current range is {:?} and current page has {} rows",
3140 range, cur_page.num_rows
3141 );
3142 while cur_page.num_rows + self.global_row_offset <= range.start {
3144 self.global_row_offset += cur_page.num_rows;
3145 self.page_idx += 1;
3146 trace!("Skipping entire page of {} rows", cur_page.num_rows);
3147 cur_page = &self.scheduler.page_schedulers[self.page_idx];
3148 }
3149
3150 let mut ranges_in_page = Vec::new();
3154 while cur_page.num_rows + self.global_row_offset > range.start {
3155 range.start = range.start.max(self.global_row_offset);
3156 let start_in_page = range.start - self.global_row_offset;
3157 let end_in_page = start_in_page + (range.end - range.start);
3158 let end_in_page = end_in_page.min(cur_page.num_rows);
3159 let last_in_range = (end_in_page + self.global_row_offset) >= range.end;
3160
3161 ranges_in_page.push(start_in_page..end_in_page);
3162 if last_in_range {
3163 self.range_idx += 1;
3164 if self.range_idx == self.ranges.len() {
3165 break;
3166 }
3167 range = self.ranges[self.range_idx].clone();
3168 } else {
3169 break;
3170 }
3171 }
3172
3173 trace!(
3174 "Scheduling {} rows across {} ranges from page with {} rows (priority={}, column_index={}, page_index={})",
3175 ranges_in_page.iter().map(|r| r.end - r.start).sum::<u64>(),
3176 ranges_in_page.len(),
3177 cur_page.num_rows,
3178 priority,
3179 self.scheduler.column_index,
3180 cur_page.page_index,
3181 );
3182
3183 self.global_row_offset += cur_page.num_rows;
3184 self.page_idx += 1;
3185
3186 let page_decoders = cur_page
3187 .scheduler
3188 .schedule_ranges(&ranges_in_page, context.io())?;
3189
3190 let cur_path = context.current_path();
3191 page_decoders
3192 .into_iter()
3193 .map(|page_load_task| {
3194 let cur_path = cur_path.clone();
3195 let page_decoder = page_load_task.decoder_fut;
3196 let unloaded_page = async move {
3197 let page_decoder = page_decoder.await?;
3198 Ok(LoadedPageShard {
3199 decoder: page_decoder,
3200 path: cur_path,
3201 })
3202 }
3203 .boxed();
3204 Ok(ScheduledScanLine {
3205 decoders: vec![MessageType::UnloadedPage(UnloadedPageShard(unloaded_page))],
3206 rows_scheduled: page_load_task.num_rows,
3207 })
3208 })
3209 .collect::<Result<Vec<_>>>()
3210 }
3211}
3212
3213#[derive(Debug)]
3214struct PageInfoAndScheduler {
3215 page_index: usize,
3216 num_rows: u64,
3217 scheduler: Box<dyn StructuralPageScheduler>,
3218}
3219
3220#[derive(Debug)]
3225pub struct StructuralPrimitiveFieldScheduler {
3226 page_schedulers: Vec<PageInfoAndScheduler>,
3227 column_index: u32,
3228}
3229
3230impl StructuralPrimitiveFieldScheduler {
3231 pub fn try_new(
3232 column_info: &ColumnInfo,
3233 decompressors: &dyn DecompressionStrategy,
3234 cache_repetition_index: bool,
3235 target_field: &Field,
3236 ) -> Result<Self> {
3237 let page_schedulers = column_info
3238 .page_infos
3239 .iter()
3240 .enumerate()
3241 .map(|(page_index, page_info)| {
3242 Self::page_info_to_scheduler(
3243 page_info,
3244 page_index,
3245 decompressors,
3246 cache_repetition_index,
3247 target_field,
3248 )
3249 })
3250 .collect::<Result<Vec<_>>>()?;
3251 Ok(Self {
3252 page_schedulers,
3253 column_index: column_info.index,
3254 })
3255 }
3256
3257 fn page_layout_to_scheduler(
3258 page_info: &PageInfo,
3259 page_layout: &PageLayout,
3260 decompressors: &dyn DecompressionStrategy,
3261 cache_repetition_index: bool,
3262 target_field: &Field,
3263 ) -> Result<Box<dyn StructuralPageScheduler>> {
3264 use pb21::page_layout::Layout;
3265 Ok(match page_layout.layout.as_ref().expect_ok()? {
3266 Layout::MiniBlockLayout(mini_block) => Box::new(MiniBlockScheduler::try_new(
3267 &page_info.buffer_offsets_and_sizes,
3268 page_info.priority,
3269 mini_block.num_items,
3270 mini_block,
3271 decompressors,
3272 )?),
3273 Layout::FullZipLayout(full_zip) => {
3274 let mut scheduler = FullZipScheduler::try_new(
3275 &page_info.buffer_offsets_and_sizes,
3276 page_info.priority,
3277 page_info.num_rows,
3278 full_zip,
3279 decompressors,
3280 )?;
3281 scheduler.enable_cache = cache_repetition_index;
3282 Box::new(scheduler)
3283 }
3284 Layout::ConstantLayout(constant_layout) => {
3285 let def_meaning = constant_layout
3286 .layers
3287 .iter()
3288 .map(|l| ProtobufUtils21::repdef_layer_to_def_interp(*l))
3289 .collect::<Vec<_>>();
3290 let has_scalar_value = constant_layout.inline_value.is_some()
3291 || page_info.buffer_offsets_and_sizes.len() == 1
3292 || page_info.buffer_offsets_and_sizes.len() == 3;
3293 if has_scalar_value {
3294 Box::new(constant::ConstantPageScheduler::try_new(
3295 page_info.buffer_offsets_and_sizes.clone(),
3296 constant_layout.inline_value.clone(),
3297 target_field.data_type(),
3298 def_meaning.into(),
3299 )?) as Box<dyn StructuralPageScheduler>
3300 } else if def_meaning.len() == 1
3301 && def_meaning[0] == DefinitionInterpretation::NullableItem
3302 {
3303 Box::new(SimpleAllNullScheduler::default()) as Box<dyn StructuralPageScheduler>
3304 } else {
3305 let rep_decompressor = constant_layout
3306 .rep_compression
3307 .as_ref()
3308 .map(|encoding| decompressors.create_block_decompressor(encoding))
3309 .transpose()?
3310 .map(Arc::from);
3311
3312 let def_decompressor = constant_layout
3313 .def_compression
3314 .as_ref()
3315 .map(|encoding| decompressors.create_block_decompressor(encoding))
3316 .transpose()?
3317 .map(Arc::from);
3318
3319 Box::new(ComplexAllNullScheduler::new(
3320 page_info.buffer_offsets_and_sizes.clone(),
3321 def_meaning.into(),
3322 rep_decompressor,
3323 def_decompressor,
3324 constant_layout.num_rep_values,
3325 constant_layout.num_def_values,
3326 )) as Box<dyn StructuralPageScheduler>
3327 }
3328 }
3329 Layout::BlobLayout(blob) => {
3330 let inner_scheduler = Self::page_layout_to_scheduler(
3331 page_info,
3332 blob.inner_layout.as_ref().expect_ok()?.as_ref(),
3333 decompressors,
3334 cache_repetition_index,
3335 target_field,
3336 )?;
3337 let def_meaning = blob
3338 .layers
3339 .iter()
3340 .map(|l| ProtobufUtils21::repdef_layer_to_def_interp(*l))
3341 .collect::<Vec<_>>();
3342 if matches!(target_field.data_type(), DataType::Struct(_)) {
3343 Box::new(BlobDescriptionPageScheduler::new(
3345 inner_scheduler,
3346 def_meaning.into(),
3347 ))
3348 } else {
3349 Box::new(BlobPageScheduler::new(
3351 inner_scheduler,
3352 page_info.priority,
3353 page_info.num_rows,
3354 def_meaning.into(),
3355 ))
3356 }
3357 }
3358 })
3359 }
3360
3361 fn page_info_to_scheduler(
3362 page_info: &PageInfo,
3363 page_index: usize,
3364 decompressors: &dyn DecompressionStrategy,
3365 cache_repetition_index: bool,
3366 target_field: &Field,
3367 ) -> Result<PageInfoAndScheduler> {
3368 let page_layout = page_info.encoding.as_structural();
3369 let scheduler = Self::page_layout_to_scheduler(
3370 page_info,
3371 page_layout,
3372 decompressors,
3373 cache_repetition_index,
3374 target_field,
3375 )?;
3376 Ok(PageInfoAndScheduler {
3377 page_index,
3378 num_rows: page_info.num_rows,
3379 scheduler,
3380 })
3381 }
3382}
3383
3384pub trait CachedPageData: Any + Send + Sync + DeepSizeOf + 'static {
3385 fn as_arc_any(self: Arc<Self>) -> Arc<dyn Any + Send + Sync + 'static>;
3386}
3387
3388pub struct NoCachedPageData;
3389
3390impl DeepSizeOf for NoCachedPageData {
3391 fn deep_size_of_children(&self, _ctx: &mut Context) -> usize {
3392 0
3393 }
3394}
3395impl CachedPageData for NoCachedPageData {
3396 fn as_arc_any(self: Arc<Self>) -> Arc<dyn Any + Send + Sync + 'static> {
3397 self
3398 }
3399}
3400
3401pub struct CachedFieldData {
3402 pages: Vec<Arc<dyn CachedPageData>>,
3403}
3404
3405impl DeepSizeOf for CachedFieldData {
3406 fn deep_size_of_children(&self, ctx: &mut Context) -> usize {
3407 self.pages.deep_size_of_children(ctx)
3408 }
3409}
3410
3411#[derive(Debug, Clone)]
3413pub struct FieldDataCacheKey {
3414 pub column_index: u32,
3415}
3416
3417impl CacheKey for FieldDataCacheKey {
3418 type ValueType = CachedFieldData;
3419
3420 fn key(&self) -> std::borrow::Cow<'_, str> {
3421 self.column_index.to_string().into()
3422 }
3423
3424 fn type_name() -> &'static str {
3425 "FieldData"
3426 }
3427}
3428
3429impl StructuralFieldScheduler for StructuralPrimitiveFieldScheduler {
3430 fn initialize<'a>(
3431 &'a mut self,
3432 _filter: &'a FilterExpression,
3433 context: &'a SchedulerContext,
3434 ) -> BoxFuture<'a, Result<()>> {
3435 let cache_key = FieldDataCacheKey {
3436 column_index: self.column_index,
3437 };
3438 let cache = context.cache().clone();
3439
3440 async move {
3441 if let Some(cached_data) = cache.get_with_key(&cache_key).await {
3442 self.page_schedulers
3443 .iter_mut()
3444 .zip(cached_data.pages.iter())
3445 .for_each(|(page_scheduler, cached_data)| {
3446 page_scheduler.scheduler.load(cached_data);
3447 });
3448 return Ok(());
3449 }
3450
3451 let page_data = self
3452 .page_schedulers
3453 .iter_mut()
3454 .map(|s| s.scheduler.initialize(context.io()))
3455 .collect::<FuturesOrdered<_>>();
3456
3457 let page_data = page_data.try_collect::<Vec<_>>().await?;
3458 let cached_data = Arc::new(CachedFieldData { pages: page_data });
3459 cache.insert_with_key(&cache_key, cached_data).await;
3460 Ok(())
3461 }
3462 .boxed()
3463 }
3464
3465 fn schedule_ranges<'a>(
3466 &'a self,
3467 ranges: &[Range<u64>],
3468 _filter: &FilterExpression,
3469 ) -> Result<Box<dyn StructuralSchedulingJob + 'a>> {
3470 let ranges = ranges.to_vec();
3471 Ok(Box::new(StructuralPrimitiveFieldSchedulingJob::new(
3472 self, ranges,
3473 )))
3474 }
3475}
3476
3477#[derive(Debug)]
3480pub struct StructuralCompositeDecodeArrayTask {
3481 tasks: Vec<Box<dyn DecodePageTask>>,
3482 should_validate: bool,
3483 data_type: DataType,
3484}
3485
3486impl StructuralCompositeDecodeArrayTask {
3487 fn restore_validity(
3488 array: Arc<dyn Array>,
3489 unraveler: &mut CompositeRepDefUnraveler,
3490 ) -> Arc<dyn Array> {
3491 let validity = unraveler.unravel_validity(array.len());
3492 let Some(validity) = validity else {
3493 return array;
3494 };
3495 if array.data_type() == &DataType::Null {
3496 return array;
3498 }
3499 assert_eq!(validity.len(), array.len());
3500 make_array(unsafe {
3503 array
3504 .to_data()
3505 .into_builder()
3506 .nulls(Some(validity))
3507 .build_unchecked()
3508 })
3509 }
3510}
3511
3512impl StructuralDecodeArrayTask for StructuralCompositeDecodeArrayTask {
3513 fn decode(self: Box<Self>) -> Result<DecodedArray> {
3514 let mut arrays = Vec::with_capacity(self.tasks.len());
3515 let mut unravelers = Vec::with_capacity(self.tasks.len());
3516 let mut data_size = 0u64;
3517 for task in self.tasks {
3518 let decoded = task.decode()?;
3519 data_size += decoded.data.data_size();
3520 unravelers.push(decoded.repdef);
3521
3522 let array = make_array(
3523 decoded
3524 .data
3525 .into_arrow(self.data_type.clone(), self.should_validate)?,
3526 );
3527
3528 arrays.push(array);
3529 }
3530 let array_refs = arrays.iter().map(|arr| arr.as_ref()).collect::<Vec<_>>();
3531 let array = arrow_select::concat::concat(&array_refs)?;
3532 let mut repdef = CompositeRepDefUnraveler::new(unravelers);
3533
3534 let array = Self::restore_validity(array, &mut repdef);
3535
3536 Ok(DecodedArray {
3537 array,
3538 repdef,
3539 data_size,
3540 })
3541 }
3542}
3543
3544#[derive(Debug)]
3545pub struct StructuralPrimitiveFieldDecoder {
3546 field: Arc<ArrowField>,
3547 page_decoders: VecDeque<Box<dyn StructuralPageDecoder>>,
3548 should_validate: bool,
3549 rows_drained_in_current: u64,
3550}
3551
3552impl StructuralPrimitiveFieldDecoder {
3553 pub fn new(field: &Arc<ArrowField>, should_validate: bool) -> Self {
3554 Self {
3555 field: field.clone(),
3556 page_decoders: VecDeque::new(),
3557 should_validate,
3558 rows_drained_in_current: 0,
3559 }
3560 }
3561}
3562
3563impl StructuralFieldDecoder for StructuralPrimitiveFieldDecoder {
3564 fn accept_page(&mut self, child: LoadedPageShard) -> Result<()> {
3565 assert!(child.path.is_empty());
3566 self.page_decoders.push_back(child.decoder);
3567 Ok(())
3568 }
3569
3570 fn drain(&mut self, num_rows: u64) -> Result<Box<dyn StructuralDecodeArrayTask>> {
3571 let mut remaining = num_rows;
3572 let mut tasks = Vec::new();
3573 while remaining > 0 {
3574 let cur_page = self.page_decoders.front_mut().unwrap();
3575 let num_in_page = cur_page.num_rows() - self.rows_drained_in_current;
3576 let to_take = num_in_page.min(remaining);
3577
3578 let task = cur_page.drain(to_take)?;
3579 tasks.push(task);
3580
3581 if to_take == num_in_page {
3582 self.page_decoders.pop_front();
3583 self.rows_drained_in_current = 0;
3584 } else {
3585 self.rows_drained_in_current += to_take;
3586 }
3587
3588 remaining -= to_take;
3589 }
3590 Ok(Box::new(StructuralCompositeDecodeArrayTask {
3591 tasks,
3592 should_validate: self.should_validate,
3593 data_type: self.field.data_type().clone(),
3594 }))
3595 }
3596
3597 fn data_type(&self) -> &DataType {
3598 self.field.data_type()
3599 }
3600}
3601
3602struct SerializedFullZip {
3604 values: LanceBuffer,
3606 repetition_index: Option<LanceBuffer>,
3608}
3609
3610const MINIBLOCK_ALIGNMENT: usize = 8;
3630
3631pub struct PrimitiveStructuralEncoder {
3658 accumulation_queue: AccumulationQueue,
3660
3661 keep_original_array: bool,
3662 support_large_chunk: bool,
3663 accumulated_repdefs: Vec<RepDefBuilder>,
3664 compression_strategy: Arc<dyn CompressionStrategy>,
3666 column_index: u32,
3667 field: Field,
3668 encoding_metadata: Arc<HashMap<String, String>>,
3669 version: LanceFileVersion,
3670}
3671
3672struct CompressedLevelsChunk {
3673 data: LanceBuffer,
3674 num_levels: u16,
3675}
3676
3677struct CompressedLevels {
3678 data: Vec<CompressedLevelsChunk>,
3679 compression: CompressiveEncoding,
3680 rep_index: Option<LanceBuffer>,
3681}
3682
3683struct SerializedMiniBlockPage {
3684 num_buffers: u64,
3685 data: LanceBuffer,
3686 metadata: LanceBuffer,
3687}
3688
3689#[derive(Debug, Clone, Copy)]
3690struct DictEncodingBudget {
3691 max_dict_entries: u32,
3692 max_encoded_size: usize,
3693}
3694
3695impl PrimitiveStructuralEncoder {
3696 pub fn try_new(
3697 options: &EncodingOptions,
3698 compression_strategy: Arc<dyn CompressionStrategy>,
3699 column_index: u32,
3700 field: Field,
3701 encoding_metadata: Arc<HashMap<String, String>>,
3702 ) -> Result<Self> {
3703 Ok(Self {
3704 accumulation_queue: AccumulationQueue::new(
3705 options.cache_bytes_per_column,
3706 column_index,
3707 options.keep_original_array,
3708 ),
3709 support_large_chunk: options.support_large_chunk(),
3710 keep_original_array: options.keep_original_array,
3711 accumulated_repdefs: Vec::new(),
3712 column_index,
3713 compression_strategy,
3714 field,
3715 encoding_metadata,
3716 version: options.version,
3717 })
3718 }
3719
3720 fn is_narrow(data_block: &DataBlock) -> bool {
3728 const MINIBLOCK_MAX_BYTE_LENGTH_PER_VALUE: u64 = 256;
3729
3730 if let Some(max_len_array) = data_block.get_stat(Stat::MaxLength) {
3731 let max_len_array = max_len_array
3732 .as_any()
3733 .downcast_ref::<PrimitiveArray<UInt64Type>>()
3734 .unwrap();
3735 if max_len_array.value(0) < MINIBLOCK_MAX_BYTE_LENGTH_PER_VALUE {
3736 return true;
3737 }
3738 }
3739 false
3740 }
3741
3742 fn prefers_miniblock(
3743 data_block: &DataBlock,
3744 encoding_metadata: &HashMap<String, String>,
3745 ) -> bool {
3746 if let Some(user_requested) = encoding_metadata.get(STRUCTURAL_ENCODING_META_KEY) {
3748 return user_requested.to_lowercase() == STRUCTURAL_ENCODING_MINIBLOCK;
3749 }
3750 Self::is_narrow(data_block)
3752 }
3753
3754 fn repdef_too_sparse_for_miniblock(
3767 repdef: &crate::repdef::SerializedRepDefs,
3768 num_values: u64,
3769 ) -> bool {
3770 if num_values == 0 {
3771 return false;
3772 }
3773 let num_levels = repdef
3774 .repetition_levels
3775 .as_ref()
3776 .map(|r| r.len() as u64)
3777 .max(repdef.definition_levels.as_ref().map(|d| d.len() as u64))
3778 .unwrap_or(0);
3779 if num_levels == 0 {
3780 return false;
3781 }
3782
3783 let bits_per_rep = repdef
3785 .repetition_levels
3786 .as_ref()
3787 .and_then(|r| r.iter().max().copied())
3788 .map(|max_val| u16::BITS - max_val.leading_zeros())
3789 .unwrap_or(0) as u64;
3790 let bits_per_def = repdef
3791 .definition_levels
3792 .as_ref()
3793 .and_then(|d| d.iter().max().copied())
3794 .map(|max_val| u16::BITS - max_val.leading_zeros())
3795 .unwrap_or(0) as u64;
3796
3797 let bits_per_level = bits_per_rep + bits_per_def;
3798 if bits_per_level == 0 {
3799 return false;
3800 }
3801
3802 const REPDEF_BUDGET_BITS: u64 = 16 * 1024 * 8;
3804 let max_levels_per_chunk = REPDEF_BUDGET_BITS / bits_per_level;
3805
3806 let levels_per_chunk =
3809 (num_levels as f64 / num_values as f64) * *miniblock::MAX_MINIBLOCK_VALUES as f64;
3810
3811 levels_per_chunk > max_levels_per_chunk as f64
3812 }
3813
3814 fn prefers_fullzip(encoding_metadata: &HashMap<String, String>) -> bool {
3815 if let Some(user_requested) = encoding_metadata.get(STRUCTURAL_ENCODING_META_KEY) {
3819 return user_requested.to_lowercase() == STRUCTURAL_ENCODING_FULLZIP;
3820 }
3821 true
3822 }
3823
3824 fn serialize_miniblocks(
3871 miniblocks: MiniBlockCompressed,
3872 rep: Option<Vec<CompressedLevelsChunk>>,
3873 def: Option<Vec<CompressedLevelsChunk>>,
3874 support_large_chunk: bool,
3875 ) -> Result<SerializedMiniBlockPage> {
3876 let bytes_rep = rep
3877 .as_ref()
3878 .map(|rep| rep.iter().map(|r| r.data.len()).sum::<usize>())
3879 .unwrap_or(0);
3880 let bytes_def = def
3881 .as_ref()
3882 .map(|def| def.iter().map(|d| d.data.len()).sum::<usize>())
3883 .unwrap_or(0);
3884 let bytes_data = miniblocks.data.iter().map(|d| d.len()).sum::<usize>();
3885 let mut num_buffers = miniblocks.data.len();
3886 if rep.is_some() {
3887 num_buffers += 1;
3888 }
3889 if def.is_some() {
3890 num_buffers += 1;
3891 }
3892 let max_extra = 9 * num_buffers;
3894 let mut data_buffer = Vec::with_capacity(bytes_rep + bytes_def + bytes_data + max_extra);
3895 let chunk_size_bytes = if support_large_chunk { 4 } else { 2 };
3896 let mut meta_buffer = Vec::with_capacity(miniblocks.chunks.len() * chunk_size_bytes);
3897
3898 let mut rep_iter = rep.map(|r| r.into_iter());
3899 let mut def_iter = def.map(|d| d.into_iter());
3900
3901 let mut buffer_offsets = vec![0; miniblocks.data.len()];
3902 for chunk in miniblocks.chunks {
3903 let start_pos = data_buffer.len();
3904 debug_assert_eq!(start_pos % MINIBLOCK_ALIGNMENT, 0);
3906
3907 let rep = rep_iter.as_mut().map(|r| r.next().unwrap());
3908 let def = def_iter.as_mut().map(|d| d.next().unwrap());
3909
3910 let num_levels = rep
3912 .as_ref()
3913 .map(|r| r.num_levels)
3914 .unwrap_or(def.as_ref().map(|d| d.num_levels).unwrap_or(0));
3915 data_buffer.extend_from_slice(&num_levels.to_le_bytes());
3916
3917 if let Some(rep) = rep.as_ref() {
3919 let bytes_rep = u16::try_from(rep.data.len()).map_err(|_| {
3920 Error::internal(format!(
3921 "Repetition buffer size ({} bytes) too large",
3922 rep.data.len()
3923 ))
3924 })?;
3925 data_buffer.extend_from_slice(&bytes_rep.to_le_bytes());
3926 }
3927 if let Some(def) = def.as_ref() {
3928 let bytes_def = u16::try_from(def.data.len()).map_err(|_| {
3929 Error::internal(format!(
3930 "Definition buffer size ({} bytes) too large",
3931 def.data.len()
3932 ))
3933 })?;
3934 data_buffer.extend_from_slice(&bytes_def.to_le_bytes());
3935 }
3936
3937 if support_large_chunk {
3938 for &buffer_size in &chunk.buffer_sizes {
3939 data_buffer.extend_from_slice(&buffer_size.to_le_bytes());
3940 }
3941 } else {
3942 for &buffer_size in &chunk.buffer_sizes {
3943 data_buffer.extend_from_slice(&(buffer_size as u16).to_le_bytes());
3944 }
3945 }
3946
3947 let add_padding = |data_buffer: &mut Vec<u8>| {
3949 let pad = pad_bytes::<MINIBLOCK_ALIGNMENT>(data_buffer.len());
3950 data_buffer.extend(iter::repeat_n(FILL_BYTE, pad));
3951 };
3952 add_padding(&mut data_buffer);
3953
3954 if let Some(rep) = rep.as_ref() {
3956 data_buffer.extend_from_slice(&rep.data);
3957 add_padding(&mut data_buffer);
3958 }
3959 if let Some(def) = def.as_ref() {
3960 data_buffer.extend_from_slice(&def.data);
3961 add_padding(&mut data_buffer);
3962 }
3963 for (buffer_size, (buffer, buffer_offset)) in chunk
3964 .buffer_sizes
3965 .iter()
3966 .zip(miniblocks.data.iter().zip(buffer_offsets.iter_mut()))
3967 {
3968 let start = *buffer_offset;
3969 let end = start + *buffer_size as usize;
3970 *buffer_offset += *buffer_size as usize;
3971 data_buffer.extend_from_slice(&buffer[start..end]);
3972 add_padding(&mut data_buffer);
3973 }
3974
3975 let chunk_bytes = data_buffer.len() - start_pos;
3976 let max_chunk_size = if support_large_chunk {
3977 4 * 1024 * 1024 * 1024 } else {
3979 32 * 1024 };
3981 assert!(chunk_bytes <= max_chunk_size);
3982 assert!(chunk_bytes > 0);
3983 assert_eq!(chunk_bytes % 8, 0);
3984 assert!(chunk.log_num_values <= 12);
3986 let divided_bytes = chunk_bytes / MINIBLOCK_ALIGNMENT;
3990 let divided_bytes_minus_one = (divided_bytes - 1) as u64;
3991
3992 let metadata = (divided_bytes_minus_one << 4) | chunk.log_num_values as u64;
3993 if support_large_chunk {
3994 meta_buffer.extend_from_slice(&(metadata as u32).to_le_bytes());
3995 } else {
3996 meta_buffer.extend_from_slice(&(metadata as u16).to_le_bytes());
3997 }
3998 }
3999
4000 let data_buffer = LanceBuffer::from(data_buffer);
4001 let metadata_buffer = LanceBuffer::from(meta_buffer);
4002
4003 Ok(SerializedMiniBlockPage {
4004 num_buffers: miniblocks.data.len() as u64,
4005 data: data_buffer,
4006 metadata: metadata_buffer,
4007 })
4008 }
4009
4010 fn compress_levels(
4015 mut levels: RepDefSlicer<'_>,
4016 num_elements: u64,
4017 compression_strategy: &dyn CompressionStrategy,
4018 chunks: &[MiniBlockChunk],
4019 max_rep: u16,
4021 ) -> Result<CompressedLevels> {
4022 let mut rep_index = if max_rep > 0 {
4023 Vec::with_capacity(chunks.len())
4024 } else {
4025 vec![]
4026 };
4027 let num_levels = levels.num_levels() as u64;
4029 let levels_buf = levels.all_levels().clone();
4030
4031 let mut fixed_width_block = FixedWidthDataBlock {
4032 data: levels_buf,
4033 bits_per_value: 16,
4034 num_values: num_levels,
4035 block_info: BlockInfo::new(),
4036 };
4037 fixed_width_block.compute_stat();
4039
4040 let levels_block = DataBlock::FixedWidth(fixed_width_block);
4041 let levels_field = Field::new_arrow("", DataType::UInt16, false)?;
4042 let (compressor, compressor_desc) =
4044 compression_strategy.create_block_compressor(&levels_field, &levels_block)?;
4045 let mut level_chunks = Vec::with_capacity(chunks.len());
4047 let mut values_counter = 0;
4048 for (chunk_idx, chunk) in chunks.iter().enumerate() {
4049 let chunk_num_values = chunk.num_values(values_counter, num_elements);
4050 debug_assert!(chunk_num_values > 0);
4051 values_counter += chunk_num_values;
4052 let chunk_levels = if chunk_idx < chunks.len() - 1 {
4053 levels.slice_next(chunk_num_values as usize)
4054 } else {
4055 levels.slice_rest()
4056 };
4057 let num_chunk_levels = (chunk_levels.len() / 2) as u64;
4058 if max_rep > 0 {
4059 let rep_values = chunk_levels.borrow_to_typed_slice::<u16>();
4069 let rep_values = rep_values.as_ref();
4070
4071 let mut num_rows = rep_values.iter().skip(1).filter(|v| **v == max_rep).count();
4074 let num_leftovers = if chunk_idx < chunks.len() - 1 {
4075 rep_values
4076 .iter()
4077 .rev()
4078 .position(|v| *v == max_rep)
4079 .map(|pos| pos + 1)
4081 .unwrap_or(rep_values.len())
4082 } else {
4083 0
4085 };
4086
4087 if chunk_idx != 0 && rep_values.first() == Some(&max_rep) {
4088 let rep_len = rep_index.len();
4092 if rep_index[rep_len - 1] != 0 {
4093 rep_index[rep_len - 2] += 1;
4095 rep_index[rep_len - 1] = 0;
4096 }
4097 }
4098
4099 if chunk_idx == chunks.len() - 1 {
4100 num_rows += 1;
4102 }
4103 rep_index.push(num_rows as u64);
4104 rep_index.push(num_leftovers as u64);
4105 }
4106 let mut chunk_fixed_width = FixedWidthDataBlock {
4107 data: chunk_levels,
4108 bits_per_value: 16,
4109 num_values: num_chunk_levels,
4110 block_info: BlockInfo::new(),
4111 };
4112 chunk_fixed_width.compute_stat();
4113 let chunk_levels_block = DataBlock::FixedWidth(chunk_fixed_width);
4114 let compressed_levels = compressor.compress(chunk_levels_block)?;
4115 level_chunks.push(CompressedLevelsChunk {
4116 data: compressed_levels,
4117 num_levels: num_chunk_levels as u16,
4118 });
4119 }
4120 debug_assert_eq!(levels.num_levels_remaining(), 0);
4121 let rep_index = if rep_index.is_empty() {
4122 None
4123 } else {
4124 Some(LanceBuffer::reinterpret_vec(rep_index))
4125 };
4126 Ok(CompressedLevels {
4127 data: level_chunks,
4128 compression: compressor_desc,
4129 rep_index,
4130 })
4131 }
4132
4133 fn encode_simple_all_null(
4134 column_idx: u32,
4135 num_rows: u64,
4136 row_number: u64,
4137 ) -> Result<EncodedPage> {
4138 let description =
4139 ProtobufUtils21::constant_layout(&[DefinitionInterpretation::NullableItem], None);
4140 Ok(EncodedPage {
4141 column_idx,
4142 data: vec![],
4143 description: PageEncoding::Structural(description),
4144 num_rows,
4145 row_number,
4146 })
4147 }
4148
4149 fn encode_complex_all_null_vals(
4150 data: &Arc<[u16]>,
4151 compression_strategy: &dyn CompressionStrategy,
4152 ) -> Result<(LanceBuffer, pb21::CompressiveEncoding)> {
4153 let buffer = LanceBuffer::reinterpret_slice(data.clone());
4154 let mut fixed_width_block = FixedWidthDataBlock {
4155 data: buffer,
4156 bits_per_value: 16,
4157 num_values: data.len() as u64,
4158 block_info: BlockInfo::new(),
4159 };
4160 fixed_width_block.compute_stat();
4161
4162 let levels_block = DataBlock::FixedWidth(fixed_width_block);
4163 let levels_field = Field::new_arrow("", DataType::UInt16, false)?;
4164 let (compressor, encoding) =
4165 compression_strategy.create_block_compressor(&levels_field, &levels_block)?;
4166 let compressed_buffer = compressor.compress(levels_block)?;
4167 Ok((compressed_buffer, encoding))
4168 }
4169
4170 fn encode_complex_all_null(
4174 column_idx: u32,
4175 repdef: crate::repdef::SerializedRepDefs,
4176 row_number: u64,
4177 num_rows: u64,
4178 version: LanceFileVersion,
4179 compression_strategy: &dyn CompressionStrategy,
4180 ) -> Result<EncodedPage> {
4181 if version.resolve() < LanceFileVersion::V2_2 {
4182 let rep_bytes = if let Some(rep) = repdef.repetition_levels.as_ref() {
4183 LanceBuffer::reinterpret_slice(rep.clone())
4184 } else {
4185 LanceBuffer::empty()
4186 };
4187
4188 let def_bytes = if let Some(def) = repdef.definition_levels.as_ref() {
4189 LanceBuffer::reinterpret_slice(def.clone())
4190 } else {
4191 LanceBuffer::empty()
4192 };
4193
4194 let description = ProtobufUtils21::constant_layout(&repdef.def_meaning, None);
4195 return Ok(EncodedPage {
4196 column_idx,
4197 data: vec![rep_bytes, def_bytes],
4198 description: PageEncoding::Structural(description),
4199 num_rows,
4200 row_number,
4201 });
4202 }
4203
4204 let (rep_bytes, rep_encoding, num_rep_values) = if let Some(rep) =
4205 repdef.repetition_levels.as_ref()
4206 {
4207 let num_values = rep.len() as u64;
4208 let (buffer, encoding) = Self::encode_complex_all_null_vals(rep, compression_strategy)?;
4209 (buffer, Some(encoding), num_values)
4210 } else {
4211 (LanceBuffer::empty(), None, 0)
4212 };
4213
4214 let (def_bytes, def_encoding, num_def_values) = if let Some(def) =
4215 repdef.definition_levels.as_ref()
4216 {
4217 let num_values = def.len() as u64;
4218 let (buffer, encoding) = Self::encode_complex_all_null_vals(def, compression_strategy)?;
4219 (buffer, Some(encoding), num_values)
4220 } else {
4221 (LanceBuffer::empty(), None, 0)
4222 };
4223
4224 let description = ProtobufUtils21::compressed_all_null_constant_layout(
4225 &repdef.def_meaning,
4226 rep_encoding,
4227 def_encoding,
4228 num_rep_values,
4229 num_def_values,
4230 );
4231 Ok(EncodedPage {
4232 column_idx,
4233 data: vec![rep_bytes, def_bytes],
4234 description: PageEncoding::Structural(description),
4235 num_rows,
4236 row_number,
4237 })
4238 }
4239
4240 fn leaf_validity(
4241 repdef: &crate::repdef::SerializedRepDefs,
4242 num_values: usize,
4243 ) -> Result<Option<BooleanBuffer>> {
4244 let rep = repdef
4245 .repetition_levels
4246 .as_ref()
4247 .map(|rep| rep.as_ref().to_vec());
4248 let def = repdef
4249 .definition_levels
4250 .as_ref()
4251 .map(|def| def.as_ref().to_vec());
4252 let mut unraveler = RepDefUnraveler::new(
4253 rep,
4254 def,
4255 repdef.def_meaning.clone().into(),
4256 num_values as u64,
4257 );
4258 if unraveler.is_all_valid() {
4259 return Ok(None);
4260 }
4261 let mut validity = BooleanBufferBuilder::new(num_values);
4262 unraveler.unravel_validity(&mut validity);
4263 Ok(Some(validity.finish()))
4264 }
4265
4266 fn is_constant_values(
4267 arrays: &[ArrayRef],
4268 scalar: &ArrayRef,
4269 validity: Option<&BooleanBuffer>,
4270 ) -> Result<bool> {
4271 debug_assert_eq!(scalar.len(), 1);
4272 debug_assert_eq!(scalar.null_count(), 0);
4273
4274 match scalar.data_type() {
4275 DataType::Boolean => {
4276 let mut global_idx = 0usize;
4277 let scalar_val = scalar.as_boolean().value(0);
4278 for arr in arrays {
4279 let bool_arr = arr.as_boolean();
4280 for i in 0..arr.len() {
4281 let is_valid = validity.map(|v| v.value(global_idx)).unwrap_or(true);
4282 global_idx += 1;
4283 if !is_valid {
4284 continue;
4285 }
4286 if bool_arr.value(i) != scalar_val {
4287 return Ok(false);
4288 }
4289 }
4290 }
4291 Ok(true)
4292 }
4293 DataType::Utf8 => Self::is_constant_utf8::<i32>(arrays, scalar, validity),
4294 DataType::LargeUtf8 => Self::is_constant_utf8::<i64>(arrays, scalar, validity),
4295 DataType::Binary => Self::is_constant_binary::<i32>(arrays, scalar, validity),
4296 DataType::LargeBinary => Self::is_constant_binary::<i64>(arrays, scalar, validity),
4297 data_type => {
4298 let mut global_idx = 0usize;
4299 let Some(byte_width) = data_type.byte_width_opt() else {
4300 return Ok(false);
4301 };
4302 let scalar_data = scalar.to_data();
4303 if scalar_data.buffers().len() != 1 || !scalar_data.child_data().is_empty() {
4304 return Ok(false);
4305 }
4306 let scalar_bytes = scalar_data.buffers()[0].as_slice();
4307 if scalar_bytes.len() != byte_width {
4308 return Ok(false);
4309 }
4310
4311 for arr in arrays {
4312 let data = arr.to_data();
4313 if data.buffers().is_empty() {
4314 return Ok(false);
4315 }
4316 let buf = data.buffers()[0].as_slice();
4317 let base = data.offset();
4318 for i in 0..arr.len() {
4319 let is_valid = validity.map(|v| v.value(global_idx)).unwrap_or(true);
4320 global_idx += 1;
4321 if !is_valid {
4322 continue;
4323 }
4324 let start = (base + i) * byte_width;
4325 if buf[start..start + byte_width] != scalar_bytes[..] {
4326 return Ok(false);
4327 }
4328 }
4329 }
4330 Ok(true)
4331 }
4332 }
4333 }
4334
4335 fn is_constant_utf8<O: arrow_array::OffsetSizeTrait>(
4336 arrays: &[ArrayRef],
4337 scalar: &ArrayRef,
4338 validity: Option<&BooleanBuffer>,
4339 ) -> Result<bool> {
4340 debug_assert_eq!(scalar.len(), 1);
4341 let scalar_val = scalar.as_string::<O>().value(0).as_bytes();
4342 let mut global_idx = 0usize;
4343 for arr in arrays {
4344 let str_arr = arr.as_string::<O>();
4345 for i in 0..arr.len() {
4346 let is_valid = validity.map(|v| v.value(global_idx)).unwrap_or(true);
4347 global_idx += 1;
4348 if !is_valid {
4349 continue;
4350 }
4351 if str_arr.value(i).as_bytes() != scalar_val {
4352 return Ok(false);
4353 }
4354 }
4355 }
4356 Ok(true)
4357 }
4358
4359 fn is_constant_binary<O: arrow_array::OffsetSizeTrait>(
4360 arrays: &[ArrayRef],
4361 scalar: &ArrayRef,
4362 validity: Option<&BooleanBuffer>,
4363 ) -> Result<bool> {
4364 debug_assert_eq!(scalar.len(), 1);
4365 let scalar_val = scalar.as_binary::<O>().value(0);
4366 let mut global_idx = 0usize;
4367 for arr in arrays {
4368 let bin_arr = arr.as_binary::<O>();
4369 for i in 0..arr.len() {
4370 let is_valid = validity.map(|v| v.value(global_idx)).unwrap_or(true);
4371 global_idx += 1;
4372 if !is_valid {
4373 continue;
4374 }
4375 if bin_arr.value(i) != scalar_val {
4376 return Ok(false);
4377 }
4378 }
4379 }
4380 Ok(true)
4381 }
4382
4383 fn find_constant_scalar(
4384 arrays: &[ArrayRef],
4385 validity: Option<&BooleanBuffer>,
4386 ) -> Result<Option<ArrayRef>> {
4387 if arrays.is_empty() {
4388 return Ok(None);
4389 }
4390
4391 let global_scalar_idx = if let Some(validity) = validity {
4392 let Some(idx) = (0..validity.len()).find(|&i| validity.value(i)) else {
4393 return Ok(None);
4394 };
4395 idx
4396 } else {
4397 0
4398 };
4399
4400 let mut idx_remaining = global_scalar_idx;
4401 let mut scalar_arr_idx = 0usize;
4402 while scalar_arr_idx < arrays.len() {
4403 let len = arrays[scalar_arr_idx].len();
4404 if idx_remaining < len {
4405 break;
4406 }
4407 idx_remaining -= len;
4408 scalar_arr_idx += 1;
4409 }
4410
4411 if scalar_arr_idx >= arrays.len() {
4412 return Ok(None);
4413 }
4414
4415 let scalar =
4416 lance_arrow::scalar::extract_scalar_value(&arrays[scalar_arr_idx], idx_remaining)?;
4417 if scalar.null_count() != 0 {
4418 return Ok(None);
4419 }
4420 if !Self::is_constant_values(arrays, &scalar, validity)? {
4421 return Ok(None);
4422 }
4423 Ok(Some(scalar))
4424 }
4425
4426 fn resolve_dict_values_compression_metadata(
4427 field_metadata: &HashMap<String, String>,
4428 env_compression: Option<String>,
4429 env_compression_level: Option<String>,
4430 ) -> HashMap<String, String> {
4431 let mut metadata = HashMap::new();
4432
4433 let compression = field_metadata
4434 .get(DICT_VALUES_COMPRESSION_META_KEY)
4435 .cloned()
4436 .or(env_compression)
4437 .unwrap_or_else(|| DEFAULT_DICT_VALUES_COMPRESSION.to_string());
4438 metadata.insert(COMPRESSION_META_KEY.to_string(), compression);
4439
4440 if let Some(compression_level) = field_metadata
4441 .get(DICT_VALUES_COMPRESSION_LEVEL_META_KEY)
4442 .cloned()
4443 .or(env_compression_level)
4444 {
4445 metadata.insert(COMPRESSION_LEVEL_META_KEY.to_string(), compression_level);
4446 }
4447
4448 metadata
4449 }
4450
4451 fn build_dict_values_compressor_field(field: &Field) -> Result<Field> {
4452 let mut dict_values_field = Field::new_arrow("", DataType::UInt16, false)?;
4457 dict_values_field.metadata = Self::resolve_dict_values_compression_metadata(
4458 &field.metadata,
4459 env::var(DICT_VALUES_COMPRESSION_ENV_VAR).ok(),
4460 env::var(DICT_VALUES_COMPRESSION_LEVEL_ENV_VAR).ok(),
4461 );
4462 Ok(dict_values_field)
4463 }
4464
4465 #[allow(clippy::too_many_arguments)]
4466 fn encode_miniblock(
4467 column_idx: u32,
4468 field: &Field,
4469 compression_strategy: &dyn CompressionStrategy,
4470 data: DataBlock,
4471 repdef: crate::repdef::SerializedRepDefs,
4472 row_number: u64,
4473 dictionary_data: Option<DataBlock>,
4474 num_rows: u64,
4475 support_large_chunk: bool,
4476 ) -> Result<EncodedPage> {
4477 if let DataBlock::AllNull(_null_block) = data {
4478 unreachable!()
4481 }
4482
4483 let num_items = data.num_values();
4484
4485 let compressor = compression_strategy.create_miniblock_compressor(field, &data)?;
4486 let (compressed_data, value_encoding) = compressor.compress(data)?;
4487
4488 let max_rep = repdef.def_meaning.iter().filter(|l| l.is_list()).count() as u16;
4489
4490 let mut compressed_rep = repdef
4491 .rep_slicer()
4492 .map(|rep_slicer| {
4493 Self::compress_levels(
4494 rep_slicer,
4495 num_items,
4496 compression_strategy,
4497 &compressed_data.chunks,
4498 max_rep,
4499 )
4500 })
4501 .transpose()?;
4502
4503 let (rep_index, rep_index_depth) =
4504 match compressed_rep.as_mut().and_then(|cr| cr.rep_index.as_mut()) {
4505 Some(rep_index) => (Some(rep_index.clone()), 1),
4506 None => (None, 0),
4507 };
4508
4509 let mut compressed_def = repdef
4510 .def_slicer()
4511 .map(|def_slicer| {
4512 Self::compress_levels(
4513 def_slicer,
4514 num_items,
4515 compression_strategy,
4516 &compressed_data.chunks,
4517 0,
4518 )
4519 })
4520 .transpose()?;
4521
4522 let rep_data = compressed_rep
4528 .as_mut()
4529 .map(|cr| std::mem::take(&mut cr.data));
4530 let def_data = compressed_def
4531 .as_mut()
4532 .map(|cd| std::mem::take(&mut cd.data));
4533
4534 let serialized =
4535 Self::serialize_miniblocks(compressed_data, rep_data, def_data, support_large_chunk)?;
4536
4537 let mut data = Vec::with_capacity(4);
4539 data.push(serialized.metadata);
4540 data.push(serialized.data);
4541
4542 if let Some(dictionary_data) = dictionary_data {
4543 let num_dictionary_items = dictionary_data.num_values();
4544 let dict_values_field = Self::build_dict_values_compressor_field(field)?;
4545
4546 let (compressor, dictionary_encoding) = compression_strategy
4547 .create_block_compressor(&dict_values_field, &dictionary_data)?;
4548 let dictionary_buffer = compressor.compress(dictionary_data)?;
4549
4550 data.push(dictionary_buffer);
4551 if let Some(rep_index) = rep_index {
4552 data.push(rep_index);
4553 }
4554
4555 let description = ProtobufUtils21::miniblock_layout(
4556 compressed_rep.map(|cr| cr.compression),
4557 compressed_def.map(|cd| cd.compression),
4558 value_encoding,
4559 rep_index_depth,
4560 serialized.num_buffers,
4561 Some((dictionary_encoding, num_dictionary_items)),
4562 &repdef.def_meaning,
4563 num_items,
4564 support_large_chunk,
4565 );
4566 Ok(EncodedPage {
4567 num_rows,
4568 column_idx,
4569 data,
4570 description: PageEncoding::Structural(description),
4571 row_number,
4572 })
4573 } else {
4574 let description = ProtobufUtils21::miniblock_layout(
4575 compressed_rep.map(|cr| cr.compression),
4576 compressed_def.map(|cd| cd.compression),
4577 value_encoding,
4578 rep_index_depth,
4579 serialized.num_buffers,
4580 None,
4581 &repdef.def_meaning,
4582 num_items,
4583 support_large_chunk,
4584 );
4585
4586 if let Some(rep_index) = rep_index {
4587 let view = rep_index.borrow_to_typed_slice::<u64>();
4588 let total = view.chunks_exact(2).map(|c| c[0]).sum::<u64>();
4589 debug_assert_eq!(total, num_rows);
4590
4591 data.push(rep_index);
4592 }
4593
4594 Ok(EncodedPage {
4595 num_rows,
4596 column_idx,
4597 data,
4598 description: PageEncoding::Structural(description),
4599 row_number,
4600 })
4601 }
4602 }
4603
4604 fn serialize_full_zip_fixed(
4606 fixed: FixedWidthDataBlock,
4607 mut repdef: ControlWordIterator,
4608 num_values: u64,
4609 ) -> SerializedFullZip {
4610 let len = fixed.data.len() + repdef.bytes_per_word() * num_values as usize;
4611 let mut zipped_data = Vec::with_capacity(len);
4612
4613 let max_rep_index_val = if repdef.has_repetition() {
4614 len as u64
4615 } else {
4616 0
4618 };
4619 let mut rep_index_builder =
4620 BytepackedIntegerEncoder::with_capacity(num_values as usize + 1, max_rep_index_val);
4621
4622 assert_eq!(
4625 fixed.bits_per_value % 8,
4626 0,
4627 "Non-byte aligned full-zip compression not yet supported"
4628 );
4629
4630 let bytes_per_value = fixed.bits_per_value as usize / 8;
4631 let mut offset = 0;
4632
4633 if bytes_per_value == 0 {
4634 while let Some(control) = repdef.append_next(&mut zipped_data) {
4636 if control.is_new_row {
4637 debug_assert!(offset <= len);
4639 unsafe { rep_index_builder.append(offset as u64) };
4641 }
4642 offset = zipped_data.len();
4643 }
4644 } else {
4645 let mut data_iter = fixed.data.chunks_exact(bytes_per_value);
4647 while let Some(control) = repdef.append_next(&mut zipped_data) {
4648 if control.is_new_row {
4649 debug_assert!(offset <= len);
4651 unsafe { rep_index_builder.append(offset as u64) };
4653 }
4654 if control.is_visible {
4655 let value = data_iter.next().unwrap();
4656 zipped_data.extend_from_slice(value);
4657 }
4658 offset = zipped_data.len();
4659 }
4660 }
4661
4662 debug_assert_eq!(zipped_data.len(), len);
4663 unsafe {
4666 rep_index_builder.append(zipped_data.len() as u64);
4667 }
4668
4669 let zipped_data = LanceBuffer::from(zipped_data);
4670 let rep_index = rep_index_builder.into_data();
4671 let rep_index = if rep_index.is_empty() {
4672 None
4673 } else {
4674 Some(LanceBuffer::from(rep_index))
4675 };
4676 SerializedFullZip {
4677 values: zipped_data,
4678 repetition_index: rep_index,
4679 }
4680 }
4681
4682 fn serialize_full_zip_variable(
4686 variable: VariableWidthBlock,
4687 mut repdef: ControlWordIterator,
4688 num_items: u64,
4689 ) -> SerializedFullZip {
4690 let bytes_per_offset = variable.bits_per_offset as usize / 8;
4691 assert_eq!(
4692 variable.bits_per_offset % 8,
4693 0,
4694 "Only byte-aligned offsets supported"
4695 );
4696 let len = variable.data.len()
4697 + repdef.bytes_per_word() * num_items as usize
4698 + bytes_per_offset * variable.num_values as usize;
4699 let mut buf = Vec::with_capacity(len);
4700
4701 let max_rep_index_val = len as u64;
4702 let mut rep_index_builder =
4703 BytepackedIntegerEncoder::with_capacity(num_items as usize + 1, max_rep_index_val);
4704
4705 match bytes_per_offset {
4707 4 => {
4708 let offs = variable.offsets.borrow_to_typed_slice::<u32>();
4709 let mut rep_offset = 0;
4710 let mut windows_iter = offs.as_ref().windows(2);
4711 while let Some(control) = repdef.append_next(&mut buf) {
4712 if control.is_new_row {
4713 debug_assert!(rep_offset <= len);
4715 unsafe { rep_index_builder.append(rep_offset as u64) };
4717 }
4718 if control.is_visible {
4719 let window = windows_iter.next().unwrap();
4720 if control.is_valid_item {
4721 buf.extend_from_slice(&(window[1] - window[0]).to_le_bytes());
4722 buf.extend_from_slice(
4723 &variable.data[window[0] as usize..window[1] as usize],
4724 );
4725 }
4726 }
4727 rep_offset = buf.len();
4728 }
4729 }
4730 8 => {
4731 let offs = variable.offsets.borrow_to_typed_slice::<u64>();
4732 let mut rep_offset = 0;
4733 let mut windows_iter = offs.as_ref().windows(2);
4734 while let Some(control) = repdef.append_next(&mut buf) {
4735 if control.is_new_row {
4736 debug_assert!(rep_offset <= len);
4738 unsafe { rep_index_builder.append(rep_offset as u64) };
4740 }
4741 if control.is_visible {
4742 let window = windows_iter.next().unwrap();
4743 if control.is_valid_item {
4744 buf.extend_from_slice(&(window[1] - window[0]).to_le_bytes());
4745 buf.extend_from_slice(
4746 &variable.data[window[0] as usize..window[1] as usize],
4747 );
4748 }
4749 }
4750 rep_offset = buf.len();
4751 }
4752 }
4753 _ => panic!("Unsupported offset size"),
4754 }
4755
4756 debug_assert!(buf.len() <= len);
4759 unsafe {
4762 rep_index_builder.append(buf.len() as u64);
4763 }
4764
4765 let zipped_data = LanceBuffer::from(buf);
4766 let rep_index = rep_index_builder.into_data();
4767 debug_assert!(!rep_index.is_empty());
4768 let rep_index = Some(LanceBuffer::from(rep_index));
4769 SerializedFullZip {
4770 values: zipped_data,
4771 repetition_index: rep_index,
4772 }
4773 }
4774
4775 fn serialize_full_zip(
4778 compressed_data: PerValueDataBlock,
4779 repdef: ControlWordIterator,
4780 num_items: u64,
4781 ) -> SerializedFullZip {
4782 match compressed_data {
4783 PerValueDataBlock::Fixed(fixed) => {
4784 Self::serialize_full_zip_fixed(fixed, repdef, num_items)
4785 }
4786 PerValueDataBlock::Variable(var) => {
4787 Self::serialize_full_zip_variable(var, repdef, num_items)
4788 }
4789 }
4790 }
4791
4792 fn encode_full_zip(
4793 column_idx: u32,
4794 field: &Field,
4795 compression_strategy: &dyn CompressionStrategy,
4796 data: DataBlock,
4797 repdef: crate::repdef::SerializedRepDefs,
4798 row_number: u64,
4799 num_lists: u64,
4800 ) -> Result<EncodedPage> {
4801 let max_rep = repdef
4802 .repetition_levels
4803 .as_ref()
4804 .map_or(0, |r| r.iter().max().copied().unwrap_or(0));
4805 let max_def = repdef
4806 .definition_levels
4807 .as_ref()
4808 .map_or(0, |d| d.iter().max().copied().unwrap_or(0));
4809
4810 let (num_items, num_visible_items) =
4814 if let Some(rep_levels) = repdef.repetition_levels.as_ref() {
4815 (rep_levels.len() as u64, data.num_values())
4818 } else {
4819 (data.num_values(), data.num_values())
4821 };
4822
4823 let max_visible_def = repdef.max_visible_level.unwrap_or(u16::MAX);
4824
4825 let repdef_iter = build_control_word_iterator(
4826 repdef.repetition_levels.as_deref(),
4827 max_rep,
4828 repdef.definition_levels.as_deref(),
4829 max_def,
4830 max_visible_def,
4831 num_items as usize,
4832 );
4833 let bits_rep = repdef_iter.bits_rep();
4834 let bits_def = repdef_iter.bits_def();
4835
4836 let compressor = compression_strategy.create_per_value(field, &data)?;
4837 let (compressed_data, value_encoding) = compressor.compress(data)?;
4838
4839 let description = match &compressed_data {
4840 PerValueDataBlock::Fixed(fixed) => ProtobufUtils21::fixed_full_zip_layout(
4841 bits_rep,
4842 bits_def,
4843 fixed.bits_per_value as u32,
4844 value_encoding,
4845 &repdef.def_meaning,
4846 num_items as u32,
4847 num_visible_items as u32,
4848 ),
4849 PerValueDataBlock::Variable(variable) => ProtobufUtils21::variable_full_zip_layout(
4850 bits_rep,
4851 bits_def,
4852 variable.bits_per_offset as u32,
4853 value_encoding,
4854 &repdef.def_meaning,
4855 num_items as u32,
4856 num_visible_items as u32,
4857 ),
4858 };
4859
4860 let zipped = Self::serialize_full_zip(compressed_data, repdef_iter, num_items);
4861
4862 let data = if let Some(repindex) = zipped.repetition_index {
4863 vec![zipped.values, repindex]
4864 } else {
4865 vec![zipped.values]
4866 };
4867
4868 Ok(EncodedPage {
4869 num_rows: num_lists,
4870 column_idx,
4871 data,
4872 description: PageEncoding::Structural(description),
4873 row_number,
4874 })
4875 }
4876
4877 fn should_dictionary_encode(
4878 data_block: &DataBlock,
4879 field: &Field,
4880 version: LanceFileVersion,
4881 ) -> Option<DictEncodingBudget> {
4882 const DEFAULT_SAMPLE_SIZE: usize = 4096;
4883 const DEFAULT_SAMPLE_UNIQUE_RATIO: f64 = 0.98;
4884
4885 match data_block {
4888 DataBlock::FixedWidth(fixed) => {
4889 if fixed.bits_per_value == 64 && version < LanceFileVersion::V2_2 {
4890 return None;
4891 }
4892 if fixed.bits_per_value != 64 && fixed.bits_per_value != 128 {
4893 return None;
4894 }
4895 if fixed.bits_per_value % 8 != 0 {
4896 return None;
4897 }
4898 }
4899 DataBlock::VariableWidth(var) => {
4900 if var.bits_per_offset != 32 && var.bits_per_offset != 64 {
4901 return None;
4902 }
4903 }
4904 _ => return None,
4905 }
4906
4907 let too_small = env::var("LANCE_ENCODING_DICT_TOO_SMALL")
4909 .ok()
4910 .and_then(|val| val.parse().ok())
4911 .unwrap_or(100);
4912 if data_block.num_values() < too_small {
4913 return None;
4914 }
4915
4916 let num_values = data_block.num_values();
4917
4918 let divisor: u64 = field
4921 .metadata
4922 .get(DICT_DIVISOR_META_KEY)
4923 .and_then(|val| val.parse().ok())
4924 .or_else(|| {
4925 env::var("LANCE_ENCODING_DICT_DIVISOR")
4926 .ok()
4927 .and_then(|val| val.parse().ok())
4928 })
4929 .unwrap_or(DEFAULT_DICT_DIVISOR);
4930
4931 let max_cardinality: u64 = env::var("LANCE_ENCODING_DICT_MAX_CARDINALITY")
4932 .ok()
4933 .and_then(|val| val.parse().ok())
4934 .unwrap_or(DEFAULT_DICT_MAX_CARDINALITY);
4935
4936 let threshold_cardinality = num_values
4937 .checked_div(divisor.max(1))
4938 .unwrap_or(0)
4939 .min(max_cardinality);
4940 if threshold_cardinality == 0 {
4941 return None;
4942 }
4943
4944 let threshold_ratio = field
4946 .metadata
4947 .get(DICT_SIZE_RATIO_META_KEY)
4948 .and_then(|val| val.parse::<f64>().ok())
4949 .or_else(|| {
4950 env::var("LANCE_ENCODING_DICT_SIZE_RATIO")
4951 .ok()
4952 .and_then(|val| val.parse().ok())
4953 })
4954 .unwrap_or(DEFAULT_DICT_SIZE_RATIO);
4955
4956 if threshold_ratio <= 0.0 || threshold_ratio > 1.0 {
4957 panic!(
4958 "Invalid parameter: dict-size-ratio is {} which is not in the range (0, 1].",
4959 threshold_ratio
4960 );
4961 }
4962
4963 let data_size = data_block.data_size();
4964 if data_size == 0 {
4965 return None;
4966 }
4967
4968 let max_encoded_size = (data_size as f64 * threshold_ratio) as u64;
4969 let max_encoded_size = usize::try_from(max_encoded_size).ok()?;
4970
4971 if Self::sample_is_near_unique(
4973 data_block,
4974 DEFAULT_SAMPLE_SIZE,
4975 DEFAULT_SAMPLE_UNIQUE_RATIO,
4976 )? {
4977 return None;
4978 }
4979
4980 let max_dict_entries = u32::try_from(threshold_cardinality.min(i32::MAX as u64)).ok()?;
4981 Some(DictEncodingBudget {
4982 max_dict_entries,
4983 max_encoded_size,
4984 })
4985 }
4986
4987 fn sample_is_near_unique(
4993 data_block: &DataBlock,
4994 max_samples: usize,
4995 unique_ratio_threshold: f64,
4996 ) -> Option<bool> {
4997 use std::collections::HashSet;
4998
4999 if unique_ratio_threshold <= 0.0 || unique_ratio_threshold > 1.0 {
5000 return None;
5001 }
5002
5003 let num_values = usize::try_from(data_block.num_values()).ok()?;
5004 if num_values == 0 {
5005 return Some(false);
5006 }
5007
5008 let sample_count = num_values.min(max_samples).max(1);
5009 let step = (num_values / sample_count).max(1);
5011
5012 match data_block {
5013 DataBlock::FixedWidth(fixed) => match fixed.bits_per_value {
5014 64 => {
5015 let values = fixed.data.borrow_to_typed_slice::<u64>();
5016 let values = values.as_ref();
5017 let mut unique: HashSet<u64> = HashSet::with_capacity(sample_count.min(1024));
5018 for idx in (0..num_values).step_by(step).take(sample_count) {
5019 unique.insert(values.get(idx).copied()?);
5020 }
5021 let ratio = unique.len() as f64 / sample_count as f64;
5022 Some(sample_count >= 1024 && ratio >= unique_ratio_threshold)
5024 }
5025 128 => {
5026 let values = fixed.data.borrow_to_typed_slice::<u128>();
5027 let values = values.as_ref();
5028 let mut unique: HashSet<u128> = HashSet::with_capacity(sample_count.min(1024));
5029 for idx in (0..num_values).step_by(step).take(sample_count) {
5030 unique.insert(values.get(idx).copied()?);
5031 }
5032 let ratio = unique.len() as f64 / sample_count as f64;
5033 Some(sample_count >= 1024 && ratio >= unique_ratio_threshold)
5034 }
5035 _ => Some(false),
5036 },
5037 DataBlock::VariableWidth(var) => {
5038 use xxhash_rust::xxh3::xxh3_64;
5039
5040 let mut unique: HashSet<u64> = HashSet::with_capacity(sample_count.min(1024));
5042 match var.bits_per_offset {
5043 32 => {
5044 let offsets_ref = var.offsets.borrow_to_typed_slice::<u32>();
5045 let offsets: &[u32] = offsets_ref.as_ref();
5046 for i in (0..num_values).step_by(step).take(sample_count) {
5047 let start = usize::try_from(*offsets.get(i)?).ok()?;
5048 let end = usize::try_from(*offsets.get(i + 1)?).ok()?;
5049 if start > end || end > var.data.len() {
5050 return None;
5051 }
5052 unique.insert(xxh3_64(&var.data[start..end]));
5053 }
5054 }
5055 64 => {
5056 let offsets_ref = var.offsets.borrow_to_typed_slice::<u64>();
5057 let offsets: &[u64] = offsets_ref.as_ref();
5058 for i in (0..num_values).step_by(step).take(sample_count) {
5059 let start = usize::try_from(*offsets.get(i)?).ok()?;
5060 let end = usize::try_from(*offsets.get(i + 1)?).ok()?;
5061 if start > end || end > var.data.len() {
5062 return None;
5063 }
5064 unique.insert(xxh3_64(&var.data[start..end]));
5065 }
5066 }
5067 _ => return Some(false),
5068 }
5069 let ratio = unique.len() as f64 / sample_count as f64;
5070 Some(sample_count >= 1024 && ratio >= unique_ratio_threshold)
5071 }
5072 _ => Some(false),
5073 }
5074 }
5075
5076 fn do_flush(
5078 &mut self,
5079 arrays: Vec<ArrayRef>,
5080 repdefs: Vec<RepDefBuilder>,
5081 row_number: u64,
5082 num_rows: u64,
5083 ) -> Result<Vec<EncodeTask>> {
5084 let column_idx = self.column_index;
5085 let compression_strategy = self.compression_strategy.clone();
5086 let field = self.field.clone();
5087 let encoding_metadata = self.encoding_metadata.clone();
5088 let support_large_chunk = self.support_large_chunk;
5089 let version = self.version;
5090 let task = spawn_cpu(move || {
5091 let num_values = arrays.iter().map(|arr| arr.len() as u64).sum();
5092 let is_simple_validity = repdefs.iter().all(|rd| rd.is_simple_validity());
5093 let has_repdef_info = repdefs.iter().any(|rd| !rd.is_empty());
5094 let repdef = RepDefBuilder::serialize(repdefs);
5095
5096 if num_values == 0 {
5097 log::debug!("Encoding column {} with {} items ({} rows) using complex-null layout", column_idx, num_values, num_rows);
5101 return Self::encode_complex_all_null(
5102 column_idx,
5103 repdef,
5104 row_number,
5105 num_rows,
5106 version,
5107 compression_strategy.as_ref(),
5108 );
5109 }
5110
5111 let leaf_validity = Self::leaf_validity(&repdef, num_values as usize)?;
5112 let all_null = leaf_validity
5113 .as_ref()
5114 .map(|validity| validity.count_set_bits() == 0)
5115 .unwrap_or(false);
5116
5117 if all_null {
5118 return if is_simple_validity {
5119 log::debug!(
5120 "Encoding column {} with {} items ({} rows) using simple-null layout",
5121 column_idx,
5122 num_values,
5123 num_rows
5124 );
5125 Self::encode_simple_all_null(column_idx, num_values, row_number)
5126 } else {
5127 log::debug!(
5128 "Encoding column {} with {} items ({} rows) using complex-null layout",
5129 column_idx,
5130 num_values,
5131 num_rows
5132 );
5133 Self::encode_complex_all_null(
5134 column_idx,
5135 repdef,
5136 row_number,
5137 num_rows,
5138 version,
5139 compression_strategy.as_ref(),
5140 )
5141 };
5142 }
5143
5144 if let DataType::Struct(fields) = &field.data_type()
5145 && fields.is_empty()
5146 {
5147 if has_repdef_info {
5148 return Err(Error::invalid_input_source(format!("Empty structs with rep/def information are not yet supported. The field {} is an empty struct that either has nulls or is in a list.", field.name).into()));
5149 }
5150 return Self::encode_simple_all_null(column_idx, num_values, row_number);
5153 }
5154
5155 let data_block = DataBlock::from_arrays(&arrays, num_values);
5156
5157 if version.resolve() >= LanceFileVersion::V2_2
5158 && let Some(scalar) = Self::find_constant_scalar(&arrays, leaf_validity.as_ref())?
5159 {
5160 log::debug!(
5161 "Encoding column {} with {} items ({} rows) using constant layout",
5162 column_idx,
5163 num_values,
5164 num_rows
5165 );
5166 return constant::encode_constant_page(
5167 column_idx,
5168 scalar,
5169 repdef,
5170 row_number,
5171 num_rows,
5172 );
5173 }
5174
5175 let requires_full_zip_packed_struct =
5176 if let DataBlock::Struct(ref struct_data_block) = data_block {
5177 struct_data_block.has_variable_width_child()
5178 } else {
5179 false
5180 };
5181
5182 if requires_full_zip_packed_struct {
5183 log::debug!(
5184 "Encoding column {} with {} items using full-zip packed struct layout",
5185 column_idx,
5186 num_values
5187 );
5188 return Self::encode_full_zip(
5189 column_idx,
5190 &field,
5191 compression_strategy.as_ref(),
5192 data_block,
5193 repdef,
5194 row_number,
5195 num_rows,
5196 );
5197 }
5198
5199 let too_sparse = Self::repdef_too_sparse_for_miniblock(&repdef, num_values);
5203
5204 if !too_sparse {
5205 if let DataBlock::Dictionary(dict) = data_block {
5206 log::debug!("Encoding column {} with {} items using dictionary encoding (already dictionary encoded)", column_idx, num_values);
5207 let (mut indices_data_block, dictionary_data_block) = dict.into_parts();
5208 indices_data_block.compute_stat();
5213 return Self::encode_miniblock(
5214 column_idx,
5215 &field,
5216 compression_strategy.as_ref(),
5217 indices_data_block,
5218 repdef,
5219 row_number,
5220 Some(dictionary_data_block),
5221 num_rows,
5222 support_large_chunk,
5223 );
5224 }
5225 } else {
5226 log::debug!(
5227 "Encoding column {} with {} items using full-zip layout \
5228 (rep/def too sparse for mini-block)",
5229 column_idx,
5230 num_values
5231 );
5232 }
5233
5234 {
5235 let dict_result = if too_sparse {
5238 None
5239 } else {
5240 Self::should_dictionary_encode(&data_block, &field, version)
5241 .and_then(|budget| {
5242 log::debug!(
5243 "Encoding column {} with {} items using dictionary encoding (mini-block layout)",
5244 column_idx,
5245 num_values
5246 );
5247 dict::dictionary_encode(
5248 &data_block,
5249 budget.max_dict_entries,
5250 budget.max_encoded_size,
5251 )
5252 })
5253 };
5254
5255 if let Some((indices_data_block, dictionary_data_block)) = dict_result {
5256 Self::encode_miniblock(
5257 column_idx,
5258 &field,
5259 compression_strategy.as_ref(),
5260 indices_data_block,
5261 repdef,
5262 row_number,
5263 Some(dictionary_data_block),
5264 num_rows,
5265 support_large_chunk,
5266 )
5267 } else if !too_sparse && Self::prefers_miniblock(&data_block, encoding_metadata.as_ref()) {
5268 log::debug!(
5269 "Encoding column {} with {} items using mini-block layout",
5270 column_idx,
5271 num_values
5272 );
5273 Self::encode_miniblock(
5274 column_idx,
5275 &field,
5276 compression_strategy.as_ref(),
5277 data_block,
5278 repdef,
5279 row_number,
5280 None,
5281 num_rows,
5282 support_large_chunk,
5283 )
5284 } else if too_sparse || Self::prefers_fullzip(encoding_metadata.as_ref()) {
5285 log::debug!(
5286 "Encoding column {} with {} items using full-zip layout",
5287 column_idx,
5288 num_values
5289 );
5290 Self::encode_full_zip(
5291 column_idx,
5292 &field,
5293 compression_strategy.as_ref(),
5294 data_block,
5295 repdef,
5296 row_number,
5297 num_rows,
5298 )
5299 } else {
5300 Err(Error::invalid_input_source(format!("Cannot determine structural encoding for field {}. This typically indicates an invalid value of the field metadata key {}", field.name, STRUCTURAL_ENCODING_META_KEY).into()))
5301 }
5302 }
5303 })
5304 .boxed();
5305 Ok(vec![task])
5306 }
5307
5308 fn extract_validity_buf(
5309 array: Arc<dyn Array>,
5310 repdef: &mut RepDefBuilder,
5311 keep_original_array: bool,
5312 ) -> Result<Arc<dyn Array>> {
5313 if let Some(validity) = array.nulls() {
5314 if keep_original_array {
5315 repdef.add_validity_bitmap(validity.clone());
5316 } else {
5317 repdef.add_validity_bitmap(deep_copy_nulls(Some(validity)).unwrap());
5318 }
5319 let data_no_nulls = array.to_data().into_builder().nulls(None).build()?;
5320 Ok(make_array(data_no_nulls))
5321 } else {
5322 repdef.add_no_null(array.len());
5323 Ok(array)
5324 }
5325 }
5326
5327 fn extract_validity(
5328 mut array: Arc<dyn Array>,
5329 repdef: &mut RepDefBuilder,
5330 keep_original_array: bool,
5331 ) -> Result<Arc<dyn Array>> {
5332 match array.data_type() {
5333 DataType::Null => {
5334 repdef.add_validity_bitmap(NullBuffer::new(BooleanBuffer::new_unset(array.len())));
5335 Ok(array)
5336 }
5337 DataType::Dictionary(_, _) => {
5338 array = dict::normalize_dict_nulls(array)?;
5339 Self::extract_validity_buf(array, repdef, keep_original_array)
5340 }
5341 _ => Self::extract_validity_buf(array, repdef, keep_original_array),
5350 }
5351 }
5352}
5353
5354impl FieldEncoder for PrimitiveStructuralEncoder {
5355 fn maybe_encode(
5357 &mut self,
5358 array: ArrayRef,
5359 _external_buffers: &mut OutOfLineBuffers,
5360 mut repdef: RepDefBuilder,
5361 row_number: u64,
5362 num_rows: u64,
5363 ) -> Result<Vec<EncodeTask>> {
5364 let array = Self::extract_validity(array, &mut repdef, self.keep_original_array)?;
5365 self.accumulated_repdefs.push(repdef);
5366
5367 if let Some((arrays, row_number, num_rows)) =
5368 self.accumulation_queue.insert(array, row_number, num_rows)
5369 {
5370 let accumulated_repdefs = std::mem::take(&mut self.accumulated_repdefs);
5371 Ok(self.do_flush(arrays, accumulated_repdefs, row_number, num_rows)?)
5372 } else {
5373 Ok(vec![])
5374 }
5375 }
5376
5377 fn flush(&mut self, _external_buffers: &mut OutOfLineBuffers) -> Result<Vec<EncodeTask>> {
5379 if let Some((arrays, row_number, num_rows)) = self.accumulation_queue.flush() {
5380 let accumulated_repdefs = std::mem::take(&mut self.accumulated_repdefs);
5381 Ok(self.do_flush(arrays, accumulated_repdefs, row_number, num_rows)?)
5382 } else {
5383 Ok(vec![])
5384 }
5385 }
5386
5387 fn num_columns(&self) -> u32 {
5388 1
5389 }
5390
5391 fn finish(
5392 &mut self,
5393 _external_buffers: &mut OutOfLineBuffers,
5394 ) -> BoxFuture<'_, Result<Vec<crate::encoder::EncodedColumn>>> {
5395 std::future::ready(Ok(vec![EncodedColumn::default()])).boxed()
5396 }
5397}
5398
5399#[cfg(test)]
5400#[allow(clippy::single_range_in_vec_init)]
5401mod tests {
5402 use super::{
5403 ChunkInstructions, DataBlock, DecodeMiniBlockTask, FixedPerValueDecompressor,
5404 FixedWidthDataBlock, FullZipCacheableState, FullZipDecodeDetails, FullZipReadSource,
5405 FullZipRepIndexDetails, FullZipScheduler, MiniBlockRepIndex, PerValueDecompressor,
5406 PreambleAction, StructuralPageScheduler, VariableFullZipDecoder,
5407 };
5408 use crate::buffer::LanceBuffer;
5409 use crate::compression::DefaultDecompressionStrategy;
5410 use crate::constants::{
5411 COMPRESSION_LEVEL_META_KEY, COMPRESSION_META_KEY, DICT_VALUES_COMPRESSION_LEVEL_META_KEY,
5412 DICT_VALUES_COMPRESSION_META_KEY, STRUCTURAL_ENCODING_META_KEY,
5413 STRUCTURAL_ENCODING_MINIBLOCK,
5414 };
5415 use crate::data::BlockInfo;
5416 use crate::decoder::PageEncoding;
5417 use crate::encodings::logical::primitive::{
5418 ChunkDrainInstructions, PrimitiveStructuralEncoder,
5419 };
5420 use crate::format::ProtobufUtils21;
5421 use crate::format::pb21;
5422 use crate::format::pb21::compressive_encoding::Compression;
5423 use crate::testing::{TestCases, check_round_trip_encoding_of_data};
5424 use crate::version::LanceFileVersion;
5425 use arrow_array::{ArrayRef, Int8Array, StringArray};
5426 use arrow_schema::DataType;
5427 use std::collections::HashMap;
5428 use std::{collections::VecDeque, sync::Arc};
5429
5430 #[test]
5431 fn test_is_narrow() {
5432 let int8_array = Int8Array::from(vec![1, 2, 3]);
5433 let array_ref: ArrayRef = Arc::new(int8_array);
5434 let block = DataBlock::from_array(array_ref);
5435
5436 assert!(PrimitiveStructuralEncoder::is_narrow(&block));
5437
5438 let string_array = StringArray::from(vec![Some("hello"), Some("world")]);
5439 let block = DataBlock::from_array(string_array);
5440 assert!(PrimitiveStructuralEncoder::is_narrow(&block));
5441
5442 let string_array = StringArray::from(vec![
5443 Some("hello world".repeat(100)),
5444 Some("world".to_string()),
5445 ]);
5446 let block = DataBlock::from_array(string_array);
5447 assert!((!PrimitiveStructuralEncoder::is_narrow(&block)));
5448 }
5449
5450 #[test]
5451 fn test_map_range() {
5452 let rep = Some(vec![1, 0, 0, 1, 0, 1, 1, 0, 0]);
5455 let def = Some(vec![0, 0, 0, 0, 0, 1, 0, 0, 0]);
5456 let max_visible_def = 0;
5457 let total_items = 8;
5458 let max_rep = 1;
5459
5460 let check = |range, expected_item_range, expected_level_range| {
5461 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5462 range,
5463 rep.as_ref(),
5464 def.as_ref(),
5465 max_rep,
5466 max_visible_def,
5467 total_items,
5468 PreambleAction::Absent,
5469 );
5470 assert_eq!(item_range, expected_item_range);
5471 assert_eq!(level_range, expected_level_range);
5472 };
5473
5474 check(0..1, 0..3, 0..3);
5475 check(1..2, 3..5, 3..5);
5476 check(2..3, 5..5, 5..6);
5477 check(3..4, 5..8, 6..9);
5478 check(0..2, 0..5, 0..5);
5479 check(1..3, 3..5, 3..6);
5480 check(2..4, 5..8, 5..9);
5481 check(0..3, 0..5, 0..6);
5482 check(1..4, 3..8, 3..9);
5483 check(0..4, 0..8, 0..9);
5484
5485 let rep = Some(vec![1, 1, 0, 1]);
5488 let def = Some(vec![1, 0, 0, 0]);
5489 let max_visible_def = 0;
5490 let total_items = 3;
5491
5492 let check = |range, expected_item_range, expected_level_range| {
5493 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5494 range,
5495 rep.as_ref(),
5496 def.as_ref(),
5497 max_rep,
5498 max_visible_def,
5499 total_items,
5500 PreambleAction::Absent,
5501 );
5502 assert_eq!(item_range, expected_item_range);
5503 assert_eq!(level_range, expected_level_range);
5504 };
5505
5506 check(0..1, 0..0, 0..1);
5507 check(1..2, 0..2, 1..3);
5508 check(2..3, 2..3, 3..4);
5509 check(0..2, 0..2, 0..3);
5510 check(1..3, 0..3, 1..4);
5511 check(0..3, 0..3, 0..4);
5512
5513 let rep = Some(vec![1, 1, 0, 1]);
5516 let def = Some(vec![0, 0, 0, 1]);
5517 let max_visible_def = 0;
5518 let total_items = 3;
5519
5520 let check = |range, expected_item_range, expected_level_range| {
5521 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5522 range,
5523 rep.as_ref(),
5524 def.as_ref(),
5525 max_rep,
5526 max_visible_def,
5527 total_items,
5528 PreambleAction::Absent,
5529 );
5530 assert_eq!(item_range, expected_item_range);
5531 assert_eq!(level_range, expected_level_range);
5532 };
5533
5534 check(0..1, 0..1, 0..1);
5535 check(1..2, 1..3, 1..3);
5536 check(2..3, 3..3, 3..4);
5537 check(0..2, 0..3, 0..3);
5538 check(1..3, 1..3, 1..4);
5539 check(0..3, 0..3, 0..4);
5540
5541 let rep = Some(vec![1, 0, 1, 0, 1, 0]);
5544 let def: Option<&[u16]> = None;
5545 let max_visible_def = 0;
5546 let total_items = 6;
5547
5548 let check = |range, expected_item_range, expected_level_range| {
5549 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5550 range,
5551 rep.as_ref(),
5552 def.as_ref(),
5553 max_rep,
5554 max_visible_def,
5555 total_items,
5556 PreambleAction::Absent,
5557 );
5558 assert_eq!(item_range, expected_item_range);
5559 assert_eq!(level_range, expected_level_range);
5560 };
5561
5562 check(0..1, 0..2, 0..2);
5563 check(1..2, 2..4, 2..4);
5564 check(2..3, 4..6, 4..6);
5565 check(0..2, 0..4, 0..4);
5566 check(1..3, 2..6, 2..6);
5567 check(0..3, 0..6, 0..6);
5568
5569 let rep: Option<&[u16]> = None;
5572 let def = Some(vec![0, 0, 1, 0]);
5573 let max_visible_def = 1;
5574 let total_items = 4;
5575
5576 let check = |range, expected_item_range, expected_level_range| {
5577 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5578 range,
5579 rep.as_ref(),
5580 def.as_ref(),
5581 max_rep,
5582 max_visible_def,
5583 total_items,
5584 PreambleAction::Absent,
5585 );
5586 assert_eq!(item_range, expected_item_range);
5587 assert_eq!(level_range, expected_level_range);
5588 };
5589
5590 check(0..1, 0..1, 0..1);
5591 check(1..2, 1..2, 1..2);
5592 check(2..3, 2..3, 2..3);
5593 check(0..2, 0..2, 0..2);
5594 check(1..3, 1..3, 1..3);
5595 check(0..3, 0..3, 0..3);
5596
5597 let rep = Some(vec![0, 1, 0, 1]);
5602 let def = Some(vec![0, 0, 0, 1]);
5603 let max_visible_def = 0;
5604 let total_items = 3;
5605
5606 let check = |range, expected_item_range, expected_level_range| {
5607 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5608 range,
5609 rep.as_ref(),
5610 def.as_ref(),
5611 max_rep,
5612 max_visible_def,
5613 total_items,
5614 PreambleAction::Take,
5615 );
5616 assert_eq!(item_range, expected_item_range);
5617 assert_eq!(level_range, expected_level_range);
5618 };
5619
5620 check(0..1, 0..3, 0..3);
5622 check(0..2, 0..3, 0..4);
5623
5624 let check = |range, expected_item_range, expected_level_range| {
5625 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5626 range,
5627 rep.as_ref(),
5628 def.as_ref(),
5629 max_rep,
5630 max_visible_def,
5631 total_items,
5632 PreambleAction::Skip,
5633 );
5634 assert_eq!(item_range, expected_item_range);
5635 assert_eq!(level_range, expected_level_range);
5636 };
5637
5638 check(0..1, 1..3, 1..3);
5639 check(1..2, 3..3, 3..4);
5640 check(0..2, 1..3, 1..4);
5641
5642 let rep = Some(vec![0, 1, 1, 0]);
5647 let def = Some(vec![0, 1, 0, 0]);
5648 let max_visible_def = 0;
5649 let total_items = 4;
5650
5651 let check = |range, expected_item_range, expected_level_range| {
5652 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5653 range,
5654 rep.as_ref(),
5655 def.as_ref(),
5656 max_rep,
5657 max_visible_def,
5658 total_items,
5659 PreambleAction::Take,
5660 );
5661 assert_eq!(item_range, expected_item_range);
5662 assert_eq!(level_range, expected_level_range);
5663 };
5664
5665 check(0..1, 0..1, 0..2);
5667 check(0..2, 0..3, 0..4);
5668
5669 let check = |range, expected_item_range, expected_level_range| {
5670 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5671 range,
5672 rep.as_ref(),
5673 def.as_ref(),
5674 max_rep,
5675 max_visible_def,
5676 total_items,
5677 PreambleAction::Skip,
5678 );
5679 assert_eq!(item_range, expected_item_range);
5680 assert_eq!(level_range, expected_level_range);
5681 };
5682
5683 check(0..1, 1..1, 1..2);
5685 check(1..2, 1..3, 2..4);
5686 check(0..2, 1..3, 1..4);
5687
5688 let rep = Some(vec![0, 1, 0, 1]);
5691 let def: Option<Vec<u16>> = None;
5692 let max_visible_def = 0;
5693 let total_items = 4;
5694
5695 let check = |range, expected_item_range, expected_level_range| {
5696 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5697 range,
5698 rep.as_ref(),
5699 def.as_ref(),
5700 max_rep,
5701 max_visible_def,
5702 total_items,
5703 PreambleAction::Take,
5704 );
5705 assert_eq!(item_range, expected_item_range);
5706 assert_eq!(level_range, expected_level_range);
5707 };
5708
5709 check(0..1, 0..3, 0..3);
5711 check(0..2, 0..4, 0..4);
5712
5713 let check = |range, expected_item_range, expected_level_range| {
5714 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5715 range,
5716 rep.as_ref(),
5717 def.as_ref(),
5718 max_rep,
5719 max_visible_def,
5720 total_items,
5721 PreambleAction::Skip,
5722 );
5723 assert_eq!(item_range, expected_item_range);
5724 assert_eq!(level_range, expected_level_range);
5725 };
5726
5727 check(0..1, 1..3, 1..3);
5728 check(1..2, 3..4, 3..4);
5729 check(0..2, 1..4, 1..4);
5730
5731 let rep = Some(vec![2, 1, 2, 0, 1, 2]);
5735 let def = Some(vec![0, 1, 2, 0, 0, 0]);
5736 let max_rep = 2;
5737 let max_visible_def = 0;
5738 let total_items = 4;
5739
5740 let check = |range, expected_item_range, expected_level_range| {
5741 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5742 range,
5743 rep.as_ref(),
5744 def.as_ref(),
5745 max_rep,
5746 max_visible_def,
5747 total_items,
5748 PreambleAction::Absent,
5749 );
5750 assert_eq!(item_range, expected_item_range);
5751 assert_eq!(level_range, expected_level_range);
5752 };
5753
5754 check(0..3, 0..4, 0..6);
5755 check(0..1, 0..1, 0..2);
5756 check(1..2, 1..3, 2..5);
5757 check(2..3, 3..4, 5..6);
5758
5759 let rep = Some(vec![0, 0, 1, 0, 1, 1]);
5761 let def = Some(vec![0, 1, 0, 0, 0, 0]);
5762 let max_rep = 1;
5763 let max_visible_def = 0;
5764 let total_items = 5;
5765
5766 let check = |range, expected_item_range, expected_level_range| {
5767 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5768 range,
5769 rep.as_ref(),
5770 def.as_ref(),
5771 max_rep,
5772 max_visible_def,
5773 total_items,
5774 PreambleAction::Take,
5775 );
5776 assert_eq!(item_range, expected_item_range);
5777 assert_eq!(level_range, expected_level_range);
5778 };
5779
5780 check(0..0, 0..1, 0..2);
5781 check(0..1, 0..3, 0..4);
5782 check(0..2, 0..4, 0..5);
5783
5784 let rep = Some(vec![0, 1, 0, 1, 0, 1, 0, 1]);
5787 let def = Some(vec![1, 0, 1, 1, 0, 0, 0, 0]);
5788 let max_rep = 1;
5789 let max_visible_def = 0;
5790 let total_items = 5;
5791
5792 let check = |range, expected_item_range, expected_level_range| {
5793 let (item_range, level_range) = DecodeMiniBlockTask::map_range(
5794 range,
5795 rep.as_ref(),
5796 def.as_ref(),
5797 max_rep,
5798 max_visible_def,
5799 total_items,
5800 PreambleAction::Skip,
5801 );
5802 assert_eq!(item_range, expected_item_range);
5803 assert_eq!(level_range, expected_level_range);
5804 };
5805
5806 check(2..3, 2..4, 5..7);
5807 }
5808
5809 #[test]
5810 fn test_slice_batch_data_and_rebase_offsets_u32() {
5811 let data = LanceBuffer::copy_slice(b"0123456789abcdefghij");
5812 let offsets = LanceBuffer::reinterpret_vec(vec![6_u32, 8_u32, 8_u32, 12_u32]);
5813
5814 let (sliced_data, normalized_offsets) =
5815 VariableFullZipDecoder::slice_batch_data_and_rebase_offsets(&data, &offsets, 32)
5816 .unwrap();
5817
5818 assert_eq!(sliced_data.as_ref(), b"6789ab");
5819 let normalized = normalized_offsets.borrow_to_typed_slice::<u32>();
5820 assert_eq!(normalized.as_ref(), &[0, 2, 2, 6]);
5821 }
5822
5823 #[test]
5824 fn test_slice_batch_data_and_rebase_offsets_u64() {
5825 let data = LanceBuffer::copy_slice(b"abcdefghijklmnopqrstuvwxyz");
5826 let offsets = LanceBuffer::reinterpret_vec(vec![10_u64, 12_u64, 16_u64, 20_u64]);
5827
5828 let (sliced_data, normalized_offsets) =
5829 VariableFullZipDecoder::slice_batch_data_and_rebase_offsets(&data, &offsets, 64)
5830 .unwrap();
5831
5832 assert_eq!(sliced_data.as_ref(), b"klmnopqrst");
5833 let normalized = normalized_offsets.borrow_to_typed_slice::<u64>();
5834 assert_eq!(normalized.as_ref(), &[0, 2, 6, 10]);
5835 }
5836
5837 #[test]
5838 fn test_slice_batch_data_and_rebase_offsets_rejects_invalid_offsets() {
5839 let data = LanceBuffer::copy_slice(b"abcd");
5840 let offsets = LanceBuffer::reinterpret_vec(vec![3_u32, 2_u32]);
5841
5842 let err = VariableFullZipDecoder::slice_batch_data_and_rebase_offsets(&data, &offsets, 32)
5843 .expect_err("offset end before start should error");
5844 assert!(err.to_string().contains("less than base"));
5845 }
5846
5847 #[test]
5848 fn test_schedule_instructions() {
5849 let rep_data: Vec<u64> = vec![5, 2, 3, 0, 4, 7, 2, 0];
5851 let rep_bytes: Vec<u8> = rep_data.iter().flat_map(|v| v.to_le_bytes()).collect();
5852 let repetition_index = MiniBlockRepIndex::decode_from_bytes(&rep_bytes, 2);
5853
5854 let check = |user_ranges, expected_instructions| {
5855 let instructions =
5856 ChunkInstructions::schedule_instructions(&repetition_index, user_ranges);
5857 assert_eq!(instructions, expected_instructions);
5858 };
5859
5860 let expected_take_all = vec![
5862 ChunkInstructions {
5863 chunk_idx: 0,
5864 preamble: PreambleAction::Absent,
5865 rows_to_skip: 0,
5866 rows_to_take: 6,
5867 take_trailer: true,
5868 },
5869 ChunkInstructions {
5870 chunk_idx: 1,
5871 preamble: PreambleAction::Take,
5872 rows_to_skip: 0,
5873 rows_to_take: 2,
5874 take_trailer: false,
5875 },
5876 ChunkInstructions {
5877 chunk_idx: 2,
5878 preamble: PreambleAction::Absent,
5879 rows_to_skip: 0,
5880 rows_to_take: 5,
5881 take_trailer: true,
5882 },
5883 ChunkInstructions {
5884 chunk_idx: 3,
5885 preamble: PreambleAction::Take,
5886 rows_to_skip: 0,
5887 rows_to_take: 1,
5888 take_trailer: false,
5889 },
5890 ];
5891
5892 check(&[0..14], expected_take_all.clone());
5894
5895 check(
5897 &[
5898 0..1,
5899 1..2,
5900 2..3,
5901 3..4,
5902 4..5,
5903 5..6,
5904 6..7,
5905 7..8,
5906 8..9,
5907 9..10,
5908 10..11,
5909 11..12,
5910 12..13,
5911 13..14,
5912 ],
5913 expected_take_all,
5914 );
5915
5916 check(
5920 &[0..1, 3..4],
5921 vec![
5922 ChunkInstructions {
5923 chunk_idx: 0,
5924 preamble: PreambleAction::Absent,
5925 rows_to_skip: 0,
5926 rows_to_take: 1,
5927 take_trailer: false,
5928 },
5929 ChunkInstructions {
5930 chunk_idx: 0,
5931 preamble: PreambleAction::Absent,
5932 rows_to_skip: 3,
5933 rows_to_take: 1,
5934 take_trailer: false,
5935 },
5936 ],
5937 );
5938
5939 check(
5941 &[5..6],
5942 vec![
5943 ChunkInstructions {
5944 chunk_idx: 0,
5945 preamble: PreambleAction::Absent,
5946 rows_to_skip: 5,
5947 rows_to_take: 1,
5948 take_trailer: true,
5949 },
5950 ChunkInstructions {
5951 chunk_idx: 1,
5952 preamble: PreambleAction::Take,
5953 rows_to_skip: 0,
5954 rows_to_take: 0,
5955 take_trailer: false,
5956 },
5957 ],
5958 );
5959
5960 check(
5962 &[7..10],
5963 vec![
5964 ChunkInstructions {
5965 chunk_idx: 1,
5966 preamble: PreambleAction::Skip,
5967 rows_to_skip: 1,
5968 rows_to_take: 1,
5969 take_trailer: false,
5970 },
5971 ChunkInstructions {
5972 chunk_idx: 2,
5973 preamble: PreambleAction::Absent,
5974 rows_to_skip: 0,
5975 rows_to_take: 2,
5976 take_trailer: false,
5977 },
5978 ],
5979 );
5980 }
5981
5982 #[test]
5983 fn test_drain_instructions() {
5984 fn drain_from_instructions(
5985 instructions: &mut VecDeque<ChunkInstructions>,
5986 mut rows_desired: u64,
5987 need_preamble: &mut bool,
5988 skip_in_chunk: &mut u64,
5989 ) -> Vec<ChunkDrainInstructions> {
5990 let mut drain_instructions = Vec::with_capacity(instructions.len());
5992 while rows_desired > 0 || *need_preamble {
5993 let (next_instructions, consumed_chunk) = instructions
5994 .front()
5995 .unwrap()
5996 .drain_from_instruction(&mut rows_desired, need_preamble, skip_in_chunk);
5997 if consumed_chunk {
5998 instructions.pop_front();
5999 }
6000 drain_instructions.push(next_instructions);
6001 }
6002 drain_instructions
6003 }
6004
6005 let rep_data: Vec<u64> = vec![5, 2, 3, 0, 4, 7, 2, 0];
6007 let rep_bytes: Vec<u8> = rep_data.iter().flat_map(|v| v.to_le_bytes()).collect();
6008 let repetition_index = MiniBlockRepIndex::decode_from_bytes(&rep_bytes, 2);
6009 let user_ranges = vec![1..7, 10..14];
6010
6011 let scheduled = ChunkInstructions::schedule_instructions(&repetition_index, &user_ranges);
6013
6014 let mut to_drain = VecDeque::from(scheduled.clone());
6015
6016 let mut need_preamble = false;
6019 let mut skip_in_chunk = 0;
6020
6021 let next_batch =
6022 drain_from_instructions(&mut to_drain, 4, &mut need_preamble, &mut skip_in_chunk);
6023
6024 assert!(!need_preamble);
6025 assert_eq!(skip_in_chunk, 4);
6026 assert_eq!(
6027 next_batch,
6028 vec![ChunkDrainInstructions {
6029 chunk_instructions: scheduled[0].clone(),
6030 rows_to_take: 4,
6031 rows_to_skip: 0,
6032 preamble_action: PreambleAction::Absent,
6033 }]
6034 );
6035
6036 let next_batch =
6037 drain_from_instructions(&mut to_drain, 4, &mut need_preamble, &mut skip_in_chunk);
6038
6039 assert!(!need_preamble);
6040 assert_eq!(skip_in_chunk, 2);
6041
6042 assert_eq!(
6043 next_batch,
6044 vec![
6045 ChunkDrainInstructions {
6046 chunk_instructions: scheduled[0].clone(),
6047 rows_to_take: 1,
6048 rows_to_skip: 4,
6049 preamble_action: PreambleAction::Absent,
6050 },
6051 ChunkDrainInstructions {
6052 chunk_instructions: scheduled[1].clone(),
6053 rows_to_take: 1,
6054 rows_to_skip: 0,
6055 preamble_action: PreambleAction::Take,
6056 },
6057 ChunkDrainInstructions {
6058 chunk_instructions: scheduled[2].clone(),
6059 rows_to_take: 2,
6060 rows_to_skip: 0,
6061 preamble_action: PreambleAction::Absent,
6062 }
6063 ]
6064 );
6065
6066 let next_batch =
6067 drain_from_instructions(&mut to_drain, 2, &mut need_preamble, &mut skip_in_chunk);
6068
6069 assert!(!need_preamble);
6070 assert_eq!(skip_in_chunk, 0);
6071
6072 assert_eq!(
6073 next_batch,
6074 vec![
6075 ChunkDrainInstructions {
6076 chunk_instructions: scheduled[2].clone(),
6077 rows_to_take: 1,
6078 rows_to_skip: 2,
6079 preamble_action: PreambleAction::Absent,
6080 },
6081 ChunkDrainInstructions {
6082 chunk_instructions: scheduled[3].clone(),
6083 rows_to_take: 1,
6084 rows_to_skip: 0,
6085 preamble_action: PreambleAction::Take,
6086 },
6087 ]
6088 );
6089
6090 let rep_data: Vec<u64> = vec![5, 2, 3, 3, 20, 0];
6092 let rep_bytes: Vec<u8> = rep_data.iter().flat_map(|v| v.to_le_bytes()).collect();
6093 let repetition_index = MiniBlockRepIndex::decode_from_bytes(&rep_bytes, 2);
6094 let user_ranges = vec![0..28];
6095
6096 let scheduled = ChunkInstructions::schedule_instructions(&repetition_index, &user_ranges);
6098
6099 let mut to_drain = VecDeque::from(scheduled.clone());
6100
6101 let mut need_preamble = false;
6104 let mut skip_in_chunk = 0;
6105
6106 let next_batch =
6107 drain_from_instructions(&mut to_drain, 7, &mut need_preamble, &mut skip_in_chunk);
6108
6109 assert_eq!(
6110 next_batch,
6111 vec![
6112 ChunkDrainInstructions {
6113 chunk_instructions: scheduled[0].clone(),
6114 rows_to_take: 6,
6115 rows_to_skip: 0,
6116 preamble_action: PreambleAction::Absent,
6117 },
6118 ChunkDrainInstructions {
6119 chunk_instructions: scheduled[1].clone(),
6120 rows_to_take: 1,
6121 rows_to_skip: 0,
6122 preamble_action: PreambleAction::Take,
6123 },
6124 ]
6125 );
6126
6127 assert!(!need_preamble);
6128 assert_eq!(skip_in_chunk, 1);
6129
6130 let next_batch =
6133 drain_from_instructions(&mut to_drain, 2, &mut need_preamble, &mut skip_in_chunk);
6134
6135 assert_eq!(
6136 next_batch,
6137 vec![
6138 ChunkDrainInstructions {
6139 chunk_instructions: scheduled[1].clone(),
6140 rows_to_take: 2,
6141 rows_to_skip: 1,
6142 preamble_action: PreambleAction::Skip,
6143 },
6144 ChunkDrainInstructions {
6145 chunk_instructions: scheduled[2].clone(),
6146 rows_to_take: 0,
6147 rows_to_skip: 0,
6148 preamble_action: PreambleAction::Take,
6149 },
6150 ]
6151 );
6152
6153 assert!(!need_preamble);
6154 assert_eq!(skip_in_chunk, 0);
6155 }
6156
6157 #[tokio::test]
6158 async fn test_fullzip_initialize_is_lazy() {
6159 use futures::{FutureExt, future::BoxFuture};
6160 use std::ops::Range;
6161 use std::sync::Mutex;
6162
6163 #[derive(Debug, Clone)]
6164 struct RecordingScheduler {
6165 data: bytes::Bytes,
6166 requests: Arc<Mutex<Vec<Vec<Range<u64>>>>>,
6167 }
6168
6169 impl RecordingScheduler {
6170 fn new(data: bytes::Bytes) -> Self {
6171 Self {
6172 data,
6173 requests: Arc::new(Mutex::new(Vec::new())),
6174 }
6175 }
6176
6177 fn requests(&self) -> Vec<Vec<Range<u64>>> {
6178 self.requests.lock().unwrap().clone()
6179 }
6180 }
6181
6182 impl crate::EncodingsIo for RecordingScheduler {
6183 fn submit_request(
6184 &self,
6185 ranges: Vec<Range<u64>>,
6186 _priority: u64,
6187 ) -> BoxFuture<'static, crate::Result<Vec<bytes::Bytes>>> {
6188 self.requests.lock().unwrap().push(ranges.clone());
6189 let data = ranges
6190 .into_iter()
6191 .map(|range| self.data.slice(range.start as usize..range.end as usize))
6192 .collect::<Vec<_>>();
6193 std::future::ready(Ok(data)).boxed()
6194 }
6195 }
6196
6197 #[derive(Debug)]
6198 struct TestFixedDecompressor;
6199
6200 impl FixedPerValueDecompressor for TestFixedDecompressor {
6201 fn decompress(
6202 &self,
6203 _data: FixedWidthDataBlock,
6204 _num_rows: u64,
6205 ) -> crate::Result<DataBlock> {
6206 unimplemented!("Test decompressor")
6207 }
6208
6209 fn bits_per_value(&self) -> u64 {
6210 32
6211 }
6212 }
6213
6214 let io = Arc::new(RecordingScheduler::new(bytes::Bytes::from(vec![
6215 0;
6216 16 * 1024
6217 ])));
6218 let mut scheduler = FullZipScheduler {
6219 data_buf_position: 0,
6220 data_buf_size: 4096,
6221 rep_index: Some(FullZipRepIndexDetails {
6222 buf_position: 1000,
6223 bytes_per_value: 4,
6224 }),
6225 priority: 0,
6226 rows_in_page: 100,
6227 bits_per_offset: 32,
6228 details: Arc::new(FullZipDecodeDetails {
6229 value_decompressor: PerValueDecompressor::Fixed(Arc::new(TestFixedDecompressor)),
6230 def_meaning: Arc::new([crate::repdef::DefinitionInterpretation::NullableItem]),
6231 ctrl_word_parser: crate::repdef::ControlWordParser::new(0, 1),
6232 max_rep: 0,
6233 max_visible_def: 0,
6234 }),
6235 cached_state: None,
6236 enable_cache: false,
6237 };
6238
6239 let io_dyn: Arc<dyn crate::EncodingsIo> = io.clone();
6240 let cached_data = scheduler.initialize(&io_dyn).await.unwrap();
6241
6242 assert!(
6243 cached_data
6244 .as_arc_any()
6245 .downcast_ref::<super::NoCachedPageData>()
6246 .is_some(),
6247 "FullZip initialize should not eagerly load repetition index data"
6248 );
6249 assert!(scheduler.cached_state.is_none());
6250 assert!(
6251 io.requests().is_empty(),
6252 "FullZip initialize should not issue any I/O"
6253 );
6254 }
6255
6256 #[tokio::test]
6257 async fn test_fullzip_read_source_slices_prefetched_page() {
6258 let page_start = 200_u64;
6259 let page_data = LanceBuffer::copy_slice(&[0, 1, 2, 3, 4, 5, 6, 7]);
6260 let source = FullZipReadSource::PrefetchedPage {
6261 base_offset: page_start,
6262 data: page_data,
6263 };
6264 let ranges = vec![
6265 page_start..(page_start + 3),
6266 (page_start + 4)..(page_start + 8),
6267 ];
6268 let mut data = source.fetch(&ranges, 0).await.unwrap();
6269 assert_eq!(data.pop_front().unwrap().as_ref(), &[0, 1, 2]);
6270 assert_eq!(data.pop_front().unwrap().as_ref(), &[4, 5, 6, 7]);
6271 }
6272
6273 #[tokio::test]
6274 async fn test_fullzip_initialize_caches_rep_index_when_enabled() {
6275 use futures::{FutureExt, future::BoxFuture};
6276 use std::ops::Range;
6277 use std::sync::Mutex;
6278
6279 #[derive(Debug, Clone)]
6280 struct RecordingScheduler {
6281 data: bytes::Bytes,
6282 requests: Arc<Mutex<Vec<Vec<Range<u64>>>>>,
6283 }
6284
6285 impl RecordingScheduler {
6286 fn new(data: bytes::Bytes) -> Self {
6287 Self {
6288 data,
6289 requests: Arc::new(Mutex::new(Vec::new())),
6290 }
6291 }
6292
6293 fn requests(&self) -> Vec<Vec<Range<u64>>> {
6294 self.requests.lock().unwrap().clone()
6295 }
6296 }
6297
6298 impl crate::EncodingsIo for RecordingScheduler {
6299 fn submit_request(
6300 &self,
6301 ranges: Vec<Range<u64>>,
6302 _priority: u64,
6303 ) -> BoxFuture<'static, crate::Result<Vec<bytes::Bytes>>> {
6304 self.requests.lock().unwrap().push(ranges.clone());
6305 let data = ranges
6306 .into_iter()
6307 .map(|range| self.data.slice(range.start as usize..range.end as usize))
6308 .collect::<Vec<_>>();
6309 std::future::ready(Ok(data)).boxed()
6310 }
6311 }
6312
6313 #[derive(Debug)]
6314 struct TestFixedDecompressor;
6315
6316 impl FixedPerValueDecompressor for TestFixedDecompressor {
6317 fn decompress(
6318 &self,
6319 _data: FixedWidthDataBlock,
6320 _num_rows: u64,
6321 ) -> crate::Result<DataBlock> {
6322 unimplemented!("Test decompressor")
6323 }
6324
6325 fn bits_per_value(&self) -> u64 {
6326 32
6327 }
6328 }
6329
6330 let rows_in_page = 100_u64;
6331 let bytes_per_value = 4_u64;
6332 let rep_start = 1000_u64;
6333 let rep_size = ((rows_in_page + 1) * bytes_per_value) as usize;
6334 let mut data = vec![0_u8; 16 * 1024];
6335 data[rep_start as usize..rep_start as usize + rep_size].fill(7);
6336 let io = Arc::new(RecordingScheduler::new(bytes::Bytes::from(data)));
6337
6338 let mut scheduler = FullZipScheduler {
6339 data_buf_position: 0,
6340 data_buf_size: 4096,
6341 rep_index: Some(FullZipRepIndexDetails {
6342 buf_position: rep_start,
6343 bytes_per_value,
6344 }),
6345 priority: 0,
6346 rows_in_page,
6347 bits_per_offset: 32,
6348 details: Arc::new(FullZipDecodeDetails {
6349 value_decompressor: PerValueDecompressor::Fixed(Arc::new(TestFixedDecompressor)),
6350 def_meaning: Arc::new([crate::repdef::DefinitionInterpretation::NullableItem]),
6351 ctrl_word_parser: crate::repdef::ControlWordParser::new(0, 1),
6352 max_rep: 0,
6353 max_visible_def: 0,
6354 }),
6355 cached_state: None,
6356 enable_cache: true,
6357 };
6358
6359 let io_dyn: Arc<dyn crate::EncodingsIo> = io.clone();
6360 let cached_data = scheduler.initialize(&io_dyn).await.unwrap();
6361 assert!(
6362 cached_data
6363 .as_arc_any()
6364 .downcast_ref::<FullZipCacheableState>()
6365 .is_some()
6366 );
6367 assert!(scheduler.cached_state.is_some());
6368 assert_eq!(
6369 io.requests(),
6370 vec![vec![
6371 rep_start..(rep_start + (rows_in_page + 1) * bytes_per_value)
6372 ]]
6373 );
6374 }
6375
6376 #[tokio::test]
6377 async fn test_fullzip_full_page_bypasses_rep_index_io() {
6378 use futures::{FutureExt, future::BoxFuture};
6379 use std::ops::Range;
6380 use std::sync::Mutex;
6381
6382 #[derive(Debug, Clone)]
6383 struct RecordingScheduler {
6384 data: bytes::Bytes,
6385 requests: Arc<Mutex<Vec<Vec<Range<u64>>>>>,
6386 }
6387
6388 impl RecordingScheduler {
6389 fn new(data: bytes::Bytes) -> Self {
6390 Self {
6391 data,
6392 requests: Arc::new(Mutex::new(Vec::new())),
6393 }
6394 }
6395
6396 fn requests(&self) -> Vec<Vec<Range<u64>>> {
6397 self.requests.lock().unwrap().clone()
6398 }
6399 }
6400
6401 impl crate::EncodingsIo for RecordingScheduler {
6402 fn submit_request(
6403 &self,
6404 ranges: Vec<Range<u64>>,
6405 _priority: u64,
6406 ) -> BoxFuture<'static, crate::Result<Vec<bytes::Bytes>>> {
6407 self.requests.lock().unwrap().push(ranges.clone());
6408 let data = ranges
6409 .into_iter()
6410 .map(|range| self.data.slice(range.start as usize..range.end as usize))
6411 .collect::<Vec<_>>();
6412 std::future::ready(Ok(data)).boxed()
6413 }
6414 }
6415
6416 #[derive(Debug)]
6417 struct TestFixedDecompressor;
6418
6419 impl FixedPerValueDecompressor for TestFixedDecompressor {
6420 fn decompress(
6421 &self,
6422 _data: FixedWidthDataBlock,
6423 _num_rows: u64,
6424 ) -> crate::Result<DataBlock> {
6425 unimplemented!("Test decompressor")
6426 }
6427
6428 fn bits_per_value(&self) -> u64 {
6429 32
6430 }
6431 }
6432
6433 let rows_in_page = 100_u64;
6434 let data_start = 256_u64;
6435 let data_size = 500_u64;
6436 let rep_start = 4096_u64;
6437 let bytes_per_value = 4_u64;
6438
6439 let mut bytes = vec![0_u8; 16 * 1024];
6440 for i in 0..=rows_in_page {
6441 let offset = (i * 5) as u32;
6442 let pos = rep_start as usize + (i * bytes_per_value) as usize;
6443 bytes[pos..pos + 4].copy_from_slice(&offset.to_le_bytes());
6444 }
6445 let io = Arc::new(RecordingScheduler::new(bytes::Bytes::from(bytes)));
6446
6447 let scheduler = FullZipScheduler {
6448 data_buf_position: data_start,
6449 data_buf_size: data_size,
6450 rep_index: Some(FullZipRepIndexDetails {
6451 buf_position: rep_start,
6452 bytes_per_value,
6453 }),
6454 priority: 0,
6455 rows_in_page,
6456 bits_per_offset: 32,
6457 details: Arc::new(FullZipDecodeDetails {
6458 value_decompressor: PerValueDecompressor::Fixed(Arc::new(TestFixedDecompressor)),
6459 def_meaning: Arc::new([crate::repdef::DefinitionInterpretation::NullableItem]),
6460 ctrl_word_parser: crate::repdef::ControlWordParser::new(0, 1),
6461 max_rep: 0,
6462 max_visible_def: 0,
6463 }),
6464 cached_state: None,
6465 enable_cache: false,
6466 };
6467
6468 let io_dyn: Arc<dyn crate::EncodingsIo> = io.clone();
6469 let tasks = scheduler
6470 .schedule_ranges_rep(
6471 &[0..rows_in_page],
6472 &io_dyn,
6473 FullZipRepIndexDetails {
6474 buf_position: rep_start,
6475 bytes_per_value,
6476 },
6477 )
6478 .unwrap();
6479
6480 let requests = io.requests();
6481 assert_eq!(requests.len(), 1);
6482 assert_eq!(requests[0], vec![data_start..(data_start + data_size)]);
6483
6484 let _ = tasks.into_iter().next().unwrap().decoder_fut.await.unwrap();
6485 let requests_after_await = io.requests();
6486 assert_eq!(
6487 requests_after_await.len(),
6488 1,
6489 "full page path should not issue rep-index I/O"
6490 );
6491 }
6492
6493 #[tokio::test]
6495 async fn test_fuzz_issue_4492_empty_rep_values() {
6496 use lance_datagen::{RowCount, Seed, array, gen_batch};
6497
6498 let seed = 1823859942947654717u64;
6499 let num_rows = 2741usize;
6500
6501 let batch_gen = gen_batch().with_seed(Seed::from(seed));
6503 let base_generator = array::rand_type(&DataType::FixedSizeBinary(32));
6504 let list_generator = array::rand_list_any(base_generator, false);
6505
6506 let batch = batch_gen
6507 .anon_col(list_generator)
6508 .into_batch_rows(RowCount::from(num_rows as u64))
6509 .unwrap();
6510
6511 let list_array = batch.column(0).clone();
6512
6513 let mut metadata = HashMap::new();
6515 metadata.insert(
6516 STRUCTURAL_ENCODING_META_KEY.to_string(),
6517 STRUCTURAL_ENCODING_MINIBLOCK.to_string(),
6518 );
6519
6520 let test_cases = TestCases::default()
6521 .with_min_file_version(LanceFileVersion::V2_1)
6522 .with_batch_size(100)
6523 .with_range(0..num_rows.min(500) as u64)
6524 .with_indices(vec![0, num_rows as u64 / 2, (num_rows - 1) as u64]);
6525
6526 check_round_trip_encoding_of_data(vec![list_array], &test_cases, metadata).await
6527 }
6528
6529 async fn test_minichunk_size_helper(
6530 string_data: Vec<Option<String>>,
6531 minichunk_size: u64,
6532 file_version: LanceFileVersion,
6533 ) {
6534 use crate::constants::MINICHUNK_SIZE_META_KEY;
6535 use crate::testing::{TestCases, check_round_trip_encoding_of_data};
6536 use arrow_array::{ArrayRef, StringArray};
6537 use std::sync::Arc;
6538
6539 let string_array: ArrayRef = Arc::new(StringArray::from(string_data));
6540
6541 let mut metadata = HashMap::new();
6542 metadata.insert(
6543 MINICHUNK_SIZE_META_KEY.to_string(),
6544 minichunk_size.to_string(),
6545 );
6546 metadata.insert(
6547 STRUCTURAL_ENCODING_META_KEY.to_string(),
6548 STRUCTURAL_ENCODING_MINIBLOCK.to_string(),
6549 );
6550
6551 let test_cases = TestCases::default()
6552 .with_min_file_version(file_version)
6553 .with_batch_size(1000);
6554
6555 check_round_trip_encoding_of_data(vec![string_array], &test_cases, metadata).await;
6556 }
6557
6558 #[tokio::test]
6559 async fn test_minichunk_size_roundtrip() {
6560 let mut string_data = Vec::new();
6562 for i in 0..100 {
6563 string_data.push(Some(format!("test_string_{}", i).repeat(50)));
6564 }
6565 test_minichunk_size_helper(string_data, 64, LanceFileVersion::V2_1).await;
6567 }
6568
6569 #[tokio::test]
6570 async fn test_minichunk_size_128kb_v2_2() {
6571 let mut string_data = Vec::new();
6573 for i in 0..10000 {
6575 string_data.push(Some(format!("test_string_{}", i).repeat(50)));
6576 }
6577 test_minichunk_size_helper(string_data, 128 * 1024, LanceFileVersion::V2_2).await;
6578 }
6579
6580 #[tokio::test]
6581 async fn test_binary_large_minichunk_size_over_max_miniblock_values() {
6582 let mut string_data = Vec::new();
6583 for i in 0..10000 {
6585 string_data.push(Some(format!("t_{}", i)));
6586 }
6587 test_minichunk_size_helper(string_data, 128 * 1024, LanceFileVersion::V2_2).await;
6588 }
6589
6590 #[tokio::test]
6591 async fn test_large_dictionary_general_compression() {
6592 use arrow_array::{ArrayRef, StringArray};
6593 use std::collections::HashMap;
6594 use std::sync::Arc;
6595
6596 let unique_values: Vec<String> = (0..100)
6599 .map(|i| format!("value_{:04}_{}", i, "x".repeat(500)))
6600 .collect();
6601
6602 let repeated_strings: Vec<_> = unique_values
6604 .iter()
6605 .cycle()
6606 .take(100_000)
6607 .map(|s| Some(s.as_str()))
6608 .collect();
6609
6610 let string_array = Arc::new(StringArray::from(repeated_strings)) as ArrayRef;
6611
6612 let test_cases = TestCases::default()
6614 .with_min_file_version(LanceFileVersion::V2_2)
6615 .with_verify_encoding(Arc::new(|cols: &[crate::encoder::EncodedColumn], _| {
6616 assert_eq!(cols.len(), 1);
6617 let col = &cols[0];
6618
6619 if let Some(PageEncoding::Structural(page_layout)) =
6621 &col.final_pages.first().map(|p| &p.description)
6622 && let Some(pb21::page_layout::Layout::MiniBlockLayout(mini_block)) =
6623 &page_layout.layout
6624 && let Some(dictionary_encoding) = &mini_block.dictionary
6625 {
6626 match dictionary_encoding.compression.as_ref() {
6627 Some(Compression::General(general)) => {
6628 let compression = general.compression.as_ref().unwrap();
6630 assert!(
6631 compression.scheme()
6632 == pb21::CompressionScheme::CompressionAlgorithmLz4
6633 || compression.scheme()
6634 == pb21::CompressionScheme::CompressionAlgorithmZstd,
6635 "Expected LZ4 or Zstd compression for large dictionary"
6636 );
6637 }
6638 _ => panic!("Expected General compression for large dictionary"),
6639 }
6640 }
6641 }));
6642
6643 check_round_trip_encoding_of_data(vec![string_array], &test_cases, HashMap::new()).await;
6644 }
6645
6646 fn dictionary_encoding_from_page(
6647 page: &crate::encoder::EncodedPage,
6648 ) -> &crate::format::pb21::CompressiveEncoding {
6649 let PageEncoding::Structural(layout) = &page.description else {
6650 panic!("Expected structural page encoding");
6651 };
6652 let pb21::page_layout::Layout::MiniBlockLayout(layout) = layout.layout.as_ref().unwrap()
6653 else {
6654 panic!("Expected mini-block layout");
6655 };
6656 layout
6657 .dictionary
6658 .as_ref()
6659 .unwrap_or_else(|| panic!("Expected dictionary encoding"))
6660 }
6661
6662 async fn encode_variable_dict_page(
6663 metadata: HashMap<String, String>,
6664 ) -> crate::encoder::EncodedPage {
6665 use arrow_array::types::Int32Type;
6666 use arrow_array::{ArrayRef, DictionaryArray, Int32Array, StringArray};
6667
6668 let values = Arc::new(StringArray::from(
6669 (0..128)
6670 .map(|i| format!("value_{i:04}_{}", "x".repeat(256)))
6671 .collect::<Vec<_>>(),
6672 )) as ArrayRef;
6673 let keys = Int32Array::from_iter_values((0..20_000).map(|i| i % 128));
6674 let dict_array =
6675 Arc::new(DictionaryArray::<Int32Type>::try_new(keys, values).unwrap()) as ArrayRef;
6676
6677 let field = arrow_schema::Field::new(
6678 "dict_col",
6679 DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
6680 false,
6681 )
6682 .with_metadata(metadata);
6683
6684 encode_first_page(field, dict_array, LanceFileVersion::V2_2).await
6685 }
6686
6687 async fn encode_auto_fixed_dict_page(
6688 metadata: HashMap<String, String>,
6689 ) -> crate::encoder::EncodedPage {
6690 use arrow_array::{ArrayRef, Decimal128Array};
6691
6692 let values = (0..20_000)
6694 .map(|i| match i % 3 {
6695 0 => 10_i128,
6696 1 => 20_i128,
6697 _ => 30_i128,
6698 })
6699 .collect::<Vec<_>>();
6700 let decimal = Decimal128Array::from_iter_values(values)
6701 .with_precision_and_scale(38, 0)
6702 .unwrap();
6703 let decimal = Arc::new(decimal) as ArrayRef;
6704
6705 let mut field_metadata = metadata;
6706 field_metadata.insert(
6708 "lance-encoding:dict-size-ratio".to_string(),
6709 "0.99".to_string(),
6710 );
6711 let field = arrow_schema::Field::new("fixed_col", DataType::Decimal128(38, 0), false)
6712 .with_metadata(field_metadata);
6713
6714 encode_first_page(field, decimal, LanceFileVersion::V2_2).await
6715 }
6716
6717 #[tokio::test]
6718 async fn test_dict_values_general_compression_default_lz4_for_variable_dict_values() {
6719 let page = encode_variable_dict_page(HashMap::new()).await;
6720 let dictionary_encoding = dictionary_encoding_from_page(&page);
6721 let Some(Compression::General(general)) = dictionary_encoding.compression.as_ref() else {
6722 panic!("Expected General compression for dictionary values");
6723 };
6724 let compression = general.compression.as_ref().unwrap();
6725 assert_eq!(
6726 compression.scheme(),
6727 pb21::CompressionScheme::CompressionAlgorithmLz4
6728 );
6729 }
6730
6731 #[tokio::test]
6732 async fn test_dict_values_general_compression_default_lz4_for_fixed_dict_values() {
6733 let page = encode_auto_fixed_dict_page(HashMap::new()).await;
6734 let dictionary_encoding = dictionary_encoding_from_page(&page);
6735 let Some(Compression::General(general)) = dictionary_encoding.compression.as_ref() else {
6736 panic!("Expected General compression for dictionary values");
6737 };
6738 let compression = general.compression.as_ref().unwrap();
6739 assert_eq!(
6740 compression.scheme(),
6741 pb21::CompressionScheme::CompressionAlgorithmLz4
6742 );
6743 }
6744
6745 #[tokio::test]
6746 async fn test_dict_values_general_compression_zstd() {
6747 let mut metadata = HashMap::new();
6748 metadata.insert(
6749 DICT_VALUES_COMPRESSION_META_KEY.to_string(),
6750 "zstd".to_string(),
6751 );
6752 let page = encode_variable_dict_page(metadata).await;
6753 let dictionary_encoding = dictionary_encoding_from_page(&page);
6754 let Some(Compression::General(general)) = dictionary_encoding.compression.as_ref() else {
6755 panic!("Expected General compression for dictionary values");
6756 };
6757 let compression = general.compression.as_ref().unwrap();
6758 assert_eq!(
6759 compression.scheme(),
6760 pb21::CompressionScheme::CompressionAlgorithmZstd
6761 );
6762 }
6763
6764 #[tokio::test]
6765 async fn test_dict_values_general_compression_none() {
6766 let mut metadata = HashMap::new();
6767 metadata.insert(
6768 DICT_VALUES_COMPRESSION_META_KEY.to_string(),
6769 "none".to_string(),
6770 );
6771 let page = encode_variable_dict_page(metadata).await;
6772 let dictionary_encoding = dictionary_encoding_from_page(&page);
6773 assert!(
6774 !matches!(
6775 dictionary_encoding.compression.as_ref(),
6776 Some(Compression::General(_))
6777 ),
6778 "Expected dictionary values to avoid General compression"
6779 );
6780 }
6781
6782 #[test]
6783 fn test_resolve_dict_values_compression_metadata_defaults_to_lz4() {
6784 let metadata = PrimitiveStructuralEncoder::resolve_dict_values_compression_metadata(
6785 &HashMap::new(),
6786 None,
6787 None,
6788 );
6789 assert_eq!(metadata.get(COMPRESSION_META_KEY), Some(&"lz4".to_string()),);
6790 assert!(!metadata.contains_key(COMPRESSION_LEVEL_META_KEY));
6791 }
6792
6793 #[test]
6794 fn test_resolve_dict_values_compression_metadata_metadata_overrides_env() {
6795 let field_metadata = HashMap::from([
6796 (
6797 DICT_VALUES_COMPRESSION_META_KEY.to_string(),
6798 "none".to_string(),
6799 ),
6800 (
6801 DICT_VALUES_COMPRESSION_LEVEL_META_KEY.to_string(),
6802 "7".to_string(),
6803 ),
6804 ]);
6805 let metadata = PrimitiveStructuralEncoder::resolve_dict_values_compression_metadata(
6806 &field_metadata,
6807 Some("zstd".to_string()),
6808 Some("3".to_string()),
6809 );
6810 assert_eq!(
6811 metadata.get(COMPRESSION_META_KEY),
6812 Some(&"none".to_string()),
6813 );
6814 assert_eq!(
6815 metadata.get(COMPRESSION_LEVEL_META_KEY),
6816 Some(&"7".to_string()),
6817 );
6818 }
6819
6820 #[test]
6821 fn test_resolve_dict_values_compression_metadata_env_fallback() {
6822 let metadata = PrimitiveStructuralEncoder::resolve_dict_values_compression_metadata(
6823 &HashMap::new(),
6824 Some("zstd".to_string()),
6825 Some("9".to_string()),
6826 );
6827 assert_eq!(
6828 metadata.get(COMPRESSION_META_KEY),
6829 Some(&"zstd".to_string()),
6830 );
6831 assert_eq!(
6832 metadata.get(COMPRESSION_LEVEL_META_KEY),
6833 Some(&"9".to_string()),
6834 );
6835 }
6836
6837 #[tokio::test]
6838 async fn test_dictionary_encode_int64() {
6839 use crate::constants::{DICT_SIZE_RATIO_META_KEY, STRUCTURAL_ENCODING_META_KEY};
6840 use crate::testing::{TestCases, check_round_trip_encoding_of_data};
6841 use crate::version::LanceFileVersion;
6842 use arrow_array::{ArrayRef, Int64Array};
6843 use std::collections::HashMap;
6844 use std::sync::Arc;
6845
6846 let values = (0..1000)
6848 .map(|i| match i % 3 {
6849 0 => 10i64,
6850 1 => 20i64,
6851 _ => 30i64,
6852 })
6853 .collect::<Vec<_>>();
6854 let array = Arc::new(Int64Array::from(values)) as ArrayRef;
6855
6856 let mut metadata = HashMap::new();
6857 metadata.insert(
6858 STRUCTURAL_ENCODING_META_KEY.to_string(),
6859 STRUCTURAL_ENCODING_MINIBLOCK.to_string(),
6860 );
6861 metadata.insert(DICT_SIZE_RATIO_META_KEY.to_string(), "0.99".to_string());
6862
6863 let test_cases = TestCases::default()
6864 .with_min_file_version(LanceFileVersion::V2_2)
6865 .with_batch_size(1000)
6866 .with_range(0..1000)
6867 .with_indices(vec![0, 1, 10, 999])
6868 .with_expected_encoding("dictionary");
6869
6870 check_round_trip_encoding_of_data(vec![array], &test_cases, metadata).await;
6871 }
6872
6873 #[tokio::test]
6874 async fn test_dictionary_encode_float64() {
6875 use crate::constants::{DICT_SIZE_RATIO_META_KEY, STRUCTURAL_ENCODING_META_KEY};
6876 use crate::testing::{TestCases, check_round_trip_encoding_of_data};
6877 use crate::version::LanceFileVersion;
6878 use arrow_array::{ArrayRef, Float64Array};
6879 use std::collections::HashMap;
6880 use std::sync::Arc;
6881
6882 let values = (0..1000)
6884 .map(|i| match i % 3 {
6885 0 => 0.1f64,
6886 1 => 0.2f64,
6887 _ => 0.3f64,
6888 })
6889 .collect::<Vec<_>>();
6890 let array = Arc::new(Float64Array::from(values)) as ArrayRef;
6891
6892 let mut metadata = HashMap::new();
6893 metadata.insert(
6894 STRUCTURAL_ENCODING_META_KEY.to_string(),
6895 STRUCTURAL_ENCODING_MINIBLOCK.to_string(),
6896 );
6897 metadata.insert(DICT_SIZE_RATIO_META_KEY.to_string(), "0.99".to_string());
6898
6899 let test_cases = TestCases::default()
6900 .with_min_file_version(LanceFileVersion::V2_2)
6901 .with_batch_size(1000)
6902 .with_range(0..1000)
6903 .with_indices(vec![0, 1, 10, 999])
6904 .with_expected_encoding("dictionary");
6905
6906 check_round_trip_encoding_of_data(vec![array], &test_cases, metadata).await;
6907 }
6908
6909 #[test]
6910 fn test_miniblock_dictionary_out_of_line_bitpacking_decode() {
6911 let rows = 10_000;
6912 let unique_values = 2_000;
6913
6914 let dictionary_encoding =
6915 ProtobufUtils21::out_of_line_bitpacking(64, ProtobufUtils21::flat(11, None));
6916 let layout = pb21::MiniBlockLayout {
6917 rep_compression: None,
6918 def_compression: None,
6919 value_compression: Some(ProtobufUtils21::flat(64, None)),
6920 dictionary: Some(dictionary_encoding),
6921 num_dictionary_items: unique_values,
6922 layers: vec![pb21::RepDefLayer::RepdefAllValidItem as i32],
6923 num_buffers: 1,
6924 repetition_index_depth: 0,
6925 num_items: rows,
6926 has_large_chunk: false,
6927 };
6928
6929 let buffer_offsets_and_sizes = vec![(0, 0), (0, 0), (0, 0)];
6930 let scheduler = super::MiniBlockScheduler::try_new(
6931 &buffer_offsets_and_sizes,
6932 0,
6933 rows,
6934 &layout,
6935 &DefaultDecompressionStrategy::default(),
6936 )
6937 .unwrap();
6938
6939 let dictionary = scheduler.dictionary.unwrap();
6940 assert_eq!(dictionary.num_dictionary_items, unique_values);
6941 assert_eq!(
6942 dictionary.dictionary_data_alignment,
6943 crate::encoder::MIN_PAGE_BUFFER_ALIGNMENT
6944 );
6945 }
6946
6947 fn create_test_fixed_data_block(
6949 num_values: u64,
6950 cardinality: u64,
6951 bits_per_value: u64,
6952 ) -> DataBlock {
6953 assert!(cardinality > 0);
6954 assert!(cardinality <= num_values);
6955 let block_info = BlockInfo::default();
6956
6957 assert_eq!(bits_per_value % 8, 0);
6958 let data = match bits_per_value {
6959 32 => {
6960 let values = (0..num_values)
6961 .map(|i| (i % cardinality) as u32)
6962 .collect::<Vec<_>>();
6963 crate::buffer::LanceBuffer::reinterpret_vec(values)
6964 }
6965 64 => {
6966 let values = (0..num_values).map(|i| i % cardinality).collect::<Vec<_>>();
6967 crate::buffer::LanceBuffer::reinterpret_vec(values)
6968 }
6969 128 => {
6970 let values = (0..num_values)
6971 .map(|i| (i % cardinality) as u128)
6972 .collect::<Vec<_>>();
6973 crate::buffer::LanceBuffer::reinterpret_vec(values)
6974 }
6975 _ => unreachable!(),
6976 };
6977 DataBlock::FixedWidth(FixedWidthDataBlock {
6978 bits_per_value,
6979 data,
6980 num_values,
6981 block_info,
6982 })
6983 }
6984
6985 fn create_test_variable_width_block(num_values: u64, cardinality: u64) -> DataBlock {
6987 use arrow_array::StringArray;
6988
6989 assert!(cardinality <= num_values && cardinality > 0);
6990
6991 let mut values = Vec::with_capacity(num_values as usize);
6992 for i in 0..num_values {
6993 values.push(format!("value_{:016}", i % cardinality));
6994 }
6995
6996 let array = StringArray::from(values);
6997 DataBlock::from_array(Arc::new(array) as ArrayRef)
6998 }
6999
7000 #[test]
7001 fn test_should_dictionary_encode() {
7002 use crate::constants::DICT_SIZE_RATIO_META_KEY;
7003 use lance_core::datatypes::Field as LanceField;
7004
7005 let block = create_test_variable_width_block(1000, 10);
7007
7008 let mut metadata = HashMap::new();
7009 metadata.insert(DICT_SIZE_RATIO_META_KEY.to_string(), "0.8".to_string());
7010 let arrow_field =
7011 arrow_schema::Field::new("test", DataType::Utf8, false).with_metadata(metadata);
7012 let field = LanceField::try_from(&arrow_field).unwrap();
7013
7014 let result = PrimitiveStructuralEncoder::should_dictionary_encode(
7015 &block,
7016 &field,
7017 LanceFileVersion::V2_1,
7018 );
7019
7020 assert!(
7021 result.is_some(),
7022 "Should use dictionary encode based on size"
7023 );
7024 }
7025
7026 #[test]
7027 fn test_should_not_dictionary_encode_unsupported_bits() {
7028 use crate::constants::DICT_SIZE_RATIO_META_KEY;
7029 use lance_core::datatypes::Field as LanceField;
7030
7031 let block = create_test_fixed_data_block(1000, 1000, 32);
7032
7033 let mut metadata = HashMap::new();
7034 metadata.insert(DICT_SIZE_RATIO_META_KEY.to_string(), "0.8".to_string());
7035 let arrow_field =
7036 arrow_schema::Field::new("test", DataType::Int32, false).with_metadata(metadata);
7037 let field = LanceField::try_from(&arrow_field).unwrap();
7038
7039 let result = PrimitiveStructuralEncoder::should_dictionary_encode(
7040 &block,
7041 &field,
7042 LanceFileVersion::V2_1,
7043 );
7044
7045 assert!(
7046 result.is_none(),
7047 "Should not use dictionary encode for unsupported bit width"
7048 );
7049 }
7050
7051 #[test]
7052 fn test_should_not_dictionary_encode_near_unique_sample() {
7053 use crate::constants::DICT_SIZE_RATIO_META_KEY;
7054 use lance_core::datatypes::Field as LanceField;
7055
7056 let num_values = 5000;
7057 let block = create_test_variable_width_block(num_values, num_values);
7058
7059 let mut metadata = HashMap::new();
7060 metadata.insert(DICT_SIZE_RATIO_META_KEY.to_string(), "1.0".to_string());
7061 let arrow_field =
7062 arrow_schema::Field::new("test", DataType::Utf8, false).with_metadata(metadata);
7063 let field = LanceField::try_from(&arrow_field).unwrap();
7064
7065 let result = PrimitiveStructuralEncoder::should_dictionary_encode(
7066 &block,
7067 &field,
7068 LanceFileVersion::V2_1,
7069 );
7070
7071 assert!(
7072 result.is_none(),
7073 "Should not probe dictionary encoding for near-unique data"
7074 );
7075 }
7076
7077 async fn encode_first_page(
7078 field: arrow_schema::Field,
7079 array: ArrayRef,
7080 version: LanceFileVersion,
7081 ) -> crate::encoder::EncodedPage {
7082 use crate::encoder::{
7083 ColumnIndexSequence, EncodingOptions, MIN_PAGE_BUFFER_ALIGNMENT, OutOfLineBuffers,
7084 default_encoding_strategy,
7085 };
7086 use crate::repdef::RepDefBuilder;
7087
7088 let lance_field = lance_core::datatypes::Field::try_from(&field).unwrap();
7089 let encoding_strategy = default_encoding_strategy(version);
7090 let mut column_index_seq = ColumnIndexSequence::default();
7091 let encoding_options = EncodingOptions {
7092 cache_bytes_per_column: 1,
7093 max_page_bytes: 32 * 1024 * 1024,
7094 keep_original_array: true,
7095 buffer_alignment: MIN_PAGE_BUFFER_ALIGNMENT,
7096 version,
7097 };
7098
7099 let mut encoder = encoding_strategy
7100 .create_field_encoder(
7101 encoding_strategy.as_ref(),
7102 &lance_field,
7103 &mut column_index_seq,
7104 &encoding_options,
7105 )
7106 .unwrap();
7107
7108 let mut external_buffers = OutOfLineBuffers::new(0, MIN_PAGE_BUFFER_ALIGNMENT);
7109 let repdef = RepDefBuilder::default();
7110 let num_rows = array.len() as u64;
7111 let mut pages = Vec::new();
7112 for task in encoder
7113 .maybe_encode(array, &mut external_buffers, repdef, 0, num_rows)
7114 .unwrap()
7115 {
7116 pages.push(task.await.unwrap());
7117 }
7118 for task in encoder.flush(&mut external_buffers).unwrap() {
7119 pages.push(task.await.unwrap());
7120 }
7121 pages.into_iter().next().unwrap()
7122 }
7123
7124 #[tokio::test]
7125 async fn test_constant_layout_out_of_line_fixed_size_binary_v2_2() {
7126 use crate::format::pb21::page_layout::Layout;
7127
7128 let val = vec![0xABu8; 33];
7129 let arr: ArrayRef = Arc::new(
7130 arrow_array::FixedSizeBinaryArray::try_from_sparse_iter_with_size(
7131 std::iter::repeat_n(Some(val.as_slice()), 256),
7132 33,
7133 )
7134 .unwrap(),
7135 );
7136 let field = arrow_schema::Field::new("c", DataType::FixedSizeBinary(33), true);
7137 let page = encode_first_page(field, arr.clone(), LanceFileVersion::V2_2).await;
7138
7139 let PageEncoding::Structural(layout) = &page.description else {
7140 panic!("Expected structural encoding");
7141 };
7142 let Layout::ConstantLayout(layout) = layout.layout.as_ref().unwrap() else {
7143 panic!("Expected constant layout in slot 2");
7144 };
7145 assert!(layout.inline_value.is_none());
7146 assert_eq!(page.data.len(), 1);
7147
7148 let test_cases = TestCases::default()
7149 .with_min_file_version(LanceFileVersion::V2_2)
7150 .with_max_file_version(LanceFileVersion::V2_2)
7151 .with_page_sizes(vec![4096]);
7152 check_round_trip_encoding_of_data(vec![arr], &test_cases, HashMap::new()).await;
7153 }
7154
7155 #[tokio::test]
7156 async fn test_constant_layout_out_of_line_utf8_v2_2() {
7157 use crate::format::pb21::page_layout::Layout;
7158
7159 let arr: ArrayRef = Arc::new(arrow_array::StringArray::from_iter_values(
7160 std::iter::repeat_n("hello", 512),
7161 ));
7162 let field = arrow_schema::Field::new("c", DataType::Utf8, true);
7163 let page = encode_first_page(field, arr.clone(), LanceFileVersion::V2_2).await;
7164
7165 let PageEncoding::Structural(layout) = &page.description else {
7166 panic!("Expected structural encoding");
7167 };
7168 let Layout::ConstantLayout(layout) = layout.layout.as_ref().unwrap() else {
7169 panic!("Expected constant layout in slot 2");
7170 };
7171 assert!(layout.inline_value.is_none());
7172 assert_eq!(page.data.len(), 1);
7173
7174 let test_cases = TestCases::default()
7175 .with_min_file_version(LanceFileVersion::V2_2)
7176 .with_max_file_version(LanceFileVersion::V2_2)
7177 .with_page_sizes(vec![4096]);
7178 check_round_trip_encoding_of_data(vec![arr], &test_cases, HashMap::new()).await;
7179 }
7180
7181 #[tokio::test]
7182 async fn test_constant_layout_nullable_item_v2_2() {
7183 use crate::format::pb21::page_layout::Layout;
7184
7185 let arr: ArrayRef = Arc::new(arrow_array::Int32Array::from(vec![
7186 Some(7),
7187 None,
7188 Some(7),
7189 None,
7190 Some(7),
7191 ]));
7192 let field = arrow_schema::Field::new("c", DataType::Int32, true);
7193 let page = encode_first_page(field, arr.clone(), LanceFileVersion::V2_2).await;
7194
7195 let PageEncoding::Structural(layout) = &page.description else {
7196 panic!("Expected structural encoding");
7197 };
7198 let Layout::ConstantLayout(layout) = layout.layout.as_ref().unwrap() else {
7199 panic!("Expected constant layout in slot 2");
7200 };
7201 assert!(layout.inline_value.is_some());
7202 assert_eq!(page.data.len(), 2);
7203
7204 let test_cases = TestCases::default()
7205 .with_min_file_version(LanceFileVersion::V2_2)
7206 .with_max_file_version(LanceFileVersion::V2_2)
7207 .with_page_sizes(vec![4096]);
7208 check_round_trip_encoding_of_data(vec![arr], &test_cases, HashMap::new()).await;
7209 }
7210
7211 #[tokio::test]
7212 async fn test_constant_layout_list_repdef_v2_2() {
7213 use crate::format::pb21::page_layout::Layout;
7214 use arrow_array::builder::{Int32Builder, ListBuilder};
7215
7216 let mut builder = ListBuilder::new(Int32Builder::new());
7217 builder.values().append_value(7);
7218 builder.values().append_null();
7219 builder.values().append_value(7);
7220 builder.append(true);
7221
7222 builder.append(true);
7223
7224 builder.values().append_value(7);
7225 builder.append(true);
7226
7227 builder.append_null();
7228
7229 let arr: ArrayRef = Arc::new(builder.finish());
7230 let field = arrow_schema::Field::new(
7231 "c",
7232 DataType::List(Arc::new(arrow_schema::Field::new(
7233 "item",
7234 DataType::Int32,
7235 true,
7236 ))),
7237 true,
7238 );
7239 let page = encode_first_page(field, arr.clone(), LanceFileVersion::V2_2).await;
7240
7241 let PageEncoding::Structural(layout) = &page.description else {
7242 panic!("Expected structural encoding");
7243 };
7244 let Layout::ConstantLayout(layout) = layout.layout.as_ref().unwrap() else {
7245 panic!("Expected constant layout in slot 2");
7246 };
7247 assert!(layout.inline_value.is_some());
7248 assert_eq!(page.data.len(), 2);
7249
7250 let test_cases = TestCases::default()
7251 .with_min_file_version(LanceFileVersion::V2_2)
7252 .with_max_file_version(LanceFileVersion::V2_2)
7253 .with_page_sizes(vec![4096]);
7254 check_round_trip_encoding_of_data(vec![arr], &test_cases, HashMap::new()).await;
7255 }
7256
7257 #[tokio::test]
7258 async fn test_constant_layout_fixed_size_list_not_used_v2_2() {
7259 use crate::format::pb21::page_layout::Layout;
7260 use arrow_array::builder::{FixedSizeListBuilder, Int32Builder};
7261
7262 let mut builder = FixedSizeListBuilder::new(Int32Builder::new(), 3);
7263 for _ in 0..64 {
7264 builder.values().append_value(1);
7265 builder.values().append_null();
7266 builder.values().append_value(3);
7267 builder.append(true);
7268 }
7269 let arr: ArrayRef = Arc::new(builder.finish());
7270 let field = arrow_schema::Field::new(
7271 "c",
7272 DataType::FixedSizeList(
7273 Arc::new(arrow_schema::Field::new("item", DataType::Int32, true)),
7274 3,
7275 ),
7276 true,
7277 );
7278 let page = encode_first_page(field, arr.clone(), LanceFileVersion::V2_2).await;
7279
7280 if let PageEncoding::Structural(layout) = &page.description {
7281 assert!(
7282 !matches!(layout.layout.as_ref().unwrap(), Layout::ConstantLayout(_)),
7283 "FixedSizeList should not use constant layout yet"
7284 );
7285 }
7286
7287 let test_cases = TestCases::default()
7288 .with_min_file_version(LanceFileVersion::V2_2)
7289 .with_max_file_version(LanceFileVersion::V2_2)
7290 .with_page_sizes(vec![4096]);
7291 check_round_trip_encoding_of_data(vec![arr], &test_cases, HashMap::new()).await;
7292 }
7293
7294 #[tokio::test]
7295 async fn test_constant_layout_not_written_before_v2_2() {
7296 use crate::format::pb21::page_layout::Layout;
7297
7298 let arr: ArrayRef = Arc::new(arrow_array::Int32Array::from(vec![7; 1024]));
7299 let field = arrow_schema::Field::new("c", DataType::Int32, true);
7300 let page = encode_first_page(field, arr.clone(), LanceFileVersion::V2_1).await;
7301
7302 let PageEncoding::Structural(layout) = &page.description else {
7303 return;
7304 };
7305 assert!(
7306 !matches!(layout.layout.as_ref().unwrap(), Layout::ConstantLayout(_)),
7307 "Should not emit constant layout before v2.2"
7308 );
7309
7310 let test_cases = TestCases::default()
7311 .with_min_file_version(LanceFileVersion::V2_1)
7312 .with_max_file_version(LanceFileVersion::V2_1)
7313 .with_page_sizes(vec![4096]);
7314 check_round_trip_encoding_of_data(vec![arr], &test_cases, HashMap::new()).await;
7315 }
7316
7317 #[tokio::test]
7318 async fn test_all_null_constant_layout_still_works_v2_2() {
7319 use crate::format::pb21::page_layout::Layout;
7320
7321 let arr: ArrayRef = Arc::new(arrow_array::Int32Array::from(vec![None, None, None]));
7322 let field = arrow_schema::Field::new("c", DataType::Int32, true);
7323 let page = encode_first_page(field, arr.clone(), LanceFileVersion::V2_2).await;
7324
7325 let PageEncoding::Structural(layout) = &page.description else {
7326 panic!("Expected structural encoding");
7327 };
7328 let Layout::ConstantLayout(layout) = layout.layout.as_ref().unwrap() else {
7329 panic!("Expected layout in slot 2");
7330 };
7331 assert!(layout.inline_value.is_none());
7332 assert_eq!(page.data.len(), 0);
7333
7334 let test_cases = TestCases::default()
7335 .with_min_file_version(LanceFileVersion::V2_2)
7336 .with_max_file_version(LanceFileVersion::V2_2)
7337 .with_page_sizes(vec![4096]);
7338 check_round_trip_encoding_of_data(vec![arr], &test_cases, HashMap::new()).await;
7339 }
7340
7341 #[test]
7342 fn test_encode_decode_complex_all_null_vals_roundtrip() {
7343 use crate::compression::{
7344 DecompressionStrategy, DefaultCompressionStrategy, DefaultDecompressionStrategy,
7345 };
7346
7347 let values: Arc<[u16]> = Arc::from((0..2048).map(|i| (i % 5) as u16).collect::<Vec<u16>>());
7348
7349 let compression_strategy = DefaultCompressionStrategy::default();
7350 let decompression_strategy = DefaultDecompressionStrategy::default();
7351
7352 let (compressed_buf, encoding) = PrimitiveStructuralEncoder::encode_complex_all_null_vals(
7353 &values,
7354 &compression_strategy,
7355 )
7356 .unwrap();
7357
7358 let decompressor = decompression_strategy
7359 .create_block_decompressor(&encoding)
7360 .unwrap();
7361 let decompressed = decompressor
7362 .decompress(compressed_buf, values.len() as u64)
7363 .unwrap();
7364 let decompressed_fixed_width = decompressed.as_fixed_width().unwrap();
7365 assert_eq!(decompressed_fixed_width.num_values, values.len() as u64);
7366 assert_eq!(decompressed_fixed_width.bits_per_value, 16);
7367 let rep_result = decompressed_fixed_width.data.borrow_to_typed_slice::<u16>();
7368 assert_eq!(rep_result.as_ref(), values.as_ref());
7369 }
7370
7371 #[tokio::test]
7372 async fn test_complex_all_null_compression_gated_by_version() {
7373 use crate::format::pb21::page_layout::Layout;
7374 use arrow_array::ListArray;
7375
7376 let list_array = ListArray::from_iter_primitive::<arrow_array::types::Int32Type, _, _>(
7377 (0..1000).map(|i| if i % 2 == 0 { None } else { Some(vec![]) }),
7378 );
7379 let arr: ArrayRef = Arc::new(list_array);
7380 let field = arrow_schema::Field::new(
7381 "c",
7382 DataType::List(Arc::new(arrow_schema::Field::new(
7383 "item",
7384 DataType::Int32,
7385 true,
7386 ))),
7387 true,
7388 );
7389
7390 let page_v21 = encode_first_page(field.clone(), arr.clone(), LanceFileVersion::V2_1).await;
7391 let PageEncoding::Structural(layout_v21) = &page_v21.description else {
7392 panic!("Expected structural encoding");
7393 };
7394 let Layout::ConstantLayout(layout_v21) = layout_v21.layout.as_ref().unwrap() else {
7395 panic!("Expected constant layout");
7396 };
7397 assert!(layout_v21.rep_compression.is_none());
7398 assert!(layout_v21.def_compression.is_none());
7399 assert_eq!(layout_v21.num_rep_values, 0);
7400 assert_eq!(layout_v21.num_def_values, 0);
7401
7402 let page_v22 = encode_first_page(field, arr, LanceFileVersion::V2_2).await;
7403 let PageEncoding::Structural(layout_v22) = &page_v22.description else {
7404 panic!("Expected structural encoding");
7405 };
7406 let Layout::ConstantLayout(layout_v22) = layout_v22.layout.as_ref().unwrap() else {
7407 panic!("Expected constant layout");
7408 };
7409 assert!(layout_v22.def_compression.is_some());
7410 assert!(layout_v22.num_def_values > 0);
7411 }
7412
7413 #[tokio::test]
7414 async fn test_complex_all_null_round_trip() {
7415 use arrow_array::ListArray;
7416
7417 let list_array = ListArray::from_iter_primitive::<arrow_array::types::Int32Type, _, _>(
7418 (0..1000).map(|i| if i % 2 == 0 { None } else { Some(vec![]) }),
7419 );
7420
7421 let test_cases = TestCases::default().with_min_file_version(LanceFileVersion::V2_2);
7422 check_round_trip_encoding_of_data(vec![Arc::new(list_array)], &test_cases, HashMap::new())
7423 .await;
7424 }
7425 fn truncated_tail_details() -> std::sync::Arc<super::FullZipDecodeDetails> {
7426 use crate::compression::VariablePerValueDecompressor;
7427 use crate::encodings::physical::binary::VariableDecoder;
7428 use crate::repdef::{ControlWordParser, DefinitionInterpretation};
7429 use std::sync::Arc;
7430 Arc::new(super::FullZipDecodeDetails {
7431 value_decompressor: super::PerValueDecompressor::Variable(Arc::new(
7432 VariableDecoder::default(),
7433 )
7434 as Arc<dyn VariablePerValueDecompressor>),
7435 def_meaning: vec![DefinitionInterpretation::NullableItem].into(),
7436 ctrl_word_parser: ControlWordParser::new(0, 0),
7437 max_rep: 0,
7438 max_visible_def: 0,
7439 })
7440 }
7441
7442 fn decode_variable_full_zip(
7443 buf: Vec<u8>,
7444 bits_per_offset: u8,
7445 ) -> lance_core::Result<super::VariableFullZipDecoder> {
7446 use std::collections::VecDeque;
7447 let mut data = VecDeque::new();
7448 data.push_back(crate::buffer::LanceBuffer::from(buf));
7449 super::VariableFullZipDecoder::new(
7450 truncated_tail_details(),
7451 data,
7452 1,
7453 bits_per_offset,
7454 bits_per_offset,
7455 )
7456 }
7457
7458 #[test]
7467 fn variable_full_zip_truncated_length_prefix_is_corrupt_file() {
7468 use lance_core::Error;
7469
7470 for (bits, buf_len) in [(32u8, 3usize), (64u8, 4usize)] {
7471 let err = decode_variable_full_zip(vec![0xAA; buf_len], bits)
7472 .expect_err("a truncated length prefix must not decode");
7473 assert!(
7474 matches!(err, Error::CorruptFile { .. }),
7475 "expected CorruptFile for a {}-bit prefix with {} byte(s), got: {:?}",
7476 bits,
7477 buf_len,
7478 err
7479 );
7480 let msg = err.to_string();
7481 assert!(
7482 msg.contains("truncated length prefix"),
7483 "error should say what is wrong, got: {msg}"
7484 );
7485 }
7486 }
7487}