1use bytes::Bytes;
21
22use super::page::{Page, PageReader};
23use crate::basic::*;
24use crate::column::reader::decoder::{
25 ColumnValueDecoder, ColumnValueDecoderImpl, DefinitionLevelDecoder, DefinitionLevelDecoderImpl,
26 RepetitionLevelDecoder, RepetitionLevelDecoderImpl,
27};
28use crate::data_type::*;
29use crate::errors::{ParquetError, Result};
30use crate::schema::types::ColumnDescPtr;
31use crate::util::bit_util::{ceil, num_required_bits, read_num_bytes};
32
33pub(crate) mod decoder;
34
35pub enum ColumnReader {
37 BoolColumnReader(ColumnReaderImpl<BoolType>),
39 Int32ColumnReader(ColumnReaderImpl<Int32Type>),
41 Int64ColumnReader(ColumnReaderImpl<Int64Type>),
43 Int96ColumnReader(ColumnReaderImpl<Int96Type>),
45 FloatColumnReader(ColumnReaderImpl<FloatType>),
47 DoubleColumnReader(ColumnReaderImpl<DoubleType>),
49 ByteArrayColumnReader(ColumnReaderImpl<ByteArrayType>),
51 FixedLenByteArrayColumnReader(ColumnReaderImpl<FixedLenByteArrayType>),
53}
54
55pub fn get_column_reader(
58 col_descr: ColumnDescPtr,
59 col_page_reader: Box<dyn PageReader>,
60) -> ColumnReader {
61 match col_descr.physical_type() {
62 Type::BOOLEAN => {
63 ColumnReader::BoolColumnReader(ColumnReaderImpl::new(col_descr, col_page_reader))
64 }
65 Type::INT32 => {
66 ColumnReader::Int32ColumnReader(ColumnReaderImpl::new(col_descr, col_page_reader))
67 }
68 Type::INT64 => {
69 ColumnReader::Int64ColumnReader(ColumnReaderImpl::new(col_descr, col_page_reader))
70 }
71 Type::INT96 => {
72 ColumnReader::Int96ColumnReader(ColumnReaderImpl::new(col_descr, col_page_reader))
73 }
74 Type::FLOAT => {
75 ColumnReader::FloatColumnReader(ColumnReaderImpl::new(col_descr, col_page_reader))
76 }
77 Type::DOUBLE => {
78 ColumnReader::DoubleColumnReader(ColumnReaderImpl::new(col_descr, col_page_reader))
79 }
80 Type::BYTE_ARRAY => {
81 ColumnReader::ByteArrayColumnReader(ColumnReaderImpl::new(col_descr, col_page_reader))
82 }
83 Type::FIXED_LEN_BYTE_ARRAY => ColumnReader::FixedLenByteArrayColumnReader(
84 ColumnReaderImpl::new(col_descr, col_page_reader),
85 ),
86 }
87}
88
89pub fn get_typed_column_reader<T: DataType>(col_reader: ColumnReader) -> ColumnReaderImpl<T> {
96 T::get_column_reader(col_reader).unwrap_or_else(|| {
97 panic!(
98 "Failed to convert column reader into a typed column reader for `{}` type",
99 T::get_physical_type()
100 )
101 })
102}
103
104pub type ColumnReaderImpl<T> = GenericColumnReader<
106 RepetitionLevelDecoderImpl,
107 DefinitionLevelDecoderImpl,
108 ColumnValueDecoderImpl<T>,
109>;
110
111pub struct GenericColumnReader<R, D, V> {
117 descr: ColumnDescPtr,
118
119 page_reader: Box<dyn PageReader>,
120
121 num_buffered_values: usize,
123
124 num_decoded_values: usize,
127
128 has_record_delimiter: bool,
130
131 def_level_decoder: Option<D>,
133
134 rep_level_decoder: Option<R>,
136
137 values_decoder: V,
139}
140
141impl<V> GenericColumnReader<RepetitionLevelDecoderImpl, DefinitionLevelDecoderImpl, V>
142where
143 V: ColumnValueDecoder,
144{
145 pub fn new(descr: ColumnDescPtr, page_reader: Box<dyn PageReader>) -> Self {
147 let values_decoder = V::new(&descr);
148
149 let def_level_decoder = (descr.max_def_level() != 0)
150 .then(|| DefinitionLevelDecoderImpl::new(descr.max_def_level()));
151
152 let rep_level_decoder = (descr.max_rep_level() != 0)
153 .then(|| RepetitionLevelDecoderImpl::new(descr.max_rep_level()));
154
155 Self::new_with_decoders(
156 descr,
157 page_reader,
158 values_decoder,
159 def_level_decoder,
160 rep_level_decoder,
161 )
162 }
163}
164
165impl<R, D, V> GenericColumnReader<R, D, V>
166where
167 R: RepetitionLevelDecoder,
168 D: DefinitionLevelDecoder,
169 V: ColumnValueDecoder,
170{
171 pub(crate) fn new_with_decoders(
172 descr: ColumnDescPtr,
173 page_reader: Box<dyn PageReader>,
174 values_decoder: V,
175 def_level_decoder: Option<D>,
176 rep_level_decoder: Option<R>,
177 ) -> Self {
178 Self {
179 descr,
180 def_level_decoder,
181 rep_level_decoder,
182 page_reader,
183 num_buffered_values: 0,
184 num_decoded_values: 0,
185 values_decoder,
186 has_record_delimiter: false,
187 }
188 }
189
190 pub fn read_records(
205 &mut self,
206 max_records: usize,
207 def_levels: Option<&mut D::Buffer>,
208 rep_levels: Option<&mut R::Buffer>,
209 values: &mut V::Buffer,
210 ) -> Result<(usize, usize, usize)> {
211 self.read_records_with_reservation(
212 max_records,
213 def_levels,
214 rep_levels,
215 values,
216 |_, _, _, _| Ok(()),
217 )
218 }
219
220 pub(crate) fn read_records_with_reservation<F>(
223 &mut self,
224 max_records: usize,
225 mut def_levels: Option<&mut D::Buffer>,
226 mut rep_levels: Option<&mut R::Buffer>,
227 values: &mut V::Buffer,
228 mut reserve_values: F,
229 ) -> Result<(usize, usize, usize)>
230 where
231 F: FnMut(&mut V::Buffer, usize, usize, Option<&D::Buffer>) -> Result<()>,
232 {
233 let mut total_records_read = 0;
234 let mut total_levels_read = 0;
235 let mut total_values_read = 0;
236
237 while total_records_read < max_records && self.has_next()? {
238 let remaining_records = max_records - total_records_read;
239 let remaining_levels = self.num_buffered_values - self.num_decoded_values;
240
241 let (records_read, levels_to_read) = match self.rep_level_decoder.as_mut() {
242 Some(reader) => {
243 let out = rep_levels
244 .as_mut()
245 .ok_or_else(|| general_err!("must specify repetition levels"))?;
246
247 let (mut records_read, levels_read) =
248 reader.read_rep_levels(out, remaining_records, remaining_levels)?;
249
250 if records_read == 0 && levels_read == 0 {
251 return Err(general_err!(
253 "Insufficient repetition levels read from column"
254 ));
255 }
256 if levels_read == remaining_levels && self.has_record_delimiter {
257 check_partial_record_fits(records_read, remaining_records)?;
258 records_read += reader.flush_partial() as usize;
259 }
260 (records_read, levels_read)
261 }
262 None => {
263 let min = remaining_records.min(remaining_levels);
264 (min, min)
265 }
266 };
267
268 let values_to_read = match self.def_level_decoder.as_mut() {
269 Some(reader) => {
270 let out = def_levels
271 .as_mut()
272 .ok_or_else(|| general_err!("must specify definition levels"))?;
273
274 let (values_read, levels_read) = reader.read_def_levels(out, levels_to_read)?;
275
276 if levels_read != levels_to_read {
277 return Err(general_err!(
278 "insufficient definition levels read from column - expected {levels_to_read}, got {levels_read}"
279 ));
280 }
281
282 values_read
283 }
284 None => levels_to_read,
285 };
286
287 let def_levels = def_levels.as_deref();
288 reserve_values(values, values_to_read, levels_to_read, def_levels)?;
289
290 let values_read = self.values_decoder.read(values, values_to_read)?;
291
292 if values_read != values_to_read {
293 return Err(general_err!(
294 "insufficient values read from column - expected: {values_to_read}, got: {values_read}",
295 ));
296 }
297
298 self.num_decoded_values += levels_to_read;
299 total_records_read += records_read;
300 total_levels_read += levels_to_read;
301 total_values_read += values_read;
302 }
303
304 Ok((total_records_read, total_values_read, total_levels_read))
305 }
306
307 pub fn skip_records(&mut self, num_records: usize) -> Result<usize> {
313 let mut remaining_records = num_records;
314 while remaining_records != 0 {
315 if self.num_buffered_values == self.num_decoded_values {
316 let Some(metadata) = self.page_reader.peek_next_page()? else {
317 return Ok(num_records - remaining_records);
318 };
319
320 if metadata.is_dict {
322 self.read_dictionary_page()?;
323 continue;
324 }
325
326 let rows = metadata.num_rows.or_else(|| {
329 self.rep_level_decoder
331 .is_none()
332 .then_some(metadata.num_levels)?
333 });
334
335 if let Some(rows) = rows
336 && rows <= remaining_records
337 {
338 self.page_reader.skip_next_page()?;
339 remaining_records -= rows;
340 continue;
341 }
342 if !self.read_new_page()? {
345 return Ok(num_records - remaining_records);
346 }
347 }
348
349 let remaining_levels = self.num_buffered_values - self.num_decoded_values;
353
354 let (records_read, rep_levels_read) = match self.rep_level_decoder.as_mut() {
355 Some(decoder) => {
356 let (mut records_read, levels_read) =
357 decoder.skip_rep_levels(remaining_records, remaining_levels)?;
358
359 if levels_read == remaining_levels && self.has_record_delimiter {
360 check_partial_record_fits(records_read, remaining_records)?;
361 records_read += decoder.flush_partial() as usize;
362 }
363
364 (records_read, levels_read)
365 }
366 None => {
367 let levels = remaining_levels.min(remaining_records);
369 (levels, levels)
370 }
371 };
372
373 self.num_decoded_values += rep_levels_read;
374 remaining_records -= records_read;
375
376 if self.num_buffered_values == self.num_decoded_values {
377 continue;
379 }
380
381 let (values_read, def_levels_read) = match self.def_level_decoder.as_mut() {
382 Some(decoder) => decoder.skip_def_levels(rep_levels_read)?,
383 None => (rep_levels_read, rep_levels_read),
384 };
385
386 if rep_levels_read != def_levels_read {
387 return Err(general_err!(
388 "levels mismatch, read {} repetition levels and {} definition levels",
389 rep_levels_read,
390 def_levels_read
391 ));
392 }
393
394 let values = self.values_decoder.skip_values(values_read)?;
395 if values != values_read {
396 return Err(general_err!(
397 "skipped {} values, expected {}",
398 values,
399 values_read
400 ));
401 }
402 }
403 Ok(num_records - remaining_records)
404 }
405
406 fn read_dictionary_page(&mut self) -> Result<()> {
409 match self.page_reader.get_next_page()? {
410 Some(Page::DictionaryPage {
411 buf,
412 num_values,
413 encoding,
414 is_sorted,
415 }) => self
416 .values_decoder
417 .set_dict(buf, num_values, encoding, is_sorted),
418 _ => Err(ParquetError::General(
419 "Invalid page. Expecting dictionary page".to_string(),
420 )),
421 }
422 }
423
424 fn read_new_page(&mut self) -> Result<bool> {
427 loop {
428 match self.page_reader.get_next_page()? {
429 None => return Ok(false),
431 Some(current_page) => {
432 match current_page {
433 Page::DictionaryPage {
435 buf,
436 num_values,
437 encoding,
438 is_sorted,
439 } => {
440 self.values_decoder
441 .set_dict(buf, num_values, encoding, is_sorted)?;
442 }
443 Page::DataPage {
445 buf,
446 num_values,
447 encoding,
448 def_level_encoding,
449 rep_level_encoding,
450 statistics: _,
451 } => {
452 self.num_buffered_values = num_values as _;
453 self.num_decoded_values = 0;
454
455 let max_rep_level = self.descr.max_rep_level();
456 let max_def_level = self.descr.max_def_level();
457
458 let mut offset = 0;
459
460 if max_rep_level > 0 {
461 let (bytes_read, level_data) = parse_v1_level(
462 max_rep_level,
463 num_values,
464 rep_level_encoding,
465 buf.slice(offset..),
466 )?;
467 offset += bytes_read;
468
469 self.has_record_delimiter =
470 self.page_reader.at_record_boundary()?;
471
472 self.rep_level_decoder
473 .as_mut()
474 .unwrap()
475 .set_data(rep_level_encoding, level_data)?;
476 }
477
478 if max_def_level > 0 {
479 let (bytes_read, level_data) = parse_v1_level(
480 max_def_level,
481 num_values,
482 def_level_encoding,
483 buf.slice(offset..),
484 )?;
485 offset += bytes_read;
486
487 self.def_level_decoder
488 .as_mut()
489 .unwrap()
490 .set_data(def_level_encoding, level_data)?;
491 }
492
493 self.values_decoder.set_data(
494 encoding,
495 buf.slice(offset..),
496 num_values as usize,
497 None,
498 )?;
499 return Ok(true);
500 }
501 Page::DataPageV2 {
503 buf,
504 num_values,
505 encoding,
506 num_nulls,
507 num_rows: _,
508 def_levels_byte_len,
509 rep_levels_byte_len,
510 is_compressed: _,
511 statistics: _,
512 } => {
513 if num_nulls > num_values {
514 return Err(general_err!(
515 "more nulls than values in page, contained {} values and {} nulls",
516 num_values,
517 num_nulls
518 ));
519 }
520
521 self.num_buffered_values = num_values as _;
522 self.num_decoded_values = 0;
523
524 if self.descr.max_rep_level() > 0 {
527 self.has_record_delimiter =
531 self.page_reader.at_record_boundary()?;
532
533 self.rep_level_decoder.as_mut().unwrap().set_data(
534 Encoding::RLE,
535 buf.slice(..rep_levels_byte_len as usize),
536 )?;
537 }
538
539 if self.descr.max_def_level() > 0 {
542 self.def_level_decoder.as_mut().unwrap().set_data(
543 Encoding::RLE,
544 buf.slice(
545 rep_levels_byte_len as usize
546 ..(rep_levels_byte_len + def_levels_byte_len) as usize,
547 ),
548 )?;
549 }
550
551 self.values_decoder.set_data(
552 encoding,
553 buf.slice((rep_levels_byte_len + def_levels_byte_len) as usize..),
554 num_values as usize,
555 Some((num_values - num_nulls) as usize),
556 )?;
557 return Ok(true);
558 }
559 }
560 }
561 }
562 }
563 }
564
565 #[inline]
569 pub(crate) fn has_next(&mut self) -> Result<bool> {
570 if self.num_buffered_values == 0 || self.num_buffered_values == self.num_decoded_values {
571 if !self.read_new_page()? {
574 Ok(false)
575 } else {
576 Ok(self.num_buffered_values != 0)
577 }
578 } else {
579 Ok(true)
580 }
581 }
582}
583
584fn check_partial_record_fits(records_read: usize, remaining_records: usize) -> Result<()> {
591 if remaining_records <= records_read {
592 return Err(general_err!(
593 "page ended after {records_read} record(s), which is already all of the \
594 {remaining_records} record(s) asked for, so there is no partial record to flush"
595 ));
596 }
597 Ok(())
598}
599
600fn parse_v1_level(
601 max_level: i16,
602 num_buffered_values: u32,
603 encoding: Encoding,
604 buf: Bytes,
605) -> Result<(usize, Bytes)> {
606 match encoding {
607 Encoding::RLE => {
608 let i32_size = std::mem::size_of::<i32>();
609 if i32_size <= buf.len() {
610 let data_size = read_num_bytes::<i32>(i32_size, buf.as_ref()) as usize;
611 let end = i32_size
612 .checked_add(data_size)
613 .ok_or(general_err!("invalid level length"))?;
614 if end <= buf.len() {
615 return Ok((end, buf.slice(i32_size..end)));
616 }
617 }
618 Err(general_err!("not enough data to read levels"))
619 }
620 #[expect(deprecated)]
621 Encoding::BIT_PACKED => {
622 let bit_width = num_required_bits(max_level as u64);
623 let num_bytes = ceil(num_buffered_values as usize * bit_width as usize, 8);
624 Ok((num_bytes, buf.slice(..num_bytes)))
625 }
626 _ => Err(general_err!("invalid level encoding: {}", encoding)),
627 }
628}
629
630#[cfg(test)]
631mod tests {
632 use super::*;
633
634 use rand::distr::uniform::SampleUniform;
635 use std::{collections::VecDeque, sync::Arc};
636
637 use crate::basic::Type as PhysicalType;
638 use crate::schema::types::{ColumnDescriptor, ColumnPath, Type as SchemaType};
639 use crate::util::test_common::page_util::InMemoryPageReader;
640 use crate::util::test_common::rand_gen::make_pages;
641
642 #[test]
643 fn test_parse_v1_level_invalid_length() {
644 let buf = Bytes::from(vec![10, 0, 0, 0]);
646 let err = parse_v1_level(1, 100, Encoding::RLE, buf).unwrap_err();
647 assert_eq!(
648 err.to_string(),
649 "Parquet error: not enough data to read levels"
650 );
651
652 let buf = Bytes::from(vec![4, 0, 0]);
654 let err = parse_v1_level(1, 100, Encoding::RLE, buf).unwrap_err();
655 assert_eq!(
656 err.to_string(),
657 "Parquet error: not enough data to read levels"
658 );
659 }
660
661 const NUM_LEVELS: usize = 128;
662 const NUM_PAGES: usize = 2;
663 const MAX_DEF_LEVEL: i16 = 5;
664 const MAX_REP_LEVEL: i16 = 5;
665
666 macro_rules! test {
668 ($test_func:ident, i32, $func:ident, $def_level:expr, $rep_level:expr,
670 $num_pages:expr, $num_levels:expr, $batch_size:expr, $min:expr, $max:expr) => {
671 test_internal!(
672 $test_func,
673 Int32Type,
674 get_test_int32_type,
675 $func,
676 $def_level,
677 $rep_level,
678 $num_pages,
679 $num_levels,
680 $batch_size,
681 $min,
682 $max
683 );
684 };
685 ($test_func:ident, i64, $func:ident, $def_level:expr, $rep_level:expr,
687 $num_pages:expr, $num_levels:expr, $batch_size:expr, $min:expr, $max:expr) => {
688 test_internal!(
689 $test_func,
690 Int64Type,
691 get_test_int64_type,
692 $func,
693 $def_level,
694 $rep_level,
695 $num_pages,
696 $num_levels,
697 $batch_size,
698 $min,
699 $max
700 );
701 };
702 }
703
704 macro_rules! test_internal {
705 ($test_func:ident, $ty:ident, $pty:ident, $func:ident, $def_level:expr,
706 $rep_level:expr, $num_pages:expr, $num_levels:expr, $batch_size:expr,
707 $min:expr, $max:expr) => {
708 #[test]
709 fn $test_func() {
710 let desc = Arc::new(ColumnDescriptor::new(
711 Arc::new($pty()),
712 $def_level,
713 $rep_level,
714 ColumnPath::new(Vec::new()),
715 ));
716 let mut tester = ColumnReaderTester::<$ty>::new();
717 tester.$func(desc, $num_pages, $num_levels, $batch_size, $min, $max);
718 }
719 };
720 }
721
722 test!(
723 test_read_plain_v1_int32,
724 i32,
725 plain_v1,
726 MAX_DEF_LEVEL,
727 MAX_REP_LEVEL,
728 NUM_PAGES,
729 NUM_LEVELS,
730 16,
731 i32::MIN,
732 i32::MAX
733 );
734 test!(
735 test_read_plain_v2_int32,
736 i32,
737 plain_v2,
738 MAX_DEF_LEVEL,
739 MAX_REP_LEVEL,
740 NUM_PAGES,
741 NUM_LEVELS,
742 16,
743 i32::MIN,
744 i32::MAX
745 );
746
747 test!(
748 test_read_plain_v1_int32_uneven,
749 i32,
750 plain_v1,
751 MAX_DEF_LEVEL,
752 MAX_REP_LEVEL,
753 NUM_PAGES,
754 NUM_LEVELS,
755 17,
756 i32::MIN,
757 i32::MAX
758 );
759 test!(
760 test_read_plain_v2_int32_uneven,
761 i32,
762 plain_v2,
763 MAX_DEF_LEVEL,
764 MAX_REP_LEVEL,
765 NUM_PAGES,
766 NUM_LEVELS,
767 17,
768 i32::MIN,
769 i32::MAX
770 );
771
772 test!(
773 test_read_plain_v1_int32_multi_page,
774 i32,
775 plain_v1,
776 MAX_DEF_LEVEL,
777 MAX_REP_LEVEL,
778 NUM_PAGES,
779 NUM_LEVELS,
780 512,
781 i32::MIN,
782 i32::MAX
783 );
784 test!(
785 test_read_plain_v2_int32_multi_page,
786 i32,
787 plain_v2,
788 MAX_DEF_LEVEL,
789 MAX_REP_LEVEL,
790 NUM_PAGES,
791 NUM_LEVELS,
792 512,
793 i32::MIN,
794 i32::MAX
795 );
796
797 test!(
799 test_read_plain_v1_int32_required_non_repeated,
800 i32,
801 plain_v1,
802 0,
803 0,
804 NUM_PAGES,
805 NUM_LEVELS,
806 16,
807 i32::MIN,
808 i32::MAX
809 );
810 test!(
811 test_read_plain_v2_int32_required_non_repeated,
812 i32,
813 plain_v2,
814 0,
815 0,
816 NUM_PAGES,
817 NUM_LEVELS,
818 16,
819 i32::MIN,
820 i32::MAX
821 );
822
823 test!(
824 test_read_plain_v1_int64,
825 i64,
826 plain_v1,
827 1,
828 1,
829 NUM_PAGES,
830 NUM_LEVELS,
831 16,
832 i64::MIN,
833 i64::MAX
834 );
835 test!(
836 test_read_plain_v2_int64,
837 i64,
838 plain_v2,
839 1,
840 1,
841 NUM_PAGES,
842 NUM_LEVELS,
843 16,
844 i64::MIN,
845 i64::MAX
846 );
847
848 test!(
849 test_read_plain_v1_int64_uneven,
850 i64,
851 plain_v1,
852 1,
853 1,
854 NUM_PAGES,
855 NUM_LEVELS,
856 17,
857 i64::MIN,
858 i64::MAX
859 );
860 test!(
861 test_read_plain_v2_int64_uneven,
862 i64,
863 plain_v2,
864 1,
865 1,
866 NUM_PAGES,
867 NUM_LEVELS,
868 17,
869 i64::MIN,
870 i64::MAX
871 );
872
873 test!(
874 test_read_plain_v1_int64_multi_page,
875 i64,
876 plain_v1,
877 1,
878 1,
879 NUM_PAGES,
880 NUM_LEVELS,
881 512,
882 i64::MIN,
883 i64::MAX
884 );
885 test!(
886 test_read_plain_v2_int64_multi_page,
887 i64,
888 plain_v2,
889 1,
890 1,
891 NUM_PAGES,
892 NUM_LEVELS,
893 512,
894 i64::MIN,
895 i64::MAX
896 );
897
898 test!(
900 test_read_plain_v1_int64_required_non_repeated,
901 i64,
902 plain_v1,
903 0,
904 0,
905 NUM_PAGES,
906 NUM_LEVELS,
907 16,
908 i64::MIN,
909 i64::MAX
910 );
911 test!(
912 test_read_plain_v2_int64_required_non_repeated,
913 i64,
914 plain_v2,
915 0,
916 0,
917 NUM_PAGES,
918 NUM_LEVELS,
919 16,
920 i64::MIN,
921 i64::MAX
922 );
923
924 test!(
925 test_read_dict_v1_int32_small,
926 i32,
927 dict_v1,
928 MAX_DEF_LEVEL,
929 MAX_REP_LEVEL,
930 2,
931 2,
932 16,
933 0,
934 3
935 );
936 test!(
937 test_read_dict_v2_int32_small,
938 i32,
939 dict_v2,
940 MAX_DEF_LEVEL,
941 MAX_REP_LEVEL,
942 2,
943 2,
944 16,
945 0,
946 3
947 );
948
949 test!(
950 test_read_dict_v1_int32,
951 i32,
952 dict_v1,
953 MAX_DEF_LEVEL,
954 MAX_REP_LEVEL,
955 NUM_PAGES,
956 NUM_LEVELS,
957 16,
958 0,
959 3
960 );
961 test!(
962 test_read_dict_v2_int32,
963 i32,
964 dict_v2,
965 MAX_DEF_LEVEL,
966 MAX_REP_LEVEL,
967 NUM_PAGES,
968 NUM_LEVELS,
969 16,
970 0,
971 3
972 );
973
974 test!(
975 test_read_dict_v1_int32_uneven,
976 i32,
977 dict_v1,
978 MAX_DEF_LEVEL,
979 MAX_REP_LEVEL,
980 NUM_PAGES,
981 NUM_LEVELS,
982 17,
983 0,
984 3
985 );
986 test!(
987 test_read_dict_v2_int32_uneven,
988 i32,
989 dict_v2,
990 MAX_DEF_LEVEL,
991 MAX_REP_LEVEL,
992 NUM_PAGES,
993 NUM_LEVELS,
994 17,
995 0,
996 3
997 );
998
999 test!(
1000 test_read_dict_v1_int32_multi_page,
1001 i32,
1002 dict_v1,
1003 MAX_DEF_LEVEL,
1004 MAX_REP_LEVEL,
1005 NUM_PAGES,
1006 NUM_LEVELS,
1007 512,
1008 0,
1009 3
1010 );
1011 test!(
1012 test_read_dict_v2_int32_multi_page,
1013 i32,
1014 dict_v2,
1015 MAX_DEF_LEVEL,
1016 MAX_REP_LEVEL,
1017 NUM_PAGES,
1018 NUM_LEVELS,
1019 512,
1020 0,
1021 3
1022 );
1023
1024 test!(
1025 test_read_dict_v1_int64,
1026 i64,
1027 dict_v1,
1028 MAX_DEF_LEVEL,
1029 MAX_REP_LEVEL,
1030 NUM_PAGES,
1031 NUM_LEVELS,
1032 16,
1033 0,
1034 3
1035 );
1036 test!(
1037 test_read_dict_v2_int64,
1038 i64,
1039 dict_v2,
1040 MAX_DEF_LEVEL,
1041 MAX_REP_LEVEL,
1042 NUM_PAGES,
1043 NUM_LEVELS,
1044 16,
1045 0,
1046 3
1047 );
1048
1049 #[test]
1050 fn test_read_batch_values_only() {
1051 test_read_batch_int32(16, 0, 0);
1052 }
1053
1054 #[test]
1055 fn test_read_batch_values_def_levels() {
1056 test_read_batch_int32(16, MAX_DEF_LEVEL, 0);
1057 }
1058
1059 #[test]
1060 fn test_read_batch_values_rep_levels() {
1061 test_read_batch_int32(16, 0, MAX_REP_LEVEL);
1062 }
1063
1064 #[test]
1065 fn test_read_batch_values_def_rep_levels() {
1066 test_read_batch_int32(128, MAX_DEF_LEVEL, MAX_REP_LEVEL);
1067 }
1068
1069 #[test]
1070 fn test_read_batch_adjust_after_buffering_page() {
1071 let primitive_type = get_test_int32_type();
1078 let desc = Arc::new(ColumnDescriptor::new(
1079 Arc::new(primitive_type),
1080 1,
1081 1,
1082 ColumnPath::new(Vec::new()),
1083 ));
1084
1085 let num_pages = 2;
1086 let num_levels = 4;
1087 let batch_size = 5;
1088
1089 let mut tester = ColumnReaderTester::<Int32Type>::new();
1090 tester.test_read_batch(
1091 desc,
1092 Encoding::RLE_DICTIONARY,
1093 num_pages,
1094 num_levels,
1095 batch_size,
1096 i32::MIN,
1097 i32::MAX,
1098 false,
1099 );
1100 }
1101
1102 fn get_test_int32_type() -> SchemaType {
1148 SchemaType::primitive_type_builder("a", PhysicalType::INT32)
1149 .with_repetition(Repetition::REQUIRED)
1150 .with_converted_type(ConvertedType::INT_32)
1151 .with_length(-1)
1152 .build()
1153 .expect("build() should be OK")
1154 }
1155
1156 fn get_test_int64_type() -> SchemaType {
1158 SchemaType::primitive_type_builder("a", PhysicalType::INT64)
1159 .with_repetition(Repetition::REQUIRED)
1160 .with_converted_type(ConvertedType::INT_64)
1161 .with_length(-1)
1162 .build()
1163 .expect("build() should be OK")
1164 }
1165
1166 fn test_read_batch_int32(batch_size: usize, max_def_level: i16, max_rep_level: i16) {
1171 let primitive_type = get_test_int32_type();
1172
1173 let desc = Arc::new(ColumnDescriptor::new(
1174 Arc::new(primitive_type),
1175 max_def_level,
1176 max_rep_level,
1177 ColumnPath::new(Vec::new()),
1178 ));
1179
1180 let mut tester = ColumnReaderTester::<Int32Type>::new();
1181 tester.test_read_batch(
1182 desc,
1183 Encoding::RLE_DICTIONARY,
1184 NUM_PAGES,
1185 NUM_LEVELS,
1186 batch_size,
1187 i32::MIN,
1188 i32::MAX,
1189 false,
1190 );
1191 }
1192
1193 struct ColumnReaderTester<T: DataType>
1194 where
1195 T::T: PartialOrd + SampleUniform + Copy,
1196 {
1197 rep_levels: Vec<i16>,
1198 def_levels: Vec<i16>,
1199 values: Vec<T::T>,
1200 }
1201
1202 impl<T: DataType> ColumnReaderTester<T>
1203 where
1204 T::T: PartialOrd + SampleUniform + Copy,
1205 {
1206 pub fn new() -> Self {
1207 Self {
1208 rep_levels: Vec::new(),
1209 def_levels: Vec::new(),
1210 values: Vec::new(),
1211 }
1212 }
1213
1214 fn plain_v1(
1216 &mut self,
1217 desc: ColumnDescPtr,
1218 num_pages: usize,
1219 num_levels: usize,
1220 batch_size: usize,
1221 min: T::T,
1222 max: T::T,
1223 ) {
1224 self.test_read_batch_general(
1225 desc,
1226 Encoding::PLAIN,
1227 num_pages,
1228 num_levels,
1229 batch_size,
1230 min,
1231 max,
1232 false,
1233 );
1234 }
1235
1236 fn plain_v2(
1238 &mut self,
1239 desc: ColumnDescPtr,
1240 num_pages: usize,
1241 num_levels: usize,
1242 batch_size: usize,
1243 min: T::T,
1244 max: T::T,
1245 ) {
1246 self.test_read_batch_general(
1247 desc,
1248 Encoding::PLAIN,
1249 num_pages,
1250 num_levels,
1251 batch_size,
1252 min,
1253 max,
1254 true,
1255 );
1256 }
1257
1258 fn dict_v1(
1260 &mut self,
1261 desc: ColumnDescPtr,
1262 num_pages: usize,
1263 num_levels: usize,
1264 batch_size: usize,
1265 min: T::T,
1266 max: T::T,
1267 ) {
1268 self.test_read_batch_general(
1269 desc,
1270 Encoding::RLE_DICTIONARY,
1271 num_pages,
1272 num_levels,
1273 batch_size,
1274 min,
1275 max,
1276 false,
1277 );
1278 }
1279
1280 fn dict_v2(
1282 &mut self,
1283 desc: ColumnDescPtr,
1284 num_pages: usize,
1285 num_levels: usize,
1286 batch_size: usize,
1287 min: T::T,
1288 max: T::T,
1289 ) {
1290 self.test_read_batch_general(
1291 desc,
1292 Encoding::RLE_DICTIONARY,
1293 num_pages,
1294 num_levels,
1295 batch_size,
1296 min,
1297 max,
1298 true,
1299 );
1300 }
1301
1302 #[expect(clippy::too_many_arguments)]
1305 fn test_read_batch_general(
1306 &mut self,
1307 desc: ColumnDescPtr,
1308 encoding: Encoding,
1309 num_pages: usize,
1310 num_levels: usize,
1311 batch_size: usize,
1312 min: T::T,
1313 max: T::T,
1314 use_v2: bool,
1315 ) {
1316 self.test_read_batch(
1317 desc, encoding, num_pages, num_levels, batch_size, min, max, use_v2,
1318 );
1319 }
1320
1321 #[expect(clippy::too_many_arguments)]
1324 fn test_read_batch(
1325 &mut self,
1326 desc: ColumnDescPtr,
1327 encoding: Encoding,
1328 num_pages: usize,
1329 num_levels: usize,
1330 batch_size: usize,
1331 min: T::T,
1332 max: T::T,
1333 use_v2: bool,
1334 ) {
1335 let mut pages = VecDeque::new();
1336 make_pages::<T>(
1337 desc.clone(),
1338 encoding,
1339 num_pages,
1340 num_levels,
1341 min,
1342 max,
1343 &mut self.def_levels,
1344 &mut self.rep_levels,
1345 &mut self.values,
1346 &mut pages,
1347 use_v2,
1348 );
1349 let max_def_level = desc.max_def_level();
1350 let max_rep_level = desc.max_rep_level();
1351 let page_reader = InMemoryPageReader::new(pages);
1352 let column_reader = get_column_reader(desc, Box::new(page_reader));
1353 let mut typed_column_reader = get_typed_column_reader::<T>(column_reader);
1354
1355 let mut values = Vec::new();
1356 let mut def_levels = Vec::new();
1357 let mut rep_levels = Vec::new();
1358
1359 let mut curr_values_read = 0;
1360 let mut curr_levels_read = 0;
1361 loop {
1362 let (_, values_read, levels_read) = typed_column_reader
1363 .read_records(
1364 batch_size,
1365 Some(&mut def_levels),
1366 Some(&mut rep_levels),
1367 &mut values,
1368 )
1369 .expect("read_batch() should be OK");
1370
1371 curr_values_read += values_read;
1372 curr_levels_read += levels_read;
1373
1374 if values_read == 0 && levels_read == 0 {
1375 break;
1376 }
1377 }
1378
1379 assert_eq!(values, self.values, "values content doesn't match");
1380
1381 if max_def_level > 0 {
1382 assert_eq!(
1383 def_levels, self.def_levels,
1384 "definition levels content doesn't match"
1385 );
1386 }
1387
1388 if max_rep_level > 0 {
1389 assert_eq!(
1390 rep_levels, self.rep_levels,
1391 "repetition levels content doesn't match"
1392 );
1393 }
1394
1395 assert!(
1396 curr_levels_read >= curr_values_read,
1397 "expected levels read to be greater than values read"
1398 );
1399 }
1400 }
1401
1402 #[test]
1420 fn test_skip_records_v2_page_skip_accounts_for_partial() {
1421 use crate::encodings::levels::LevelEncoder;
1422
1423 let max_rep_level: i16 = 1;
1424 let max_def_level: i16 = 1;
1425
1426 let primitive_type = SchemaType::primitive_type_builder("element", PhysicalType::INT32)
1428 .with_repetition(Repetition::REQUIRED)
1429 .build()
1430 .unwrap();
1431 let desc = Arc::new(ColumnDescriptor::new(
1432 Arc::new(primitive_type),
1433 max_def_level,
1434 max_rep_level,
1435 ColumnPath::new(vec!["list".to_string(), "element".to_string()]),
1436 ));
1437
1438 let make_v2_page =
1440 |rep_levels: &[i16], def_levels: &[i16], values: &[i32], num_rows: u32| -> Page {
1441 let mut rep_enc = LevelEncoder::v2_streaming(max_rep_level);
1442 rep_enc.put_with_observer(rep_levels, |_, _| {});
1443 let rep_bytes = rep_enc.consume();
1444
1445 let mut def_enc = LevelEncoder::v2_streaming(max_def_level);
1446 def_enc.put_with_observer(def_levels, |_, _| {});
1447 let def_bytes = def_enc.consume();
1448
1449 let val_bytes: Vec<u8> = values.iter().flat_map(|v| v.to_le_bytes()).collect();
1450
1451 let mut buf = Vec::new();
1452 buf.extend_from_slice(&rep_bytes);
1453 buf.extend_from_slice(&def_bytes);
1454 buf.extend_from_slice(&val_bytes);
1455
1456 Page::DataPageV2 {
1457 buf: Bytes::from(buf),
1458 num_values: rep_levels.len() as u32,
1459 encoding: Encoding::PLAIN,
1460 num_nulls: 0,
1461 num_rows,
1462 def_levels_byte_len: def_bytes.len() as u32,
1463 rep_levels_byte_len: rep_bytes.len() as u32,
1464 is_compressed: false,
1465 statistics: None,
1466 }
1467 };
1468
1469 let page1 = make_v2_page(&[0, 1, 0, 1], &[1, 1, 1, 1], &[10, 20, 30, 40], 2);
1475
1476 let page2 = make_v2_page(&[0, 1, 0, 1], &[1, 1, 1, 1], &[50, 60, 70, 80], 2);
1478
1479 let page3 = make_v2_page(&[0, 1], &[1, 1], &[90, 100], 1);
1481
1482 let pages = VecDeque::from(vec![page1, page2, page3]);
1484 let page_reader = InMemoryPageReader::new(pages);
1485 let column_reader = get_column_reader(desc, Box::new(page_reader));
1486 let mut typed_reader = get_typed_column_reader::<Int32Type>(column_reader);
1487
1488 let skipped = typed_reader.skip_records(1).unwrap();
1494 assert_eq!(skipped, 1);
1495
1496 let skipped = typed_reader.skip_records(2).unwrap();
1511 assert_eq!(skipped, 2);
1512
1513 let mut values = Vec::new();
1515 let mut def_levels = Vec::new();
1516 let mut rep_levels = Vec::new();
1517
1518 let (records, values_read, levels_read) = typed_reader
1519 .read_records(1, Some(&mut def_levels), Some(&mut rep_levels), &mut values)
1520 .unwrap();
1521
1522 assert_eq!(records, 1, "should read exactly 1 record");
1526 assert_eq!(levels_read, 2, "should read 2 levels for the record");
1527 assert_eq!(values_read, 2, "should read 2 non-null values");
1528 assert_eq!(values, vec![70, 80], "should contain 4th record's values");
1529 assert_eq!(rep_levels, vec![0, 1], "rep levels for a 2-element list");
1530 assert_eq!(def_levels, vec![1, 1], "def levels (all non-null)");
1531 }
1532}