1use crate::basic::Type as PhysicalType;
19use crate::column::reader::{ColumnReader, ColumnReaderImpl, get_typed_column_reader};
20use crate::data_type::*;
21use crate::errors::{ParquetError, Result};
22use crate::record::api::Field;
23use crate::schema::types::ColumnDescPtr;
24
25macro_rules! triplet_enum_func {
29 ($self:ident, $func:ident, $( $token:tt ),*) => ({
30 match *$self {
31 TripletIter::BoolTripletIter($($token)* typed) => typed.$func(),
32 TripletIter::Int32TripletIter($($token)* typed) => typed.$func(),
33 TripletIter::Int64TripletIter($($token)* typed) => typed.$func(),
34 TripletIter::Int96TripletIter($($token)* typed) => typed.$func(),
35 TripletIter::FloatTripletIter($($token)* typed) => typed.$func(),
36 TripletIter::DoubleTripletIter($($token)* typed) => typed.$func(),
37 TripletIter::ByteArrayTripletIter($($token)* typed) => typed.$func(),
38 TripletIter::FixedLenByteArrayTripletIter($($token)* typed) => typed.$func()
39 }
40 });
41}
42
43#[expect(clippy::enum_variant_names)]
46pub enum TripletIter {
47 BoolTripletIter(TypedTripletIter<BoolType>),
48 Int32TripletIter(TypedTripletIter<Int32Type>),
49 Int64TripletIter(TypedTripletIter<Int64Type>),
50 Int96TripletIter(TypedTripletIter<Int96Type>),
51 FloatTripletIter(TypedTripletIter<FloatType>),
52 DoubleTripletIter(TypedTripletIter<DoubleType>),
53 ByteArrayTripletIter(TypedTripletIter<ByteArrayType>),
54 FixedLenByteArrayTripletIter(TypedTripletIter<FixedLenByteArrayType>),
55}
56
57impl TripletIter {
58 pub fn new(descr: ColumnDescPtr, reader: ColumnReader, batch_size: usize) -> Self {
60 match descr.physical_type() {
61 PhysicalType::BOOLEAN => {
62 TripletIter::BoolTripletIter(TypedTripletIter::new(descr, batch_size, reader))
63 }
64 PhysicalType::INT32 => {
65 TripletIter::Int32TripletIter(TypedTripletIter::new(descr, batch_size, reader))
66 }
67 PhysicalType::INT64 => {
68 TripletIter::Int64TripletIter(TypedTripletIter::new(descr, batch_size, reader))
69 }
70 PhysicalType::INT96 => {
71 TripletIter::Int96TripletIter(TypedTripletIter::new(descr, batch_size, reader))
72 }
73 PhysicalType::FLOAT => {
74 TripletIter::FloatTripletIter(TypedTripletIter::new(descr, batch_size, reader))
75 }
76 PhysicalType::DOUBLE => {
77 TripletIter::DoubleTripletIter(TypedTripletIter::new(descr, batch_size, reader))
78 }
79 PhysicalType::BYTE_ARRAY => {
80 TripletIter::ByteArrayTripletIter(TypedTripletIter::new(descr, batch_size, reader))
81 }
82 PhysicalType::FIXED_LEN_BYTE_ARRAY => TripletIter::FixedLenByteArrayTripletIter(
83 TypedTripletIter::new(descr, batch_size, reader),
84 ),
85 }
86 }
87
88 #[inline]
91 pub fn read_next(&mut self) -> Result<bool> {
92 triplet_enum_func!(self, read_next, ref, mut)
93 }
94
95 #[inline]
100 pub fn has_next(&self) -> bool {
101 triplet_enum_func!(self, has_next, ref)
102 }
103
104 #[inline]
106 pub fn current_def_level(&self) -> i16 {
107 triplet_enum_func!(self, current_def_level, ref)
108 }
109
110 #[inline]
112 pub fn max_def_level(&self) -> i16 {
113 triplet_enum_func!(self, max_def_level, ref)
114 }
115
116 #[inline]
118 pub fn current_rep_level(&self) -> i16 {
119 triplet_enum_func!(self, current_rep_level, ref)
120 }
121
122 #[inline]
124 pub fn max_rep_level(&self) -> i16 {
125 triplet_enum_func!(self, max_rep_level, ref)
126 }
127
128 #[inline]
132 pub fn is_null(&self) -> bool {
133 self.current_def_level() < self.max_def_level()
134 }
135
136 pub fn current_value(&self) -> Result<Field> {
138 if self.is_null() {
139 return Ok(Field::Null);
140 }
141 let field = match *self {
142 TripletIter::BoolTripletIter(ref typed) => {
143 Field::convert_bool(typed.column_descr(), *typed.current_value())
144 }
145 TripletIter::Int32TripletIter(ref typed) => {
146 Field::convert_int32(typed.column_descr(), *typed.current_value())
147 }
148 TripletIter::Int64TripletIter(ref typed) => {
149 Field::convert_int64(typed.column_descr(), *typed.current_value())
150 }
151 TripletIter::Int96TripletIter(ref typed) => {
152 Field::convert_int96(typed.column_descr(), *typed.current_value())
153 }
154 TripletIter::FloatTripletIter(ref typed) => {
155 Field::convert_float(typed.column_descr(), *typed.current_value())
156 }
157 TripletIter::DoubleTripletIter(ref typed) => {
158 Field::convert_double(typed.column_descr(), *typed.current_value())
159 }
160 TripletIter::ByteArrayTripletIter(ref typed) => {
161 Field::convert_byte_array(typed.column_descr(), typed.current_value().clone())?
162 }
163 TripletIter::FixedLenByteArrayTripletIter(ref typed) => Field::convert_byte_array(
164 typed.column_descr(),
165 typed.current_value().clone().into(),
166 )?,
167 };
168 Ok(field)
169 }
170}
171
172pub struct TypedTripletIter<T: DataType> {
175 reader: ColumnReaderImpl<T>,
176 column_descr: ColumnDescPtr,
177 batch_size: usize,
178 max_def_level: i16,
180 max_rep_level: i16,
181 values: Vec<T::T>,
183 def_levels: Option<Vec<i16>>,
184 rep_levels: Option<Vec<i16>>,
185 curr_triplet_index: usize,
187 triplets_left: usize,
189 has_next: bool,
191}
192
193impl<T: DataType> TypedTripletIter<T> {
194 fn new(descr: ColumnDescPtr, batch_size: usize, column_reader: ColumnReader) -> Self {
197 assert!(
198 batch_size > 0,
199 "Expected positive batch size, found: {batch_size}"
200 );
201
202 let max_def_level = descr.max_def_level();
203 let max_rep_level = descr.max_rep_level();
204
205 let def_levels = if max_def_level == 0 {
206 None
207 } else {
208 Some(vec![0; batch_size])
209 };
210 let rep_levels = if max_rep_level == 0 {
211 None
212 } else {
213 Some(vec![0; batch_size])
214 };
215
216 Self {
217 reader: get_typed_column_reader(column_reader),
218 column_descr: descr,
219 batch_size,
220 max_def_level,
221 max_rep_level,
222 values: vec![T::T::default(); batch_size],
223 def_levels,
224 rep_levels,
225 curr_triplet_index: 0,
226 triplets_left: 0,
227 has_next: false,
228 }
229 }
230
231 #[inline]
233 pub fn column_descr(&self) -> &ColumnDescPtr {
234 &self.column_descr
235 }
236
237 #[inline]
239 fn max_def_level(&self) -> i16 {
240 self.max_def_level
241 }
242
243 #[inline]
245 fn max_rep_level(&self) -> i16 {
246 self.max_rep_level
247 }
248
249 #[inline]
252 fn current_value(&self) -> &T::T {
253 assert!(
254 self.current_def_level() == self.max_def_level(),
255 "Cannot extract value, max definition level: {}, current level: {}",
256 self.max_def_level(),
257 self.current_def_level()
258 );
259 &self.values[self.curr_triplet_index]
260 }
261
262 #[inline]
265 fn current_def_level(&self) -> i16 {
266 if !self.has_next {
267 return 0;
268 }
269 match self.def_levels {
270 Some(ref vec) => vec[self.curr_triplet_index],
271 None => self.max_def_level,
272 }
273 }
274
275 #[inline]
278 fn current_rep_level(&self) -> i16 {
279 if !self.has_next {
280 return 0;
281 }
282 match self.rep_levels {
283 Some(ref vec) => vec[self.curr_triplet_index],
284 None => self.max_rep_level,
285 }
286 }
287
288 #[inline]
291 fn has_next(&self) -> bool {
292 self.has_next
293 }
294
295 fn read_next(&mut self) -> Result<bool> {
298 self.curr_triplet_index += 1;
299
300 while self.curr_triplet_index >= self.triplets_left {
304 let (records_read, values_read, levels_read) = {
305 self.values.clear();
306 if let Some(x) = &mut self.def_levels {
307 x.clear()
308 }
309 if let Some(x) = &mut self.rep_levels {
310 x.clear()
311 }
312
313 self.reader.read_records(
315 self.batch_size,
316 self.def_levels.as_mut(),
317 self.rep_levels.as_mut(),
318 &mut self.values,
319 )?
320 };
321
322 if records_read == 0 && values_read == 0 && levels_read == 0 {
324 self.curr_triplet_index = 0;
325 self.has_next = false;
326 return Ok(false);
327 }
328
329 if levels_read == 0 || values_read == levels_read {
331 self.curr_triplet_index = 0;
334 self.triplets_left = values_read;
335 } else if values_read < levels_read {
336 let mut idx = values_read;
344 let def_levels = self.def_levels.as_ref().unwrap();
345 self.values.resize(levels_read, T::T::default());
346 for i in 0..levels_read {
347 if def_levels[levels_read - i - 1] == self.max_def_level {
348 idx -= 1; self.values.swap(levels_read - i - 1, idx);
350 }
351 }
352 self.curr_triplet_index = 0;
353 self.triplets_left = levels_read;
354 } else {
355 return Err(general_err!(
356 "Spacing of values/levels is wrong, values_read: {}, levels_read: {}",
357 values_read,
358 levels_read
359 ));
360 }
361 }
362
363 self.has_next = true;
364 Ok(true)
365 }
366}
367
368#[cfg(test)]
369mod tests {
370 use super::*;
371
372 use crate::file::reader::{FileReader, SerializedFileReader};
373 use crate::schema::types::ColumnPath;
374 use crate::util::test_common::file_util::get_test_file;
375
376 #[test]
377 #[should_panic(expected = "Expected positive batch size, found: 0")]
378 fn test_triplet_zero_batch_size() {
379 let column_path = ColumnPath::from(vec!["b_struct".to_string(), "b_c_int".to_string()]);
380 test_column_in_file("nulls.snappy.parquet", 0, &column_path, &[], &[], &[]);
381 }
382
383 #[test]
384 fn test_triplet_null_column() {
385 let path = vec!["b_struct", "b_c_int"];
386 let values = vec![];
387 let def_levels = vec![1, 1, 1, 1, 1, 1, 1, 1];
388 let rep_levels = vec![0, 0, 0, 0, 0, 0, 0, 0];
389 test_triplet_iter(
390 "nulls.snappy.parquet",
391 path,
392 &values,
393 &def_levels,
394 &rep_levels,
395 );
396 }
397
398 #[test]
399 #[cfg_attr(miri, ignore)] fn test_triplet_required_column() {
401 let path = vec!["ID"];
402 let values = vec![Field::Long(8)];
403 let def_levels = vec![0];
404 let rep_levels = vec![0];
405 test_triplet_iter(
406 "nonnullable.impala.parquet",
407 path,
408 &values,
409 &def_levels,
410 &rep_levels,
411 );
412 }
413
414 #[test]
415 #[cfg_attr(miri, ignore)] fn test_triplet_optional_column() {
417 let path = vec!["nested_struct", "A"];
418 let values = vec![Field::Int(1), Field::Int(7)];
419 let def_levels = vec![2, 1, 1, 1, 1, 0, 2];
420 let rep_levels = vec![0, 0, 0, 0, 0, 0, 0];
421 test_triplet_iter(
422 "nullable.impala.parquet",
423 path,
424 &values,
425 &def_levels,
426 &rep_levels,
427 );
428 }
429
430 #[test]
431 #[cfg_attr(miri, ignore)] fn test_triplet_optional_list_column() {
433 let path = vec!["a", "list", "element", "list", "element", "list", "element"];
434 let values = vec![
435 Field::Str("a".to_string()),
436 Field::Str("b".to_string()),
437 Field::Str("c".to_string()),
438 Field::Str("d".to_string()),
439 Field::Str("a".to_string()),
440 Field::Str("b".to_string()),
441 Field::Str("c".to_string()),
442 Field::Str("d".to_string()),
443 Field::Str("e".to_string()),
444 Field::Str("a".to_string()),
445 Field::Str("b".to_string()),
446 Field::Str("c".to_string()),
447 Field::Str("d".to_string()),
448 Field::Str("e".to_string()),
449 Field::Str("f".to_string()),
450 ];
451 let def_levels = vec![7, 7, 7, 4, 7, 7, 7, 7, 7, 4, 7, 7, 7, 7, 7, 7, 4, 7];
452 let rep_levels = vec![0, 3, 2, 1, 2, 0, 3, 2, 3, 1, 2, 0, 3, 2, 3, 2, 1, 2];
453 test_triplet_iter(
454 "nested_lists.snappy.parquet",
455 path,
456 &values,
457 &def_levels,
458 &rep_levels,
459 );
460 }
461
462 #[test]
463 #[cfg_attr(miri, ignore)] fn test_triplet_optional_map_column() {
465 let path = vec!["a", "key_value", "value", "key_value", "key"];
466 let values = vec![
467 Field::Int(1),
468 Field::Int(2),
469 Field::Int(1),
470 Field::Int(1),
471 Field::Int(3),
472 Field::Int(4),
473 Field::Int(5),
474 ];
475 let def_levels = vec![4, 4, 4, 2, 3, 4, 4, 4, 4];
476 let rep_levels = vec![0, 2, 0, 0, 0, 0, 0, 2, 2];
477 test_triplet_iter(
478 "nested_maps.snappy.parquet",
479 path,
480 &values,
481 &def_levels,
482 &rep_levels,
483 );
484 }
485
486 fn test_triplet_iter(
488 file_name: &str,
489 column_path: Vec<&str>,
490 expected_values: &[Field],
491 expected_def_levels: &[i16],
492 expected_rep_levels: &[i16],
493 ) {
494 let path: Vec<String> = column_path.iter().map(|x| x.to_string()).collect();
496 let column_path = ColumnPath::from(path);
497
498 let batch_sizes = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 128, 256];
499 for batch_size in batch_sizes {
500 test_column_in_file(
501 file_name,
502 batch_size,
503 &column_path,
504 expected_values,
505 expected_def_levels,
506 expected_rep_levels,
507 );
508 }
509 }
510
511 fn test_column_in_file(
513 file_name: &str,
514 batch_size: usize,
515 column_path: &ColumnPath,
516 expected_values: &[Field],
517 expected_def_levels: &[i16],
518 expected_rep_levels: &[i16],
519 ) {
520 let file = get_test_file(file_name);
521 let file_reader = SerializedFileReader::new(file).unwrap();
522 let metadata = file_reader.metadata();
523 let file_metadata = metadata.file_metadata();
525 let schema = file_metadata.schema_descr();
526 let row_group_reader = file_reader.get_row_group(0).unwrap();
528
529 for i in 0..schema.num_columns() {
530 let descr = schema.column(i);
531 if descr.path() == column_path {
532 let reader = row_group_reader.get_column_reader(i).unwrap();
533 test_triplet_column(
534 descr,
535 reader,
536 batch_size,
537 expected_values,
538 expected_def_levels,
539 expected_rep_levels,
540 );
541 }
542 }
543 }
544
545 fn test_triplet_column(
547 descr: ColumnDescPtr,
548 reader: ColumnReader,
549 batch_size: usize,
550 expected_values: &[Field],
551 expected_def_levels: &[i16],
552 expected_rep_levels: &[i16],
553 ) {
554 let mut iter = TripletIter::new(descr.clone(), reader, batch_size);
555 let mut values: Vec<Field> = Vec::new();
556 let mut def_levels: Vec<i16> = Vec::new();
557 let mut rep_levels: Vec<i16> = Vec::new();
558
559 assert_eq!(iter.max_def_level(), descr.max_def_level());
560 assert_eq!(iter.max_rep_level(), descr.max_rep_level());
561
562 while matches!(iter.read_next(), Ok(true)) {
563 assert!(iter.has_next());
564 if !iter.is_null() {
565 values.push(iter.current_value().unwrap());
566 }
567 def_levels.push(iter.current_def_level());
568 rep_levels.push(iter.current_rep_level());
569 }
570
571 assert_eq!(values, expected_values);
572 assert_eq!(def_levels, expected_def_levels);
573 assert_eq!(rep_levels, expected_rep_levels);
574 }
575
576 fn open_triplet_iter(file_name: &str, path: &[&str], batch_size: usize) -> TripletIter {
577 let column_path = ColumnPath::from(path.iter().map(|x| x.to_string()).collect::<Vec<_>>());
578 let file = get_test_file(file_name);
579 let file_reader = SerializedFileReader::new(file).unwrap();
580 let metadata = file_reader.metadata();
581 let schema = metadata.file_metadata().schema_descr();
582 let row_group_reader = file_reader.get_row_group(0).unwrap();
583 for i in 0..schema.num_columns() {
584 let descr = schema.column(i);
585 if descr.path() == &column_path {
586 let reader = row_group_reader.get_column_reader(i).unwrap();
587 return TripletIter::new(descr.clone(), reader, batch_size);
588 }
589 }
590 panic!("Column {column_path:?} not found in {file_name}");
591 }
592
593 #[test]
594 fn test_current_def_level_safe_after_exhaustion() {
595 let mut iter = open_triplet_iter("nulls.snappy.parquet", &["b_struct", "b_c_int"], 256);
596 while matches!(iter.read_next(), Ok(true)) {}
597 assert!(!iter.has_next());
598 assert_eq!(iter.current_def_level(), 0);
599 }
600
601 #[test]
602 fn test_current_rep_level_safe_after_exhaustion() {
603 let mut iter = open_triplet_iter(
604 "nested_lists.snappy.parquet",
605 &["a", "list", "element", "list", "element", "list", "element"],
606 256,
607 );
608 while matches!(iter.read_next(), Ok(true)) {}
609 assert!(!iter.has_next());
610 assert_eq!(iter.current_rep_level(), 0);
611 }
612}