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> {
94 T::get_column_reader(col_reader).unwrap_or_else(|| {
95 panic!(
96 "Failed to convert column reader into a typed column reader for `{}` type",
97 T::get_physical_type()
98 )
99 })
100}
101
102pub type ColumnReaderImpl<T> = GenericColumnReader<
104 RepetitionLevelDecoderImpl,
105 DefinitionLevelDecoderImpl,
106 ColumnValueDecoderImpl<T>,
107>;
108
109pub struct GenericColumnReader<R, D, V> {
115 descr: ColumnDescPtr,
116
117 page_reader: Box<dyn PageReader>,
118
119 num_buffered_values: usize,
121
122 num_decoded_values: usize,
125
126 has_record_delimiter: bool,
128
129 def_level_decoder: Option<D>,
131
132 rep_level_decoder: Option<R>,
134
135 values_decoder: V,
137}
138
139impl<V> GenericColumnReader<RepetitionLevelDecoderImpl, DefinitionLevelDecoderImpl, V>
140where
141 V: ColumnValueDecoder,
142{
143 pub fn new(descr: ColumnDescPtr, page_reader: Box<dyn PageReader>) -> Self {
145 let values_decoder = V::new(&descr);
146
147 let def_level_decoder = (descr.max_def_level() != 0)
148 .then(|| DefinitionLevelDecoderImpl::new(descr.max_def_level()));
149
150 let rep_level_decoder = (descr.max_rep_level() != 0)
151 .then(|| RepetitionLevelDecoderImpl::new(descr.max_rep_level()));
152
153 Self::new_with_decoders(
154 descr,
155 page_reader,
156 values_decoder,
157 def_level_decoder,
158 rep_level_decoder,
159 )
160 }
161}
162
163impl<R, D, V> GenericColumnReader<R, D, V>
164where
165 R: RepetitionLevelDecoder,
166 D: DefinitionLevelDecoder,
167 V: ColumnValueDecoder,
168{
169 pub(crate) fn new_with_decoders(
170 descr: ColumnDescPtr,
171 page_reader: Box<dyn PageReader>,
172 values_decoder: V,
173 def_level_decoder: Option<D>,
174 rep_level_decoder: Option<R>,
175 ) -> Self {
176 Self {
177 descr,
178 def_level_decoder,
179 rep_level_decoder,
180 page_reader,
181 num_buffered_values: 0,
182 num_decoded_values: 0,
183 values_decoder,
184 has_record_delimiter: false,
185 }
186 }
187
188 pub fn read_records(
203 &mut self,
204 max_records: usize,
205 def_levels: Option<&mut D::Buffer>,
206 rep_levels: Option<&mut R::Buffer>,
207 values: &mut V::Buffer,
208 ) -> Result<(usize, usize, usize)> {
209 self.read_records_with_reservation(
210 max_records,
211 def_levels,
212 rep_levels,
213 values,
214 |_, _, _, _| Ok(()),
215 )
216 }
217
218 pub(crate) fn read_records_with_reservation<F>(
221 &mut self,
222 max_records: usize,
223 mut def_levels: Option<&mut D::Buffer>,
224 mut rep_levels: Option<&mut R::Buffer>,
225 values: &mut V::Buffer,
226 mut reserve_values: F,
227 ) -> Result<(usize, usize, usize)>
228 where
229 F: FnMut(&mut V::Buffer, usize, usize, Option<&D::Buffer>) -> Result<()>,
230 {
231 let mut total_records_read = 0;
232 let mut total_levels_read = 0;
233 let mut total_values_read = 0;
234
235 while total_records_read < max_records && self.has_next()? {
236 let remaining_records = max_records - total_records_read;
237 let remaining_levels = self.num_buffered_values - self.num_decoded_values;
238
239 let (records_read, levels_to_read) = match self.rep_level_decoder.as_mut() {
240 Some(reader) => {
241 let out = rep_levels
242 .as_mut()
243 .ok_or_else(|| general_err!("must specify repetition levels"))?;
244
245 let (mut records_read, levels_read) =
246 reader.read_rep_levels(out, remaining_records, remaining_levels)?;
247
248 if records_read == 0 && levels_read == 0 {
249 return Err(general_err!(
251 "Insufficient repetition levels read from column"
252 ));
253 }
254 if levels_read == remaining_levels && self.has_record_delimiter {
255 assert!(records_read < remaining_records); 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 metadata = match self.page_reader.peek_next_page()? {
317 None => return Ok(num_records - remaining_records),
318 Some(metadata) => metadata,
319 };
320
321 if metadata.is_dict {
323 self.read_dictionary_page()?;
324 continue;
325 }
326
327 let rows = metadata.num_rows.or_else(|| {
330 self.rep_level_decoder
332 .is_none()
333 .then_some(metadata.num_levels)?
334 });
335
336 if let Some(rows) = rows {
337 if rows <= remaining_records {
338 self.page_reader.skip_next_page()?;
339 remaining_records -= rows;
340 continue;
341 }
342 }
343 if !self.read_new_page()? {
346 return Ok(num_records - remaining_records);
347 }
348 }
349
350 let remaining_levels = self.num_buffered_values - self.num_decoded_values;
354
355 let (records_read, rep_levels_read) = match self.rep_level_decoder.as_mut() {
356 Some(decoder) => {
357 let (mut records_read, levels_read) =
358 decoder.skip_rep_levels(remaining_records, remaining_levels)?;
359
360 if levels_read == remaining_levels && self.has_record_delimiter {
361 assert!(records_read < remaining_records); records_read += decoder.flush_partial() as usize;
365 }
366
367 (records_read, levels_read)
368 }
369 None => {
370 let levels = remaining_levels.min(remaining_records);
372 (levels, levels)
373 }
374 };
375
376 self.num_decoded_values += rep_levels_read;
377 remaining_records -= records_read;
378
379 if self.num_buffered_values == self.num_decoded_values {
380 continue;
382 }
383
384 let (values_read, def_levels_read) = match self.def_level_decoder.as_mut() {
385 Some(decoder) => decoder.skip_def_levels(rep_levels_read)?,
386 None => (rep_levels_read, rep_levels_read),
387 };
388
389 if rep_levels_read != def_levels_read {
390 return Err(general_err!(
391 "levels mismatch, read {} repetition levels and {} definition levels",
392 rep_levels_read,
393 def_levels_read
394 ));
395 }
396
397 let values = self.values_decoder.skip_values(values_read)?;
398 if values != values_read {
399 return Err(general_err!(
400 "skipped {} values, expected {}",
401 values,
402 values_read
403 ));
404 }
405 }
406 Ok(num_records - remaining_records)
407 }
408
409 fn read_dictionary_page(&mut self) -> Result<()> {
412 match self.page_reader.get_next_page()? {
413 Some(Page::DictionaryPage {
414 buf,
415 num_values,
416 encoding,
417 is_sorted,
418 }) => self
419 .values_decoder
420 .set_dict(buf, num_values, encoding, is_sorted),
421 _ => Err(ParquetError::General(
422 "Invalid page. Expecting dictionary page".to_string(),
423 )),
424 }
425 }
426
427 fn read_new_page(&mut self) -> Result<bool> {
430 loop {
431 match self.page_reader.get_next_page()? {
432 None => return Ok(false),
434 Some(current_page) => {
435 match current_page {
436 Page::DictionaryPage {
438 buf,
439 num_values,
440 encoding,
441 is_sorted,
442 } => {
443 self.values_decoder
444 .set_dict(buf, num_values, encoding, is_sorted)?;
445 continue;
446 }
447 Page::DataPage {
449 buf,
450 num_values,
451 encoding,
452 def_level_encoding,
453 rep_level_encoding,
454 statistics: _,
455 } => {
456 self.num_buffered_values = num_values as _;
457 self.num_decoded_values = 0;
458
459 let max_rep_level = self.descr.max_rep_level();
460 let max_def_level = self.descr.max_def_level();
461
462 let mut offset = 0;
463
464 if max_rep_level > 0 {
465 let (bytes_read, level_data) = parse_v1_level(
466 max_rep_level,
467 num_values,
468 rep_level_encoding,
469 buf.slice(offset..),
470 )?;
471 offset += bytes_read;
472
473 self.has_record_delimiter =
474 self.page_reader.at_record_boundary()?;
475
476 self.rep_level_decoder
477 .as_mut()
478 .unwrap()
479 .set_data(rep_level_encoding, level_data)?;
480 }
481
482 if max_def_level > 0 {
483 let (bytes_read, level_data) = parse_v1_level(
484 max_def_level,
485 num_values,
486 def_level_encoding,
487 buf.slice(offset..),
488 )?;
489 offset += bytes_read;
490
491 self.def_level_decoder
492 .as_mut()
493 .unwrap()
494 .set_data(def_level_encoding, level_data)?;
495 }
496
497 self.values_decoder.set_data(
498 encoding,
499 buf.slice(offset..),
500 num_values as usize,
501 None,
502 )?;
503 return Ok(true);
504 }
505 Page::DataPageV2 {
507 buf,
508 num_values,
509 encoding,
510 num_nulls,
511 num_rows: _,
512 def_levels_byte_len,
513 rep_levels_byte_len,
514 is_compressed: _,
515 statistics: _,
516 } => {
517 if num_nulls > num_values {
518 return Err(general_err!(
519 "more nulls than values in page, contained {} values and {} nulls",
520 num_values,
521 num_nulls
522 ));
523 }
524
525 self.num_buffered_values = num_values as _;
526 self.num_decoded_values = 0;
527
528 if self.descr.max_rep_level() > 0 {
531 self.has_record_delimiter =
535 self.page_reader.at_record_boundary()?;
536
537 self.rep_level_decoder.as_mut().unwrap().set_data(
538 Encoding::RLE,
539 buf.slice(..rep_levels_byte_len as usize),
540 )?;
541 }
542
543 if self.descr.max_def_level() > 0 {
546 self.def_level_decoder.as_mut().unwrap().set_data(
547 Encoding::RLE,
548 buf.slice(
549 rep_levels_byte_len as usize
550 ..(rep_levels_byte_len + def_levels_byte_len) as usize,
551 ),
552 )?;
553 }
554
555 self.values_decoder.set_data(
556 encoding,
557 buf.slice((rep_levels_byte_len + def_levels_byte_len) as usize..),
558 num_values as usize,
559 Some((num_values - num_nulls) as usize),
560 )?;
561 return Ok(true);
562 }
563 };
564 }
565 }
566 }
567 }
568
569 #[inline]
573 pub(crate) fn has_next(&mut self) -> Result<bool> {
574 if self.num_buffered_values == 0 || self.num_buffered_values == self.num_decoded_values {
575 if !self.read_new_page()? {
578 Ok(false)
579 } else {
580 Ok(self.num_buffered_values != 0)
581 }
582 } else {
583 Ok(true)
584 }
585 }
586}
587
588fn parse_v1_level(
589 max_level: i16,
590 num_buffered_values: u32,
591 encoding: Encoding,
592 buf: Bytes,
593) -> Result<(usize, Bytes)> {
594 match encoding {
595 Encoding::RLE => {
596 let i32_size = std::mem::size_of::<i32>();
597 if i32_size <= buf.len() {
598 let data_size = read_num_bytes::<i32>(i32_size, buf.as_ref()) as usize;
599 let end = i32_size
600 .checked_add(data_size)
601 .ok_or(general_err!("invalid level length"))?;
602 if end <= buf.len() {
603 return Ok((end, buf.slice(i32_size..end)));
604 }
605 }
606 Err(general_err!("not enough data to read levels"))
607 }
608 #[allow(deprecated)]
609 Encoding::BIT_PACKED => {
610 let bit_width = num_required_bits(max_level as u64);
611 let num_bytes = ceil(num_buffered_values as usize * bit_width as usize, 8);
612 Ok((num_bytes, buf.slice(..num_bytes)))
613 }
614 _ => Err(general_err!("invalid level encoding: {}", encoding)),
615 }
616}
617
618#[cfg(test)]
619mod tests {
620 use super::*;
621
622 use rand::distr::uniform::SampleUniform;
623 use std::{collections::VecDeque, sync::Arc};
624
625 use crate::basic::Type as PhysicalType;
626 use crate::schema::types::{ColumnDescriptor, ColumnPath, Type as SchemaType};
627 use crate::util::test_common::page_util::InMemoryPageReader;
628 use crate::util::test_common::rand_gen::make_pages;
629
630 #[test]
631 fn test_parse_v1_level_invalid_length() {
632 let buf = Bytes::from(vec![10, 0, 0, 0]);
634 let err = parse_v1_level(1, 100, Encoding::RLE, buf).unwrap_err();
635 assert_eq!(
636 err.to_string(),
637 "Parquet error: not enough data to read levels"
638 );
639
640 let buf = Bytes::from(vec![4, 0, 0]);
642 let err = parse_v1_level(1, 100, Encoding::RLE, buf).unwrap_err();
643 assert_eq!(
644 err.to_string(),
645 "Parquet error: not enough data to read levels"
646 );
647 }
648
649 const NUM_LEVELS: usize = 128;
650 const NUM_PAGES: usize = 2;
651 const MAX_DEF_LEVEL: i16 = 5;
652 const MAX_REP_LEVEL: i16 = 5;
653
654 macro_rules! test {
656 ($test_func:ident, i32, $func:ident, $def_level:expr, $rep_level:expr,
658 $num_pages:expr, $num_levels:expr, $batch_size:expr, $min:expr, $max:expr) => {
659 test_internal!(
660 $test_func,
661 Int32Type,
662 get_test_int32_type,
663 $func,
664 $def_level,
665 $rep_level,
666 $num_pages,
667 $num_levels,
668 $batch_size,
669 $min,
670 $max
671 );
672 };
673 ($test_func:ident, i64, $func:ident, $def_level:expr, $rep_level:expr,
675 $num_pages:expr, $num_levels:expr, $batch_size:expr, $min:expr, $max:expr) => {
676 test_internal!(
677 $test_func,
678 Int64Type,
679 get_test_int64_type,
680 $func,
681 $def_level,
682 $rep_level,
683 $num_pages,
684 $num_levels,
685 $batch_size,
686 $min,
687 $max
688 );
689 };
690 }
691
692 macro_rules! test_internal {
693 ($test_func:ident, $ty:ident, $pty:ident, $func:ident, $def_level:expr,
694 $rep_level:expr, $num_pages:expr, $num_levels:expr, $batch_size:expr,
695 $min:expr, $max:expr) => {
696 #[test]
697 fn $test_func() {
698 let desc = Arc::new(ColumnDescriptor::new(
699 Arc::new($pty()),
700 $def_level,
701 $rep_level,
702 ColumnPath::new(Vec::new()),
703 ));
704 let mut tester = ColumnReaderTester::<$ty>::new();
705 tester.$func(desc, $num_pages, $num_levels, $batch_size, $min, $max);
706 }
707 };
708 }
709
710 test!(
711 test_read_plain_v1_int32,
712 i32,
713 plain_v1,
714 MAX_DEF_LEVEL,
715 MAX_REP_LEVEL,
716 NUM_PAGES,
717 NUM_LEVELS,
718 16,
719 i32::MIN,
720 i32::MAX
721 );
722 test!(
723 test_read_plain_v2_int32,
724 i32,
725 plain_v2,
726 MAX_DEF_LEVEL,
727 MAX_REP_LEVEL,
728 NUM_PAGES,
729 NUM_LEVELS,
730 16,
731 i32::MIN,
732 i32::MAX
733 );
734
735 test!(
736 test_read_plain_v1_int32_uneven,
737 i32,
738 plain_v1,
739 MAX_DEF_LEVEL,
740 MAX_REP_LEVEL,
741 NUM_PAGES,
742 NUM_LEVELS,
743 17,
744 i32::MIN,
745 i32::MAX
746 );
747 test!(
748 test_read_plain_v2_int32_uneven,
749 i32,
750 plain_v2,
751 MAX_DEF_LEVEL,
752 MAX_REP_LEVEL,
753 NUM_PAGES,
754 NUM_LEVELS,
755 17,
756 i32::MIN,
757 i32::MAX
758 );
759
760 test!(
761 test_read_plain_v1_int32_multi_page,
762 i32,
763 plain_v1,
764 MAX_DEF_LEVEL,
765 MAX_REP_LEVEL,
766 NUM_PAGES,
767 NUM_LEVELS,
768 512,
769 i32::MIN,
770 i32::MAX
771 );
772 test!(
773 test_read_plain_v2_int32_multi_page,
774 i32,
775 plain_v2,
776 MAX_DEF_LEVEL,
777 MAX_REP_LEVEL,
778 NUM_PAGES,
779 NUM_LEVELS,
780 512,
781 i32::MIN,
782 i32::MAX
783 );
784
785 test!(
787 test_read_plain_v1_int32_required_non_repeated,
788 i32,
789 plain_v1,
790 0,
791 0,
792 NUM_PAGES,
793 NUM_LEVELS,
794 16,
795 i32::MIN,
796 i32::MAX
797 );
798 test!(
799 test_read_plain_v2_int32_required_non_repeated,
800 i32,
801 plain_v2,
802 0,
803 0,
804 NUM_PAGES,
805 NUM_LEVELS,
806 16,
807 i32::MIN,
808 i32::MAX
809 );
810
811 test!(
812 test_read_plain_v1_int64,
813 i64,
814 plain_v1,
815 1,
816 1,
817 NUM_PAGES,
818 NUM_LEVELS,
819 16,
820 i64::MIN,
821 i64::MAX
822 );
823 test!(
824 test_read_plain_v2_int64,
825 i64,
826 plain_v2,
827 1,
828 1,
829 NUM_PAGES,
830 NUM_LEVELS,
831 16,
832 i64::MIN,
833 i64::MAX
834 );
835
836 test!(
837 test_read_plain_v1_int64_uneven,
838 i64,
839 plain_v1,
840 1,
841 1,
842 NUM_PAGES,
843 NUM_LEVELS,
844 17,
845 i64::MIN,
846 i64::MAX
847 );
848 test!(
849 test_read_plain_v2_int64_uneven,
850 i64,
851 plain_v2,
852 1,
853 1,
854 NUM_PAGES,
855 NUM_LEVELS,
856 17,
857 i64::MIN,
858 i64::MAX
859 );
860
861 test!(
862 test_read_plain_v1_int64_multi_page,
863 i64,
864 plain_v1,
865 1,
866 1,
867 NUM_PAGES,
868 NUM_LEVELS,
869 512,
870 i64::MIN,
871 i64::MAX
872 );
873 test!(
874 test_read_plain_v2_int64_multi_page,
875 i64,
876 plain_v2,
877 1,
878 1,
879 NUM_PAGES,
880 NUM_LEVELS,
881 512,
882 i64::MIN,
883 i64::MAX
884 );
885
886 test!(
888 test_read_plain_v1_int64_required_non_repeated,
889 i64,
890 plain_v1,
891 0,
892 0,
893 NUM_PAGES,
894 NUM_LEVELS,
895 16,
896 i64::MIN,
897 i64::MAX
898 );
899 test!(
900 test_read_plain_v2_int64_required_non_repeated,
901 i64,
902 plain_v2,
903 0,
904 0,
905 NUM_PAGES,
906 NUM_LEVELS,
907 16,
908 i64::MIN,
909 i64::MAX
910 );
911
912 test!(
913 test_read_dict_v1_int32_small,
914 i32,
915 dict_v1,
916 MAX_DEF_LEVEL,
917 MAX_REP_LEVEL,
918 2,
919 2,
920 16,
921 0,
922 3
923 );
924 test!(
925 test_read_dict_v2_int32_small,
926 i32,
927 dict_v2,
928 MAX_DEF_LEVEL,
929 MAX_REP_LEVEL,
930 2,
931 2,
932 16,
933 0,
934 3
935 );
936
937 test!(
938 test_read_dict_v1_int32,
939 i32,
940 dict_v1,
941 MAX_DEF_LEVEL,
942 MAX_REP_LEVEL,
943 NUM_PAGES,
944 NUM_LEVELS,
945 16,
946 0,
947 3
948 );
949 test!(
950 test_read_dict_v2_int32,
951 i32,
952 dict_v2,
953 MAX_DEF_LEVEL,
954 MAX_REP_LEVEL,
955 NUM_PAGES,
956 NUM_LEVELS,
957 16,
958 0,
959 3
960 );
961
962 test!(
963 test_read_dict_v1_int32_uneven,
964 i32,
965 dict_v1,
966 MAX_DEF_LEVEL,
967 MAX_REP_LEVEL,
968 NUM_PAGES,
969 NUM_LEVELS,
970 17,
971 0,
972 3
973 );
974 test!(
975 test_read_dict_v2_int32_uneven,
976 i32,
977 dict_v2,
978 MAX_DEF_LEVEL,
979 MAX_REP_LEVEL,
980 NUM_PAGES,
981 NUM_LEVELS,
982 17,
983 0,
984 3
985 );
986
987 test!(
988 test_read_dict_v1_int32_multi_page,
989 i32,
990 dict_v1,
991 MAX_DEF_LEVEL,
992 MAX_REP_LEVEL,
993 NUM_PAGES,
994 NUM_LEVELS,
995 512,
996 0,
997 3
998 );
999 test!(
1000 test_read_dict_v2_int32_multi_page,
1001 i32,
1002 dict_v2,
1003 MAX_DEF_LEVEL,
1004 MAX_REP_LEVEL,
1005 NUM_PAGES,
1006 NUM_LEVELS,
1007 512,
1008 0,
1009 3
1010 );
1011
1012 test!(
1013 test_read_dict_v1_int64,
1014 i64,
1015 dict_v1,
1016 MAX_DEF_LEVEL,
1017 MAX_REP_LEVEL,
1018 NUM_PAGES,
1019 NUM_LEVELS,
1020 16,
1021 0,
1022 3
1023 );
1024 test!(
1025 test_read_dict_v2_int64,
1026 i64,
1027 dict_v2,
1028 MAX_DEF_LEVEL,
1029 MAX_REP_LEVEL,
1030 NUM_PAGES,
1031 NUM_LEVELS,
1032 16,
1033 0,
1034 3
1035 );
1036
1037 #[test]
1038 fn test_read_batch_values_only() {
1039 test_read_batch_int32(16, 0, 0);
1040 }
1041
1042 #[test]
1043 fn test_read_batch_values_def_levels() {
1044 test_read_batch_int32(16, MAX_DEF_LEVEL, 0);
1045 }
1046
1047 #[test]
1048 fn test_read_batch_values_rep_levels() {
1049 test_read_batch_int32(16, 0, MAX_REP_LEVEL);
1050 }
1051
1052 #[test]
1053 fn test_read_batch_values_def_rep_levels() {
1054 test_read_batch_int32(128, MAX_DEF_LEVEL, MAX_REP_LEVEL);
1055 }
1056
1057 #[test]
1058 fn test_read_batch_adjust_after_buffering_page() {
1059 let primitive_type = get_test_int32_type();
1066 let desc = Arc::new(ColumnDescriptor::new(
1067 Arc::new(primitive_type),
1068 1,
1069 1,
1070 ColumnPath::new(Vec::new()),
1071 ));
1072
1073 let num_pages = 2;
1074 let num_levels = 4;
1075 let batch_size = 5;
1076
1077 let mut tester = ColumnReaderTester::<Int32Type>::new();
1078 tester.test_read_batch(
1079 desc,
1080 Encoding::RLE_DICTIONARY,
1081 num_pages,
1082 num_levels,
1083 batch_size,
1084 i32::MIN,
1085 i32::MAX,
1086 false,
1087 );
1088 }
1089
1090 fn get_test_int32_type() -> SchemaType {
1136 SchemaType::primitive_type_builder("a", PhysicalType::INT32)
1137 .with_repetition(Repetition::REQUIRED)
1138 .with_converted_type(ConvertedType::INT_32)
1139 .with_length(-1)
1140 .build()
1141 .expect("build() should be OK")
1142 }
1143
1144 fn get_test_int64_type() -> SchemaType {
1146 SchemaType::primitive_type_builder("a", PhysicalType::INT64)
1147 .with_repetition(Repetition::REQUIRED)
1148 .with_converted_type(ConvertedType::INT_64)
1149 .with_length(-1)
1150 .build()
1151 .expect("build() should be OK")
1152 }
1153
1154 fn test_read_batch_int32(batch_size: usize, max_def_level: i16, max_rep_level: i16) {
1159 let primitive_type = get_test_int32_type();
1160
1161 let desc = Arc::new(ColumnDescriptor::new(
1162 Arc::new(primitive_type),
1163 max_def_level,
1164 max_rep_level,
1165 ColumnPath::new(Vec::new()),
1166 ));
1167
1168 let mut tester = ColumnReaderTester::<Int32Type>::new();
1169 tester.test_read_batch(
1170 desc,
1171 Encoding::RLE_DICTIONARY,
1172 NUM_PAGES,
1173 NUM_LEVELS,
1174 batch_size,
1175 i32::MIN,
1176 i32::MAX,
1177 false,
1178 );
1179 }
1180
1181 struct ColumnReaderTester<T: DataType>
1182 where
1183 T::T: PartialOrd + SampleUniform + Copy,
1184 {
1185 rep_levels: Vec<i16>,
1186 def_levels: Vec<i16>,
1187 values: Vec<T::T>,
1188 }
1189
1190 impl<T: DataType> ColumnReaderTester<T>
1191 where
1192 T::T: PartialOrd + SampleUniform + Copy,
1193 {
1194 pub fn new() -> Self {
1195 Self {
1196 rep_levels: Vec::new(),
1197 def_levels: Vec::new(),
1198 values: Vec::new(),
1199 }
1200 }
1201
1202 fn plain_v1(
1204 &mut self,
1205 desc: ColumnDescPtr,
1206 num_pages: usize,
1207 num_levels: usize,
1208 batch_size: usize,
1209 min: T::T,
1210 max: T::T,
1211 ) {
1212 self.test_read_batch_general(
1213 desc,
1214 Encoding::PLAIN,
1215 num_pages,
1216 num_levels,
1217 batch_size,
1218 min,
1219 max,
1220 false,
1221 );
1222 }
1223
1224 fn plain_v2(
1226 &mut self,
1227 desc: ColumnDescPtr,
1228 num_pages: usize,
1229 num_levels: usize,
1230 batch_size: usize,
1231 min: T::T,
1232 max: T::T,
1233 ) {
1234 self.test_read_batch_general(
1235 desc,
1236 Encoding::PLAIN,
1237 num_pages,
1238 num_levels,
1239 batch_size,
1240 min,
1241 max,
1242 true,
1243 );
1244 }
1245
1246 fn dict_v1(
1248 &mut self,
1249 desc: ColumnDescPtr,
1250 num_pages: usize,
1251 num_levels: usize,
1252 batch_size: usize,
1253 min: T::T,
1254 max: T::T,
1255 ) {
1256 self.test_read_batch_general(
1257 desc,
1258 Encoding::RLE_DICTIONARY,
1259 num_pages,
1260 num_levels,
1261 batch_size,
1262 min,
1263 max,
1264 false,
1265 );
1266 }
1267
1268 fn dict_v2(
1270 &mut self,
1271 desc: ColumnDescPtr,
1272 num_pages: usize,
1273 num_levels: usize,
1274 batch_size: usize,
1275 min: T::T,
1276 max: T::T,
1277 ) {
1278 self.test_read_batch_general(
1279 desc,
1280 Encoding::RLE_DICTIONARY,
1281 num_pages,
1282 num_levels,
1283 batch_size,
1284 min,
1285 max,
1286 true,
1287 );
1288 }
1289
1290 #[allow(clippy::too_many_arguments)]
1293 fn test_read_batch_general(
1294 &mut self,
1295 desc: ColumnDescPtr,
1296 encoding: Encoding,
1297 num_pages: usize,
1298 num_levels: usize,
1299 batch_size: usize,
1300 min: T::T,
1301 max: T::T,
1302 use_v2: bool,
1303 ) {
1304 self.test_read_batch(
1305 desc, encoding, num_pages, num_levels, batch_size, min, max, use_v2,
1306 );
1307 }
1308
1309 #[allow(clippy::too_many_arguments)]
1312 fn test_read_batch(
1313 &mut self,
1314 desc: ColumnDescPtr,
1315 encoding: Encoding,
1316 num_pages: usize,
1317 num_levels: usize,
1318 batch_size: usize,
1319 min: T::T,
1320 max: T::T,
1321 use_v2: bool,
1322 ) {
1323 let mut pages = VecDeque::new();
1324 make_pages::<T>(
1325 desc.clone(),
1326 encoding,
1327 num_pages,
1328 num_levels,
1329 min,
1330 max,
1331 &mut self.def_levels,
1332 &mut self.rep_levels,
1333 &mut self.values,
1334 &mut pages,
1335 use_v2,
1336 );
1337 let max_def_level = desc.max_def_level();
1338 let max_rep_level = desc.max_rep_level();
1339 let page_reader = InMemoryPageReader::new(pages);
1340 let column_reader: ColumnReader = get_column_reader(desc, Box::new(page_reader));
1341 let mut typed_column_reader = get_typed_column_reader::<T>(column_reader);
1342
1343 let mut values = Vec::new();
1344 let mut def_levels = Vec::new();
1345 let mut rep_levels = Vec::new();
1346
1347 let mut curr_values_read = 0;
1348 let mut curr_levels_read = 0;
1349 loop {
1350 let (_, values_read, levels_read) = typed_column_reader
1351 .read_records(
1352 batch_size,
1353 Some(&mut def_levels),
1354 Some(&mut rep_levels),
1355 &mut values,
1356 )
1357 .expect("read_batch() should be OK");
1358
1359 curr_values_read += values_read;
1360 curr_levels_read += levels_read;
1361
1362 if values_read == 0 && levels_read == 0 {
1363 break;
1364 }
1365 }
1366
1367 assert_eq!(values, self.values, "values content doesn't match");
1368
1369 if max_def_level > 0 {
1370 assert_eq!(
1371 def_levels, self.def_levels,
1372 "definition levels content doesn't match"
1373 );
1374 }
1375
1376 if max_rep_level > 0 {
1377 assert_eq!(
1378 rep_levels, self.rep_levels,
1379 "repetition levels content doesn't match"
1380 );
1381 }
1382
1383 assert!(
1384 curr_levels_read >= curr_values_read,
1385 "expected levels read to be greater than values read"
1386 );
1387 }
1388 }
1389
1390 #[test]
1408 fn test_skip_records_v2_page_skip_accounts_for_partial() {
1409 use crate::encodings::levels::LevelEncoder;
1410
1411 let max_rep_level: i16 = 1;
1412 let max_def_level: i16 = 1;
1413
1414 let primitive_type = SchemaType::primitive_type_builder("element", PhysicalType::INT32)
1416 .with_repetition(Repetition::REQUIRED)
1417 .build()
1418 .unwrap();
1419 let desc = Arc::new(ColumnDescriptor::new(
1420 Arc::new(primitive_type),
1421 max_def_level,
1422 max_rep_level,
1423 ColumnPath::new(vec!["list".to_string(), "element".to_string()]),
1424 ));
1425
1426 let make_v2_page =
1428 |rep_levels: &[i16], def_levels: &[i16], values: &[i32], num_rows: u32| -> Page {
1429 let mut rep_enc = LevelEncoder::v2_streaming(max_rep_level);
1430 rep_enc.put_with_observer(rep_levels, |_, _| {});
1431 let rep_bytes = rep_enc.consume();
1432
1433 let mut def_enc = LevelEncoder::v2_streaming(max_def_level);
1434 def_enc.put_with_observer(def_levels, |_, _| {});
1435 let def_bytes = def_enc.consume();
1436
1437 let val_bytes: Vec<u8> = values.iter().flat_map(|v| v.to_le_bytes()).collect();
1438
1439 let mut buf = Vec::new();
1440 buf.extend_from_slice(&rep_bytes);
1441 buf.extend_from_slice(&def_bytes);
1442 buf.extend_from_slice(&val_bytes);
1443
1444 Page::DataPageV2 {
1445 buf: Bytes::from(buf),
1446 num_values: rep_levels.len() as u32,
1447 encoding: Encoding::PLAIN,
1448 num_nulls: 0,
1449 num_rows,
1450 def_levels_byte_len: def_bytes.len() as u32,
1451 rep_levels_byte_len: rep_bytes.len() as u32,
1452 is_compressed: false,
1453 statistics: None,
1454 }
1455 };
1456
1457 let page1 = make_v2_page(&[0, 1, 0, 1], &[1, 1, 1, 1], &[10, 20, 30, 40], 2);
1463
1464 let page2 = make_v2_page(&[0, 1, 0, 1], &[1, 1, 1, 1], &[50, 60, 70, 80], 2);
1466
1467 let page3 = make_v2_page(&[0, 1], &[1, 1], &[90, 100], 1);
1469
1470 let pages = VecDeque::from(vec![page1, page2, page3]);
1472 let page_reader = InMemoryPageReader::new(pages);
1473 let column_reader: ColumnReader = get_column_reader(desc, Box::new(page_reader));
1474 let mut typed_reader = get_typed_column_reader::<Int32Type>(column_reader);
1475
1476 let skipped = typed_reader.skip_records(1).unwrap();
1482 assert_eq!(skipped, 1);
1483
1484 let skipped = typed_reader.skip_records(2).unwrap();
1499 assert_eq!(skipped, 2);
1500
1501 let mut values = Vec::new();
1503 let mut def_levels = Vec::new();
1504 let mut rep_levels = Vec::new();
1505
1506 let (records, values_read, levels_read) = typed_reader
1507 .read_records(1, Some(&mut def_levels), Some(&mut rep_levels), &mut values)
1508 .unwrap();
1509
1510 assert_eq!(records, 1, "should read exactly 1 record");
1514 assert_eq!(levels_read, 2, "should read 2 levels for the record");
1515 assert_eq!(values_read, 2, "should read 2 non-null values");
1516 assert_eq!(values, vec![70, 80], "should contain 4th record's values");
1517 assert_eq!(rep_levels, vec![0, 1], "rep levels for a 2-element list");
1518 assert_eq!(def_levels, vec![1, 1], "def levels (all non-null)");
1519 }
1520}