1use crate::column::chunker::CdcChunk;
44use crate::column::writer::LevelDataRef;
45use crate::errors::{ParquetError, Result};
46use arrow_array::cast::AsArray;
47use arrow_array::types::RunEndIndexType;
48use arrow_array::{Array, ArrayRef, Int32Array, OffsetSizeTrait, RunArray, downcast_run_array};
49use arrow_buffer::bit_iterator::BitIndexIterator;
50use arrow_buffer::{NullBuffer, OffsetBuffer, ScalarBuffer};
51use arrow_schema::{DataType, Field};
52use std::ops::Range;
53use std::sync::Arc;
54
55fn expand_ree_array(array: &ArrayRef) -> Result<ArrayRef> {
60 downcast_run_array!(
61 array => expand_typed_ree(array),
62 _ => unreachable!("expand_ree_array called on non-REE array"),
63 )
64}
65
66fn expand_typed_ree<R: RunEndIndexType>(run_array: &RunArray<R>) -> Result<ArrayRef> {
67 let run_ends = run_array.run_ends();
68 let values = run_array.values();
69 let len = run_array.len();
70 let indices: Int32Array = (0..len)
71 .map(|i| run_ends.get_physical_index(i) as i32)
72 .collect();
73 arrow_select::take::take(values.as_ref(), &indices, None)
74 .map_err(|e| arrow_err!("Failed to expand REE array: {}", e))
75}
76
77pub(crate) fn calculate_array_levels(array: &ArrayRef, field: &Field) -> Result<Vec<ArrayLevels>> {
80 let mut builder = LevelInfoBuilder::try_new(field, Default::default(), array)?;
81 builder.write(0..array.len());
82 Ok(builder.finish())
83}
84
85fn is_leaf(data_type: &DataType) -> bool {
88 matches!(
89 data_type,
90 DataType::Null
91 | DataType::Boolean
92 | DataType::Int8
93 | DataType::Int16
94 | DataType::Int32
95 | DataType::Int64
96 | DataType::UInt8
97 | DataType::UInt16
98 | DataType::UInt32
99 | DataType::UInt64
100 | DataType::Float16
101 | DataType::Float32
102 | DataType::Float64
103 | DataType::Utf8
104 | DataType::Utf8View
105 | DataType::LargeUtf8
106 | DataType::Timestamp(_, _)
107 | DataType::Date32
108 | DataType::Date64
109 | DataType::Time32(_)
110 | DataType::Time64(_)
111 | DataType::Duration(_)
112 | DataType::Interval(_)
113 | DataType::Binary
114 | DataType::LargeBinary
115 | DataType::BinaryView
116 | DataType::Decimal32(_, _)
117 | DataType::Decimal64(_, _)
118 | DataType::Decimal128(_, _)
119 | DataType::Decimal256(_, _)
120 | DataType::FixedSizeBinary(_)
121 )
122}
123
124#[derive(Debug, Default, Clone, Copy)]
126struct LevelContext {
127 rep_level: i16,
129 def_level: i16,
131}
132
133#[derive(Debug)]
135enum LevelInfoBuilder {
136 Primitive(ArrayLevels),
138 List(
140 Box<LevelInfoBuilder>, LevelContext, OffsetBuffer<i32>, Option<NullBuffer>, bool, ),
146 LargeList(
148 Box<LevelInfoBuilder>, LevelContext, OffsetBuffer<i64>, Option<NullBuffer>, bool, ),
154 FixedSizeList(
156 Box<LevelInfoBuilder>, LevelContext, usize, Option<NullBuffer>, ),
161 ListView(
163 Box<LevelInfoBuilder>, LevelContext, ScalarBuffer<i32>, ScalarBuffer<i32>, Option<NullBuffer>, ),
169 LargeListView(
171 Box<LevelInfoBuilder>, LevelContext, ScalarBuffer<i64>, ScalarBuffer<i64>, Option<NullBuffer>, ),
177 Struct(Vec<LevelInfoBuilder>, LevelContext, Option<NullBuffer>),
179}
180
181const BULK_FILL_MIN_LEN: usize = 64;
187
188impl LevelInfoBuilder {
189 fn try_new(field: &Field, parent_ctx: LevelContext, array: &ArrayRef) -> Result<Self> {
191 if !Self::types_compatible(field.data_type(), array.data_type()) {
192 return Err(arrow_err!(format!(
193 "Incompatible type. Field '{}' has type {}, array has type {}",
194 field.name(),
195 field.data_type(),
196 array.data_type(),
197 )));
198 }
199
200 let is_nullable = field.is_nullable();
201
202 match array.data_type() {
203 d if is_leaf(d) => {
204 let levels = ArrayLevels::new(parent_ctx, is_nullable, array.clone());
205 Ok(Self::Primitive(levels))
206 }
207 DataType::Dictionary(_, v) if is_leaf(v.as_ref()) => {
208 let levels = ArrayLevels::new(parent_ctx, is_nullable, array.clone());
209 Ok(Self::Primitive(levels))
210 }
211 DataType::RunEndEncoded(_, value_field) => {
212 let flat = expand_ree_array(array)?;
213 let flat_field = Field::new(
214 field.name(),
215 value_field.data_type().clone(),
216 field.is_nullable(),
217 );
218 Self::try_new(&flat_field, parent_ctx, &flat)
219 }
220 DataType::Struct(children) => {
221 let array = array.as_struct();
222 let def_level = match is_nullable {
223 true => parent_ctx.def_level + 1,
224 false => parent_ctx.def_level,
225 };
226
227 let ctx = LevelContext {
228 rep_level: parent_ctx.rep_level,
229 def_level,
230 };
231
232 let children = children
233 .iter()
234 .zip(array.columns())
235 .map(|(f, a)| Self::try_new(f, ctx, a))
236 .collect::<Result<_>>()?;
237
238 Ok(Self::Struct(children, ctx, array.nulls().cloned()))
239 }
240 DataType::List(child)
241 | DataType::LargeList(child)
242 | DataType::Map(child, _)
243 | DataType::FixedSizeList(child, _)
244 | DataType::ListView(child)
245 | DataType::LargeListView(child) => {
246 let def_level = match is_nullable {
247 true => parent_ctx.def_level + 2,
248 false => parent_ctx.def_level + 1,
249 };
250
251 let ctx = LevelContext {
252 rep_level: parent_ctx.rep_level + 1,
253 def_level,
254 };
255
256 Ok(match field.data_type() {
257 DataType::List(_) => {
258 let list = array.as_list();
259 let child = Self::try_new(child.as_ref(), ctx, list.values())?;
260 let is_last = child.child_has_no_nested_rep();
261 let offsets = list.offsets().clone();
262 Self::List(
263 Box::new(child),
264 ctx,
265 offsets,
266 list.nulls().cloned(),
267 is_last,
268 )
269 }
270 DataType::LargeList(_) => {
271 let list = array.as_list();
272 let child = Self::try_new(child.as_ref(), ctx, list.values())?;
273 let is_last = child.child_has_no_nested_rep();
274 let offsets = list.offsets().clone();
275 let nulls = list.nulls().cloned();
276 Self::LargeList(Box::new(child), ctx, offsets, nulls, is_last)
277 }
278 DataType::Map(_, _) => {
279 let map = array.as_map();
280 let entries = Arc::new(map.entries().clone()) as ArrayRef;
281 let child = Self::try_new(child.as_ref(), ctx, &entries)?;
282 let is_last = child.child_has_no_nested_rep();
283 let offsets = map.offsets().clone();
284 Self::List(Box::new(child), ctx, offsets, map.nulls().cloned(), is_last)
285 }
286 DataType::FixedSizeList(_, size) => {
287 let list = array.as_fixed_size_list();
288 let child = Self::try_new(child.as_ref(), ctx, list.values())?;
289 let nulls = list.nulls().cloned();
290 Self::FixedSizeList(Box::new(child), ctx, *size as _, nulls)
291 }
292 DataType::ListView(_) => {
293 let list = array.as_list_view();
294 let child = Self::try_new(child.as_ref(), ctx, list.values())?;
295 let offsets = list.offsets().clone();
296 let sizes = list.sizes().clone();
297 let nulls = list.nulls().cloned();
298 Self::ListView(Box::new(child), ctx, offsets, sizes, nulls)
299 }
300 DataType::LargeListView(_) => {
301 let list = array.as_list_view();
302 let child = Self::try_new(child.as_ref(), ctx, list.values())?;
303 let offsets = list.offsets().clone();
304 let sizes = list.sizes().clone();
305 let nulls = list.nulls().cloned();
306 Self::LargeListView(Box::new(child), ctx, offsets, sizes, nulls)
307 }
308 _ => unreachable!(),
309 })
310 }
311 d => Err(nyi_err!("Datatype {} is not yet supported", d)),
312 }
313 }
314
315 fn finish(self) -> Vec<ArrayLevels> {
318 match self {
319 LevelInfoBuilder::Primitive(v) => vec![v],
320 LevelInfoBuilder::List(v, _, _, _, _)
321 | LevelInfoBuilder::LargeList(v, _, _, _, _)
322 | LevelInfoBuilder::FixedSizeList(v, _, _, _)
323 | LevelInfoBuilder::ListView(v, _, _, _, _)
324 | LevelInfoBuilder::LargeListView(v, _, _, _, _) => v.finish(),
325 LevelInfoBuilder::Struct(v, _, _) => v.into_iter().flat_map(|l| l.finish()).collect(),
326 }
327 }
328
329 fn write(&mut self, range: Range<usize>) {
331 match self {
332 LevelInfoBuilder::Primitive(info) => Self::write_leaf(info, range),
333 LevelInfoBuilder::List(child, ctx, offsets, nulls, is_last) => {
334 Self::write_list(child, ctx, offsets, nulls.as_ref(), range, *is_last)
335 }
336 LevelInfoBuilder::LargeList(child, ctx, offsets, nulls, is_last) => {
337 Self::write_list(child, ctx, offsets, nulls.as_ref(), range, *is_last)
338 }
339 LevelInfoBuilder::FixedSizeList(child, ctx, size, nulls) => {
340 Self::write_fixed_size_list(child, ctx, *size, nulls.as_ref(), range)
341 }
342 LevelInfoBuilder::ListView(child, ctx, offsets, sizes, nulls) => {
343 Self::write_list_view(child, ctx, offsets, sizes, nulls.as_ref(), range)
344 }
345 LevelInfoBuilder::LargeListView(child, ctx, offsets, sizes, nulls) => {
346 Self::write_list_view(child, ctx, offsets, sizes, nulls.as_ref(), range)
347 }
348 LevelInfoBuilder::Struct(children, ctx, nulls) => {
349 Self::write_struct(children, ctx, nulls.as_ref(), range)
350 }
351 }
352 }
353
354 fn child_has_no_nested_rep(&self) -> bool {
358 match self {
359 LevelInfoBuilder::Primitive(_) => true,
360 LevelInfoBuilder::Struct(children, _, _) => {
361 children.iter().all(|c| c.child_has_no_nested_rep())
362 }
363 _ => false,
364 }
365 }
366
367 fn write_list<O: OffsetSizeTrait>(
371 child: &mut LevelInfoBuilder,
372 ctx: &LevelContext,
373 offsets: &[O],
374 nulls: Option<&NullBuffer>,
375 range: Range<usize>,
376 is_last_level: bool,
377 ) {
378 if nulls.is_some_and(|nulls| nulls.null_count() == nulls.len()) {
380 let count = range.end - range.start;
381 child.visit_leaves(|leaf| {
382 leaf.extend_uniform_levels(ctx.def_level - 2, ctx.rep_level - 1, count);
383 });
384 return;
385 }
386
387 if is_last_level {
390 Self::write_list_direct(child, ctx, offsets, nulls, range);
391 } else {
392 Self::write_list_scan(child, ctx, offsets, nulls, range);
393 }
394 }
395
396 fn write_list_direct<O: OffsetSizeTrait>(
400 child: &mut LevelInfoBuilder,
401 ctx: &LevelContext,
402 offsets: &[O],
403 nulls: Option<&NullBuffer>,
404 range: Range<usize>,
405 ) {
406 let list_start_rep = ctx.rep_level - 1;
407
408 let emit_non_empty_run = |child: &mut LevelInfoBuilder, run_offsets: &[O]| {
409 debug_assert!(run_offsets.len() >= 2);
410 let values_start = run_offsets[0].as_usize();
411 let values_end = run_offsets[run_offsets.len() - 1].as_usize();
412 debug_assert!(values_end > values_start);
413
414 child.write(values_start..values_end);
415
416 child.visit_leaves(|leaf| {
421 debug_assert!(leaf.max_rep_level == ctx.rep_level);
422 let rep_levels = leaf.rep_levels.materialize_mut().unwrap();
423 let batch_len = values_end - values_start;
424 let batch_base = rep_levels.len() - batch_len;
425 for slot_offset in run_offsets.iter().take(run_offsets.len() - 1) {
426 let pos = batch_base + (slot_offset.as_usize() - values_start);
427 rep_levels[pos] = list_start_rep;
428 }
429 });
430 };
431
432 Self::write_list_impl(child, ctx, offsets, nulls, range, emit_non_empty_run);
433 }
434
435 fn write_list_scan<O: OffsetSizeTrait>(
442 child: &mut LevelInfoBuilder,
443 ctx: &LevelContext,
444 offsets: &[O],
445 nulls: Option<&NullBuffer>,
446 range: Range<usize>,
447 ) {
448 let list_start_rep = ctx.rep_level - 1;
449
450 let emit_non_empty_run = |child: &mut LevelInfoBuilder, run_offsets: &[O]| {
451 debug_assert!(run_offsets.len() >= 2);
452 let values_start = run_offsets[0].as_usize();
453 let values_end = run_offsets[run_offsets.len() - 1].as_usize();
454 debug_assert!(values_end > values_start);
455
456 child.write(values_start..values_end);
457
458 child.visit_leaves(|leaf| {
459 let rep_levels = leaf.rep_levels.materialize_mut().unwrap();
460
461 if leaf.max_rep_level == ctx.rep_level {
462 let batch_len = values_end - values_start;
466 let batch_base = rep_levels.len() - batch_len;
467 for slot_offset in run_offsets.iter().take(run_offsets.len() - 1) {
468 let pos = batch_base + (slot_offset.as_usize() - values_start);
469 rep_levels[pos] = list_start_rep;
470 }
471 } else {
472 let mut slot_bounds = run_offsets[..run_offsets.len() - 1].iter().rev();
475 let mut next_stamp_at = values_end - slot_bounds.next().unwrap().as_usize();
476 let mut seen = 0usize;
477
478 for rep in rep_levels.iter_mut().rev() {
479 debug_assert!(*rep >= ctx.rep_level);
482 if *rep <= ctx.rep_level {
488 seen += 1;
489 if seen == next_stamp_at {
490 *rep = list_start_rep;
491 match slot_bounds.next() {
492 Some(offset) => next_stamp_at = values_end - offset.as_usize(),
493 None => break,
494 }
495 }
496 }
497 }
498 }
499 });
500 };
501
502 Self::write_list_impl(child, ctx, offsets, nulls, range, emit_non_empty_run);
503 }
504
505 fn write_list_impl<O: OffsetSizeTrait>(
509 child: &mut LevelInfoBuilder,
510 ctx: &LevelContext,
511 offsets: &[O],
512 nulls: Option<&NullBuffer>,
513 range: Range<usize>,
514 mut emit_non_empty_run: impl FnMut(&mut LevelInfoBuilder, &[O]),
515 ) {
516 let null_offset = range.start;
517 let offsets = &offsets[range.start..range.end + 1];
518 let list_start_rep = ctx.rep_level - 1;
519
520 let emit_nulls = |child: &mut LevelInfoBuilder, count: usize| {
521 child.visit_leaves(|leaf| {
522 leaf.append_rep_level_run(list_start_rep, count);
523 leaf.append_def_level_run(ctx.def_level - 2, count);
524 });
525 };
526
527 let emit_empties = |child: &mut LevelInfoBuilder, count: usize| {
528 child.visit_leaves(|leaf| {
529 leaf.append_rep_level_run(list_start_rep, count);
530 leaf.append_def_level_run(ctx.def_level - 1, count);
531 });
532 };
533
534 #[derive(Clone, Copy, PartialEq)]
536 enum SlotKind {
537 Null,
538 Empty,
539 NonEmpty,
540 }
541
542 let num_slots = offsets.len() - 1;
543 if num_slots == 0 {
544 return;
545 }
546
547 macro_rules! classify {
548 ($i:expr, $nulls:expr) => {
549 if !$nulls.is_valid($i + null_offset) {
550 SlotKind::Null
551 } else if offsets[$i] == offsets[$i + 1] {
552 SlotKind::Empty
553 } else {
554 SlotKind::NonEmpty
555 }
556 };
557 }
558
559 macro_rules! flush_run {
560 ($kind:expr, $start:expr, $end:expr) => {
561 match $kind {
562 SlotKind::Null => emit_nulls(child, $end - $start),
563 SlotKind::Empty => emit_empties(child, $end - $start),
564 SlotKind::NonEmpty => emit_non_empty_run(child, &offsets[$start..$end + 1]),
565 }
566 };
567 }
568
569 match nulls {
570 Some(nulls) if nulls.null_count() > 0 => {
573 let mut run_kind = classify!(0, nulls);
574 let mut run_start: usize = 0;
575 for i in 1..num_slots {
576 let kind = classify!(i, nulls);
577 if kind != run_kind {
578 flush_run!(run_kind, run_start, i);
579 run_kind = kind;
580 run_start = i;
581 }
582 }
583 flush_run!(run_kind, run_start, num_slots);
584 }
585 _ => {
586 let mut run_kind = if offsets[0] == offsets[1] {
587 SlotKind::Empty
588 } else {
589 SlotKind::NonEmpty
590 };
591 let mut run_start: usize = 0;
592 for i in 1..num_slots {
593 let kind = if offsets[i] == offsets[i + 1] {
594 SlotKind::Empty
595 } else {
596 SlotKind::NonEmpty
597 };
598 if kind != run_kind {
599 flush_run!(run_kind, run_start, i);
600 run_kind = kind;
601 run_start = i;
602 }
603 }
604 flush_run!(run_kind, run_start, num_slots);
605 }
606 }
607 }
608
609 fn write_list_view<O: OffsetSizeTrait>(
611 child: &mut LevelInfoBuilder,
612 ctx: &LevelContext,
613 offsets: &[O],
614 sizes: &[O],
615 nulls: Option<&NullBuffer>,
616 range: Range<usize>,
617 ) {
618 let offsets = &offsets[range.start..range.end];
619 let sizes = &sizes[range.start..range.end];
620
621 let write_non_null_slice =
622 |child: &mut LevelInfoBuilder, start_idx: usize, end_idx: usize| {
623 child.write(start_idx..end_idx);
624 child.visit_leaves(|leaf| {
625 let rep_levels = leaf.rep_levels.materialize_mut().unwrap();
626 let mut rev = rep_levels.iter_mut().rev();
627 let mut remaining = end_idx - start_idx;
628
629 loop {
630 let next = rev.next().unwrap();
631 if *next > ctx.rep_level {
632 continue;
634 }
635
636 remaining -= 1;
637 if remaining == 0 {
638 *next = ctx.rep_level - 1;
639 break;
640 }
641 }
642 })
643 };
644
645 let write_empty_slice = |child: &mut LevelInfoBuilder| {
646 child.visit_leaves(|leaf| {
647 leaf.append_rep_level_run(ctx.rep_level - 1, 1);
648 leaf.append_def_level_run(ctx.def_level - 1, 1);
649 })
650 };
651
652 let write_null_slice = |child: &mut LevelInfoBuilder| {
653 child.visit_leaves(|leaf| {
654 leaf.append_rep_level_run(ctx.rep_level - 1, 1);
655 leaf.append_def_level_run(ctx.def_level - 2, 1);
656 })
657 };
658
659 match nulls {
660 Some(nulls) => {
661 let null_offset = range.start;
662 for (idx, (offset, size)) in offsets.iter().zip(sizes.iter()).enumerate() {
664 let is_valid = nulls.is_valid(idx + null_offset);
665 let start_idx = offset.as_usize();
666 let size = size.as_usize();
667 let end_idx = start_idx + size;
668 if !is_valid {
669 write_null_slice(child)
670 } else if size == 0 {
671 write_empty_slice(child)
672 } else {
673 write_non_null_slice(child, start_idx, end_idx)
674 }
675 }
676 }
677 None => {
678 for (offset, size) in offsets.iter().zip(sizes.iter()) {
679 let start_idx = offset.as_usize();
680 let size = size.as_usize();
681 let end_idx = start_idx + size;
682 if size == 0 {
683 write_empty_slice(child)
684 } else {
685 write_non_null_slice(child, start_idx, end_idx)
686 }
687 }
688 }
689 }
690 }
691
692 fn write_struct(
694 children: &mut [LevelInfoBuilder],
695 ctx: &LevelContext,
696 nulls: Option<&NullBuffer>,
697 range: Range<usize>,
698 ) {
699 let write_null = |children: &mut [LevelInfoBuilder], range: Range<usize>| {
700 let len = range.end - range.start;
701 for child in children {
702 child.visit_leaves(|info| {
703 info.extend_uniform_levels(ctx.def_level - 1, ctx.rep_level, len);
704 })
705 }
706 };
707
708 if nulls.is_some_and(|nulls| nulls.null_count() == nulls.len()) {
710 write_null(children, range);
711 return;
712 }
713
714 let write_non_null = |children: &mut [LevelInfoBuilder], range: Range<usize>| {
715 for child in children {
716 child.write(range.clone())
717 }
718 };
719
720 match nulls {
721 Some(validity) => {
722 let mut last_non_null_idx = None;
723 let mut last_null_idx = None;
724
725 for i in range.clone() {
727 match validity.is_valid(i) {
728 true => {
729 if let Some(last_idx) = last_null_idx.take() {
730 write_null(children, last_idx..i)
731 }
732 last_non_null_idx.get_or_insert(i);
733 }
734 false => {
735 if let Some(last_idx) = last_non_null_idx.take() {
736 write_non_null(children, last_idx..i)
737 }
738 last_null_idx.get_or_insert(i);
739 }
740 }
741 }
742
743 if let Some(last_idx) = last_null_idx.take() {
744 write_null(children, last_idx..range.end)
745 }
746
747 if let Some(last_idx) = last_non_null_idx.take() {
748 write_non_null(children, last_idx..range.end)
749 }
750 }
751 None => write_non_null(children, range),
752 }
753 }
754
755 fn write_fixed_size_list(
757 child: &mut LevelInfoBuilder,
758 ctx: &LevelContext,
759 fixed_size: usize,
760 nulls: Option<&NullBuffer>,
761 range: Range<usize>,
762 ) {
763 if nulls.is_some_and(|nulls| nulls.null_count() == nulls.len()) {
765 let count = range.end - range.start;
766 child.visit_leaves(|leaf| {
767 leaf.extend_uniform_levels(ctx.def_level - 2, ctx.rep_level - 1, count);
768 });
769 return;
770 }
771
772 let write_non_null = |child: &mut LevelInfoBuilder, start_idx: usize, end_idx: usize| {
773 let values_start = start_idx * fixed_size;
774 let values_end = end_idx * fixed_size;
775 child.write(values_start..values_end);
776
777 child.visit_leaves(|leaf| {
778 let rep_levels = leaf.rep_levels.materialize_mut().unwrap();
779
780 let row_indices = (0..fixed_size)
781 .rev()
782 .cycle()
783 .take(values_end - values_start);
784
785 rep_levels
787 .iter_mut()
788 .rev()
789 .filter(|&&mut r| r == ctx.rep_level)
791 .zip(row_indices)
792 .for_each(|(r, idx)| {
793 if idx == 0 {
794 *r = ctx.rep_level - 1;
795 }
796 });
797 })
798 };
799
800 let write_empty = |child: &mut LevelInfoBuilder, start_idx: usize, end_idx: usize| {
802 let len = end_idx - start_idx;
803 child.visit_leaves(|leaf| {
804 leaf.append_rep_level_run(ctx.rep_level - 1, len);
805 leaf.append_def_level_run(ctx.def_level - 1, len);
806 })
807 };
808
809 let write_rows = |child: &mut LevelInfoBuilder, start_idx: usize, end_idx: usize| {
810 if fixed_size > 0 {
811 write_non_null(child, start_idx, end_idx)
812 } else {
813 write_empty(child, start_idx, end_idx)
814 }
815 };
816
817 match nulls {
818 Some(nulls) => {
819 let mut start_idx = None;
820 for idx in range.clone() {
821 if nulls.is_valid(idx) {
822 start_idx.get_or_insert(idx);
824 } else {
825 if let Some(start) = start_idx.take() {
827 write_rows(child, start, idx);
828 }
829 child.visit_leaves(|leaf| {
831 leaf.append_rep_level_run(ctx.rep_level - 1, 1);
832 leaf.append_def_level_run(ctx.def_level - 2, 1);
833 })
834 }
835 }
836 if let Some(start) = start_idx.take() {
838 write_rows(child, start, range.end);
839 }
840 }
841 None => write_rows(child, range.start, range.end),
843 }
844 }
845
846 fn write_leaf(info: &mut ArrayLevels, range: Range<usize>) {
848 let len = range.end - range.start;
849
850 if let Some(nulls) = &info.logical_nulls {
852 if !matches!(info.def_levels, LevelData::Absent) && nulls.null_count() == nulls.len() {
853 info.extend_uniform_levels(info.max_def_level - 1, info.max_rep_level, len);
854 return;
855 }
856 }
857
858 if matches!(info.def_levels, LevelData::Absent) {
859 info.non_null_indices.extend(range.clone());
860 } else {
861 let max_def_level = info.max_def_level;
862 match &info.logical_nulls {
863 Some(nulls) => {
864 assert!(range.end <= nulls.len());
865 if len >= BULK_FILL_MIN_LEN && nulls.null_count() * 2 >= nulls.len() {
870 let range_nulls = nulls.slice(range.start, len);
871 let valid_in_range = len - range_nulls.null_count();
872 let null_def_level = max_def_level - 1;
873 let buf = info
874 .def_levels
875 .materialize_mut()
876 .expect("definition levels present");
877 let base = buf.len();
878 buf.resize(base + len, null_def_level);
879 for i in range_nulls.valid_indices() {
880 buf[base + i] = max_def_level;
881 }
882 info.non_null_indices.reserve(valid_in_range);
883 info.non_null_indices
884 .extend(range_nulls.valid_indices().map(|i| i + range.start));
885 } else {
886 let bits = nulls.inner();
887 info.def_levels.extend_from_iter(range.clone().map(|i| {
888 let valid = unsafe { bits.value_unchecked(i) };
890 max_def_level - (!valid as i16)
891 }));
892 info.non_null_indices.reserve(len);
893 info.non_null_indices.extend(
894 BitIndexIterator::new(bits.inner(), bits.offset() + range.start, len)
895 .map(|i| i + range.start),
896 );
897 }
898 }
899 None => {
900 info.append_def_level_run(max_def_level, len);
901 info.non_null_indices.reserve(len);
902 info.non_null_indices.extend(range.clone());
903 }
904 }
905 }
906
907 if !matches!(info.rep_levels, LevelData::Absent) {
908 info.append_rep_level_run(info.max_rep_level, len);
909 }
910 }
911
912 fn visit_leaves(&mut self, visit: impl Fn(&mut ArrayLevels) + Copy) {
914 match self {
915 LevelInfoBuilder::Primitive(info) => visit(info),
916 LevelInfoBuilder::List(c, _, _, _, _)
917 | LevelInfoBuilder::LargeList(c, _, _, _, _)
918 | LevelInfoBuilder::FixedSizeList(c, _, _, _)
919 | LevelInfoBuilder::ListView(c, _, _, _, _)
920 | LevelInfoBuilder::LargeListView(c, _, _, _, _) => c.visit_leaves(visit),
921 LevelInfoBuilder::Struct(children, _, _) => {
922 for c in children {
923 c.visit_leaves(visit)
924 }
925 }
926 }
927 }
928
929 fn types_compatible(a: &DataType, b: &DataType) -> bool {
935 if a.equals_datatype(b) {
937 return true;
938 }
939
940 let (a, b) = match (a, b) {
942 (DataType::Dictionary(_, va), DataType::Dictionary(_, vb)) => {
943 (va.as_ref(), vb.as_ref())
944 }
945 (DataType::Dictionary(_, v), b) => (v.as_ref(), b),
946 (a, DataType::Dictionary(_, v)) => (a, v.as_ref()),
947 _ => (a, b),
948 };
949
950 if a == b {
953 return true;
954 }
955
956 match a {
959 DataType::Utf8 => matches!(b, DataType::LargeUtf8 | DataType::Utf8View),
961 DataType::Utf8View => matches!(b, DataType::LargeUtf8 | DataType::Utf8),
962 DataType::LargeUtf8 => matches!(b, DataType::Utf8 | DataType::Utf8View),
963
964 DataType::Binary => matches!(b, DataType::LargeBinary | DataType::BinaryView),
966 DataType::BinaryView => matches!(b, DataType::LargeBinary | DataType::Binary),
967 DataType::LargeBinary => matches!(b, DataType::Binary | DataType::BinaryView),
968
969 _ => false,
971 }
972 }
973}
974
975#[derive(Debug, Clone)]
978pub(crate) enum LevelData {
979 Absent,
980 Materialized(Vec<i16>),
981 Uniform { value: i16, count: usize },
982}
983
984impl PartialEq for LevelData {
987 fn eq(&self, other: &Self) -> bool {
988 match (self, other) {
989 (Self::Absent, Self::Absent) => true,
990 (Self::Materialized(a), Self::Materialized(b)) => a == b,
991 (Self::Uniform { value: v, count: n }, Self::Materialized(b))
992 | (Self::Materialized(b), Self::Uniform { value: v, count: n }) => {
993 b.len() == *n && b.iter().all(|x| x == v)
994 }
995 (
996 Self::Uniform {
997 value: v1,
998 count: n1,
999 },
1000 Self::Uniform {
1001 value: v2,
1002 count: n2,
1003 },
1004 ) => v1 == v2 && n1 == n2,
1005 _ => false,
1006 }
1007 }
1008}
1009
1010impl Eq for LevelData {}
1011
1012impl LevelData {
1013 fn new(present: bool) -> Self {
1014 match present {
1015 true => Self::Materialized(Vec::new()),
1016 false => Self::Absent,
1017 }
1018 }
1019
1020 pub(crate) fn as_ref(&self) -> LevelDataRef<'_> {
1021 match self {
1022 Self::Absent => LevelDataRef::Absent,
1023 Self::Materialized(values) => LevelDataRef::Materialized(values),
1024 Self::Uniform { value, count } => LevelDataRef::Uniform {
1025 value: *value,
1026 count: *count,
1027 },
1028 }
1029 }
1030
1031 pub(crate) fn slice(&self, offset: usize, len: usize) -> Self {
1032 match self {
1033 Self::Absent => Self::Absent,
1034 Self::Materialized(values) => Self::Materialized(values[offset..offset + len].to_vec()),
1035 Self::Uniform { value, .. } => Self::Uniform {
1036 value: *value,
1037 count: len,
1038 },
1039 }
1040 }
1041
1042 fn append_run(&mut self, value: i16, count: usize) {
1043 if count == 0 {
1044 return;
1045 }
1046
1047 match self {
1048 Self::Absent => {}
1051 Self::Materialized(values) if values.is_empty() => {
1054 *self = Self::Uniform { value, count };
1055 }
1056 Self::Materialized(values) => values.extend(std::iter::repeat_n(value, count)),
1058 Self::Uniform {
1061 value: uniform_value,
1062 count: uniform_count,
1063 } if *uniform_value == value => {
1064 *uniform_count += count;
1065 }
1066 Self::Uniform { .. } => {
1069 let values = self.materialize_mut().unwrap();
1070 values.extend(std::iter::repeat_n(value, count));
1071 }
1072 }
1073 }
1074
1075 fn extend_from_iter<I>(&mut self, iter: I)
1076 where
1077 I: IntoIterator<Item = i16>,
1078 {
1079 if let Some(values) = self.materialize_mut() {
1080 values.extend(iter);
1081 }
1082 }
1083
1084 fn materialize_mut(&mut self) -> Option<&mut Vec<i16>> {
1087 match self {
1088 Self::Absent => None,
1089 Self::Materialized(values) => Some(values),
1090 Self::Uniform { value, count } => {
1091 let values = vec![*value; *count];
1092 *self = Self::Materialized(values);
1093 match self {
1094 Self::Materialized(values) => Some(values),
1095 _ => unreachable!(),
1096 }
1097 }
1098 }
1099 }
1100}
1101
1102#[derive(Debug, Clone)]
1103pub(crate) struct ArrayLevels {
1104 def_levels: LevelData,
1108
1109 rep_levels: LevelData,
1113
1114 non_null_indices: Vec<usize>,
1117
1118 max_def_level: i16,
1120
1121 max_rep_level: i16,
1123
1124 array: ArrayRef,
1126
1127 logical_nulls: Option<NullBuffer>,
1129}
1130
1131impl PartialEq for ArrayLevels {
1132 fn eq(&self, other: &Self) -> bool {
1133 self.def_levels == other.def_levels
1134 && self.rep_levels == other.rep_levels
1135 && self.non_null_indices == other.non_null_indices
1136 && self.max_def_level == other.max_def_level
1137 && self.max_rep_level == other.max_rep_level
1138 && self.array.as_ref() == other.array.as_ref()
1139 && self.logical_nulls.as_ref() == other.logical_nulls.as_ref()
1140 }
1141}
1142impl Eq for ArrayLevels {}
1143
1144impl ArrayLevels {
1145 fn new(ctx: LevelContext, is_nullable: bool, array: ArrayRef) -> Self {
1146 let max_rep_level = ctx.rep_level;
1147 let max_def_level = match is_nullable {
1148 true => ctx.def_level + 1,
1149 false => ctx.def_level,
1150 };
1151
1152 let logical_nulls = array.logical_nulls();
1153
1154 Self {
1155 def_levels: LevelData::new(max_def_level != 0),
1156 rep_levels: LevelData::new(max_rep_level != 0),
1157 non_null_indices: vec![],
1158 max_def_level,
1159 max_rep_level,
1160 array,
1161 logical_nulls,
1162 }
1163 }
1164
1165 pub fn array(&self) -> &ArrayRef {
1166 &self.array
1167 }
1168
1169 pub(crate) fn def_level_data(&self) -> &LevelData {
1170 &self.def_levels
1171 }
1172
1173 pub(crate) fn rep_level_data(&self) -> &LevelData {
1174 &self.rep_levels
1175 }
1176
1177 pub fn non_null_indices(&self) -> &[usize] {
1178 &self.non_null_indices
1179 }
1180
1181 pub(crate) fn slice_for_chunk(&self, chunk: &CdcChunk) -> Self {
1187 let def_levels = self.def_levels.slice(chunk.level_offset, chunk.num_levels);
1188 let rep_levels = self.rep_levels.slice(chunk.level_offset, chunk.num_levels);
1189
1190 let nni = &self.non_null_indices[chunk.value_offset..chunk.value_offset + chunk.num_values];
1192 let start = nni.first().copied().unwrap_or(0);
1197 let end = nni.last().map_or(0, |&i| i + 1);
1198 let non_null_indices = nni.iter().map(|&idx| idx - start).collect();
1200 let array = self.array.slice(start, end - start);
1202 let logical_nulls = array.logical_nulls();
1203
1204 Self {
1205 def_levels,
1206 rep_levels,
1207 non_null_indices,
1208 max_def_level: self.max_def_level,
1209 max_rep_level: self.max_rep_level,
1210 array,
1211 logical_nulls,
1212 }
1213 }
1214
1215 fn extend_uniform_levels(&mut self, def_val: i16, rep_val: i16, count: usize) {
1217 self.def_levels.append_run(def_val, count);
1218 self.rep_levels.append_run(rep_val, count);
1219 }
1220
1221 fn append_def_level_run(&mut self, value: i16, count: usize) {
1222 self.def_levels.append_run(value, count);
1223 }
1224
1225 fn append_rep_level_run(&mut self, value: i16, count: usize) {
1226 self.rep_levels.append_run(value, count);
1227 }
1228}
1229
1230#[cfg(test)]
1231mod tests {
1232 use super::*;
1233 use crate::column::chunker::CdcChunk;
1234
1235 use arrow_array::builder::*;
1236 use arrow_array::types::Int32Type;
1237 use arrow_array::*;
1238 use arrow_buffer::{Buffer, ToByteSlice};
1239 use arrow_cast::display::array_value_to_string;
1240 use arrow_data::{ArrayData, ArrayDataBuilder};
1241 use arrow_schema::{Fields, Schema};
1242
1243 #[test]
1244 fn test_calculate_array_levels_twitter_example() {
1245 let leaf_type = Field::new_list_field(DataType::Int32, false);
1249 let inner_type = DataType::List(Arc::new(leaf_type));
1250 let inner_field = Field::new("l2", inner_type.clone(), false);
1251 let outer_type = DataType::List(Arc::new(inner_field));
1252 let outer_field = Field::new("l1", outer_type.clone(), false);
1253
1254 let primitives = Int32Array::from_iter(0..10);
1255
1256 let offsets = Buffer::from_iter([0_i32, 3, 7, 8, 10]);
1258 let inner_list = ArrayDataBuilder::new(inner_type)
1259 .len(4)
1260 .add_buffer(offsets)
1261 .add_child_data(primitives.to_data())
1262 .build()
1263 .unwrap();
1264
1265 let offsets = Buffer::from_iter([0_i32, 2, 4]);
1266 let outer_list = ArrayDataBuilder::new(outer_type)
1267 .len(2)
1268 .add_buffer(offsets)
1269 .add_child_data(inner_list)
1270 .build()
1271 .unwrap();
1272 let outer_list = make_array(outer_list);
1273
1274 let levels = calculate_array_levels(&outer_list, &outer_field).unwrap();
1275 assert_eq!(levels.len(), 1);
1276
1277 let expected = ArrayLevels {
1278 def_levels: LevelData::Materialized(vec![2; 10]),
1279 rep_levels: LevelData::Materialized(vec![0, 2, 2, 1, 2, 2, 2, 0, 1, 2]),
1280 non_null_indices: vec![0, 1, 2, 3, 4, 5, 6, 7, 8, 9],
1281 max_def_level: 2,
1282 max_rep_level: 2,
1283 array: Arc::new(primitives),
1284 logical_nulls: None,
1285 };
1286 assert_eq!(&levels[0], &expected);
1287 }
1288
1289 #[test]
1290 fn test_calculate_one_level_1() {
1291 let array = Arc::new(Int32Array::from_iter(0..10)) as ArrayRef;
1293 let field = Field::new_list_field(DataType::Int32, false);
1294
1295 let levels = calculate_array_levels(&array, &field).unwrap();
1296 assert_eq!(levels.len(), 1);
1297
1298 let expected_levels = ArrayLevels {
1299 def_levels: LevelData::Absent,
1300 rep_levels: LevelData::Absent,
1301 non_null_indices: (0..10).collect(),
1302 max_def_level: 0,
1303 max_rep_level: 0,
1304 array,
1305 logical_nulls: None,
1306 };
1307 assert_eq!(&levels[0], &expected_levels);
1308 }
1309
1310 #[test]
1311 fn test_calculate_one_level_2() {
1312 let array = Arc::new(Int32Array::from_iter([
1314 Some(0),
1315 None,
1316 Some(0),
1317 Some(0),
1318 None,
1319 ])) as ArrayRef;
1320 let field = Field::new_list_field(DataType::Int32, true);
1321
1322 let levels = calculate_array_levels(&array, &field).unwrap();
1323 assert_eq!(levels.len(), 1);
1324
1325 let logical_nulls = array.logical_nulls();
1326 let expected_levels = ArrayLevels {
1327 def_levels: LevelData::Materialized(vec![1, 0, 1, 1, 0]),
1328 rep_levels: LevelData::Absent,
1329 non_null_indices: vec![0, 2, 3],
1330 max_def_level: 1,
1331 max_rep_level: 0,
1332 array,
1333 logical_nulls,
1334 };
1335 assert_eq!(&levels[0], &expected_levels);
1336 }
1337
1338 #[test]
1339 fn test_calculate_array_levels_1() {
1340 let leaf_field = Field::new_list_field(DataType::Int32, false);
1341 let list_type = DataType::List(Arc::new(leaf_field));
1342
1343 let leaf_array = Int32Array::from_iter(0..5);
1347 let offsets = Buffer::from_iter(0_i32..6);
1349 let list = ArrayDataBuilder::new(list_type.clone())
1350 .len(5)
1351 .add_buffer(offsets)
1352 .add_child_data(leaf_array.to_data())
1353 .build()
1354 .unwrap();
1355 let list = make_array(list);
1356
1357 let list_field = Field::new("list", list_type.clone(), false);
1358 let levels = calculate_array_levels(&list, &list_field).unwrap();
1359 assert_eq!(levels.len(), 1);
1360
1361 let expected_levels = ArrayLevels {
1362 def_levels: LevelData::Materialized(vec![1; 5]),
1363 rep_levels: LevelData::Materialized(vec![0; 5]),
1364 non_null_indices: (0..5).collect(),
1365 max_def_level: 1,
1366 max_rep_level: 1,
1367 array: Arc::new(leaf_array),
1368 logical_nulls: None,
1369 };
1370 assert_eq!(&levels[0], &expected_levels);
1371
1372 let leaf_array = Int32Array::from_iter([0, 0, 2, 2, 3, 3, 3, 3, 4, 4, 4]);
1381 let offsets = Buffer::from_iter([0_i32, 2, 2, 4, 8, 11]);
1382 let list = ArrayDataBuilder::new(list_type.clone())
1383 .len(5)
1384 .add_buffer(offsets)
1385 .add_child_data(leaf_array.to_data())
1386 .null_bit_buffer(Some(Buffer::from([0b00011101])))
1387 .build()
1388 .unwrap();
1389 let list = make_array(list);
1390
1391 let list_field = Field::new("list", list_type, true);
1392 let levels = calculate_array_levels(&list, &list_field).unwrap();
1393 assert_eq!(levels.len(), 1);
1394
1395 let expected_levels = ArrayLevels {
1396 def_levels: LevelData::Materialized(vec![2, 2, 0, 2, 2, 2, 2, 2, 2, 2, 2, 2]),
1397 rep_levels: LevelData::Materialized(vec![0, 1, 0, 0, 1, 0, 1, 1, 1, 0, 1, 1]),
1398 non_null_indices: (0..11).collect(),
1399 max_def_level: 2,
1400 max_rep_level: 1,
1401 array: Arc::new(leaf_array),
1402 logical_nulls: None,
1403 };
1404 assert_eq!(&levels[0], &expected_levels);
1405 }
1406
1407 #[test]
1408 fn test_calculate_array_levels_2() {
1409 let leaf = Int32Array::from_iter(0..11);
1423 let leaf_field = Field::new("leaf", DataType::Int32, false);
1424
1425 let list_type = DataType::List(Arc::new(leaf_field));
1426 let list = ArrayData::builder(list_type.clone())
1427 .len(5)
1428 .add_child_data(leaf.to_data())
1429 .add_buffer(Buffer::from_iter([0_i32, 2, 2, 4, 8, 11]))
1430 .build()
1431 .unwrap();
1432
1433 let list = make_array(list);
1434 let list_field = Arc::new(Field::new("list", list_type, true));
1435
1436 let struct_array =
1437 StructArray::from((vec![(list_field, list)], Buffer::from([0b00011010])));
1438 let array = Arc::new(struct_array) as ArrayRef;
1439
1440 let struct_field = Field::new("struct", array.data_type().clone(), true);
1441
1442 let levels = calculate_array_levels(&array, &struct_field).unwrap();
1443 assert_eq!(levels.len(), 1);
1444
1445 let expected_levels = ArrayLevels {
1446 def_levels: LevelData::Materialized(vec![0, 2, 0, 3, 3, 3, 3, 3, 3, 3]),
1447 rep_levels: LevelData::Materialized(vec![0, 0, 0, 0, 1, 1, 1, 0, 1, 1]),
1448 non_null_indices: (4..11).collect(),
1449 max_def_level: 3,
1450 max_rep_level: 1,
1451 array: Arc::new(leaf),
1452 logical_nulls: None,
1453 };
1454
1455 assert_eq!(&levels[0], &expected_levels);
1456
1457 let leaf = Int32Array::from_iter(100..122);
1466 let leaf_field = Field::new("leaf", DataType::Int32, true);
1467
1468 let l1_type = DataType::List(Arc::new(leaf_field));
1469 let offsets = Buffer::from_iter([0_i32, 2, 4, 6, 8, 10, 12, 14, 16, 18, 20, 22]);
1470 let l1 = ArrayData::builder(l1_type.clone())
1471 .len(11)
1472 .add_child_data(leaf.to_data())
1473 .add_buffer(offsets)
1474 .build()
1475 .unwrap();
1476
1477 let l1_field = Field::new("l1", l1_type, true);
1478 let l2_type = DataType::List(Arc::new(l1_field));
1479 let l2 = ArrayData::builder(l2_type)
1480 .len(5)
1481 .add_child_data(l1)
1482 .add_buffer(Buffer::from_iter([0, 2, 2, 4, 8, 11]))
1483 .build()
1484 .unwrap();
1485
1486 let l2 = make_array(l2);
1487 let l2_field = Field::new("l2", l2.data_type().clone(), true);
1488
1489 let levels = calculate_array_levels(&l2, &l2_field).unwrap();
1490 assert_eq!(levels.len(), 1);
1491
1492 let expected_levels = ArrayLevels {
1493 def_levels: LevelData::Materialized(vec![
1494 5, 5, 5, 5, 1, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5,
1495 ]),
1496 rep_levels: LevelData::Materialized(vec![
1497 0, 2, 1, 2, 0, 0, 2, 1, 2, 0, 2, 1, 2, 1, 2, 1, 2, 0, 2, 1, 2, 1, 2,
1498 ]),
1499 non_null_indices: (0..22).collect(),
1500 max_def_level: 5,
1501 max_rep_level: 2,
1502 array: Arc::new(leaf),
1503 logical_nulls: None,
1504 };
1505
1506 assert_eq!(&levels[0], &expected_levels);
1507 }
1508
1509 #[test]
1510 fn test_calculate_array_levels_nested_list() {
1511 let leaf_field = Field::new("leaf", DataType::Int32, false);
1512 let list_type = DataType::List(Arc::new(leaf_field));
1513
1514 let leaf = Int32Array::from_iter([0; 4]);
1522 let list = ArrayData::builder(list_type.clone())
1523 .len(4)
1524 .add_buffer(Buffer::from_iter(0_i32..5))
1525 .add_child_data(leaf.to_data())
1526 .build()
1527 .unwrap();
1528 let list = make_array(list);
1529
1530 let list_field = Field::new("list", list_type.clone(), false);
1531 let levels = calculate_array_levels(&list, &list_field).unwrap();
1532 assert_eq!(levels.len(), 1);
1533
1534 let expected_levels = ArrayLevels {
1535 def_levels: LevelData::Materialized(vec![1; 4]),
1536 rep_levels: LevelData::Materialized(vec![0; 4]),
1537 non_null_indices: (0..4).collect(),
1538 max_def_level: 1,
1539 max_rep_level: 1,
1540 array: Arc::new(leaf),
1541 logical_nulls: None,
1542 };
1543 assert_eq!(&levels[0], &expected_levels);
1544
1545 let leaf = Int32Array::from_iter(0..8);
1550 let list = ArrayData::builder(list_type.clone())
1551 .len(4)
1552 .add_buffer(Buffer::from_iter([0_i32, 0, 3, 5, 7]))
1553 .null_bit_buffer(Some(Buffer::from([0b00001110])))
1554 .add_child_data(leaf.to_data())
1555 .build()
1556 .unwrap();
1557 let list = make_array(list);
1558 let list_field = Arc::new(Field::new("list", list_type, true));
1559
1560 let struct_array = StructArray::from(vec![(list_field, list)]);
1561 let array = Arc::new(struct_array) as ArrayRef;
1562
1563 let struct_field = Field::new("struct", array.data_type().clone(), true);
1564 let levels = calculate_array_levels(&array, &struct_field).unwrap();
1565 assert_eq!(levels.len(), 1);
1566
1567 let expected_levels = ArrayLevels {
1568 def_levels: LevelData::Materialized(vec![1, 3, 3, 3, 3, 3, 3, 3]),
1569 rep_levels: LevelData::Materialized(vec![0, 0, 1, 1, 0, 1, 0, 1]),
1570 non_null_indices: (0..7).collect(),
1571 max_def_level: 3,
1572 max_rep_level: 1,
1573 array: Arc::new(leaf),
1574 logical_nulls: None,
1575 };
1576 assert_eq!(&levels[0], &expected_levels);
1577
1578 let leaf = Int32Array::from_iter(201..216);
1586 let leaf_field = Field::new("leaf", DataType::Int32, false);
1587 let list_1_type = DataType::List(Arc::new(leaf_field));
1588 let list_1 = ArrayData::builder(list_1_type.clone())
1589 .len(7)
1590 .add_buffer(Buffer::from_iter([0_i32, 1, 3, 3, 6, 10, 10, 15]))
1591 .add_child_data(leaf.to_data())
1592 .build()
1593 .unwrap();
1594
1595 let list_1_field = Field::new("l1", list_1_type, true);
1596 let list_2_type = DataType::List(Arc::new(list_1_field));
1597 let list_2 = ArrayData::builder(list_2_type.clone())
1598 .len(4)
1599 .add_buffer(Buffer::from_iter([0_i32, 0, 3, 5, 7]))
1600 .null_bit_buffer(Some(Buffer::from([0b00001110])))
1601 .add_child_data(list_1)
1602 .build()
1603 .unwrap();
1604
1605 let list_2 = make_array(list_2);
1606 let list_2_field = Arc::new(Field::new("list_2", list_2_type, true));
1607
1608 let struct_array =
1609 StructArray::from((vec![(list_2_field, list_2)], Buffer::from([0b00001111])));
1610 let struct_field = Field::new("struct", struct_array.data_type().clone(), true);
1611
1612 let array = Arc::new(struct_array) as ArrayRef;
1613 let levels = calculate_array_levels(&array, &struct_field).unwrap();
1614 assert_eq!(levels.len(), 1);
1615
1616 let expected_levels = ArrayLevels {
1617 def_levels: LevelData::Materialized(vec![
1618 1, 5, 5, 5, 4, 5, 5, 5, 5, 5, 5, 5, 4, 5, 5, 5, 5, 5,
1619 ]),
1620 rep_levels: LevelData::Materialized(vec![
1621 0, 0, 1, 2, 1, 0, 2, 2, 1, 2, 2, 2, 0, 1, 2, 2, 2, 2,
1622 ]),
1623 non_null_indices: (0..15).collect(),
1624 max_def_level: 5,
1625 max_rep_level: 2,
1626 array: Arc::new(leaf),
1627 logical_nulls: None,
1628 };
1629 assert_eq!(&levels[0], &expected_levels);
1630 }
1631
1632 #[test]
1633 fn test_calculate_nested_struct_levels() {
1634 let c = Int32Array::from_iter([Some(1), None, Some(3), None, Some(5), Some(6)]);
1644 let leaf = Arc::new(c) as ArrayRef;
1645 let c_field = Arc::new(Field::new("c", DataType::Int32, true));
1646 let b = StructArray::from(((vec![(c_field, leaf.clone())]), Buffer::from([0b00110111])));
1647
1648 let b_field = Arc::new(Field::new("b", b.data_type().clone(), true));
1649 let a = StructArray::from((
1650 (vec![(b_field, Arc::new(b) as ArrayRef)]),
1651 Buffer::from([0b00101111]),
1652 ));
1653
1654 let a_field = Field::new("a", a.data_type().clone(), true);
1655 let a_array = Arc::new(a) as ArrayRef;
1656
1657 let levels = calculate_array_levels(&a_array, &a_field).unwrap();
1658 assert_eq!(levels.len(), 1);
1659
1660 let logical_nulls = leaf.logical_nulls();
1661 let expected_levels = ArrayLevels {
1662 def_levels: LevelData::Materialized(vec![3, 2, 3, 1, 0, 3]),
1663 rep_levels: LevelData::Absent,
1664 non_null_indices: vec![0, 2, 5],
1665 max_def_level: 3,
1666 max_rep_level: 0,
1667 array: leaf,
1668 logical_nulls,
1669 };
1670 assert_eq!(&levels[0], &expected_levels);
1671 }
1672
1673 #[test]
1674 fn list_single_column() {
1675 let a_values = Int32Array::from(vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10]);
1678 let a_value_offsets = arrow::buffer::Buffer::from_iter([0_i32, 1, 3, 3, 6, 10]);
1679 let a_list_type = DataType::List(Arc::new(Field::new_list_field(DataType::Int32, true)));
1680 let a_list_data = ArrayData::builder(a_list_type.clone())
1681 .len(5)
1682 .add_buffer(a_value_offsets)
1683 .null_bit_buffer(Some(Buffer::from([0b00011011])))
1684 .add_child_data(a_values.to_data())
1685 .build()
1686 .unwrap();
1687
1688 assert_eq!(a_list_data.null_count(), 1);
1689
1690 let a = ListArray::from(a_list_data);
1691
1692 let item_field = Field::new_list_field(a_list_type, true);
1693 let mut builder = levels(&item_field, a);
1694 builder.write(2..4);
1695 let levels = builder.finish();
1696
1697 assert_eq!(levels.len(), 1);
1698
1699 let list_level = &levels[0];
1700
1701 let expected_level = ArrayLevels {
1702 def_levels: LevelData::Materialized(vec![0, 3, 3, 3]),
1703 rep_levels: LevelData::Materialized(vec![0, 0, 1, 1]),
1704 non_null_indices: vec![3, 4, 5],
1705 max_def_level: 3,
1706 max_rep_level: 1,
1707 array: Arc::new(a_values),
1708 logical_nulls: None,
1709 };
1710 assert_eq!(list_level, &expected_level);
1711 }
1712
1713 #[test]
1714 fn mixed_struct_list() {
1715 let struct_field_d = Arc::new(Field::new("d", DataType::Float64, true));
1719 let struct_field_f = Arc::new(Field::new("f", DataType::Float32, true));
1720 let struct_field_g = Arc::new(Field::new(
1721 "g",
1722 DataType::List(Arc::new(Field::new("items", DataType::Int16, false))),
1723 false,
1724 ));
1725 let struct_field_e = Arc::new(Field::new(
1726 "e",
1727 DataType::Struct(vec![struct_field_f.clone(), struct_field_g.clone()].into()),
1728 true,
1729 ));
1730 let schema = Schema::new(vec![
1731 Field::new("a", DataType::Int32, false),
1732 Field::new("b", DataType::Int32, true),
1733 Field::new(
1734 "c",
1735 DataType::Struct(vec![struct_field_d.clone(), struct_field_e.clone()].into()),
1736 true, ),
1738 ]);
1739
1740 let a = Int32Array::from(vec![1, 2, 3, 4, 5]);
1742 let b = Int32Array::from(vec![Some(1), None, None, Some(4), Some(5)]);
1743 let d = Float64Array::from(vec![None, None, None, Some(1.0), None]);
1744 let f = Float32Array::from(vec![Some(0.0), None, Some(333.3), None, Some(5.25)]);
1745
1746 let g_value = Int16Array::from(vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10]);
1747
1748 let g_value_offsets = arrow::buffer::Buffer::from([0, 1, 3, 3, 6, 10].to_byte_slice());
1751
1752 let g_list_data = ArrayData::builder(struct_field_g.data_type().clone())
1754 .len(5)
1755 .add_buffer(g_value_offsets)
1756 .add_child_data(g_value.into_data())
1757 .build()
1758 .unwrap();
1759 let g = ListArray::from(g_list_data);
1760
1761 let e = StructArray::from(vec![
1762 (struct_field_f, Arc::new(f.clone()) as ArrayRef),
1763 (struct_field_g, Arc::new(g) as ArrayRef),
1764 ]);
1765
1766 let c = StructArray::from(vec![
1767 (struct_field_d, Arc::new(d.clone()) as ArrayRef),
1768 (struct_field_e, Arc::new(e) as ArrayRef),
1769 ]);
1770
1771 let batch = RecordBatch::try_new(
1773 Arc::new(schema),
1774 vec![Arc::new(a.clone()), Arc::new(b.clone()), Arc::new(c)],
1775 )
1776 .unwrap();
1777
1778 let mut levels = vec![];
1781 batch
1782 .columns()
1783 .iter()
1784 .zip(batch.schema().fields())
1785 .for_each(|(array, field)| {
1786 let mut array_levels = calculate_array_levels(array, field).unwrap();
1787 levels.append(&mut array_levels);
1788 });
1789 assert_eq!(levels.len(), 5);
1790
1791 let list_level = &levels[0];
1793
1794 let expected_level = ArrayLevels {
1795 def_levels: LevelData::Absent,
1796 rep_levels: LevelData::Absent,
1797 non_null_indices: vec![0, 1, 2, 3, 4],
1798 max_def_level: 0,
1799 max_rep_level: 0,
1800 array: Arc::new(a),
1801 logical_nulls: None,
1802 };
1803 assert_eq!(list_level, &expected_level);
1804
1805 let list_level = levels.get(1).unwrap();
1807
1808 let b_logical_nulls = b.logical_nulls();
1809 let expected_level = ArrayLevels {
1810 def_levels: LevelData::Materialized(vec![1, 0, 0, 1, 1]),
1811 rep_levels: LevelData::Absent,
1812 non_null_indices: vec![0, 3, 4],
1813 max_def_level: 1,
1814 max_rep_level: 0,
1815 array: Arc::new(b),
1816 logical_nulls: b_logical_nulls,
1817 };
1818 assert_eq!(list_level, &expected_level);
1819
1820 let list_level = levels.get(2).unwrap();
1822
1823 let d_logical_nulls = d.logical_nulls();
1824 let expected_level = ArrayLevels {
1825 def_levels: LevelData::Materialized(vec![1, 1, 1, 2, 1]),
1826 rep_levels: LevelData::Absent,
1827 non_null_indices: vec![3],
1828 max_def_level: 2,
1829 max_rep_level: 0,
1830 array: Arc::new(d),
1831 logical_nulls: d_logical_nulls,
1832 };
1833 assert_eq!(list_level, &expected_level);
1834
1835 let list_level = levels.get(3).unwrap();
1837
1838 let f_logical_nulls = f.logical_nulls();
1839 let expected_level = ArrayLevels {
1840 def_levels: LevelData::Materialized(vec![3, 2, 3, 2, 3]),
1841 rep_levels: LevelData::Absent,
1842 non_null_indices: vec![0, 2, 4],
1843 max_def_level: 3,
1844 max_rep_level: 0,
1845 array: Arc::new(f),
1846 logical_nulls: f_logical_nulls,
1847 };
1848 assert_eq!(list_level, &expected_level);
1849 }
1850
1851 #[test]
1852 fn test_null_vs_nonnull_struct() {
1853 let offset_field = Arc::new(Field::new("offset", DataType::Int32, true));
1855 let schema = Schema::new(vec![Field::new(
1856 "some_nested_object",
1857 DataType::Struct(vec![offset_field.clone()].into()),
1858 false,
1859 )]);
1860
1861 let offset = Int32Array::from(vec![1, 2, 3, 4, 5]);
1863
1864 let some_nested_object =
1865 StructArray::from(vec![(offset_field, Arc::new(offset) as ArrayRef)]);
1866
1867 let batch =
1869 RecordBatch::try_new(Arc::new(schema), vec![Arc::new(some_nested_object)]).unwrap();
1870
1871 let struct_null_level =
1872 calculate_array_levels(batch.column(0), batch.schema().field(0)).unwrap();
1873
1874 let offset_field = Arc::new(Field::new("offset", DataType::Int32, true));
1877 let schema = Schema::new(vec![Field::new(
1878 "some_nested_object",
1879 DataType::Struct(vec![offset_field.clone()].into()),
1880 true,
1881 )]);
1882
1883 let offset = Int32Array::from(vec![1, 2, 3, 4, 5]);
1885
1886 let some_nested_object =
1887 StructArray::from(vec![(offset_field, Arc::new(offset) as ArrayRef)]);
1888
1889 let batch =
1891 RecordBatch::try_new(Arc::new(schema), vec![Arc::new(some_nested_object)]).unwrap();
1892
1893 let struct_non_null_level =
1894 calculate_array_levels(batch.column(0), batch.schema().field(0)).unwrap();
1895
1896 if struct_non_null_level == struct_null_level {
1898 panic!("Levels should not be equal, to reflect the difference in struct nullness");
1899 }
1900 }
1901
1902 #[test]
1903 fn test_map_array() {
1904 let json_content = r#"
1906 {"stocks":{"long": "$AAA", "short": "$BBB"}}
1907 {"stocks":{"long": "$CCC", "short": null}}
1908 {"stocks":{"hedged": "$YYY", "long": null, "short": "$D"}}
1909 "#;
1910 let entries_struct_type = DataType::Struct(Fields::from(vec![
1911 Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
1912 Field::new(Field::MAP_VALUE_FIELD_DEFAULT_NAME, DataType::Utf8, true),
1913 ]));
1914 let stocks_field = Field::new(
1915 "stocks",
1916 DataType::Map(
1917 Arc::new(Field::new(
1918 Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
1919 entries_struct_type,
1920 false,
1921 )),
1922 false,
1923 ),
1924 false,
1926 );
1927 let schema = Arc::new(Schema::new(vec![stocks_field]));
1928 let builder = arrow::json::ReaderBuilder::new(schema).with_batch_size(64);
1929 let mut reader = builder.build(std::io::Cursor::new(json_content)).unwrap();
1930
1931 let batch = reader.next().unwrap().unwrap();
1932
1933 let mut levels = vec![];
1935 batch
1936 .columns()
1937 .iter()
1938 .zip(batch.schema().fields())
1939 .for_each(|(array, field)| {
1940 let mut array_levels = calculate_array_levels(array, field).unwrap();
1941 levels.append(&mut array_levels);
1942 });
1943 assert_eq!(levels.len(), 2);
1944
1945 let map = batch.column(0).as_map();
1946 let map_keys_logical_nulls = map.keys().logical_nulls();
1947
1948 let list_level = &levels[0];
1950
1951 let expected_level = ArrayLevels {
1952 def_levels: LevelData::Materialized(vec![1; 7]),
1953 rep_levels: LevelData::Materialized(vec![0, 1, 0, 1, 0, 1, 1]),
1954 non_null_indices: vec![0, 1, 2, 3, 4, 5, 6],
1955 max_def_level: 1,
1956 max_rep_level: 1,
1957 array: map.keys().clone(),
1958 logical_nulls: map_keys_logical_nulls,
1959 };
1960 assert_eq!(list_level, &expected_level);
1961
1962 let list_level = levels.get(1).unwrap();
1964 let map_values_logical_nulls = map.values().logical_nulls();
1965
1966 let expected_level = ArrayLevels {
1967 def_levels: LevelData::Materialized(vec![2, 2, 2, 1, 2, 1, 2]),
1968 rep_levels: LevelData::Materialized(vec![0, 1, 0, 1, 0, 1, 1]),
1969 non_null_indices: vec![0, 1, 2, 4, 6],
1970 max_def_level: 2,
1971 max_rep_level: 1,
1972 array: map.values().clone(),
1973 logical_nulls: map_values_logical_nulls,
1974 };
1975 assert_eq!(list_level, &expected_level);
1976 }
1977
1978 #[test]
1979 fn test_list_of_struct() {
1980 let int_field = Field::new("a", DataType::Int32, true);
1982 let fields = Fields::from([Arc::new(int_field)]);
1983 let item_field = Field::new_list_field(DataType::Struct(fields.clone()), true);
1984 let list_field = Field::new("list", DataType::List(Arc::new(item_field)), true);
1985
1986 let int_builder = Int32Builder::with_capacity(10);
1987 let struct_builder = StructBuilder::new(fields, vec![Box::new(int_builder)]);
1988 let mut list_builder = ListBuilder::new(struct_builder);
1989
1990 let values = list_builder.values();
1994 values
1995 .field_builder::<Int32Builder>(0)
1996 .unwrap()
1997 .append_value(1);
1998 values.append(true);
1999 list_builder.append(true);
2000
2001 list_builder.append(true);
2003
2004 list_builder.append(false);
2006
2007 let values = list_builder.values();
2009 values
2010 .field_builder::<Int32Builder>(0)
2011 .unwrap()
2012 .append_null();
2013 values.append(false);
2014 values
2015 .field_builder::<Int32Builder>(0)
2016 .unwrap()
2017 .append_null();
2018 values.append(false);
2019 list_builder.append(true);
2020
2021 let values = list_builder.values();
2023 values
2024 .field_builder::<Int32Builder>(0)
2025 .unwrap()
2026 .append_null();
2027 values.append(true);
2028 list_builder.append(true);
2029
2030 let values = list_builder.values();
2032 values
2033 .field_builder::<Int32Builder>(0)
2034 .unwrap()
2035 .append_value(2);
2036 values.append(true);
2037 list_builder.append(true);
2038
2039 let array = Arc::new(list_builder.finish());
2040
2041 let values = array.values().as_struct().column(0).clone();
2042 let values_len = values.len();
2043 assert_eq!(values_len, 5);
2044
2045 let schema = Arc::new(Schema::new(vec![list_field]));
2046
2047 let rb = RecordBatch::try_new(schema, vec![array]).unwrap();
2048
2049 let levels = calculate_array_levels(rb.column(0), rb.schema().field(0)).unwrap();
2050 let list_level = &levels[0];
2051
2052 let logical_nulls = values.logical_nulls();
2053 let expected_level = ArrayLevels {
2054 def_levels: LevelData::Materialized(vec![4, 1, 0, 2, 2, 3, 4]),
2055 rep_levels: LevelData::Materialized(vec![0, 0, 0, 0, 1, 0, 0]),
2056 non_null_indices: vec![0, 4],
2057 max_def_level: 4,
2058 max_rep_level: 1,
2059 array: values,
2060 logical_nulls,
2061 };
2062
2063 assert_eq!(list_level, &expected_level);
2064 }
2065
2066 #[test]
2067 fn test_struct_mask_list() {
2068 let inner = ListArray::from_iter_primitive::<Int32Type, _, _>(vec![
2070 Some(vec![Some(1), Some(2)]),
2071 Some(vec![None]),
2072 Some(vec![]),
2073 Some(vec![Some(3), None]), Some(vec![Some(4), Some(5)]),
2075 None, None,
2077 ]);
2078 let values = inner.values().clone();
2079
2080 assert_eq!(inner.values().len(), 7);
2082
2083 let field = Arc::new(Field::new("list", inner.data_type().clone(), true));
2084 let array = Arc::new(inner) as ArrayRef;
2085 let nulls = Buffer::from([0b01010111]);
2086 let struct_a = StructArray::from((vec![(field, array)], nulls));
2087
2088 let field = Field::new("struct", struct_a.data_type().clone(), true);
2089 let array = Arc::new(struct_a) as ArrayRef;
2090 let levels = calculate_array_levels(&array, &field).unwrap();
2091
2092 assert_eq!(levels.len(), 1);
2093
2094 let logical_nulls = values.logical_nulls();
2095 let expected_level = ArrayLevels {
2096 def_levels: LevelData::Materialized(vec![4, 4, 3, 2, 0, 4, 4, 0, 1]),
2097 rep_levels: LevelData::Materialized(vec![0, 1, 0, 0, 0, 0, 1, 0, 0]),
2098 non_null_indices: vec![0, 1, 5, 6],
2099 max_def_level: 4,
2100 max_rep_level: 1,
2101 array: values,
2102 logical_nulls,
2103 };
2104
2105 assert_eq!(&levels[0], &expected_level);
2106 }
2107
2108 #[test]
2109 fn test_list_mask_struct() {
2110 let a1 = ListArray::from_iter_primitive::<Int32Type, _, _>(vec![
2114 Some(vec![None]), Some(vec![]), Some(vec![Some(3), None]),
2117 Some(vec![Some(4), Some(5), None, Some(6)]), None,
2119 None,
2120 ]);
2121 let a1_values = a1.values().clone();
2122 let a1 = Arc::new(a1) as ArrayRef;
2123
2124 let a2 = Arc::new(Int32Array::from_iter(vec![
2125 Some(1), Some(2), None,
2128 Some(4), Some(5),
2130 None,
2131 ])) as ArrayRef;
2132 let a2_values = a2.clone();
2133
2134 let field_a1 = Arc::new(Field::new("list", a1.data_type().clone(), true));
2135 let field_a2 = Arc::new(Field::new("integers", a2.data_type().clone(), true));
2136
2137 let nulls = Buffer::from([0b00110111]);
2138 let struct_a = Arc::new(StructArray::from((
2139 vec![(field_a1, a1), (field_a2, a2)],
2140 nulls,
2141 ))) as ArrayRef;
2142
2143 let offsets = Buffer::from_iter([0_i32, 0, 2, 2, 3, 5, 5]);
2144 let nulls = Buffer::from([0b00111100]);
2145
2146 let list_type = DataType::List(Arc::new(Field::new(
2147 "struct",
2148 struct_a.data_type().clone(),
2149 true,
2150 )));
2151
2152 let data = ArrayDataBuilder::new(list_type.clone())
2153 .len(6)
2154 .null_bit_buffer(Some(nulls))
2155 .add_buffer(offsets)
2156 .add_child_data(struct_a.into_data())
2157 .build()
2158 .unwrap();
2159
2160 let list = make_array(data);
2161 let list_field = Field::new("col", list_type, true);
2162
2163 let expected = vec![
2164 r#""#.to_string(),
2165 r#""#.to_string(),
2166 r#"[]"#.to_string(),
2167 r#"[{list: [3, ], integers: }]"#.to_string(),
2168 r#"[, {list: , integers: 5}]"#.to_string(),
2169 r#"[]"#.to_string(),
2170 ];
2171
2172 let actual: Vec<_> = (0..6)
2173 .map(|x| array_value_to_string(&list, x).unwrap())
2174 .collect();
2175 assert_eq!(actual, expected);
2176
2177 let levels = calculate_array_levels(&list, &list_field).unwrap();
2178
2179 assert_eq!(levels.len(), 2);
2180
2181 let a1_logical_nulls = a1_values.logical_nulls();
2182 let expected_level = ArrayLevels {
2183 def_levels: LevelData::Materialized(vec![0, 0, 1, 6, 5, 2, 3, 1]),
2184 rep_levels: LevelData::Materialized(vec![0, 0, 0, 0, 2, 0, 1, 0]),
2185 non_null_indices: vec![1],
2186 max_def_level: 6,
2187 max_rep_level: 2,
2188 array: a1_values,
2189 logical_nulls: a1_logical_nulls,
2190 };
2191
2192 assert_eq!(&levels[0], &expected_level);
2193
2194 let a2_logical_nulls = a2_values.logical_nulls();
2195 let expected_level = ArrayLevels {
2196 def_levels: LevelData::Materialized(vec![0, 0, 1, 3, 2, 4, 1]),
2197 rep_levels: LevelData::Materialized(vec![0, 0, 0, 0, 0, 1, 0]),
2198 non_null_indices: vec![4],
2199 max_def_level: 4,
2200 max_rep_level: 1,
2201 array: a2_values,
2202 logical_nulls: a2_logical_nulls,
2203 };
2204
2205 assert_eq!(&levels[1], &expected_level);
2206 }
2207
2208 #[test]
2209 fn test_fixed_size_list() {
2210 let mut builder = FixedSizeListBuilder::new(Int32Builder::new(), 2);
2212 builder.values().append_slice(&[1, 2]);
2213 builder.append(true);
2214 builder.values().append_slice(&[3, 4]);
2215 builder.append(false);
2216 builder.values().append_slice(&[5, 6]);
2217 builder.append(false);
2218 builder.values().append_slice(&[7, 8]);
2219 builder.append(true);
2220 builder.values().append_slice(&[9, 10]);
2221 builder.append(false);
2222 let a = builder.finish();
2223 let values = a.values().clone();
2224
2225 let item_field = Field::new_list_field(a.data_type().clone(), true);
2226 let mut builder = levels(&item_field, a);
2227 builder.write(1..4);
2228 let levels = builder.finish();
2229
2230 assert_eq!(levels.len(), 1);
2231
2232 let list_level = &levels[0];
2233
2234 let logical_nulls = values.logical_nulls();
2235 let expected_level = ArrayLevels {
2236 def_levels: LevelData::Materialized(vec![0, 0, 3, 3]),
2237 rep_levels: LevelData::Materialized(vec![0, 0, 0, 1]),
2238 non_null_indices: vec![6, 7],
2239 max_def_level: 3,
2240 max_rep_level: 1,
2241 array: values,
2242 logical_nulls,
2243 };
2244 assert_eq!(list_level, &expected_level);
2245 }
2246
2247 #[test]
2248 fn test_fixed_size_list_of_struct() {
2249 let field_a = Field::new("a", DataType::Int32, true);
2251 let field_b = Field::new("b", DataType::Int64, false);
2252 let fields = Fields::from([Arc::new(field_a), Arc::new(field_b)]);
2253 let item_field = Field::new_list_field(DataType::Struct(fields.clone()), true);
2254 let list_field = Field::new(
2255 "list",
2256 DataType::FixedSizeList(Arc::new(item_field), 2),
2257 true,
2258 );
2259
2260 let builder_a = Int32Builder::with_capacity(10);
2261 let builder_b = Int64Builder::with_capacity(10);
2262 let struct_builder =
2263 StructBuilder::new(fields, vec![Box::new(builder_a), Box::new(builder_b)]);
2264 let mut list_builder = FixedSizeListBuilder::new(struct_builder, 2);
2265
2266 let values = list_builder.values();
2275 values
2277 .field_builder::<Int32Builder>(0)
2278 .unwrap()
2279 .append_value(1);
2280 values
2281 .field_builder::<Int64Builder>(1)
2282 .unwrap()
2283 .append_value(2);
2284 values.append(true);
2285 values
2287 .field_builder::<Int32Builder>(0)
2288 .unwrap()
2289 .append_null();
2290 values
2291 .field_builder::<Int64Builder>(1)
2292 .unwrap()
2293 .append_value(0);
2294 values.append(false);
2295 list_builder.append(true);
2296
2297 let values = list_builder.values();
2299 values
2301 .field_builder::<Int32Builder>(0)
2302 .unwrap()
2303 .append_null();
2304 values
2305 .field_builder::<Int64Builder>(1)
2306 .unwrap()
2307 .append_value(0);
2308 values.append(false);
2309 values
2311 .field_builder::<Int32Builder>(0)
2312 .unwrap()
2313 .append_null();
2314 values
2315 .field_builder::<Int64Builder>(1)
2316 .unwrap()
2317 .append_value(0);
2318 values.append(false);
2319 list_builder.append(false);
2320
2321 let values = list_builder.values();
2323 values
2325 .field_builder::<Int32Builder>(0)
2326 .unwrap()
2327 .append_null();
2328 values
2329 .field_builder::<Int64Builder>(1)
2330 .unwrap()
2331 .append_value(0);
2332 values.append(false);
2333 values
2335 .field_builder::<Int32Builder>(0)
2336 .unwrap()
2337 .append_null();
2338 values
2339 .field_builder::<Int64Builder>(1)
2340 .unwrap()
2341 .append_value(0);
2342 values.append(false);
2343 list_builder.append(true);
2344
2345 let values = list_builder.values();
2347 values
2349 .field_builder::<Int32Builder>(0)
2350 .unwrap()
2351 .append_null();
2352 values
2353 .field_builder::<Int64Builder>(1)
2354 .unwrap()
2355 .append_value(3);
2356 values.append(true);
2357 values
2359 .field_builder::<Int32Builder>(0)
2360 .unwrap()
2361 .append_value(2);
2362 values
2363 .field_builder::<Int64Builder>(1)
2364 .unwrap()
2365 .append_value(4);
2366 values.append(true);
2367 list_builder.append(true);
2368
2369 let array = Arc::new(list_builder.finish());
2370
2371 assert_eq!(array.values().len(), 8);
2372 assert_eq!(array.len(), 4);
2373
2374 let struct_values = array.values().as_struct();
2375 let values_a = struct_values.column(0).clone();
2376 let values_b = struct_values.column(1).clone();
2377
2378 let schema = Arc::new(Schema::new(vec![list_field]));
2379 let rb = RecordBatch::try_new(schema, vec![array]).unwrap();
2380
2381 let levels = calculate_array_levels(rb.column(0), rb.schema().field(0)).unwrap();
2382 let a_levels = &levels[0];
2383 let b_levels = &levels[1];
2384
2385 let values_a_logical_nulls = values_a.logical_nulls();
2387 let expected_a = ArrayLevels {
2388 def_levels: LevelData::Materialized(vec![4, 2, 0, 2, 2, 3, 4]),
2389 rep_levels: LevelData::Materialized(vec![0, 1, 0, 0, 1, 0, 1]),
2390 non_null_indices: vec![0, 7],
2391 max_def_level: 4,
2392 max_rep_level: 1,
2393 array: values_a,
2394 logical_nulls: values_a_logical_nulls,
2395 };
2396 let values_b_logical_nulls = values_b.logical_nulls();
2398 let expected_b = ArrayLevels {
2399 def_levels: LevelData::Materialized(vec![3, 2, 0, 2, 2, 3, 3]),
2400 rep_levels: LevelData::Materialized(vec![0, 1, 0, 0, 1, 0, 1]),
2401 non_null_indices: vec![0, 6, 7],
2402 max_def_level: 3,
2403 max_rep_level: 1,
2404 array: values_b,
2405 logical_nulls: values_b_logical_nulls,
2406 };
2407
2408 assert_eq!(a_levels, &expected_a);
2409 assert_eq!(b_levels, &expected_b);
2410 }
2411
2412 #[test]
2413 fn test_fixed_size_list_empty() {
2414 let mut builder = FixedSizeListBuilder::new(Int32Builder::new(), 0);
2415 builder.append(true);
2416 builder.append(false);
2417 builder.append(true);
2418 let array = builder.finish();
2419 let values = array.values().clone();
2420
2421 let item_field = Field::new_list_field(array.data_type().clone(), true);
2422 let mut builder = levels(&item_field, array);
2423 builder.write(0..3);
2424 let levels = builder.finish();
2425
2426 assert_eq!(levels.len(), 1);
2427
2428 let list_level = &levels[0];
2429
2430 let logical_nulls = values.logical_nulls();
2431 let expected_level = ArrayLevels {
2432 def_levels: LevelData::Materialized(vec![1, 0, 1]),
2433 rep_levels: LevelData::Materialized(vec![0, 0, 0]),
2434 non_null_indices: vec![],
2435 max_def_level: 3,
2436 max_rep_level: 1,
2437 array: values,
2438 logical_nulls,
2439 };
2440 assert_eq!(list_level, &expected_level);
2441 }
2442
2443 #[test]
2444 fn test_fixed_size_list_of_var_lists() {
2445 let mut builder = FixedSizeListBuilder::new(ListBuilder::new(Int32Builder::new()), 2);
2447 builder.values().append_value([Some(1), None, Some(3)]);
2448 builder.values().append_null();
2449 builder.append(true);
2450 builder.values().append_value([Some(4)]);
2451 builder.values().append_value([]);
2452 builder.append(true);
2453 builder.values().append_value([Some(5), Some(6)]);
2454 builder.values().append_value([None, None]);
2455 builder.append(true);
2456 builder.values().append_null();
2457 builder.values().append_null();
2458 builder.append(false);
2459 let a = builder.finish();
2460 let values = a.values().as_list::<i32>().values().clone();
2461
2462 let item_field = Field::new_list_field(a.data_type().clone(), true);
2463 let mut builder = levels(&item_field, a);
2464 builder.write(0..4);
2465 let levels = builder.finish();
2466
2467 let logical_nulls = values.logical_nulls();
2468 let expected_level = ArrayLevels {
2469 def_levels: LevelData::Materialized(vec![5, 4, 5, 2, 5, 3, 5, 5, 4, 4, 0]),
2470 rep_levels: LevelData::Materialized(vec![0, 2, 2, 1, 0, 1, 0, 2, 1, 2, 0]),
2471 non_null_indices: vec![0, 2, 3, 4, 5],
2472 max_def_level: 5,
2473 max_rep_level: 2,
2474 array: values,
2475 logical_nulls,
2476 };
2477
2478 assert_eq!(levels[0], expected_level);
2479 }
2480
2481 #[test]
2482 fn test_null_dictionary_values() {
2483 let values = Int32Array::new(
2484 vec![1, 2, 3, 4].into(),
2485 Some(NullBuffer::from(vec![true, false, true, true])),
2486 );
2487 let keys = Int32Array::new(
2488 vec![1, 54, 2, 0].into(),
2489 Some(NullBuffer::from(vec![true, false, true, true])),
2490 );
2491 let dict = DictionaryArray::new(keys, Arc::new(values));
2493
2494 let item_field = Field::new_list_field(dict.data_type().clone(), true);
2495
2496 let mut builder = levels(&item_field, dict.clone());
2497 builder.write(0..4);
2498 let levels = builder.finish();
2499
2500 let logical_nulls = dict.logical_nulls();
2501 let expected_level = ArrayLevels {
2502 def_levels: LevelData::Materialized(vec![0, 0, 1, 1]),
2503 rep_levels: LevelData::Absent,
2504 non_null_indices: vec![2, 3],
2505 max_def_level: 1,
2506 max_rep_level: 0,
2507 array: Arc::new(dict),
2508 logical_nulls,
2509 };
2510 assert_eq!(levels[0], expected_level);
2511 }
2512
2513 #[test]
2514 fn mismatched_types() {
2515 let array = Arc::new(Int32Array::from_iter(0..10)) as ArrayRef;
2516 let field = Field::new_list_field(DataType::Float64, false);
2517
2518 let err = LevelInfoBuilder::try_new(&field, Default::default(), &array)
2519 .unwrap_err()
2520 .to_string();
2521
2522 assert_eq!(
2523 err,
2524 "Arrow: Incompatible type. Field 'item' has type Float64, array has type Int32",
2525 );
2526 }
2527
2528 fn levels<T: Array + 'static>(field: &Field, array: T) -> LevelInfoBuilder {
2529 let v = Arc::new(array) as ArrayRef;
2530 LevelInfoBuilder::try_new(field, Default::default(), &v).unwrap()
2531 }
2532
2533 #[test]
2534 fn test_slice_for_chunk_flat() {
2535 let array: ArrayRef = Arc::new(Int32Array::from(vec![1, 2, 3, 4, 5, 6]));
2540 let logical_nulls = array.logical_nulls();
2541 let levels = ArrayLevels {
2542 def_levels: LevelData::Absent,
2543 rep_levels: LevelData::Absent,
2544 non_null_indices: vec![0, 1, 2, 3, 4, 5],
2545 max_def_level: 0,
2546 max_rep_level: 0,
2547 array,
2548 logical_nulls,
2549 };
2550 let sliced = levels.slice_for_chunk(&CdcChunk {
2551 level_offset: 0,
2552 num_levels: 0,
2553 value_offset: 2,
2554 num_values: 3,
2555 });
2556 assert!(matches!(sliced.def_levels, LevelData::Absent));
2557 assert!(matches!(sliced.rep_levels, LevelData::Absent));
2558 assert_eq!(sliced.non_null_indices, vec![0, 1, 2]);
2559 assert_eq!(sliced.array.len(), 3);
2560
2561 let array: ArrayRef = Arc::new(Int32Array::from(vec![
2567 Some(1),
2568 None,
2569 Some(3),
2570 None,
2571 Some(5),
2572 Some(6),
2573 ]));
2574 let logical_nulls = array.logical_nulls();
2575 let levels = ArrayLevels {
2576 def_levels: LevelData::Materialized(vec![1, 0, 1, 0, 1, 1]),
2577 rep_levels: LevelData::Absent,
2578 non_null_indices: vec![0, 2, 4, 5],
2579 max_def_level: 1,
2580 max_rep_level: 0,
2581 array,
2582 logical_nulls,
2583 };
2584 let sliced = levels.slice_for_chunk(&CdcChunk {
2585 level_offset: 1,
2586 num_levels: 3,
2587 value_offset: 1,
2588 num_values: 1,
2589 });
2590 assert_eq!(sliced.def_levels, LevelData::Materialized(vec![0, 1, 0]));
2591 assert!(matches!(sliced.rep_levels, LevelData::Absent));
2592 assert_eq!(sliced.non_null_indices, vec![0]); assert_eq!(sliced.array.len(), 1);
2594 }
2595
2596 #[test]
2597 fn test_slice_for_chunk_nested_with_nulls() {
2598 let array: ArrayRef = Arc::new(Int32Array::from(vec![
2617 Some(1), None, None, Some(2), None, None, None, None, Some(4), Some(5), ]));
2628 let logical_nulls = array.logical_nulls();
2629 let levels = ArrayLevels {
2630 def_levels: LevelData::Materialized(vec![3, 0, 3, 2, 0, 3, 3]),
2631 rep_levels: LevelData::Materialized(vec![0, 0, 0, 1, 0, 0, 1]),
2632 non_null_indices: vec![0, 3, 8, 9],
2633 max_def_level: 3,
2634 max_rep_level: 1,
2635 array,
2636 logical_nulls,
2637 };
2638
2639 let chunk0 = levels.slice_for_chunk(&CdcChunk {
2641 level_offset: 0,
2642 num_levels: 2,
2643 value_offset: 0,
2644 num_values: 1,
2645 });
2646 assert_eq!(chunk0.non_null_indices, vec![0]);
2647 assert_eq!(chunk0.array.len(), 1);
2648
2649 let chunk1 = levels.slice_for_chunk(&CdcChunk {
2651 level_offset: 2,
2652 num_levels: 3,
2653 value_offset: 1,
2654 num_values: 1,
2655 });
2656 assert_eq!(chunk1.non_null_indices, vec![0]);
2657 assert_eq!(chunk1.array.len(), 1);
2658
2659 let chunk2 = levels.slice_for_chunk(&CdcChunk {
2661 level_offset: 5,
2662 num_levels: 2,
2663 value_offset: 2,
2664 num_values: 2,
2665 });
2666 assert_eq!(chunk2.non_null_indices, vec![0, 1]);
2667 assert_eq!(chunk2.array.len(), 2);
2668 }
2669
2670 #[test]
2671 fn test_slice_for_chunk_all_null() {
2672 let array: ArrayRef = Arc::new(Int32Array::from(vec![Some(1), None, None, Some(4)]));
2674 let logical_nulls = array.logical_nulls();
2675 let levels = ArrayLevels {
2676 def_levels: LevelData::Materialized(vec![1, 0, 0, 1]),
2677 rep_levels: LevelData::Absent,
2678 non_null_indices: vec![0, 3],
2679 max_def_level: 1,
2680 max_rep_level: 0,
2681 array,
2682 logical_nulls,
2683 };
2684 let sliced = levels.slice_for_chunk(&CdcChunk {
2686 level_offset: 1,
2687 num_levels: 2,
2688 value_offset: 1,
2689 num_values: 0,
2690 });
2691 assert_eq!(sliced.def_levels, LevelData::Materialized(vec![0, 0]));
2692 assert_eq!(sliced.non_null_indices, Vec::<usize>::new());
2693 assert_eq!(sliced.array.len(), 0);
2694 }
2695
2696 #[test]
2697 fn test_all_null_list() {
2698 let item_field = Arc::new(Field::new_list_field(DataType::Int32, true));
2704 let list = ListArray::new_null(item_field, 4);
2705 let values = list.values().clone();
2706 let field = Field::new("list", list.data_type().clone(), true);
2707 let array = Arc::new(list) as ArrayRef;
2708
2709 let levels = calculate_array_levels(&array, &field).unwrap();
2710 assert_eq!(levels.len(), 1);
2711
2712 let logical_nulls = values.logical_nulls();
2713 let expected = ArrayLevels {
2714 def_levels: LevelData::Uniform { value: 0, count: 4 },
2715 rep_levels: LevelData::Uniform { value: 0, count: 4 },
2716 non_null_indices: vec![],
2717 max_def_level: 3,
2718 max_rep_level: 1,
2719 array: values,
2720 logical_nulls,
2721 };
2722 assert_eq!(&levels[0], &expected);
2723 }
2724
2725 #[test]
2726 fn test_all_null_fixed_size_list() {
2727 let item_field = Arc::new(Field::new_list_field(DataType::Int32, true));
2733 let list = FixedSizeListArray::new_null(item_field, 2, 3);
2734 let values = list.values().clone();
2735 let field = Field::new("list", list.data_type().clone(), true);
2736 let array = Arc::new(list) as ArrayRef;
2737
2738 let levels = calculate_array_levels(&array, &field).unwrap();
2739 assert_eq!(levels.len(), 1);
2740
2741 let logical_nulls = values.logical_nulls();
2742 let expected = ArrayLevels {
2743 def_levels: LevelData::Uniform { value: 0, count: 3 },
2744 rep_levels: LevelData::Uniform { value: 0, count: 3 },
2745 non_null_indices: vec![],
2746 max_def_level: 3,
2747 max_rep_level: 1,
2748 array: values,
2749 logical_nulls,
2750 };
2751 assert_eq!(&levels[0], &expected);
2752 }
2753
2754 #[test]
2755 fn test_all_null_struct() {
2756 let c = Int32Array::from(vec![None::<i32>; 4]);
2763 let leaf = Arc::new(c) as ArrayRef;
2764 let c_field = Arc::new(Field::new("c", DataType::Int32, true));
2765 let a = StructArray::from((vec![(c_field, leaf.clone())], Buffer::from([0b00000000])));
2766 let a_field = Field::new("a", a.data_type().clone(), true);
2767 let a_array = Arc::new(a) as ArrayRef;
2768
2769 let levels = calculate_array_levels(&a_array, &a_field).unwrap();
2770 assert_eq!(levels.len(), 1);
2771
2772 let expected = ArrayLevels {
2773 def_levels: LevelData::Uniform { value: 0, count: 4 },
2774 rep_levels: LevelData::Absent,
2775 non_null_indices: vec![],
2776 max_def_level: 2,
2777 max_rep_level: 0,
2778 array: leaf,
2779 logical_nulls: Some(NullBuffer::new_null(4)),
2780 };
2781 assert_eq!(&levels[0], &expected);
2782 }
2783
2784 #[test]
2785 fn test_all_null_nested_struct() {
2786 let c = Int32Array::from(vec![None::<i32>; 3]);
2792 let leaf = Arc::new(c) as ArrayRef;
2793 let c_field = Arc::new(Field::new("c", DataType::Int32, true));
2794 let b = StructArray::from((vec![(c_field, leaf.clone())], Buffer::from([0b00000000])));
2795 let b_field = Arc::new(Field::new("b", b.data_type().clone(), true));
2796 let a = StructArray::from((
2797 vec![(b_field, Arc::new(b) as ArrayRef)],
2798 Buffer::from([0b00000000]),
2799 ));
2800 let a_field = Field::new("a", a.data_type().clone(), true);
2801 let a_array = Arc::new(a) as ArrayRef;
2802
2803 let levels = calculate_array_levels(&a_array, &a_field).unwrap();
2804 assert_eq!(levels.len(), 1);
2805
2806 let expected = ArrayLevels {
2807 def_levels: LevelData::Uniform { value: 0, count: 3 },
2808 rep_levels: LevelData::Absent,
2809 non_null_indices: vec![],
2810 max_def_level: 3,
2811 max_rep_level: 0,
2812 array: leaf,
2813 logical_nulls: Some(NullBuffer::new_null(3)),
2814 };
2815 assert_eq!(&levels[0], &expected);
2816 }
2817
2818 #[test]
2819 fn test_all_null_struct_multiple_children() {
2820 let c1 = Arc::new(Int32Array::from(vec![None::<i32>; 2])) as ArrayRef;
2826 let c2 = Arc::new(Int32Array::from(vec![None::<i32>; 2])) as ArrayRef;
2827 let c1_field = Arc::new(Field::new("c1", DataType::Int32, true));
2828 let c2_field = Arc::new(Field::new("c2", DataType::Int32, true));
2829 let a = StructArray::from((
2830 vec![(c1_field, c1.clone()), (c2_field, c2.clone())],
2831 Buffer::from([0b00000000]),
2832 ));
2833 let a_field = Field::new("a", a.data_type().clone(), true);
2834 let a_array = Arc::new(a) as ArrayRef;
2835
2836 let levels = calculate_array_levels(&a_array, &a_field).unwrap();
2837 assert_eq!(levels.len(), 2);
2838
2839 for (i, leaf) in [c1, c2].into_iter().enumerate() {
2840 let expected = ArrayLevels {
2841 def_levels: LevelData::Uniform { value: 0, count: 2 },
2842 rep_levels: LevelData::Absent,
2843 non_null_indices: vec![],
2844 max_def_level: 2,
2845 max_rep_level: 0,
2846 array: leaf,
2847 logical_nulls: Some(NullBuffer::new_null(2)),
2848 };
2849 assert_eq!(&levels[i], &expected, "leaf {i} mismatch");
2850 }
2851 }
2852}