Skip to main content

parquet/arrow/record_reader/
definition_levels.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_array::builder::BooleanBufferBuilder;
19use arrow_buffer::Buffer;
20use arrow_buffer::bit_chunk_iterator::UnalignedBitChunk;
21use bytes::Bytes;
22
23use crate::arrow::buffer::bit_util::count_set_bits;
24use crate::basic::Encoding;
25use crate::column::reader::decoder::{
26    ColumnLevelDecoder, DefinitionLevelDecoder, DefinitionLevelDecoderImpl,
27};
28use crate::errors::{ParquetError, Result};
29use crate::schema::types::ColumnDescPtr;
30
31enum BufferInner {
32    /// Compute levels and null mask
33    Full {
34        levels: Vec<i16>,
35        nulls: BooleanBufferBuilder,
36        max_level: i16,
37    },
38    /// Only compute null bitmask - requires max level to be 1
39    ///
40    /// This is an optimisation for the common case of a nullable scalar column, as decoding
41    /// the definition level data is only required when decoding nested structures
42    ///
43    Mask { nulls: BooleanBufferBuilder },
44}
45
46pub struct DefinitionLevelBuffer {
47    inner: BufferInner,
48
49    /// The length of this buffer
50    ///
51    /// Note: `buffer` and `builder` may contain more elements
52    len: usize,
53}
54
55impl DefinitionLevelBuffer {
56    pub fn new(desc: &ColumnDescPtr, null_mask_only: bool) -> Self {
57        let inner = match null_mask_only {
58            true => {
59                assert_eq!(
60                    desc.max_def_level(),
61                    1,
62                    "max definition level must be 1 to only compute null bitmask"
63                );
64
65                assert_eq!(
66                    desc.max_rep_level(),
67                    0,
68                    "max repetition level must be 0 to only compute null bitmask"
69                );
70
71                BufferInner::Mask {
72                    nulls: BooleanBufferBuilder::new(0),
73                }
74            }
75            false => BufferInner::Full {
76                levels: Vec::new(),
77                nulls: BooleanBufferBuilder::new(0),
78                max_level: desc.max_def_level(),
79            },
80        };
81
82        Self { inner, len: 0 }
83    }
84
85    /// Returns the built level data
86    pub fn consume_levels(&mut self) -> Option<Vec<i16>> {
87        match &mut self.inner {
88            BufferInner::Full { levels, .. } => Some(std::mem::take(levels)),
89            BufferInner::Mask { .. } => None,
90        }
91    }
92
93    /// Returns the built null bitmask, or None if all values are valid
94    pub fn consume_bitmask(&mut self) -> Option<Buffer> {
95        self.len = 0;
96        let nulls = match &mut self.inner {
97            BufferInner::Full { nulls, .. } => nulls,
98            BufferInner::Mask { nulls } => nulls,
99        };
100
101        // Always call finish() to reset the builder state for the next batch.
102        let buffer = nulls.finish().into_inner();
103
104        // If no bitmap was constructed, return None
105        if buffer.is_empty() {
106            return None;
107        }
108
109        Some(buffer)
110    }
111
112    pub fn nulls(&self) -> &BooleanBufferBuilder {
113        match &self.inner {
114            BufferInner::Full { nulls, .. } => nulls,
115            BufferInner::Mask { nulls } => nulls,
116        }
117    }
118
119    /// Returns the raw definition levels accumulated so far, if available.
120    /// Only available when the buffer is in Full mode (nested columns).
121    pub fn levels(&self) -> Option<&[i16]> {
122        match &self.inner {
123            BufferInner::Full { levels, .. } => Some(levels.as_slice()),
124            BufferInner::Mask { .. } => None,
125        }
126    }
127}
128
129/// Build a filtered validity bitmap from definition/repetition levels.
130///
131/// For each level entry where `d >= include_threshold` (when set) and
132/// `r <= max_rep` (when `rep_filter` is provided), appends one bit to
133/// `bitmap`: set when `d >= value_level`, unset otherwise.
134///
135/// Returns the number of bits appended.
136///
137/// This is the shared implementation behind the compact bitmap (leaf
138/// readers with selective padding) and the struct validity bitmap.
139/// Processing uses 64-level chunks with [`compress`] for word-at-a-time
140/// packing.
141pub(crate) fn build_filtered_validity_bitmap(
142    def_levels: &[i16],
143    rep_filter: Option<(&[i16], i16)>,
144    include_threshold: Option<i16>,
145    value_level: i16,
146    bitmap: &mut BooleanBufferBuilder,
147) -> usize {
148    // Fast path: no filtering — every def level produces a bit.
149    if include_threshold.is_none() && rep_filter.is_none() {
150        let chunks = def_levels.chunks_exact(u64::BITS as usize);
151        let remainder = chunks.remainder();
152        for chunk in chunks {
153            let mut word: u64 = 0;
154            for (i, &d) in chunk.iter().enumerate() {
155                word |= ((d >= value_level) as u64) << i;
156            }
157            bitmap.append_word(word, u64::BITS as usize);
158        }
159        for &d in remainder {
160            bitmap.append(d >= value_level);
161        }
162        return def_levels.len();
163    }
164
165    // Filtered path: build include mask + value mask per chunk, compress.
166    let mut item_count: usize = 0;
167    let chunks = def_levels.chunks_exact(u64::BITS as usize);
168    let remainder_offset = def_levels.len() - chunks.remainder().len();
169
170    for (chunk_idx, chunk) in chunks.enumerate() {
171        let base = chunk_idx * u64::BITS as usize;
172        let mut include_mask: u64 = 0;
173        let mut value_mask: u64 = 0;
174        for (i, &d) in chunk.iter().enumerate() {
175            let mut include = true;
176            if let Some(threshold) = include_threshold {
177                if d < threshold {
178                    include = false;
179                }
180            }
181            if include {
182                if let Some((reps, max_rep)) = rep_filter {
183                    if reps[base + i] > max_rep {
184                        include = false;
185                    }
186                }
187            }
188            include_mask |= (include as u64) << i;
189            value_mask |= ((d >= value_level) as u64) << i;
190        }
191
192        let count = include_mask.count_ones() as usize;
193        let compact = compress(value_mask, include_mask);
194        bitmap.append_word(compact, count);
195        item_count += count;
196    }
197
198    for idx in remainder_offset..def_levels.len() {
199        let d = def_levels[idx];
200        if let Some(threshold) = include_threshold {
201            if d < threshold {
202                continue;
203            }
204        }
205        if let Some((reps, max_rep)) = rep_filter {
206            if reps[idx] > max_rep {
207                continue;
208            }
209        }
210        bitmap.append(d >= value_level);
211        item_count += 1;
212    }
213
214    item_count
215}
216
217use crate::util::bit_util::compress;
218
219enum MaybePacked {
220    Packed(PackedDecoder),
221    Fallback(DefinitionLevelDecoderImpl),
222}
223
224pub struct DefinitionLevelBufferDecoder {
225    max_level: i16,
226    decoder: MaybePacked,
227}
228
229impl DefinitionLevelBufferDecoder {
230    pub fn new(max_level: i16, packed: bool) -> Self {
231        let decoder = match packed {
232            true => MaybePacked::Packed(PackedDecoder::new()),
233            false => MaybePacked::Fallback(DefinitionLevelDecoderImpl::new(max_level)),
234        };
235
236        Self { max_level, decoder }
237    }
238}
239
240impl ColumnLevelDecoder for DefinitionLevelBufferDecoder {
241    type Buffer = DefinitionLevelBuffer;
242
243    fn set_data(&mut self, encoding: Encoding, data: Bytes) -> Result<()> {
244        match &mut self.decoder {
245            MaybePacked::Packed(d) => d.set_data(encoding, data),
246            MaybePacked::Fallback(d) => d.set_data(encoding, data)?,
247        };
248        Ok(())
249    }
250}
251
252impl DefinitionLevelDecoder for DefinitionLevelBufferDecoder {
253    fn read_def_levels(
254        &mut self,
255        writer: &mut Self::Buffer,
256        num_levels: usize,
257    ) -> Result<(usize, usize)> {
258        match (&mut writer.inner, &mut self.decoder) {
259            (
260                BufferInner::Full {
261                    levels,
262                    nulls,
263                    max_level,
264                },
265                MaybePacked::Fallback(decoder),
266            ) => {
267                assert_eq!(self.max_level, *max_level);
268
269                let start = levels.len();
270                let (values_read, levels_read) = decoder.read_def_levels(levels, num_levels)?;
271
272                // Safety: slice iterator has a trusted length
273                unsafe {
274                    nulls
275                        .extend_trusted_len(levels[start..].iter().map(|level| level == max_level));
276                }
277
278                Ok((values_read, levels_read))
279            }
280            (BufferInner::Mask { nulls }, MaybePacked::Packed(decoder)) => {
281                assert_eq!(self.max_level, 1);
282
283                // Fast path: if all requested levels are valid (max definition level),
284                // we can skip RLE decoding and just append all-ones to the bitmap.
285                // This is faster than decoding RLE data.
286                if let Some(count) = decoder.try_consume_all_valid(num_levels)? {
287                    nulls.append_n(count, true);
288                    return Ok((count, count)); // values_read == levels_read when all valid
289                }
290
291                // Normal path: decode RLE data into the bitmap
292                let start = nulls.len();
293                let levels_read = decoder.read(nulls, num_levels)?;
294
295                let values_read = count_set_bits(nulls.as_slice(), start..start + levels_read);
296                Ok((values_read, levels_read))
297            }
298            _ => unreachable!("inconsistent null mask"),
299        }
300    }
301
302    fn skip_def_levels(&mut self, num_levels: usize) -> Result<(usize, usize)> {
303        match &mut self.decoder {
304            MaybePacked::Fallback(decoder) => decoder.skip_def_levels(num_levels),
305            MaybePacked::Packed(decoder) => decoder.skip(num_levels),
306        }
307    }
308}
309
310/// An optimized decoder for decoding [RLE] and [BIT_PACKED] data with a bit width of 1
311/// directly into a bitmask
312///
313/// This is significantly faster than decoding the data into `[i16]` and then computing
314/// a bitmask from this, as not only can it skip this buffer allocation and construction,
315/// but it can exploit properties of the encoded data to reduce work further
316///
317/// In particular:
318///
319/// * Packed runs are already bitmask encoded and can simply be appended
320/// * Runs of 1 or 0 bits can be efficiently appended with byte (or larger) operations
321///
322/// [RLE]: https://github.com/apache/parquet-format/blob/master/Encodings.md#run-length-encoding--bit-packing-hybrid-rle--3
323/// [BIT_PACKED]: https://github.com/apache/parquet-format/blob/master/Encodings.md#bit-packed-deprecated-bit_packed--4
324struct PackedDecoder {
325    data: Bytes,
326    data_offset: usize,
327    rle_left: usize,
328    rle_value: bool,
329    packed_count: usize,
330    packed_offset: usize,
331}
332
333impl PackedDecoder {
334    fn next_rle_block(&mut self) -> Result<()> {
335        let indicator_value = self.decode_header()?;
336        if indicator_value & 1 == 1 {
337            let len = (indicator_value >> 1) as usize;
338            self.packed_count = len * 8;
339            self.packed_offset = 0;
340        } else {
341            self.rle_left = (indicator_value >> 1) as usize;
342            let byte = *self.data.as_ref().get(self.data_offset).ok_or_else(|| {
343                ParquetError::EOF(
344                    "unexpected end of file whilst decoding definition levels rle value".into(),
345                )
346            })?;
347
348            self.data_offset += 1;
349            self.rle_value = byte != 0;
350        }
351        Ok(())
352    }
353
354    /// Decodes a VLQ encoded little endian integer and returns it
355    fn decode_header(&mut self) -> Result<i64> {
356        let mut offset = 0;
357        let mut v: i64 = 0;
358        while offset < 10 {
359            let byte = *self
360                .data
361                .as_ref()
362                .get(self.data_offset + offset)
363                .ok_or_else(|| {
364                    ParquetError::EOF(
365                        "unexpected end of file whilst decoding definition levels rle header"
366                            .into(),
367                    )
368                })?;
369
370            v |= ((byte & 0x7F) as i64) << (offset * 7);
371            offset += 1;
372            if byte & 0x80 == 0 {
373                self.data_offset += offset;
374                return Ok(v);
375            }
376        }
377        Err(general_err!("too many bytes for VLQ"))
378    }
379}
380
381impl PackedDecoder {
382    fn new() -> Self {
383        Self {
384            data: Bytes::from(vec![]),
385            data_offset: 0,
386            rle_left: 0,
387            rle_value: false,
388            packed_count: 0,
389            packed_offset: 0,
390        }
391    }
392
393    fn set_data(&mut self, encoding: Encoding, data: Bytes) {
394        self.rle_left = 0;
395        self.rle_value = false;
396        self.packed_offset = 0;
397        self.packed_count = match encoding {
398            Encoding::RLE => 0,
399            #[allow(deprecated)]
400            Encoding::BIT_PACKED => data.len() * 8,
401            _ => unreachable!("invalid level encoding: {}", encoding),
402        };
403        self.data = data;
404        self.data_offset = 0;
405    }
406
407    /// Try to consume `len` levels if all are valid (max definition level).
408    ///
409    /// Returns `Ok(Some(count))` if successfully consumed `count` all-valid levels.
410    /// Returns `Ok(None)` if there are any nulls or packed data that prevents fast path.
411    ///
412    /// Note: On `None`, the decoder state may have advanced to the next RLE block,
413    /// but only if `rle_left` was zero (i.e., the block would have been loaded
414    /// on the next read anyway).
415    fn try_consume_all_valid(&mut self, len: usize) -> Result<Option<usize>> {
416        // If no active run and no packed data pending, try to parse the next RLE block
417        if self.rle_left == 0 && self.packed_count == self.packed_offset {
418            if self.data_offset < self.data.len() {
419                self.next_rle_block()?;
420            } else {
421                // No more data available
422                return Ok(None);
423            }
424        }
425
426        // Fast path only works when we have an active RLE run of true values
427        // that covers the entire requested length.
428        if self.rle_left >= len && self.rle_value {
429            self.rle_left -= len;
430            return Ok(Some(len));
431        }
432
433        // Any other case (null run, packed data, or insufficient length)
434        // falls back to normal path
435        Ok(None)
436    }
437
438    fn read(&mut self, buffer: &mut BooleanBufferBuilder, len: usize) -> Result<usize> {
439        let mut read = 0;
440        while read != len {
441            if self.rle_left != 0 {
442                let to_read = self.rle_left.min(len - read);
443                buffer.append_n(to_read, self.rle_value);
444                self.rle_left -= to_read;
445                read += to_read;
446            } else if self.packed_count != self.packed_offset {
447                let to_read = (self.packed_count - self.packed_offset).min(len - read);
448                let offset = self.data_offset * 8 + self.packed_offset;
449                buffer.append_packed_range(offset..offset + to_read, self.data.as_ref());
450                self.packed_offset += to_read;
451                read += to_read;
452
453                if self.packed_offset == self.packed_count {
454                    self.data_offset += self.packed_count / 8;
455                }
456            } else if self.data_offset == self.data.len() {
457                break;
458            } else {
459                self.next_rle_block()?
460            }
461        }
462        Ok(read)
463    }
464
465    /// Skips `level_num` definition levels
466    ///
467    /// Returns the number of values skipped and the number of levels skipped
468    fn skip(&mut self, level_num: usize) -> Result<(usize, usize)> {
469        let mut skipped_value = 0;
470        let mut skipped_level = 0;
471        while skipped_level != level_num {
472            if self.rle_left != 0 {
473                let to_skip = self.rle_left.min(level_num - skipped_level);
474                self.rle_left -= to_skip;
475                skipped_level += to_skip;
476                if self.rle_value {
477                    skipped_value += to_skip;
478                }
479            } else if self.packed_count != self.packed_offset {
480                let to_skip =
481                    (self.packed_count - self.packed_offset).min(level_num - skipped_level);
482                let offset = self.data_offset * 8 + self.packed_offset;
483                let bit_chunk = UnalignedBitChunk::new(self.data.as_ref(), offset, to_skip);
484                skipped_value += bit_chunk.count_ones();
485                self.packed_offset += to_skip;
486                skipped_level += to_skip;
487                if self.packed_offset == self.packed_count {
488                    self.data_offset += self.packed_count / 8;
489                }
490            } else if self.data_offset == self.data.len() {
491                break;
492            } else {
493                self.next_rle_block()?
494            }
495        }
496        Ok((skipped_value, skipped_level))
497    }
498}
499
500#[cfg(test)]
501mod tests {
502    use super::*;
503
504    use crate::encodings::rle::RleEncoder;
505    use rand::{Rng, rng};
506
507    #[test]
508    fn test_build_validity_bitmap_unfiltered_word_chunk() {
509        // 65 levels forces the unfiltered path to process one full u64 word
510        // with append_word, plus a remainder bit.
511        let def_levels = (0..65)
512            .map(|i| if i % 3 == 0 { 2 } else { 1 })
513            .collect::<Vec<_>>();
514        let mut bitmap = BooleanBufferBuilder::new(0);
515
516        assert_eq!(
517            build_filtered_validity_bitmap(&def_levels, None, None, 2, &mut bitmap),
518            def_levels.len()
519        );
520
521        let bitmap = bitmap.finish();
522        for (idx, def) in def_levels.iter().enumerate() {
523            assert_eq!(bitmap.value(idx), *def >= 2);
524        }
525    }
526
527    #[test]
528    fn test_packed_decoder() {
529        let mut rng = rng();
530        let len: usize = rng.random_range(512..1024);
531
532        let mut expected = BooleanBufferBuilder::new(len);
533        let mut encoder = RleEncoder::new(1, 1024);
534        for _ in 0..len {
535            let bool = rng.random_bool(0.8);
536            encoder.put(bool as u64);
537            expected.append(bool);
538        }
539        assert_eq!(expected.len(), len);
540
541        let encoded = encoder.consume();
542        let mut decoder = PackedDecoder::new();
543        decoder.set_data(Encoding::RLE, encoded.into());
544
545        // Decode data in random length intervals
546        let mut decoded = BooleanBufferBuilder::new(len);
547        loop {
548            let remaining = len - decoded.len();
549            if remaining == 0 {
550                break;
551            }
552
553            let to_read = rng.random_range(1..=remaining);
554            decoder.read(&mut decoded, to_read).unwrap();
555        }
556
557        assert_eq!(decoded.len(), len);
558        assert_eq!(decoded.as_slice(), expected.as_slice());
559    }
560
561    #[test]
562    fn test_packed_decoder_skip() {
563        let mut rng = rng();
564        let len: usize = rng.random_range(512..1024);
565
566        let mut expected = BooleanBufferBuilder::new(len);
567        let mut encoder = RleEncoder::new(1, 1024);
568
569        let mut total_value = 0;
570        for _ in 0..len {
571            let bool = rng.random_bool(0.8);
572            encoder.put(bool as u64);
573            expected.append(bool);
574            if bool {
575                total_value += 1;
576            }
577        }
578        assert_eq!(expected.len(), len);
579
580        let encoded = encoder.consume();
581        let mut decoder = PackedDecoder::new();
582        decoder.set_data(Encoding::RLE, encoded.into());
583
584        let mut skip_value = 0;
585        let mut read_value = 0;
586        let mut skip_level = 0;
587        let mut read_level = 0;
588
589        loop {
590            let offset = skip_level + read_level;
591            let remaining_levels = len - offset;
592            if remaining_levels == 0 {
593                break;
594            }
595            let to_read_or_skip_level = rng.random_range(1..=remaining_levels);
596            if rng.random_bool(0.5) {
597                let (skip_val_num, skip_level_num) = decoder.skip(to_read_or_skip_level).unwrap();
598                skip_value += skip_val_num;
599                skip_level += skip_level_num
600            } else {
601                let mut decoded = BooleanBufferBuilder::new(to_read_or_skip_level);
602                let read_level_num = decoder.read(&mut decoded, to_read_or_skip_level).unwrap();
603                read_level += read_level_num;
604                for i in 0..read_level_num {
605                    assert!(!decoded.is_empty());
606                    //check each read bit
607                    let read_bit = decoded.get_bit(i);
608                    if read_bit {
609                        read_value += 1;
610                    }
611                    let expect_bit = expected.get_bit(i + offset);
612                    assert_eq!(read_bit, expect_bit);
613                }
614            }
615        }
616        assert_eq!(read_level + skip_level, len);
617        assert_eq!(read_value + skip_value, total_value);
618    }
619
620    #[test]
621    fn test_try_consume_all_valid() {
622        // Test with all-valid data (all 1s) - single RLE run
623        let len = 100;
624        let mut encoder = RleEncoder::new(1, 1024);
625        for _ in 0..len {
626            encoder.put(1); // all valid
627        }
628        let encoded = encoder.consume();
629        let mut decoder = PackedDecoder::new();
630        decoder.set_data(Encoding::RLE, encoded.into());
631
632        // try_consume_all_valid now parses the RLE block itself, no need to read first
633        let result = decoder.try_consume_all_valid(len).unwrap();
634        assert_eq!(result, Some(len));
635
636        // Test with all-null data (all 0s)
637        let mut encoder = RleEncoder::new(1, 1024);
638        for _ in 0..len {
639            encoder.put(0); // all null
640        }
641        let encoded = encoder.consume();
642        let mut decoder = PackedDecoder::new();
643        decoder.set_data(Encoding::RLE, encoded.into());
644
645        // Should return None because rle_value is false (all nulls)
646        let result = decoder.try_consume_all_valid(len).unwrap();
647        assert_eq!(result, None);
648
649        // Test when requesting more than available in current RLE run
650        let mut encoder = RleEncoder::new(1, 1024);
651        for _ in 0..10 {
652            encoder.put(1); // small run of valid
653        }
654        for _ in 0..10 {
655            encoder.put(0); // followed by nulls
656        }
657        let encoded = encoder.consume();
658        let mut decoder = PackedDecoder::new();
659        decoder.set_data(Encoding::RLE, encoded.into());
660
661        // Request more than the valid run - should return None
662        // (because we don't look ahead to next block)
663        let result = decoder.try_consume_all_valid(20).unwrap();
664        assert_eq!(result, None);
665
666        // Reset decoder and try requesting within the run
667        decoder.set_data(Encoding::RLE, {
668            let mut encoder = RleEncoder::new(1, 1024);
669            for _ in 0..10 {
670                encoder.put(1);
671            }
672            for _ in 0..10 {
673                encoder.put(0);
674            }
675            encoder.consume().into()
676        });
677
678        let result = decoder.try_consume_all_valid(5).unwrap();
679        assert_eq!(result, Some(5));
680
681        // After skipping 5, we should have 5 left in the valid RLE run
682        let result = decoder.try_consume_all_valid(5).unwrap();
683        assert_eq!(result, Some(5));
684
685        // Now the valid run is exhausted, next call should parse the null run and return None
686        let result = decoder.try_consume_all_valid(5).unwrap();
687        assert_eq!(result, None);
688    }
689}