Skip to main content

parquet/arrow/record_reader/
mod.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18use 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
39/// A `RecordReader` is a stateful column reader that delimits semantic records.
40pub type RecordReader<T> = GenericRecordReader<Vec<<T as DataType>::T>, ColumnValueDecoderImpl<T>>;
41
42pub(crate) type ColumnReader<CV> =
43    GenericColumnReader<RepetitionLevelDecoderImpl, DefinitionLevelBufferDecoder, CV>;
44
45/// A generic stateful column reader that delimits semantic records
46///
47/// This type is hidden from the docs, and relies on private traits with no
48/// public implementations. As such this type signature may be changed without
49/// breaking downstream users as it can only be constructed through type aliases
50pub struct GenericRecordReader<V, CV> {
51    column_desc: ColumnDescPtr,
52
53    /// Values buffer, lazily initialized on first read to avoid
54    /// allocating a buffer that may never be used (e.g., after the last batch)
55    values: Option<V>,
56    def_levels: Option<DefinitionLevelBuffer>,
57    rep_levels: Option<Vec<i16>>,
58    column_reader: Option<ColumnReader<CV>>,
59    /// Number of buffered levels / null-padded values
60    num_values: usize,
61    /// Number of buffered records
62    num_records: usize,
63    /// Capacity hint for pre-allocating buffers based on batch size
64    capacity_hint: usize,
65    /// Number of values in the values buffer (may differ from num_values when
66    /// padding_threshold is set, since parent-level padding is excluded).
67    values_written: usize,
68    /// Definition-level threshold used for selective null padding.
69    ///
70    /// With full padding (`None`), the leaf values buffer has one slot for each
71    /// decoded definition level. This includes placeholders for null or empty
72    /// parent lists, which parent `ListArrayReader`s later have to filter out
73    /// before computing offsets.
74    ///
75    /// With selective padding (`Some(threshold)`), the threshold is the nearest
76    /// enclosing list/map definition level. Entries with `def < threshold`
77    /// describe a null/empty parent and are skipped entirely. Entries with
78    /// `def >= threshold` belong to an actual child item slot: real values are
79    /// copied, and item-level nulls are padded. The companion `compact_bitmap`
80    /// has the same compact length and becomes the leaf null bitmap.
81    padding_threshold: Option<i16>,
82    /// Compact bitmap accumulated during selective padding. Each bit
83    /// corresponds to an item-level entry (def >= threshold): set when the
84    /// value is real (def >= max_def), unset for item-level nulls. Used both
85    /// as the valid_mask for `pad_nulls` (via `as_slice()`) and as the null
86    /// bitmap consumed by the leaf reader (via `consume_compact_bitmap`).
87    compact_bitmap: Option<BooleanBufferBuilder>,
88}
89
90impl<V, CV> GenericRecordReader<V, CV>
91where
92    V: ValuesBuffer,
93    CV: ColumnValueDecoder<Buffer = V>,
94{
95    /// Create a new [`GenericRecordReader`]
96    ///
97    /// The capacity is used to pre-allocate internal buffers for full-padding
98    /// reads, avoiding reallocations when reading fragmented row selections.
99    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, // Lazily initialized on first read
107            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    /// Set the current page reader.
121    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    /// Try to read `num_records` of column data into internal buffer.
143    ///
144    /// # Returns
145    ///
146    /// Number of actual records read.
147    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    /// Try to skip the next `num_records` rows
165    ///
166    /// # Returns
167    ///
168    /// Number of records skipped
169    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    /// Returns number of records stored in buffer.
177    #[allow(unused)]
178    pub fn num_records(&self) -> usize {
179        self.num_records
180    }
181
182    /// Return number of values stored in buffer.
183    /// If the parquet column is not repeated, it should be equals to `num_records`,
184    /// otherwise it should be larger than or equal to `num_records`.
185    pub fn num_values(&self) -> usize {
186        self.num_values
187    }
188
189    /// Returns definition level data.
190    /// The implementation has side effects. It will create a new buffer to hold those
191    /// definition level values that have already been read into memory but not counted
192    /// as record values, e.g. those from `self.num_values` to `self.values_written`.
193    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    /// Return repetition level data.
198    /// The side effect is similar to `consume_def_levels`.
199    pub fn consume_rep_levels(&mut self) -> Option<Vec<i16>> {
200        self.rep_levels.as_mut().map(std::mem::take)
201    }
202
203    /// Returns currently stored buffer data.
204    /// The side effect is similar to `consume_def_levels`.
205    pub fn consume_record_data(&mut self) -> V {
206        // Take the buffer, leaving None. The next read will lazily allocate a new buffer.
207        // This avoids allocating a buffer that may never be used (e.g., after the last batch).
208        self.values.take().unwrap_or_else(|| V::with_capacity(0))
209    }
210
211    /// Reset state of record reader.
212    /// Should be called after consuming data, e.g. `consume_rep_levels`,
213    /// `consume_rep_levels`, `consume_record_data` and `consume_compact_bitmap`.
214    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    /// Returns the maximum definition level for the column being read.
222    pub fn max_def_level(&self) -> i16 {
223        self.column_desc.max_def_level()
224    }
225
226    /// Set the padding threshold. When set, `pad_nulls` only pads entries
227    /// where `def >= threshold` (item-level nulls within non-null lists),
228    /// skipping list-level padding entries (def < threshold).
229    pub fn set_padding_threshold(&mut self, threshold: i16) {
230        self.padding_threshold = Some(threshold);
231    }
232
233    /// Returns the number of values in the values buffer.
234    /// When padding_threshold is None, this equals `num_values` (full padding).
235    /// When padding_threshold is set, this is the item_count (selective padding).
236    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    /// Consume the compact null bitmap built during selective padding.
245    /// Returns the full bitmap when not using selective padding.
246    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    /// Returns bitmap data for nullable columns.
260    /// For non-nullable columns, the bitmap is discarded.
261    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        // While we always consume the bitmask here, we only want to return
268        // the bitmask for nullable arrays. (Marking nulls on a non-nullable
269        // array may fail validations, even if those nulls are masked off at
270        // a higher level.)
271        if self.column_desc.self_type().is_optional() {
272            mask
273        } else {
274            None
275        }
276    }
277
278    /// Try to read one batch of data returning the number of records read
279    fn read_one_batch(&mut self, batch_size: usize) -> Result<usize> {
280        if batch_size == 0 {
281            return Ok(0);
282        }
283        // Update capacity hint to the largest batch size seen.
284        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            // Pad values to item_count positions (only if there are gaps).
356            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            // Full padding: pad all null positions.
378            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
397/// Returns true if we do not need to unpack the nullability for this column, this is
398/// only possible if the max definition level is 1, and corresponds to nulls at the
399/// leaf level, as opposed to a nullable parent nested type
400fn 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        // Construct column schema
424        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        // Construct record reader
435        let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
436
437        // First page
438
439        // Records data:
440        // test_schema
441        //   leaf: 4
442        // test_schema
443        //   leaf: 7
444        // test_schema
445        //   leaf: 6
446        // test_schema
447        //   left: 3
448        // test_schema
449        //   left: 2
450        {
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        // Second page
467
468        // Records data:
469        // test_schema
470        //   leaf: 8
471        // test_schema
472        //   leaf: 9
473        {
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        // Construct column schema
527        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        // Construct record reader
541        let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
542
543        // First page
544
545        // Records data:
546        // test_schema
547        //   test_struct
548        // test_schema
549        //   test_struct
550        //     left: 7
551        // test_schema
552        // test_schema
553        //   test_struct
554        //     leaf: 6
555        // test_schema
556        //   test_struct
557        //     leaf: 6
558        {
559            let values = [7, 6, 3];
560            //empty, non-empty, empty, non-empty, non-empty
561            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        // Second page
578
579        // Records data:
580        // test_schema
581        // test_schema
582        //   test_struct
583        //     left: 8
584        {
585            let values = [8];
586            //empty, non-empty
587            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        // Verify result def levels
601        assert_eq!(
602            Some(vec![1i16, 2i16, 0i16, 2i16, 2i16, 0i16, 2i16]),
603            record_reader.consume_def_levels()
604        );
605
606        // Verify bitmap
607        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        // Verify result record data
612        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        // Only validate valid values are equal
618        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        // Construct column schema
629        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        // Construct record reader
643        let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
644
645        // First page
646
647        // Records data:
648        // test_schema
649        //   test_struct
650        //     leaf: 4
651        // test_schema
652        // test_schema
653        //   test_struct
654        //   test_struct
655        //     leaf: 7
656        //     leaf: 6
657        //     leaf: 3
658        //   test_struct
659        //     leaf: 2
660        {
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        // Second page
682
683        // Records data:
684        // test_schema
685        //   test_struct
686        //     leaf: 8
687        //     leaf: 9
688        {
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        // Verify result def levels
707        assert_eq!(
708            Some(vec![2i16, 0i16, 1i16, 2i16, 2i16, 2i16, 2i16, 2i16, 2i16]),
709            record_reader.consume_def_levels()
710        );
711
712        // Verify bitmap
713        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        // Verify result record data
718        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        // Only validate valid values are equal
723        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        // Construct column schema
823        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        // Construct record reader
835        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        // Construct column schema
863        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        // Construct column schema
913        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        // Construct record reader
924        let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
925
926        // First page
927
928        // Records data:
929        // test_schema
930        //   leaf: 4
931        // test_schema
932        //   leaf: 7
933        // test_schema
934        //   leaf: 6
935        // test_schema
936        //   left: 3
937        // test_schema
938        //   left: 2
939        {
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        // Second page
956
957        // Records data:
958        // test_schema
959        //   leaf: 8
960        // test_schema
961        //   leaf: 9
962        {
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        // Construct column schema
984        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        // Construct record reader
998        let mut record_reader = RecordReader::<Int32Type>::new(desc.clone(), DEFAULT_BATCH_SIZE);
999
1000        // First page
1001
1002        // Records data:
1003        // test_schema
1004        //   test_struct
1005        // test_schema
1006        //   test_struct
1007        //     leaf: 7
1008        // test_schema
1009        // test_schema
1010        //   test_struct
1011        //     leaf: 6
1012        // test_schema
1013        //   test_struct
1014        //     leaf: 6
1015        {
1016            let values = [7, 6, 3];
1017            //empty, non-empty, empty, non-empty, non-empty
1018            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        // Second page
1035
1036        // Records data:
1037        // test_schema
1038        // test_schema
1039        //   test_struct
1040        //     left: 8
1041        {
1042            let values = [8];
1043            //empty, non-empty
1044            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        // Verify result def levels
1059        assert_eq!(
1060            Some(vec![0i16, 2i16, 2i16]),
1061            record_reader.consume_def_levels()
1062        );
1063
1064        // Verify bitmap
1065        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        // Verify result record data
1070        let actual = record_reader.consume_record_data();
1071
1072        let expected = &[0, 6, 3];
1073        assert_eq!(actual.len(), expected.len());
1074
1075        // Only validate valid values are equal
1076        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}