1use arrow_buffer::{BooleanBufferBuilder, Buffer};
19
20use crate::arrow::record_reader::{
21 buffer::ValuesBuffer,
22 definition_levels::{DefinitionLevelBuffer, DefinitionLevelBufferDecoder},
23};
24use crate::column::reader::decoder::RepetitionLevelDecoderImpl;
25use crate::column::{
26 page::PageReader,
27 reader::{
28 GenericColumnReader,
29 decoder::{ColumnValueDecoder, ColumnValueDecoderImpl},
30 },
31};
32use crate::data_type::DataType;
33use crate::errors::{ParquetError, Result};
34use crate::schema::types::ColumnDescPtr;
35
36pub(crate) mod buffer;
37pub(crate) mod definition_levels;
38
39pub type RecordReader<T> = GenericRecordReader<Vec<<T as DataType>::T>, ColumnValueDecoderImpl<T>>;
41
42pub(crate) type ColumnReader<CV> =
43 GenericColumnReader<RepetitionLevelDecoderImpl, DefinitionLevelBufferDecoder, CV>;
44
45pub struct GenericRecordReader<V, CV> {
51 column_desc: ColumnDescPtr,
52
53 values: Option<V>,
56 def_levels: Option<DefinitionLevelBuffer>,
57 rep_levels: Option<Vec<i16>>,
58 column_reader: Option<ColumnReader<CV>>,
59 num_values: usize,
61 num_records: usize,
63 capacity_hint: usize,
65 values_written: usize,
68 padding_threshold: Option<i16>,
82 compact_bitmap: Option<BooleanBufferBuilder>,
88}
89
90impl<V, CV> GenericRecordReader<V, CV>
91where
92 V: ValuesBuffer,
93 CV: ColumnValueDecoder<Buffer = V>,
94{
95 pub fn new(desc: ColumnDescPtr, capacity: usize) -> Self {
100 let def_levels = (desc.max_def_level() > 0)
101 .then(|| DefinitionLevelBuffer::new(&desc, packed_null_mask(&desc)));
102
103 let rep_levels = (desc.max_rep_level() > 0).then(Vec::new);
104
105 Self {
106 values: None, def_levels,
108 rep_levels,
109 column_reader: None,
110 column_desc: desc,
111 num_values: 0,
112 num_records: 0,
113 capacity_hint: capacity,
114 values_written: 0,
115 padding_threshold: None,
116 compact_bitmap: None,
117 }
118 }
119
120 pub fn set_page_reader(&mut self, page_reader: Box<dyn PageReader>) -> Result<()> {
122 let descr = &self.column_desc;
123 let values_decoder = CV::new(descr);
124
125 let def_level_decoder = (descr.max_def_level() != 0).then(|| {
126 DefinitionLevelBufferDecoder::new(descr.max_def_level(), packed_null_mask(descr))
127 });
128
129 let rep_level_decoder = (descr.max_rep_level() != 0)
130 .then(|| RepetitionLevelDecoderImpl::new(descr.max_rep_level()));
131
132 self.column_reader = Some(GenericColumnReader::new_with_decoders(
133 self.column_desc.clone(),
134 page_reader,
135 values_decoder,
136 def_level_decoder,
137 rep_level_decoder,
138 ));
139 Ok(())
140 }
141
142 pub fn read_records(&mut self, num_records: usize) -> Result<usize> {
148 if self.column_reader.is_none() {
149 return Ok(0);
150 }
151
152 let mut records_read = 0;
153
154 loop {
155 let records_to_read = num_records - records_read;
156 records_read += self.read_one_batch(records_to_read)?;
157 if records_read == num_records || !self.column_reader.as_mut().unwrap().has_next()? {
158 break;
159 }
160 }
161 Ok(records_read)
162 }
163
164 pub fn skip_records(&mut self, num_records: usize) -> Result<usize> {
170 match self.column_reader.as_mut() {
171 Some(reader) => reader.skip_records(num_records),
172 None => Ok(0),
173 }
174 }
175
176 #[allow(unused)]
178 pub fn num_records(&self) -> usize {
179 self.num_records
180 }
181
182 pub fn num_values(&self) -> usize {
186 self.num_values
187 }
188
189 pub fn consume_def_levels(&mut self) -> Option<Vec<i16>> {
194 self.def_levels.as_mut().and_then(|x| x.consume_levels())
195 }
196
197 pub fn consume_rep_levels(&mut self) -> Option<Vec<i16>> {
200 self.rep_levels.as_mut().map(std::mem::take)
201 }
202
203 pub fn consume_record_data(&mut self) -> V {
206 self.values.take().unwrap_or_else(|| V::with_capacity(0))
209 }
210
211 pub fn reset(&mut self) {
215 self.num_values = 0;
216 self.num_records = 0;
217 self.values_written = 0;
218 self.compact_bitmap = None;
219 }
220
221 pub fn max_def_level(&self) -> i16 {
223 self.column_desc.max_def_level()
224 }
225
226 pub fn set_padding_threshold(&mut self, threshold: i16) {
230 self.padding_threshold = Some(threshold);
231 }
232
233 pub fn values_written(&self) -> usize {
237 if self.padding_threshold.is_some() {
238 self.values_written
239 } else {
240 self.num_values
241 }
242 }
243
244 pub fn consume_compact_bitmap(&mut self) -> Option<Buffer> {
247 if self.padding_threshold.is_some() {
248 if let Some(levels) = self.def_levels.as_mut() {
249 levels.consume_bitmask();
250 }
251 self.compact_bitmap
252 .as_mut()
253 .map(|b| b.finish().into_inner())
254 } else {
255 self.consume_bitmap()
256 }
257 }
258
259 pub fn consume_bitmap(&mut self) -> Option<Buffer> {
262 let mask = self
263 .def_levels
264 .as_mut()
265 .and_then(|levels| levels.consume_bitmask());
266
267 if self.column_desc.self_type().is_optional() {
272 mask
273 } else {
274 None
275 }
276 }
277
278 fn read_one_batch(&mut self, batch_size: usize) -> Result<usize> {
280 if batch_size == 0 {
281 return Ok(0);
282 }
283 if batch_size > self.capacity_hint {
285 self.capacity_hint = batch_size;
286 }
287
288 let padding_threshold = self.padding_threshold;
289 let max_def = self.column_desc.max_def_level();
290 let mut physical_values_to_read = 0usize;
291 let mut output_slots_to_read = 0usize;
292
293 let values = self.values.get_or_insert_with(|| V::with_capacity(0));
294 let compact_bitmap = &mut self.compact_bitmap;
295
296 if padding_threshold.is_none() {
297 values.reserve_exact(self.capacity_hint.saturating_sub(self.num_values));
298 }
299
300 let (records_read, values_read, levels_read) = self
301 .column_reader
302 .as_mut()
303 .unwrap()
304 .read_records_with_reservation(
305 batch_size,
306 self.def_levels.as_mut(),
307 self.rep_levels.as_mut(),
308 values,
309 |values, values_to_read, levels_to_read, def_levels| {
310 let output_slots = if let Some(threshold) = padding_threshold {
311 let def_levels = def_levels.ok_or_else(|| {
312 general_err!(
313 "Definition levels should exist when data is less than levels!"
314 )
315 })?;
316 let all_levels = def_levels.levels().ok_or_else(|| {
317 general_err!(
318 "Raw definition levels must be available for selective padding"
319 )
320 })?;
321 let batch_levels = &all_levels[all_levels.len() - levels_to_read..];
322 let bitmap =
323 compact_bitmap.get_or_insert_with(|| BooleanBufferBuilder::new(0));
324
325 definition_levels::build_filtered_validity_bitmap(
326 batch_levels,
327 None,
328 Some(threshold),
329 max_def,
330 bitmap,
331 )
332 } else {
333 levels_to_read
334 };
335
336 let additional = output_slots_to_read
337 .saturating_add(output_slots)
338 .saturating_sub(physical_values_to_read);
339 values.reserve_exact(additional);
340
341 output_slots_to_read += output_slots;
342 physical_values_to_read += values_to_read;
343 Ok(())
344 },
345 )?;
346
347 if self.padding_threshold.is_some() {
348 debug_assert_eq!(
349 physical_values_to_read, values_read,
350 "reservation accounting must match decoded values"
351 );
352 let item_count = output_slots_to_read;
353 let bitmap = compact_bitmap.get_or_insert_with(|| BooleanBufferBuilder::new(0));
354
355 if values_read < item_count {
357 values.reserve_exact(item_count - values_read);
358 values.pad_nulls(
359 self.values_written,
360 values_read,
361 item_count,
362 bitmap.as_slice(),
363 )?;
364 }
365
366 debug_assert_eq!(
367 bitmap.len(),
368 self.values_written + item_count,
369 "compact bitmap length must equal total items written"
370 );
371 self.values_written += item_count;
372 } else if values_read < levels_read {
373 debug_assert_eq!(
374 output_slots_to_read, levels_read,
375 "reservation accounting must match decoded levels"
376 );
377 let def_levels = self.def_levels.as_ref().ok_or_else(|| {
379 general_err!("Definition levels should exist when data is less than levels!")
380 })?;
381
382 values.reserve_exact(levels_read - values_read);
383 values.pad_nulls(
384 self.num_values,
385 values_read,
386 levels_read,
387 def_levels.nulls().as_slice(),
388 )?;
389 }
390
391 self.num_records += records_read;
392 self.num_values += levels_read;
393 Ok(records_read)
394 }
395}
396
397fn packed_null_mask(descr: &ColumnDescPtr) -> bool {
401 descr.max_def_level() == 1 && descr.max_rep_level() == 0 && descr.self_type().is_optional()
402}
403
404#[cfg(test)]
405mod tests {
406 use std::sync::Arc;
407
408 use arrow::buffer::Buffer;
409
410 use crate::arrow::arrow_reader::DEFAULT_BATCH_SIZE;
411 use crate::basic::Encoding;
412 use crate::data_type::Int32Type;
413 use crate::schema::parser::parse_message_type;
414 use crate::schema::types::SchemaDescriptor;
415 use crate::util::test_common::page_util::{
416 DataPageBuilder, DataPageBuilderImpl, InMemoryPageReader,
417 };
418
419 use super::RecordReader;
420
421 #[test]
422 fn test_read_required_records() {
423 let message_type = "
425 message test_schema {
426 REQUIRED INT32 leaf;
427 }
428 ";
429 let desc = parse_message_type(message_type)
430 .map(|t| SchemaDescriptor::new(Arc::new(t)))
431 .map(|s| s.column(0))
432 .unwrap();
433
434 let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
436
437 {
451 let values = [4, 7, 6, 3, 2];
452 let mut pb = DataPageBuilderImpl::new(desc.clone(), 5, true);
453 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
454 let page = pb.consume();
455
456 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
457 record_reader.set_page_reader(page_reader).unwrap();
458 assert_eq!(2, record_reader.read_records(2).unwrap());
459 assert_eq!(2, record_reader.num_records());
460 assert_eq!(2, record_reader.num_values());
461 assert_eq!(3, record_reader.read_records(3).unwrap());
462 assert_eq!(5, record_reader.num_records());
463 assert_eq!(5, record_reader.num_values());
464 }
465
466 {
474 let values = [8, 9];
475 let mut pb = DataPageBuilderImpl::new(desc, 2, true);
476 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
477 let page = pb.consume();
478
479 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
480 record_reader.set_page_reader(page_reader).unwrap();
481 assert_eq!(2, record_reader.read_records(10).unwrap());
482 assert_eq!(7, record_reader.num_records());
483 assert_eq!(7, record_reader.num_values());
484 }
485
486 assert_eq!(record_reader.consume_record_data(), &[4, 7, 6, 3, 2, 8, 9]);
487 assert_eq!(None, record_reader.consume_def_levels());
488 assert_eq!(None, record_reader.consume_bitmap());
489 }
490
491 #[test]
492 fn test_capacity_hint_preserved_across_fragmented_reads() {
493 let message_type = "
494 message test_schema {
495 REQUIRED INT32 leaf;
496 }
497 ";
498 let desc = parse_message_type(message_type)
499 .map(|t| SchemaDescriptor::new(Arc::new(t)))
500 .map(|s| s.column(0))
501 .unwrap();
502
503 let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), 8);
504 let values = [1, 2, 3, 4];
505 let mut pb = DataPageBuilderImpl::new(desc, 4, true);
506 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
507 let page = pb.consume();
508
509 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
510 record_reader.set_page_reader(page_reader).unwrap();
511
512 assert_eq!(1, record_reader.read_records(1).unwrap());
513 assert_eq!(1, record_reader.read_records(1).unwrap());
514
515 let record_data = record_reader.consume_record_data();
516 assert_eq!(record_data, &[1, 2]);
517 assert!(
518 record_data.capacity() >= 8,
519 "capacity hint should survive fragmented reads, got {}",
520 record_data.capacity()
521 );
522 }
523
524 #[test]
525 fn test_read_optional_records() {
526 let message_type = "
528 message test_schema {
529 OPTIONAL Group test_struct {
530 OPTIONAL INT32 leaf;
531 }
532 }
533 ";
534
535 let desc = parse_message_type(message_type)
536 .map(|t| SchemaDescriptor::new(Arc::new(t)))
537 .map(|s| s.column(0))
538 .unwrap();
539
540 let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
542
543 {
559 let values = [7, 6, 3];
560 let def_levels = [1i16, 2i16, 0i16, 2i16, 2i16];
562 let mut pb = DataPageBuilderImpl::new(desc.clone(), 5, true);
563 pb.add_def_levels(2, &def_levels);
564 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
565 let page = pb.consume();
566
567 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
568 record_reader.set_page_reader(page_reader).unwrap();
569 assert_eq!(2, record_reader.read_records(2).unwrap());
570 assert_eq!(2, record_reader.num_records());
571 assert_eq!(2, record_reader.num_values());
572 assert_eq!(3, record_reader.read_records(3).unwrap());
573 assert_eq!(5, record_reader.num_records());
574 assert_eq!(5, record_reader.num_values());
575 }
576
577 {
585 let values = [8];
586 let def_levels = [0i16, 2i16];
588 let mut pb = DataPageBuilderImpl::new(desc, 2, true);
589 pb.add_def_levels(2, &def_levels);
590 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
591 let page = pb.consume();
592
593 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
594 record_reader.set_page_reader(page_reader).unwrap();
595 assert_eq!(2, record_reader.read_records(10).unwrap());
596 assert_eq!(7, record_reader.num_records());
597 assert_eq!(7, record_reader.num_values());
598 }
599
600 assert_eq!(
602 Some(vec![1i16, 2i16, 0i16, 2i16, 2i16, 0i16, 2i16]),
603 record_reader.consume_def_levels()
604 );
605
606 let expected_valid = &[false, true, false, true, true, false, true];
608 let expected_buffer = Buffer::from_iter(expected_valid.iter().cloned());
609 assert_eq!(Some(expected_buffer), record_reader.consume_bitmap());
610
611 let actual = record_reader.consume_record_data();
613
614 let expected = &[0, 7, 0, 6, 3, 0, 8];
615 assert_eq!(actual.len(), expected.len());
616
617 let iter = expected_valid.iter().zip(&actual).zip(expected);
619 for ((valid, actual), expected) in iter {
620 if *valid {
621 assert_eq!(actual, expected)
622 }
623 }
624 }
625
626 #[test]
627 fn test_read_repeated_records() {
628 let message_type = "
630 message test_schema {
631 REPEATED Group test_struct {
632 REPEATED INT32 leaf;
633 }
634 }
635 ";
636
637 let desc = parse_message_type(message_type)
638 .map(|t| SchemaDescriptor::new(Arc::new(t)))
639 .map(|s| s.column(0))
640 .unwrap();
641
642 let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
644
645 {
661 let values = [4, 7, 6, 3, 2];
662 let def_levels = [2i16, 0i16, 1i16, 2i16, 2i16, 2i16, 2i16];
663 let rep_levels = [0i16, 0i16, 0i16, 1i16, 2i16, 2i16, 1i16];
664 let mut pb = DataPageBuilderImpl::new(desc.clone(), 7, true);
665 pb.add_rep_levels(2, &rep_levels);
666 pb.add_def_levels(2, &def_levels);
667 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
668 let page = pb.consume();
669
670 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
671 record_reader.set_page_reader(page_reader).unwrap();
672
673 assert_eq!(1, record_reader.read_records(1).unwrap());
674 assert_eq!(1, record_reader.num_records());
675 assert_eq!(1, record_reader.num_values());
676 assert_eq!(2, record_reader.read_records(3).unwrap());
677 assert_eq!(3, record_reader.num_records());
678 assert_eq!(7, record_reader.num_values());
679 }
680
681 {
689 let values = [8, 9];
690 let def_levels = [2i16, 2i16];
691 let rep_levels = [0i16, 2i16];
692 let mut pb = DataPageBuilderImpl::new(desc, 2, true);
693 pb.add_rep_levels(2, &rep_levels);
694 pb.add_def_levels(2, &def_levels);
695 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
696 let page = pb.consume();
697
698 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
699 record_reader.set_page_reader(page_reader).unwrap();
700
701 assert_eq!(1, record_reader.read_records(10).unwrap());
702 assert_eq!(4, record_reader.num_records());
703 assert_eq!(9, record_reader.num_values());
704 }
705
706 assert_eq!(
708 Some(vec![2i16, 0i16, 1i16, 2i16, 2i16, 2i16, 2i16, 2i16, 2i16]),
709 record_reader.consume_def_levels()
710 );
711
712 let expected_valid = &[true, false, false, true, true, true, true, true, true];
714 let expected_buffer = Buffer::from_iter(expected_valid.iter().cloned());
715 assert_eq!(Some(expected_buffer), record_reader.consume_bitmap());
716
717 let actual = record_reader.consume_record_data();
719 let expected = &[4, 0, 0, 7, 6, 3, 2, 8, 9];
720 assert_eq!(actual.len(), expected.len());
721
722 let iter = expected_valid.iter().zip(&actual).zip(expected);
724 for ((valid, actual), expected) in iter {
725 if *valid {
726 assert_eq!(actual, expected)
727 }
728 }
729 }
730
731 #[test]
732 fn test_selective_padding_consumes_full_bitmap() {
733 let message_type = "
734 message test_schema {
735 OPTIONAL GROUP my_list (LIST) {
736 REPEATED GROUP list {
737 OPTIONAL INT32 element;
738 }
739 }
740 }
741 ";
742
743 let desc = parse_message_type(message_type)
744 .map(|t| SchemaDescriptor::new(Arc::new(t)))
745 .map(|s| s.column(0))
746 .unwrap();
747
748 let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
749 record_reader.set_padding_threshold(2);
750
751 let values = [10, 20];
752 let def_levels = [3i16, 2i16, 0i16, 2i16, 3i16];
753 let rep_levels = [0i16, 1i16, 0i16, 0i16, 1i16];
754 let mut pb = DataPageBuilderImpl::new(desc, 5, true);
755 pb.add_rep_levels(1, &rep_levels);
756 pb.add_def_levels(3, &def_levels);
757 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
758
759 let page_reader = Box::new(InMemoryPageReader::new(vec![pb.consume()]));
760 record_reader.set_page_reader(page_reader).unwrap();
761 assert_eq!(3, record_reader.read_records(3).unwrap());
762
763 let expected_compact = Buffer::from_iter([true, false, false, true]);
764 assert_eq!(
765 Some(expected_compact),
766 record_reader.consume_compact_bitmap()
767 );
768 assert_eq!(None, record_reader.consume_bitmap());
769 }
770
771 #[test]
772 fn test_selective_padding_reserves_compact_capacity() {
773 let message_type = "
774 message test_schema {
775 OPTIONAL GROUP my_list (LIST) {
776 REPEATED GROUP list {
777 OPTIONAL INT32 element;
778 }
779 }
780 }
781 ";
782
783 let desc = parse_message_type(message_type)
784 .map(|t| SchemaDescriptor::new(Arc::new(t)))
785 .map(|s| s.column(0))
786 .unwrap();
787
788 let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
789 record_reader.set_padding_threshold(2);
790
791 let values = [10];
792 let mut def_levels = vec![0i16; 30];
793 def_levels.extend([3i16, 2i16]);
794 let rep_levels = vec![0i16; def_levels.len()];
795
796 let mut pb = DataPageBuilderImpl::new(desc, def_levels.len() as u32, true);
797 pb.add_rep_levels(1, &rep_levels);
798 pb.add_def_levels(3, &def_levels);
799 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
800
801 let page_reader = Box::new(InMemoryPageReader::new(vec![pb.consume()]));
802 record_reader.set_page_reader(page_reader).unwrap();
803 assert_eq!(def_levels.len(), record_reader.read_records(128).unwrap());
804
805 let expected_compact = Buffer::from_iter([true, false]);
806 assert_eq!(
807 Some(expected_compact),
808 record_reader.consume_compact_bitmap()
809 );
810
811 let actual = record_reader.consume_record_data();
812 assert_eq!(actual.len(), 2);
813 assert_eq!(actual[0], 10);
814 assert!(
815 actual.capacity() < def_levels.len(),
816 "selective padding should not reserve one child slot per list-level NULL"
817 );
818 }
819
820 #[test]
821 fn test_read_more_than_one_batch() {
822 let message_type = "
824 message test_schema {
825 REPEATED INT32 leaf;
826 }
827 ";
828
829 let desc = parse_message_type(message_type)
830 .map(|t| SchemaDescriptor::new(Arc::new(t)))
831 .map(|s| s.column(0))
832 .unwrap();
833
834 let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
836
837 {
838 let values = [100; 5000];
839 let def_levels = [1i16; 5000];
840 let mut rep_levels = [1i16; 5000];
841 for idx in 0..1000 {
842 rep_levels[idx * 5] = 0i16;
843 }
844
845 let mut pb = DataPageBuilderImpl::new(desc, 5000, true);
846 pb.add_rep_levels(1, &rep_levels);
847 pb.add_def_levels(1, &def_levels);
848 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
849 let page = pb.consume();
850
851 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
852 record_reader.set_page_reader(page_reader).unwrap();
853
854 assert_eq!(1000, record_reader.read_records(1000).unwrap());
855 assert_eq!(1000, record_reader.num_records());
856 assert_eq!(5000, record_reader.num_values());
857 }
858 }
859
860 #[test]
861 fn test_row_group_boundary() {
862 let message_type = "
864 message test_schema {
865 REPEATED Group test_struct {
866 REPEATED INT32 leaf;
867 }
868 }
869 ";
870
871 let desc = parse_message_type(message_type)
872 .map(|t| SchemaDescriptor::new(Arc::new(t)))
873 .map(|s| s.column(0))
874 .unwrap();
875
876 let values = [1, 2, 3];
877 let def_levels = [1i16, 0i16, 1i16, 2i16, 2i16, 1i16, 2i16];
878 let rep_levels = [0i16, 0i16, 0i16, 1i16, 2i16, 0i16, 1i16];
879 let mut pb = DataPageBuilderImpl::new(desc.clone(), 7, true);
880 pb.add_rep_levels(2, &rep_levels);
881 pb.add_def_levels(2, &def_levels);
882 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
883 let page = pb.consume();
884
885 let mut record_reader = RecordReader::<Int32Type>::new(desc, DEFAULT_BATCH_SIZE);
886 let page_reader = Box::new(InMemoryPageReader::new(vec![page.clone()]));
887 record_reader.set_page_reader(page_reader).unwrap();
888 assert_eq!(record_reader.read_records(4).unwrap(), 4);
889 assert_eq!(record_reader.num_records(), 4);
890 assert_eq!(record_reader.num_values(), 7);
891
892 assert_eq!(record_reader.read_records(4).unwrap(), 0);
893 assert_eq!(record_reader.num_records(), 4);
894 assert_eq!(record_reader.num_values(), 7);
895
896 record_reader.read_records(4).unwrap();
897
898 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
899 record_reader.set_page_reader(page_reader).unwrap();
900
901 assert_eq!(record_reader.read_records(4).unwrap(), 4);
902 assert_eq!(record_reader.num_records(), 8);
903 assert_eq!(record_reader.num_values(), 14);
904
905 assert_eq!(record_reader.read_records(4).unwrap(), 0);
906 assert_eq!(record_reader.num_records(), 8);
907 assert_eq!(record_reader.num_values(), 14);
908 }
909
910 #[test]
911 fn test_skip_required_records() {
912 let message_type = "
914 message test_schema {
915 REQUIRED INT32 leaf;
916 }
917 ";
918 let desc = parse_message_type(message_type)
919 .map(|t| SchemaDescriptor::new(Arc::new(t)))
920 .map(|s| s.column(0))
921 .unwrap();
922
923 let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
925
926 {
940 let values = [4, 7, 6, 3, 2];
941 let mut pb = DataPageBuilderImpl::new(desc.clone(), 5, true);
942 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
943 let page = pb.consume();
944
945 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
946 record_reader.set_page_reader(page_reader).unwrap();
947 assert_eq!(2, record_reader.skip_records(2).unwrap());
948 assert_eq!(0, record_reader.num_records());
949 assert_eq!(0, record_reader.num_values());
950 assert_eq!(3, record_reader.read_records(3).unwrap());
951 assert_eq!(3, record_reader.num_records());
952 assert_eq!(3, record_reader.num_values());
953 }
954
955 {
963 let values = [8, 9];
964 let mut pb = DataPageBuilderImpl::new(desc, 2, true);
965 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
966 let page = pb.consume();
967
968 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
969 record_reader.set_page_reader(page_reader).unwrap();
970 assert_eq!(2, record_reader.skip_records(10).unwrap());
971 assert_eq!(3, record_reader.num_records());
972 assert_eq!(3, record_reader.num_values());
973 assert_eq!(0, record_reader.read_records(10).unwrap());
974 }
975
976 assert_eq!(record_reader.consume_record_data(), &[6, 3, 2]);
977 assert_eq!(None, record_reader.consume_def_levels());
978 assert_eq!(None, record_reader.consume_bitmap());
979 }
980
981 #[test]
982 fn test_skip_optional_records() {
983 let message_type = "
985 message test_schema {
986 OPTIONAL Group test_struct {
987 OPTIONAL INT32 leaf;
988 }
989 }
990 ";
991
992 let desc = parse_message_type(message_type)
993 .map(|t| SchemaDescriptor::new(Arc::new(t)))
994 .map(|s| s.column(0))
995 .unwrap();
996
997 let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
999
1000 {
1016 let values = [7, 6, 3];
1017 let def_levels = [1i16, 2i16, 0i16, 2i16, 2i16];
1019 let mut pb = DataPageBuilderImpl::new(desc.clone(), 5, true);
1020 pb.add_def_levels(2, &def_levels);
1021 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
1022 let page = pb.consume();
1023
1024 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
1025 record_reader.set_page_reader(page_reader).unwrap();
1026 assert_eq!(2, record_reader.skip_records(2).unwrap());
1027 assert_eq!(0, record_reader.num_records());
1028 assert_eq!(0, record_reader.num_values());
1029 assert_eq!(3, record_reader.read_records(3).unwrap());
1030 assert_eq!(3, record_reader.num_records());
1031 assert_eq!(3, record_reader.num_values());
1032 }
1033
1034 {
1042 let values = [8];
1043 let def_levels = [0i16, 2i16];
1045 let mut pb = DataPageBuilderImpl::new(desc, 2, true);
1046 pb.add_def_levels(2, &def_levels);
1047 pb.add_values::<Int32Type>(Encoding::PLAIN, &values);
1048 let page = pb.consume();
1049
1050 let page_reader = Box::new(InMemoryPageReader::new(vec![page]));
1051 record_reader.set_page_reader(page_reader).unwrap();
1052 assert_eq!(2, record_reader.skip_records(10).unwrap());
1053 assert_eq!(3, record_reader.num_records());
1054 assert_eq!(3, record_reader.num_values());
1055 assert_eq!(0, record_reader.read_records(10).unwrap());
1056 }
1057
1058 assert_eq!(
1060 Some(vec![0i16, 2i16, 2i16]),
1061 record_reader.consume_def_levels()
1062 );
1063
1064 let expected_valid = &[false, true, true];
1066 let expected_buffer = Buffer::from_iter(expected_valid.iter().cloned());
1067 assert_eq!(Some(expected_buffer), record_reader.consume_bitmap());
1068
1069 let actual = record_reader.consume_record_data();
1071
1072 let expected = &[0, 6, 3];
1073 assert_eq!(actual.len(), expected.len());
1074
1075 let iter = expected_valid.iter().zip(&actual).zip(expected);
1077 for ((valid, actual), expected) in iter {
1078 if *valid {
1079 assert_eq!(actual, expected)
1080 }
1081 }
1082 }
1083}