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 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 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 if let Some(count) = decoder.try_consume_all_valid(num_levels)? {
287 nulls.append_n(count, true);
288 return Ok((count, count)); }
290
291 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
310struct 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 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 fn try_consume_all_valid(&mut self, len: usize) -> Result<Option<usize>> {
416 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 return Ok(None);
423 }
424 }
425
426 if self.rle_left >= len && self.rle_value {
429 self.rle_left -= len;
430 return Ok(Some(len));
431 }
432
433 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 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 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 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 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 let len = 100;
624 let mut encoder = RleEncoder::new(1, 1024);
625 for _ in 0..len {
626 encoder.put(1); }
628 let encoded = encoder.consume();
629 let mut decoder = PackedDecoder::new();
630 decoder.set_data(Encoding::RLE, encoded.into());
631
632 let result = decoder.try_consume_all_valid(len).unwrap();
634 assert_eq!(result, Some(len));
635
636 let mut encoder = RleEncoder::new(1, 1024);
638 for _ in 0..len {
639 encoder.put(0); }
641 let encoded = encoder.consume();
642 let mut decoder = PackedDecoder::new();
643 decoder.set_data(Encoding::RLE, encoded.into());
644
645 let result = decoder.try_consume_all_valid(len).unwrap();
647 assert_eq!(result, None);
648
649 let mut encoder = RleEncoder::new(1, 1024);
651 for _ in 0..10 {
652 encoder.put(1); }
654 for _ in 0..10 {
655 encoder.put(0); }
657 let encoded = encoder.consume();
658 let mut decoder = PackedDecoder::new();
659 decoder.set_data(Encoding::RLE, encoded.into());
660
661 let result = decoder.try_consume_all_valid(20).unwrap();
664 assert_eq!(result, None);
665
666 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 let result = decoder.try_consume_all_valid(5).unwrap();
683 assert_eq!(result, Some(5));
684
685 let result = decoder.try_consume_all_valid(5).unwrap();
687 assert_eq!(result, None);
688 }
689}