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 include = include_threshold.is_none_or(|threshold| threshold <= d)
176                && rep_filter.is_none_or(|(reps, max_rep)| reps[base + i] <= max_rep);
177            include_mask |= (include as u64) << i;
178            value_mask |= ((d >= value_level) as u64) << i;
179        }
180
181        let count = include_mask.count_ones() as usize;
182        let compact = compress(value_mask, include_mask);
183        bitmap.append_word(compact, count);
184        item_count += count;
185    }
186
187    for idx in remainder_offset..def_levels.len() {
188        let d = def_levels[idx];
189        if let Some(threshold) = include_threshold
190            && d < threshold
191        {
192            continue;
193        }
194        if let Some((reps, max_rep)) = rep_filter
195            && reps[idx] > max_rep
196        {
197            continue;
198        }
199        bitmap.append(d >= value_level);
200        item_count += 1;
201    }
202
203    item_count
204}
205
206use crate::util::bit_util::compress;
207
208enum MaybePacked {
209    Packed(PackedDecoder),
210    Fallback(DefinitionLevelDecoderImpl),
211}
212
213pub struct DefinitionLevelBufferDecoder {
214    max_level: i16,
215    decoder: MaybePacked,
216}
217
218impl DefinitionLevelBufferDecoder {
219    pub fn new(max_level: i16, packed: bool) -> Self {
220        let decoder = match packed {
221            true => MaybePacked::Packed(PackedDecoder::new()),
222            false => MaybePacked::Fallback(DefinitionLevelDecoderImpl::new(max_level)),
223        };
224
225        Self { max_level, decoder }
226    }
227}
228
229impl ColumnLevelDecoder for DefinitionLevelBufferDecoder {
230    type Buffer = DefinitionLevelBuffer;
231
232    fn set_data(&mut self, encoding: Encoding, data: Bytes) -> Result<()> {
233        match &mut self.decoder {
234            MaybePacked::Packed(d) => d.set_data(encoding, data),
235            MaybePacked::Fallback(d) => d.set_data(encoding, data)?,
236        }
237        Ok(())
238    }
239}
240
241impl DefinitionLevelDecoder for DefinitionLevelBufferDecoder {
242    fn read_def_levels(
243        &mut self,
244        writer: &mut Self::Buffer,
245        num_levels: usize,
246    ) -> Result<(usize, usize)> {
247        match (&mut writer.inner, &mut self.decoder) {
248            (
249                BufferInner::Full {
250                    levels,
251                    nulls,
252                    max_level,
253                },
254                MaybePacked::Fallback(decoder),
255            ) => {
256                assert_eq!(self.max_level, *max_level);
257
258                let start = levels.len();
259                let (values_read, levels_read) = decoder.read_def_levels(levels, num_levels)?;
260
261                // Safety: slice iterator has a trusted length
262                unsafe {
263                    nulls
264                        .extend_trusted_len(levels[start..].iter().map(|level| level == max_level));
265                }
266
267                Ok((values_read, levels_read))
268            }
269            (BufferInner::Mask { nulls }, MaybePacked::Packed(decoder)) => {
270                assert_eq!(self.max_level, 1);
271
272                // Fast path: if all requested levels are valid (max definition level),
273                // we can skip RLE decoding and just append all-ones to the bitmap.
274                // This is faster than decoding RLE data.
275                if let Some(count) = decoder.try_consume_all_valid(num_levels)? {
276                    nulls.append_n(count, true);
277                    return Ok((count, count)); // values_read == levels_read when all valid
278                }
279
280                // Normal path: decode RLE data into the bitmap
281                let start = nulls.len();
282                let levels_read = decoder.read(nulls, num_levels)?;
283
284                let values_read = count_set_bits(nulls.as_slice(), start..start + levels_read);
285                Ok((values_read, levels_read))
286            }
287            _ => unreachable!("inconsistent null mask"),
288        }
289    }
290
291    fn skip_def_levels(&mut self, num_levels: usize) -> Result<(usize, usize)> {
292        match &mut self.decoder {
293            MaybePacked::Fallback(decoder) => decoder.skip_def_levels(num_levels),
294            MaybePacked::Packed(decoder) => decoder.skip(num_levels),
295        }
296    }
297}
298
299/// An optimized decoder for decoding [RLE] and [BIT_PACKED] data with a bit width of 1
300/// directly into a bitmask
301///
302/// This is significantly faster than decoding the data into `[i16]` and then computing
303/// a bitmask from this, as not only can it skip this buffer allocation and construction,
304/// but it can exploit properties of the encoded data to reduce work further
305///
306/// In particular:
307///
308/// * Packed runs are already bitmask encoded and can simply be appended
309/// * Runs of 1 or 0 bits can be efficiently appended with byte (or larger) operations
310///
311/// [RLE]: https://github.com/apache/parquet-format/blob/master/Encodings.md#run-length-encoding--bit-packing-hybrid-rle--3
312/// [BIT_PACKED]: https://github.com/apache/parquet-format/blob/master/Encodings.md#bit-packed-deprecated-bit_packed--4
313struct PackedDecoder {
314    data: Bytes,
315    data_offset: usize,
316    rle_left: usize,
317    rle_value: bool,
318    packed_count: usize,
319    packed_offset: usize,
320}
321
322impl PackedDecoder {
323    fn next_rle_block(&mut self) -> Result<()> {
324        let indicator_value = self.decode_header()?;
325        if indicator_value & 1 == 1 {
326            let len = (indicator_value >> 1) as usize;
327            self.packed_count = len * 8;
328            self.packed_offset = 0;
329        } else {
330            self.rle_left = (indicator_value >> 1) as usize;
331            let byte = *self.data.as_ref().get(self.data_offset).ok_or_else(|| {
332                ParquetError::EOF(
333                    "unexpected end of file whilst decoding definition levels rle value".into(),
334                )
335            })?;
336
337            self.data_offset += 1;
338            self.rle_value = byte != 0;
339        }
340        Ok(())
341    }
342
343    /// Decodes a VLQ encoded little endian integer and returns it
344    fn decode_header(&mut self) -> Result<i64> {
345        let mut offset = 0;
346        let mut v: i64 = 0;
347        while offset < 10 {
348            let byte = *self
349                .data
350                .as_ref()
351                .get(self.data_offset + offset)
352                .ok_or_else(|| {
353                    ParquetError::EOF(
354                        "unexpected end of file whilst decoding definition levels rle header"
355                            .into(),
356                    )
357                })?;
358
359            v |= ((byte & 0x7F) as i64) << (offset * 7);
360            offset += 1;
361            if byte & 0x80 == 0 {
362                self.data_offset += offset;
363                return Ok(v);
364            }
365        }
366        Err(general_err!("too many bytes for VLQ"))
367    }
368}
369
370impl PackedDecoder {
371    fn new() -> Self {
372        Self {
373            data: Bytes::from(vec![]),
374            data_offset: 0,
375            rle_left: 0,
376            rle_value: false,
377            packed_count: 0,
378            packed_offset: 0,
379        }
380    }
381
382    fn set_data(&mut self, encoding: Encoding, data: Bytes) {
383        self.rle_left = 0;
384        self.rle_value = false;
385        self.packed_offset = 0;
386        self.packed_count = match encoding {
387            Encoding::RLE => 0,
388            #[expect(deprecated)]
389            Encoding::BIT_PACKED => data.len() * 8,
390            _ => unreachable!("invalid level encoding: {}", encoding),
391        };
392        self.data = data;
393        self.data_offset = 0;
394    }
395
396    /// Try to consume `len` levels if all are valid (max definition level).
397    ///
398    /// Returns `Ok(Some(count))` if successfully consumed `count` all-valid levels.
399    /// Returns `Ok(None)` if there are any nulls or packed data that prevents fast path.
400    ///
401    /// Note: On `None`, the decoder state may have advanced to the next RLE block,
402    /// but only if `rle_left` was zero (i.e., the block would have been loaded
403    /// on the next read anyway).
404    fn try_consume_all_valid(&mut self, len: usize) -> Result<Option<usize>> {
405        // If no active run and no packed data pending, try to parse the next RLE block
406        if self.rle_left == 0 && self.packed_count == self.packed_offset {
407            if self.data_offset < self.data.len() {
408                self.next_rle_block()?;
409            } else {
410                // No more data available
411                return Ok(None);
412            }
413        }
414
415        // Fast path only works when we have an active RLE run of true values
416        // that covers the entire requested length.
417        if self.rle_left >= len && self.rle_value {
418            self.rle_left -= len;
419            return Ok(Some(len));
420        }
421
422        // Any other case (null run, packed data, or insufficient length)
423        // falls back to normal path
424        Ok(None)
425    }
426
427    fn read(&mut self, buffer: &mut BooleanBufferBuilder, len: usize) -> Result<usize> {
428        let mut read = 0;
429        while read != len {
430            if self.rle_left != 0 {
431                let to_read = self.rle_left.min(len - read);
432                buffer.append_n(to_read, self.rle_value);
433                self.rle_left -= to_read;
434                read += to_read;
435            } else if self.packed_count != self.packed_offset {
436                let to_read = (self.packed_count - self.packed_offset).min(len - read);
437                let offset = self.data_offset * 8 + self.packed_offset;
438                buffer.append_packed_range(offset..offset + to_read, self.data.as_ref());
439                self.packed_offset += to_read;
440                read += to_read;
441
442                if self.packed_offset == self.packed_count {
443                    self.data_offset += self.packed_count / 8;
444                }
445            } else if self.data_offset == self.data.len() {
446                break;
447            } else {
448                self.next_rle_block()?
449            }
450        }
451        Ok(read)
452    }
453
454    /// Skips `level_num` definition levels
455    ///
456    /// Returns the number of values skipped and the number of levels skipped
457    fn skip(&mut self, level_num: usize) -> Result<(usize, usize)> {
458        let mut skipped_value = 0;
459        let mut skipped_level = 0;
460        while skipped_level != level_num {
461            if self.rle_left != 0 {
462                let to_skip = self.rle_left.min(level_num - skipped_level);
463                self.rle_left -= to_skip;
464                skipped_level += to_skip;
465                if self.rle_value {
466                    skipped_value += to_skip;
467                }
468            } else if self.packed_count != self.packed_offset {
469                let to_skip =
470                    (self.packed_count - self.packed_offset).min(level_num - skipped_level);
471                let offset = self.data_offset * 8 + self.packed_offset;
472                let bit_chunk = UnalignedBitChunk::new(self.data.as_ref(), offset, to_skip);
473                skipped_value += bit_chunk.count_ones();
474                self.packed_offset += to_skip;
475                skipped_level += to_skip;
476                if self.packed_offset == self.packed_count {
477                    self.data_offset += self.packed_count / 8;
478                }
479            } else if self.data_offset == self.data.len() {
480                break;
481            } else {
482                self.next_rle_block()?
483            }
484        }
485        Ok((skipped_value, skipped_level))
486    }
487}
488
489#[cfg(test)]
490mod tests {
491    use super::*;
492
493    use crate::encodings::rle::RleEncoder;
494    use rand::{RngExt, rng};
495
496    #[test]
497    fn test_build_validity_bitmap_unfiltered_word_chunk() {
498        // 65 levels forces the unfiltered path to process one full u64 word
499        // with append_word, plus a remainder bit.
500        let def_levels = (0..65)
501            .map(|i| if i % 3 == 0 { 2 } else { 1 })
502            .collect::<Vec<_>>();
503        let mut bitmap = BooleanBufferBuilder::new(0);
504
505        assert_eq!(
506            build_filtered_validity_bitmap(&def_levels, None, None, 2, &mut bitmap),
507            def_levels.len()
508        );
509
510        let bitmap = bitmap.finish();
511        for (idx, def) in def_levels.iter().enumerate() {
512            assert_eq!(bitmap.value(idx), *def >= 2);
513        }
514    }
515
516    #[test]
517    fn test_packed_decoder() {
518        let mut rng = rng();
519        let len: usize = rng.random_range(512..1024);
520
521        let mut expected = BooleanBufferBuilder::new(len);
522        let mut encoder = RleEncoder::new(1, 1024);
523        for _ in 0..len {
524            let bool = rng.random_bool(0.8);
525            encoder.put(bool as u64);
526            expected.append(bool);
527        }
528        assert_eq!(expected.len(), len);
529
530        let encoded = encoder.consume();
531        let mut decoder = PackedDecoder::new();
532        decoder.set_data(Encoding::RLE, encoded.into());
533
534        // Decode data in random length intervals
535        let mut decoded = BooleanBufferBuilder::new(len);
536        loop {
537            let remaining = len - decoded.len();
538            if remaining == 0 {
539                break;
540            }
541
542            let to_read = rng.random_range(1..=remaining);
543            decoder.read(&mut decoded, to_read).unwrap();
544        }
545
546        assert_eq!(decoded.len(), len);
547        assert_eq!(decoded.as_slice(), expected.as_slice());
548    }
549
550    #[test]
551    fn test_packed_decoder_skip() {
552        let mut rng = rng();
553        let len: usize = rng.random_range(512..1024);
554
555        let mut expected = BooleanBufferBuilder::new(len);
556        let mut encoder = RleEncoder::new(1, 1024);
557
558        let mut total_value = 0;
559        for _ in 0..len {
560            let bool = rng.random_bool(0.8);
561            encoder.put(bool as u64);
562            expected.append(bool);
563            if bool {
564                total_value += 1;
565            }
566        }
567        assert_eq!(expected.len(), len);
568
569        let encoded = encoder.consume();
570        let mut decoder = PackedDecoder::new();
571        decoder.set_data(Encoding::RLE, encoded.into());
572
573        let mut skip_value = 0;
574        let mut read_value = 0;
575        let mut skip_level = 0;
576        let mut read_level = 0;
577
578        loop {
579            let offset = skip_level + read_level;
580            let remaining_levels = len - offset;
581            if remaining_levels == 0 {
582                break;
583            }
584            let to_read_or_skip_level = rng.random_range(1..=remaining_levels);
585            if rng.random_bool(0.5) {
586                let (skip_val_num, skip_level_num) = decoder.skip(to_read_or_skip_level).unwrap();
587                skip_value += skip_val_num;
588                skip_level += skip_level_num
589            } else {
590                let mut decoded = BooleanBufferBuilder::new(to_read_or_skip_level);
591                let read_level_num = decoder.read(&mut decoded, to_read_or_skip_level).unwrap();
592                read_level += read_level_num;
593                for i in 0..read_level_num {
594                    assert!(!decoded.is_empty());
595                    //check each read bit
596                    let read_bit = decoded.get_bit(i);
597                    if read_bit {
598                        read_value += 1;
599                    }
600                    let expect_bit = expected.get_bit(i + offset);
601                    assert_eq!(read_bit, expect_bit);
602                }
603            }
604        }
605        assert_eq!(read_level + skip_level, len);
606        assert_eq!(read_value + skip_value, total_value);
607    }
608
609    #[test]
610    fn test_try_consume_all_valid() {
611        // Test with all-valid data (all 1s) - single RLE run
612        let len = 100;
613        let mut encoder = RleEncoder::new(1, 1024);
614        for _ in 0..len {
615            encoder.put(1); // all valid
616        }
617        let encoded = encoder.consume();
618        let mut decoder = PackedDecoder::new();
619        decoder.set_data(Encoding::RLE, encoded.into());
620
621        // try_consume_all_valid now parses the RLE block itself, no need to read first
622        let result = decoder.try_consume_all_valid(len).unwrap();
623        assert_eq!(result, Some(len));
624
625        // Test with all-null data (all 0s)
626        let mut encoder = RleEncoder::new(1, 1024);
627        for _ in 0..len {
628            encoder.put(0); // all null
629        }
630        let encoded = encoder.consume();
631        let mut decoder = PackedDecoder::new();
632        decoder.set_data(Encoding::RLE, encoded.into());
633
634        // Should return None because rle_value is false (all nulls)
635        let result = decoder.try_consume_all_valid(len).unwrap();
636        assert_eq!(result, None);
637
638        // Test when requesting more than available in current RLE run
639        let mut encoder = RleEncoder::new(1, 1024);
640        for _ in 0..10 {
641            encoder.put(1); // small run of valid
642        }
643        for _ in 0..10 {
644            encoder.put(0); // followed by nulls
645        }
646        let encoded = encoder.consume();
647        let mut decoder = PackedDecoder::new();
648        decoder.set_data(Encoding::RLE, encoded.into());
649
650        // Request more than the valid run - should return None
651        // (because we don't look ahead to next block)
652        let result = decoder.try_consume_all_valid(20).unwrap();
653        assert_eq!(result, None);
654
655        // Reset decoder and try requesting within the run
656        decoder.set_data(Encoding::RLE, {
657            let mut encoder = RleEncoder::new(1, 1024);
658            for _ in 0..10 {
659                encoder.put(1);
660            }
661            for _ in 0..10 {
662                encoder.put(0);
663            }
664            encoder.consume().into()
665        });
666
667        let result = decoder.try_consume_all_valid(5).unwrap();
668        assert_eq!(result, Some(5));
669
670        // After skipping 5, we should have 5 left in the valid RLE run
671        let result = decoder.try_consume_all_valid(5).unwrap();
672        assert_eq!(result, Some(5));
673
674        // Now the valid run is exhausted, next call should parse the null run and return None
675        let result = decoder.try_consume_all_valid(5).unwrap();
676        assert_eq!(result, None);
677    }
678}