parquet/arrow/record_reader/
definition_levels.rs1use 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 Full {
34 levels: Vec<i16>,
35 nulls: BooleanBufferBuilder,
36 max_level: i16,
37 },
38 Mask { nulls: BooleanBufferBuilder },
44}
45
46pub struct DefinitionLevelBuffer {
47 inner: BufferInner,
48
49 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 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 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 let buffer = nulls.finish().into_inner();
103
104 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 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
129pub(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 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 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 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 if let Some(count) = decoder.try_consume_all_valid(num_levels)? {
276 nulls.append_n(count, true);
277 return Ok((count, count)); }
279
280 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
299struct 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 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 fn try_consume_all_valid(&mut self, len: usize) -> Result<Option<usize>> {
405 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 return Ok(None);
412 }
413 }
414
415 if self.rle_left >= len && self.rle_value {
418 self.rle_left -= len;
419 return Ok(Some(len));
420 }
421
422 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 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 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 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 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 let len = 100;
613 let mut encoder = RleEncoder::new(1, 1024);
614 for _ in 0..len {
615 encoder.put(1); }
617 let encoded = encoder.consume();
618 let mut decoder = PackedDecoder::new();
619 decoder.set_data(Encoding::RLE, encoded.into());
620
621 let result = decoder.try_consume_all_valid(len).unwrap();
623 assert_eq!(result, Some(len));
624
625 let mut encoder = RleEncoder::new(1, 1024);
627 for _ in 0..len {
628 encoder.put(0); }
630 let encoded = encoder.consume();
631 let mut decoder = PackedDecoder::new();
632 decoder.set_data(Encoding::RLE, encoded.into());
633
634 let result = decoder.try_consume_all_valid(len).unwrap();
636 assert_eq!(result, None);
637
638 let mut encoder = RleEncoder::new(1, 1024);
640 for _ in 0..10 {
641 encoder.put(1); }
643 for _ in 0..10 {
644 encoder.put(0); }
646 let encoded = encoder.consume();
647 let mut decoder = PackedDecoder::new();
648 decoder.set_data(Encoding::RLE, encoded.into());
649
650 let result = decoder.try_consume_all_valid(20).unwrap();
653 assert_eq!(result, None);
654
655 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 let result = decoder.try_consume_all_valid(5).unwrap();
672 assert_eq!(result, Some(5));
673
674 let result = decoder.try_consume_all_valid(5).unwrap();
676 assert_eq!(result, None);
677 }
678}