1use crate::variant_array::{ShreddedVariantFieldArray, StructArrayBuilder};
21use crate::variant_to_arrow::{
22 ArrayVariantToArrowRowBuilder, PrimitiveVariantToArrowRowBuilder,
23 make_primitive_variant_to_arrow_row_builder,
24};
25use crate::{VariantArray, VariantValueArrayBuilder};
26use arrow::array::{ArrayRef, BinaryViewArray, NullBufferBuilder};
27use arrow::buffer::NullBuffer;
28use arrow::compute::CastOptions;
29use arrow::datatypes::{DataType, Field, FieldRef, Fields, TimeUnit};
30use arrow::error::{ArrowError, Result};
31use indexmap::IndexMap;
32use parquet_variant::{Variant, VariantBuilderExt, VariantPath, VariantPathElement};
33use std::collections::BTreeMap;
34use std::sync::Arc;
35
36pub fn shred_variant(array: &VariantArray, as_type: &DataType) -> Result<VariantArray> {
71 shred_variant_with_options(array, as_type, &CastOptions::default())
72}
73
74pub(crate) fn shred_variant_with_options(
75 array: &VariantArray,
76 as_type: &DataType,
77 cast_options: &CastOptions,
78) -> Result<VariantArray> {
79 if array.typed_value_column().is_some() {
80 return Err(ArrowError::InvalidArgumentError(
81 "Input is already shredded".to_string(),
82 ));
83 }
84
85 let mut builder = make_variant_to_shredded_variant_arrow_row_builder(
86 as_type,
87 cast_options,
88 array.len(),
89 NullValue::TopLevelVariant,
90 )?;
91 for i in 0..array.len() {
92 if array.is_null(i) {
93 builder.append_null()?;
94 } else {
95 builder.append_value(array.value(i))?;
96 }
97 }
98 let (value, typed_value, nulls) = builder.finish()?;
99 Ok(VariantArray::from_parts(
100 array.metadata_column().clone(),
101 Arc::new(value),
102 Some(typed_value),
103 nulls,
104 ))
105}
106
107#[derive(Debug, Clone, Copy, PartialEq, Eq)]
115pub(crate) enum NullValue {
116 TopLevelVariant,
117 ObjectField,
118 ArrayElement,
119}
120
121impl NullValue {
122 fn append_to(
123 self,
124 nulls: &mut NullBufferBuilder,
125 value_builder: &mut VariantValueArrayBuilder,
126 ) {
127 match self {
128 Self::TopLevelVariant => nulls.append_null(),
129 Self::ObjectField | Self::ArrayElement => nulls.append_non_null(),
130 }
131 match self {
132 Self::TopLevelVariant | Self::ObjectField => value_builder.append_null(),
133 Self::ArrayElement => value_builder.append_value(Variant::Null),
134 }
135 }
136}
137
138pub(crate) fn make_variant_to_shredded_variant_arrow_row_builder<'a>(
139 data_type: &'a DataType,
140 cast_options: &'a CastOptions,
141 capacity: usize,
142 null_value: NullValue,
143) -> Result<VariantToShreddedVariantRowBuilder<'a>> {
144 let builder = match data_type {
145 DataType::Struct(fields) => {
146 let typed_value_builder = VariantToShreddedObjectVariantRowBuilder::try_new(
147 fields,
148 cast_options,
149 capacity,
150 null_value,
151 )?;
152 VariantToShreddedVariantRowBuilder::Object(typed_value_builder)
153 }
154 DataType::List(_)
155 | DataType::LargeList(_)
156 | DataType::ListView(_)
157 | DataType::LargeListView(_)
158 | DataType::FixedSizeList(..) => {
159 let typed_value_builder = VariantToShreddedArrayVariantRowBuilder::try_new(
160 data_type,
161 cast_options,
162 capacity,
163 null_value,
164 )?;
165 VariantToShreddedVariantRowBuilder::Array(typed_value_builder)
166 }
167 DataType::Boolean
170 | DataType::Int8
171 | DataType::Int16
172 | DataType::Int32
173 | DataType::Int64
174 | DataType::Float32
175 | DataType::Float64
176 | DataType::Decimal32(..)
177 | DataType::Decimal64(..)
178 | DataType::Decimal128(..)
179 | DataType::Date32
180 | DataType::Time64(TimeUnit::Microsecond)
181 | DataType::Timestamp(TimeUnit::Microsecond | TimeUnit::Nanosecond, _)
182 | DataType::Binary
183 | DataType::BinaryView
184 | DataType::LargeBinary
185 | DataType::Utf8
186 | DataType::Utf8View
187 | DataType::LargeUtf8
188 | DataType::FixedSizeBinary(16) => {
190 let builder =
191 make_primitive_variant_to_arrow_row_builder(data_type, cast_options, capacity)?;
192 let typed_value_builder =
193 VariantToShreddedPrimitiveVariantRowBuilder::new(builder, capacity, null_value);
194 VariantToShreddedVariantRowBuilder::Primitive(typed_value_builder)
195 }
196 DataType::FixedSizeBinary(_) => {
197 return Err(ArrowError::InvalidArgumentError(format!("{data_type} is not a valid variant shredding type. Only FixedSizeBinary(16) for UUID is supported.")))
198 }
199 _ => {
200 return Err(ArrowError::InvalidArgumentError(format!("{data_type} is not a valid variant shredding type")))
201 }
202 };
203 Ok(builder)
204}
205
206pub(crate) enum VariantToShreddedVariantRowBuilder<'a> {
207 Primitive(VariantToShreddedPrimitiveVariantRowBuilder<'a>),
208 Array(VariantToShreddedArrayVariantRowBuilder<'a>),
209 Object(VariantToShreddedObjectVariantRowBuilder<'a>),
210}
211
212impl<'a> VariantToShreddedVariantRowBuilder<'a> {
213 pub fn append_null(&mut self) -> Result<()> {
214 use VariantToShreddedVariantRowBuilder::*;
215 match self {
216 Primitive(b) => b.append_null(),
217 Array(b) => b.append_null(),
218 Object(b) => b.append_null(),
219 }
220 }
221
222 pub fn append_value(&mut self, value: Variant<'_, '_>) -> Result<bool> {
223 use VariantToShreddedVariantRowBuilder::*;
224 match self {
225 Primitive(b) => b.append_value(value),
226 Array(b) => b.append_value(value),
227 Object(b) => b.append_value(value),
228 }
229 }
230
231 pub fn finish(self) -> Result<(BinaryViewArray, ArrayRef, Option<NullBuffer>)> {
232 use VariantToShreddedVariantRowBuilder::*;
233 match self {
234 Primitive(b) => b.finish(),
235 Array(b) => b.finish(),
236 Object(b) => b.finish(),
237 }
238 }
239}
240
241pub(crate) struct VariantToShreddedPrimitiveVariantRowBuilder<'a> {
243 value_builder: VariantValueArrayBuilder,
244 typed_value_builder: PrimitiveVariantToArrowRowBuilder<'a>,
245 nulls: NullBufferBuilder,
246 null_value: NullValue,
247}
248
249impl<'a> VariantToShreddedPrimitiveVariantRowBuilder<'a> {
250 pub(crate) fn new(
251 typed_value_builder: PrimitiveVariantToArrowRowBuilder<'a>,
252 capacity: usize,
253 null_value: NullValue,
254 ) -> Self {
255 Self {
256 value_builder: VariantValueArrayBuilder::new(capacity),
257 typed_value_builder,
258 nulls: NullBufferBuilder::new(capacity),
259 null_value,
260 }
261 }
262
263 fn append_null(&mut self) -> Result<()> {
264 self.null_value
265 .append_to(&mut self.nulls, &mut self.value_builder);
266 self.typed_value_builder.append_null()
267 }
268
269 fn append_value(&mut self, value: Variant<'_, '_>) -> Result<bool> {
270 self.nulls.append_non_null();
271 if self.typed_value_builder.append_value(&value)? {
272 self.value_builder.append_null();
273 } else {
274 self.value_builder.append_value(value);
275 }
276 Ok(true)
277 }
278
279 fn finish(mut self) -> Result<(BinaryViewArray, ArrayRef, Option<NullBuffer>)> {
280 Ok((
281 self.value_builder.build()?,
282 self.typed_value_builder.finish()?,
283 self.nulls.finish(),
284 ))
285 }
286}
287
288pub(crate) struct VariantToShreddedArrayVariantRowBuilder<'a> {
289 value_builder: VariantValueArrayBuilder,
290 typed_value_builder: ArrayVariantToArrowRowBuilder<'a>,
291 nulls: NullBufferBuilder,
292 null_value: NullValue,
293}
294
295impl<'a> VariantToShreddedArrayVariantRowBuilder<'a> {
296 fn try_new(
297 data_type: &'a DataType,
298 cast_options: &'a CastOptions,
299 capacity: usize,
300 null_value: NullValue,
301 ) -> Result<Self> {
302 Ok(Self {
303 value_builder: VariantValueArrayBuilder::new(capacity),
304 typed_value_builder: ArrayVariantToArrowRowBuilder::try_new(
305 data_type,
306 cast_options,
307 capacity,
308 true,
309 )?,
310 nulls: NullBufferBuilder::new(capacity),
311 null_value,
312 })
313 }
314
315 fn append_null(&mut self) -> Result<()> {
316 self.null_value
317 .append_to(&mut self.nulls, &mut self.value_builder);
318 self.typed_value_builder.append_null()?;
319 Ok(())
320 }
321
322 fn append_value(&mut self, variant: Variant<'_, '_>) -> Result<bool> {
323 match variant {
326 Variant::List(list) => {
327 self.nulls.append_non_null();
328 self.value_builder.append_null();
329
330 self.typed_value_builder
332 .append_value(&Variant::List(list))?;
333 Ok(true)
334 }
335 other => {
336 self.nulls.append_non_null();
337 self.value_builder.append_value(other);
338 self.typed_value_builder.append_null()?;
339 Ok(false)
340 }
341 }
342 }
343
344 fn finish(mut self) -> Result<(BinaryViewArray, ArrayRef, Option<NullBuffer>)> {
345 Ok((
346 self.value_builder.build()?,
347 self.typed_value_builder.finish()?,
348 self.nulls.finish(),
349 ))
350 }
351}
352
353pub(crate) struct VariantToShreddedObjectVariantRowBuilder<'a> {
354 value_builder: VariantValueArrayBuilder,
355 typed_value_builders: IndexMap<&'a str, VariantToShreddedVariantRowBuilder<'a>>,
356 typed_value_nulls: NullBufferBuilder,
357 nulls: NullBufferBuilder,
358 null_value: NullValue,
359}
360
361impl<'a> VariantToShreddedObjectVariantRowBuilder<'a> {
362 fn try_new(
363 fields: &'a Fields,
364 cast_options: &'a CastOptions,
365 capacity: usize,
366 null_value: NullValue,
367 ) -> Result<Self> {
368 let typed_value_builders = fields.iter().map(|field| {
369 let builder = make_variant_to_shredded_variant_arrow_row_builder(
370 field.data_type(),
371 cast_options,
372 capacity,
373 NullValue::ObjectField,
374 )?;
375 Ok((field.name().as_str(), builder))
376 });
377 Ok(Self {
378 value_builder: VariantValueArrayBuilder::new(capacity),
379 typed_value_builders: typed_value_builders.collect::<Result<_>>()?,
380 typed_value_nulls: NullBufferBuilder::new(capacity),
381 nulls: NullBufferBuilder::new(capacity),
382 null_value,
383 })
384 }
385
386 fn append_null(&mut self) -> Result<()> {
387 self.null_value
388 .append_to(&mut self.nulls, &mut self.value_builder);
389 self.typed_value_nulls.append_null();
390 for (_, typed_value_builder) in &mut self.typed_value_builders {
391 typed_value_builder.append_null()?;
392 }
393 Ok(())
394 }
395
396 fn append_value(&mut self, value: Variant<'_, '_>) -> Result<bool> {
397 let Variant::Object(ref obj) = value else {
398 self.nulls.append_non_null();
400 self.value_builder.append_value(value);
401 self.typed_value_nulls.append_null();
402 for (_, typed_value_builder) in &mut self.typed_value_builders {
403 typed_value_builder.append_null()?;
404 }
405 return Ok(false);
406 };
407
408 let mut builder = self.value_builder.builder_ext(value.metadata());
410 let mut object_builder = builder.try_new_object()?;
411 let mut seen = std::collections::HashSet::new();
412 let mut partially_shredded = false;
413 for (field_name, value) in obj.iter() {
414 match self.typed_value_builders.get_mut(field_name) {
415 Some(typed_value_builder) => {
416 typed_value_builder.append_value(value)?;
417 seen.insert(field_name);
418 }
419 None => {
420 object_builder.insert_bytes(field_name, value);
421 partially_shredded = true;
422 }
423 }
424 }
425
426 for (field_name, typed_value_builder) in &mut self.typed_value_builders {
428 if !seen.contains(field_name) {
429 typed_value_builder.append_null()?;
430 }
431 }
432
433 if partially_shredded {
435 object_builder.finish();
436 } else {
437 drop(object_builder);
438 self.value_builder.append_null();
439 }
440
441 self.typed_value_nulls.append_non_null();
442 self.nulls.append_non_null();
443 Ok(true)
444 }
445
446 fn finish(mut self) -> Result<(BinaryViewArray, ArrayRef, Option<NullBuffer>)> {
447 let mut builder = StructArrayBuilder::new();
448 for (field_name, typed_value_builder) in self.typed_value_builders {
449 let (value, typed_value, nulls) = typed_value_builder.finish()?;
450 let array =
451 ShreddedVariantFieldArray::from_parts(Arc::new(value), Some(typed_value), nulls);
452 builder = builder.with_field(field_name, ArrayRef::from(array), false);
453 }
454 if let Some(nulls) = self.typed_value_nulls.finish() {
455 builder = builder.with_nulls(nulls);
456 }
457 Ok((
458 self.value_builder.build()?,
459 Arc::new(builder.build()),
460 self.nulls.finish(),
461 ))
462 }
463}
464
465#[derive(Clone)]
467pub struct ShreddingField {
468 data_type: DataType,
469 nullable: bool,
470}
471
472impl ShreddingField {
473 fn new(data_type: DataType, nullable: bool) -> Self {
474 Self {
475 data_type,
476 nullable,
477 }
478 }
479
480 fn null() -> Self {
481 Self::new(DataType::Null, true)
482 }
483}
484
485pub trait IntoShreddingField {
487 fn into_shredding_field(self) -> ShreddingField;
488}
489
490impl IntoShreddingField for FieldRef {
491 fn into_shredding_field(self) -> ShreddingField {
492 ShreddingField::new(self.data_type().clone(), self.is_nullable())
493 }
494}
495
496impl IntoShreddingField for &DataType {
497 fn into_shredding_field(self) -> ShreddingField {
498 ShreddingField::new(self.clone(), true)
499 }
500}
501
502impl IntoShreddingField for DataType {
503 fn into_shredding_field(self) -> ShreddingField {
504 ShreddingField::new(self, true)
505 }
506}
507
508impl IntoShreddingField for (&DataType, bool) {
509 fn into_shredding_field(self) -> ShreddingField {
510 ShreddingField::new(self.0.clone(), self.1)
511 }
512}
513
514impl IntoShreddingField for (DataType, bool) {
515 fn into_shredding_field(self) -> ShreddingField {
516 ShreddingField::new(self.0, self.1)
517 }
518}
519
520#[derive(Default, Clone)]
561pub struct ShreddedSchemaBuilder {
562 root: VariantSchemaNode,
563}
564
565impl ShreddedSchemaBuilder {
566 pub fn new() -> Self {
568 Self::default()
569 }
570
571 pub fn with_path<'a, P, F>(mut self, path: P, field: F) -> Result<Self>
583 where
584 P: TryInto<VariantPath<'a>>,
585 P::Error: std::fmt::Debug,
586 F: IntoShreddingField,
587 {
588 let path: VariantPath<'a> = path
589 .try_into()
590 .map_err(|e| ArrowError::InvalidArgumentError(format!("{:?}", e)))?;
591 self.root.insert_path(&path, field.into_shredding_field());
592 Ok(self)
593 }
594
595 pub fn build(self) -> DataType {
597 let shredding_type = self.root.to_shredding_type();
598 match shredding_type {
599 Some(shredding_type) => shredding_type,
600 None => DataType::Null,
601 }
602 }
603}
604
605#[derive(Clone)]
607enum VariantSchemaNode {
608 Leaf(ShreddingField),
610 Struct(BTreeMap<String, VariantSchemaNode>),
612}
613
614impl Default for VariantSchemaNode {
615 fn default() -> Self {
616 Self::Leaf(ShreddingField::null())
617 }
618}
619
620impl VariantSchemaNode {
621 fn insert_path(&mut self, path: &VariantPath<'_>, field: ShreddingField) {
623 self.insert_path_elements(path, field);
624 }
625
626 fn insert_path_elements(&mut self, segments: &[VariantPathElement<'_>], field: ShreddingField) {
627 let Some((head, tail)) = segments.split_first() else {
628 *self = Self::Leaf(field);
629 return;
630 };
631
632 match head {
633 VariantPathElement::Field { name } => {
634 let children = match self {
636 Self::Struct(children) => children,
637 _ => {
638 *self = Self::Struct(BTreeMap::new());
639 match self {
640 Self::Struct(children) => children,
641 _ => unreachable!(),
642 }
643 }
644 };
645
646 children
647 .entry(name.to_string())
648 .or_default()
649 .insert_path_elements(tail, field);
650 }
651 VariantPathElement::Index { .. } => {
652 unreachable!("List paths are not supported yet");
654 }
655 }
656 }
657
658 fn to_shredding_type(&self) -> Option<DataType> {
662 match self {
663 Self::Leaf(field) => Some(field.data_type.clone()),
664 Self::Struct(children) => {
665 let child_fields: Vec<_> = children
666 .iter()
667 .filter_map(|(name, child)| child.to_shredding_field(name))
668 .collect();
669 if child_fields.is_empty() {
670 None
671 } else {
672 Some(DataType::Struct(Fields::from(child_fields)))
673 }
674 }
675 }
676 }
677
678 fn to_shredding_field(&self, name: &str) -> Option<FieldRef> {
679 match self {
680 Self::Leaf(field) => Some(Arc::new(Field::new(
681 name,
682 field.data_type.clone(),
683 field.nullable,
684 ))),
685 Self::Struct(_) => self
686 .to_shredding_type()
687 .map(|data_type| Arc::new(Field::new(name, data_type, true))),
688 }
689 }
690}
691
692#[cfg(test)]
693mod tests {
694 use super::*;
695 use crate::VariantArrayBuilder;
696 use crate::variant_array::{all_null_value_column, binary_array_value, variant_from_arrays_at};
697 use arrow::array::{
698 Array, BinaryViewArray, FixedSizeBinaryArray, FixedSizeListArray, Float64Array,
699 GenericListArray, GenericListViewArray, Int64Array, LargeBinaryArray, LargeStringArray,
700 ListArray, ListLikeArray, OffsetSizeTrait, PrimitiveArray, StringArray, StructArray,
701 };
702 use arrow::datatypes::{
703 ArrowPrimitiveType, DataType, Field, Fields, Int64Type, TimeUnit, UnionFields, UnionMode,
704 };
705 use parquet_variant::{
706 BuilderSpecificState, EMPTY_VARIANT_METADATA_BYTES, ObjectBuilder, ReadOnlyMetadataBuilder,
707 Variant, VariantBuilder, VariantPath, VariantPathElement,
708 };
709 use std::sync::Arc;
710 use uuid::Uuid;
711
712 const NULL_VALUES: [NullValue; 3] = [
713 NullValue::TopLevelVariant,
714 NullValue::ObjectField,
715 NullValue::ArrayElement,
716 ];
717
718 #[derive(Clone)]
719 enum VariantValue<'a> {
720 Value(Variant<'a, 'a>),
721 List(Vec<VariantValue<'a>>),
722 Object(Vec<(&'a str, VariantValue<'a>)>),
723 Null,
724 }
725
726 impl<'a, T> From<T> for VariantValue<'a>
727 where
728 T: Into<Variant<'a, 'a>>,
729 {
730 fn from(value: T) -> Self {
731 Self::Value(value.into())
732 }
733 }
734
735 #[derive(Clone)]
736 enum VariantRow<'a> {
737 Value(VariantValue<'a>),
738 List(Vec<VariantValue<'a>>),
739 Object(Vec<(&'a str, VariantValue<'a>)>),
740 Null,
741 }
742
743 fn build_variant_array(rows: Vec<VariantRow<'static>>) -> VariantArray {
744 let mut builder = VariantArrayBuilder::new(rows.len());
745
746 fn append_variant_value<B: VariantBuilderExt>(builder: &mut B, value: VariantValue) {
747 match value {
748 VariantValue::Value(v) => builder.append_value(v),
749 VariantValue::List(values) => {
750 let mut list = builder.new_list();
751 for v in values {
752 append_variant_value(&mut list, v);
753 }
754 list.finish();
755 }
756 VariantValue::Object(fields) => {
757 let mut object = builder.new_object();
758 for (name, value) in fields {
759 append_variant_field(&mut object, name, value);
760 }
761 object.finish();
762 }
763 VariantValue::Null => builder.append_null(),
764 }
765 }
766
767 fn append_variant_field<'a, S: BuilderSpecificState>(
768 object: &mut ObjectBuilder<'_, S>,
769 name: &'a str,
770 value: VariantValue<'a>,
771 ) {
772 match value {
773 VariantValue::Value(v) => {
774 object.insert(name, v);
775 }
776 VariantValue::List(values) => {
777 let mut list = object.new_list(name);
778 for v in values {
779 append_variant_value(&mut list, v);
780 }
781 list.finish();
782 }
783 VariantValue::Object(fields) => {
784 let mut nested = object.new_object(name);
785 for (field_name, v) in fields {
786 append_variant_field(&mut nested, field_name, v);
787 }
788 nested.finish();
789 }
790 VariantValue::Null => {
791 object.insert(name, Variant::Null);
792 }
793 }
794 }
795
796 rows.into_iter().for_each(|row| match row {
797 VariantRow::Value(value) => append_variant_value(&mut builder, value),
798 VariantRow::List(values) => {
799 let mut list = builder.new_list();
800 for value in values {
801 append_variant_value(&mut list, value);
802 }
803 list.finish();
804 }
805 VariantRow::Object(fields) => {
806 let mut object = builder.new_object();
807 for (name, value) in fields {
808 append_variant_field(&mut object, name, value);
809 }
810 object.finish();
811 }
812 VariantRow::Null => builder.append_null(),
813 });
814 builder.build()
815 }
816
817 trait TestListLikeArray: ListLikeArray {
818 type OffsetSize: OffsetSizeTrait;
819 fn value_offsets(&self) -> Option<&[Self::OffsetSize]>;
820 fn value_size(&self, index: usize) -> Self::OffsetSize;
821 }
822
823 impl<O: OffsetSizeTrait> TestListLikeArray for GenericListArray<O> {
824 type OffsetSize = O;
825
826 fn value_offsets(&self) -> Option<&[Self::OffsetSize]> {
827 Some(GenericListArray::value_offsets(self))
828 }
829
830 fn value_size(&self, index: usize) -> Self::OffsetSize {
831 GenericListArray::value_length(self, index)
832 }
833 }
834
835 impl<O: OffsetSizeTrait> TestListLikeArray for GenericListViewArray<O> {
836 type OffsetSize = O;
837
838 fn value_offsets(&self) -> Option<&[Self::OffsetSize]> {
839 Some(GenericListViewArray::value_offsets(self))
840 }
841
842 fn value_size(&self, index: usize) -> Self::OffsetSize {
843 GenericListViewArray::value_size(self, index)
844 }
845 }
846
847 fn downcast_list_like_array<O: OffsetSizeTrait>(
848 array: &VariantArray,
849 ) -> &dyn TestListLikeArray<OffsetSize = O> {
850 let typed_value = array.typed_value_column().unwrap();
851 if let Some(list) = typed_value.as_any().downcast_ref::<GenericListArray<O>>() {
852 list
853 } else if let Some(list_view) = typed_value
854 .as_any()
855 .downcast_ref::<GenericListViewArray<O>>()
856 {
857 list_view
858 } else {
859 panic!(
860 "Expected list-like typed_value with matching offset type, got {}",
861 typed_value.data_type()
862 );
863 }
864 }
865
866 fn assert_list_structure<O: OffsetSizeTrait>(
867 array: &VariantArray,
868 expected_len: usize,
869 expected_offsets: &[O],
870 expected_sizes: &[Option<O>],
871 expected_fallbacks: &[Option<Variant<'static, 'static>>],
872 ) {
873 assert_eq!(array.len(), expected_len);
874
875 let fallback_value = array.value_column();
876 let fallback_metadata = array.metadata_column();
877 let array = downcast_list_like_array::<O>(array);
878
879 assert_eq!(
880 array.value_offsets().unwrap(),
881 expected_offsets,
882 "list offsets mismatch"
883 );
884 assert_eq!(
885 array.len(),
886 expected_sizes.len(),
887 "expected_sizes should match array length"
888 );
889 assert_eq!(
890 array.len(),
891 expected_fallbacks.len(),
892 "expected_fallbacks should match array length"
893 );
894 assert_eq!(
895 array.len(),
896 fallback_value.len(),
897 "fallbacks value field should match array length"
898 );
899
900 for (idx, (expected_size, expected_fallback)) in expected_sizes
902 .iter()
903 .zip(expected_fallbacks.iter())
904 .enumerate()
905 {
906 match expected_size {
907 Some(len) => {
908 assert!(array.is_valid(idx));
910 assert_eq!(array.value_size(idx), *len);
911 assert!(fallback_value.is_null(idx));
912 }
913 None => {
914 assert!(array.is_null(idx));
916 assert_eq!(array.value_size(idx), O::zero());
917 match expected_fallback {
918 Some(expected_variant) => {
919 assert!(fallback_value.is_valid(idx));
920 let metadata_bytes =
921 binary_array_value(fallback_metadata.as_ref(), idx).unwrap();
922 let metadata_bytes =
923 if fallback_metadata.is_valid(idx) && !metadata_bytes.is_empty() {
924 metadata_bytes
925 } else {
926 EMPTY_VARIANT_METADATA_BYTES
927 };
928 assert_eq!(
929 Variant::new(
930 metadata_bytes,
931 binary_array_value(fallback_value.as_ref(), idx).unwrap()
932 ),
933 expected_variant.clone()
934 );
935 }
936 None => {
937 assert!(fallback_value.is_null(idx));
938 }
939 }
940 }
941 }
942 }
943 }
944
945 fn assert_list_structure_and_elements<T: ArrowPrimitiveType, O: OffsetSizeTrait>(
946 array: &VariantArray,
947 expected_len: usize,
948 expected_offsets: &[O],
949 expected_sizes: &[Option<O>],
950 expected_fallbacks: &[Option<Variant<'static, 'static>>],
951 expected_shredded_elements: (&[Option<T::Native>], &[Option<Variant<'static, 'static>>]),
952 ) {
953 assert_list_structure(
954 array,
955 expected_len,
956 expected_offsets,
957 expected_sizes,
958 expected_fallbacks,
959 );
960 let array = downcast_list_like_array::<O>(array);
961
962 let (expected_values, expected_fallbacks) = expected_shredded_elements;
964 assert_eq!(
965 expected_values.len(),
966 expected_fallbacks.len(),
967 "expected_values and expected_fallbacks should be aligned"
968 );
969
970 let element_array = ShreddedVariantFieldArray::try_new(array.values().as_ref()).unwrap();
972 let element_values = element_array
973 .typed_value_column()
974 .unwrap()
975 .as_any()
976 .downcast_ref::<PrimitiveArray<T>>()
977 .unwrap();
978 assert_eq!(element_values.len(), expected_values.len());
979 for (idx, expected_value) in expected_values.iter().enumerate() {
980 match expected_value {
981 Some(value) => {
982 assert!(element_values.is_valid(idx));
983 assert_eq!(element_values.value(idx), *value);
984 }
985 None => assert!(element_values.is_null(idx)),
986 }
987 }
988
989 let element_fallbacks = element_array.value_column();
991 assert_eq!(element_fallbacks.len(), expected_fallbacks.len());
992 for (idx, expected_fallback) in expected_fallbacks.iter().enumerate() {
993 match expected_fallback {
994 Some(expected_variant) => {
995 assert!(element_fallbacks.is_valid(idx));
996 assert_eq!(
997 Variant::new(
998 EMPTY_VARIANT_METADATA_BYTES,
999 binary_array_value(element_fallbacks.as_ref(), idx).unwrap()
1000 ),
1001 expected_variant.clone()
1002 );
1003 }
1004 None => assert!(element_fallbacks.is_null(idx)),
1005 }
1006 }
1007 }
1008
1009 fn assert_append_null_mode_value_and_struct_nulls(
1010 mode: NullValue,
1011 value: &BinaryViewArray,
1012 nulls: Option<&arrow::buffer::NullBuffer>,
1013 ) {
1014 if mode == NullValue::TopLevelVariant {
1015 assert!(nulls.is_some_and(|n| n.is_null(0)));
1016 } else {
1017 assert!(nulls.is_none());
1018 }
1019
1020 if mode == NullValue::ArrayElement {
1021 assert!(value.is_valid(0));
1022 assert_eq!(
1023 Variant::new(EMPTY_VARIANT_METADATA_BYTES, value.value(0)),
1024 Variant::Null
1025 );
1026 } else {
1027 assert!(value.is_null(0));
1028 }
1029 }
1030
1031 #[test]
1032 fn test_append_null_mode_semantics_primitive_builder() {
1033 let cast_options = arrow::compute::CastOptions::default();
1034
1035 for mode in NULL_VALUES {
1036 let mut primitive_builder = make_variant_to_shredded_variant_arrow_row_builder(
1037 &DataType::Int64,
1038 &cast_options,
1039 1,
1040 mode,
1041 )
1042 .unwrap();
1043 primitive_builder.append_null().unwrap();
1044 let (primitive_value, primitive_typed_value, primitive_nulls) =
1045 primitive_builder.finish().unwrap();
1046 let primitive_typed_value = primitive_typed_value
1047 .as_any()
1048 .downcast_ref::<Int64Array>()
1049 .unwrap();
1050
1051 assert!(primitive_typed_value.is_null(0));
1052 assert_append_null_mode_value_and_struct_nulls(
1053 mode,
1054 &primitive_value,
1055 primitive_nulls.as_ref(),
1056 );
1057 }
1058 }
1059
1060 #[test]
1061 fn test_append_null_mode_semantics_array_builder() {
1062 let cast_options = arrow::compute::CastOptions::default();
1063 let list_type = DataType::List(Arc::new(Field::new("item", DataType::Int64, true)));
1064
1065 for mode in NULL_VALUES {
1066 let mut array_builder = make_variant_to_shredded_variant_arrow_row_builder(
1067 &list_type,
1068 &cast_options,
1069 1,
1070 mode,
1071 )
1072 .unwrap();
1073 array_builder.append_null().unwrap();
1074 let (value, typed_value, nulls) = array_builder.finish().unwrap();
1075
1076 assert_append_null_mode_value_and_struct_nulls(mode, &value, nulls.as_ref());
1077
1078 let typed_value = typed_value.as_any().downcast_ref::<ListArray>().unwrap();
1079 assert_eq!(typed_value.len(), 1);
1080 assert!(typed_value.is_null(0));
1081 assert_eq!(typed_value.values().len(), 0);
1082 }
1083 }
1084
1085 #[test]
1086 fn test_append_null_mode_semantics_object_builder() {
1087 let cast_options = arrow::compute::CastOptions::default();
1088 let object_type = DataType::Struct(Fields::from(vec![
1089 Field::new("id", DataType::Int64, true),
1090 Field::new("name", DataType::Utf8, true),
1091 ]));
1092
1093 for mode in NULL_VALUES {
1094 let mut object_builder = make_variant_to_shredded_variant_arrow_row_builder(
1095 &object_type,
1096 &cast_options,
1097 1,
1098 mode,
1099 )
1100 .unwrap();
1101 object_builder.append_null().unwrap();
1102 let (value, typed_value, nulls) = object_builder.finish().unwrap();
1103
1104 assert_append_null_mode_value_and_struct_nulls(mode, &value, nulls.as_ref());
1105
1106 let typed_struct = typed_value
1107 .as_any()
1108 .downcast_ref::<arrow::array::StructArray>()
1109 .unwrap();
1110 assert_eq!(typed_struct.len(), 1);
1111 assert!(typed_struct.is_null(0));
1112
1113 for field_name in ["id", "name"] {
1114 let field = ShreddedVariantFieldArray::try_new(
1115 typed_struct.column_by_name(field_name).unwrap(),
1116 )
1117 .unwrap();
1118 assert!(field.value_column().is_null(0));
1119 assert!(field.typed_value_column().unwrap().is_null(0));
1120 }
1121 }
1122 }
1123
1124 #[test]
1125 fn test_already_shredded_input_error() {
1126 let temp_array = VariantArray::from_iter(vec![Some(Variant::from("test"))]);
1129 let metadata = temp_array.metadata_column().clone();
1130 let value = temp_array.value_column().clone();
1131 let typed_value = Arc::new(Int64Array::from(vec![42])) as ArrayRef;
1132
1133 let shredded_array = VariantArray::from_parts(metadata, value, Some(typed_value), None);
1134
1135 let result = shred_variant(&shredded_array, &DataType::Int64);
1136 assert!(matches!(
1137 result.unwrap_err(),
1138 ArrowError::InvalidArgumentError(_)
1139 ));
1140 }
1141
1142 #[test]
1143 fn test_all_null_input() {
1144 let metadata = Arc::new(BinaryViewArray::from_iter_values([
1146 EMPTY_VARIANT_METADATA_BYTES,
1147 ]));
1148 let all_null_array =
1149 VariantArray::from_parts(metadata, all_null_value_column(1), None, None);
1150 let result = shred_variant(&all_null_array, &DataType::Int64).unwrap();
1151
1152 assert!(result.typed_value_column().unwrap().is_null(0));
1155 assert_eq!(result.value(0), Variant::Null);
1156 }
1157
1158 #[test]
1159 fn test_invalid_fixed_size_binary_shredding() {
1160 let mock_uuid_1 = Uuid::new_v4();
1161
1162 let input = VariantArray::from_iter([Some(Variant::from(mock_uuid_1)), None]);
1163
1164 let err = shred_variant(&input, &DataType::FixedSizeBinary(17)).unwrap_err();
1166
1167 assert_eq!(
1168 err.to_string(),
1169 "Invalid argument error: FixedSizeBinary(17) is not a valid variant shredding type. Only FixedSizeBinary(16) for UUID is supported."
1170 );
1171 }
1172
1173 #[test]
1174 fn test_uuid_shredding() {
1175 let mock_uuid_1 = Uuid::new_v4();
1176 let mock_uuid_2 = Uuid::new_v4();
1177
1178 let input = VariantArray::from_iter([
1179 Some(Variant::from(mock_uuid_1)),
1180 None,
1181 Some(Variant::from(false)),
1182 Some(Variant::from(mock_uuid_2)),
1183 ]);
1184
1185 let variant_array = shred_variant(&input, &DataType::FixedSizeBinary(16)).unwrap();
1186
1187 let typed_value_field = variant_array.inner().field_by_name("typed_value").unwrap();
1188
1189 assert!(typed_value_field.has_valid_extension_type::<arrow_schema::extension::Uuid>());
1190
1191 let uuids = variant_array
1193 .typed_value_column()
1194 .unwrap()
1195 .as_any()
1196 .downcast_ref::<FixedSizeBinaryArray>()
1197 .unwrap();
1198
1199 assert_eq!(uuids.len(), 4);
1200
1201 assert!(!uuids.is_null(0));
1202
1203 let got_uuid_1: &[u8] = uuids.value(0);
1204 assert_eq!(got_uuid_1, mock_uuid_1.as_bytes());
1205
1206 assert!(uuids.is_null(1));
1207 assert!(uuids.is_null(2));
1208
1209 assert!(!uuids.is_null(3));
1210
1211 let got_uuid_2: &[u8] = uuids.value(3);
1212 assert_eq!(got_uuid_2, mock_uuid_2.as_bytes());
1213 }
1214
1215 #[test]
1216 fn test_uuid_nested_shredding() {
1217 let mock_uuid = Uuid::new_v4();
1218 let input = build_variant_array(vec![VariantRow::Object(vec![(
1219 "id",
1220 VariantValue::from(mock_uuid),
1221 )])]);
1222 let target = ShreddedSchemaBuilder::default()
1223 .with_path("id", DataType::FixedSizeBinary(16))
1224 .unwrap()
1225 .build();
1226
1227 let result = shred_variant(&input, &target).unwrap();
1228
1229 let typed_value = result.typed_value_column().unwrap();
1230 let typed_struct = typed_value.as_any().downcast_ref::<StructArray>().unwrap();
1231 let id =
1232 ShreddedVariantFieldArray::try_new(typed_struct.column_by_name("id").unwrap()).unwrap();
1233
1234 let leaf = id.inner().field_by_name("typed_value").unwrap();
1236
1237 assert_eq!(leaf.data_type(), &DataType::FixedSizeBinary(16));
1238 assert!(leaf.has_valid_extension_type::<arrow_schema::extension::Uuid>());
1239 }
1240
1241 #[test]
1242 fn test_primitive_shredding_comprehensive() {
1243 let input = VariantArray::from_iter(vec![
1245 Some(Variant::from(42i64)), Some(Variant::from("hello")), Some(Variant::from(100i64)), None, Some(Variant::Null), Some(Variant::from(3i8)), ]);
1252
1253 let result = shred_variant(&input, &DataType::Int64).unwrap();
1254
1255 let metadata_field = result.metadata_column();
1257 let value_field = result.value_column();
1258 let typed_value_field = result
1259 .typed_value_column()
1260 .unwrap()
1261 .as_any()
1262 .downcast_ref::<Int64Array>()
1263 .unwrap();
1264
1265 assert_eq!(result.len(), 6);
1267
1268 assert!(!result.is_null(0));
1270 assert!(value_field.is_null(0)); assert!(!typed_value_field.is_null(0));
1272 assert_eq!(typed_value_field.value(0), 42);
1273
1274 assert!(!result.is_null(1));
1276 assert!(!value_field.is_null(1)); assert!(typed_value_field.is_null(1)); assert_eq!(
1279 variant_from_arrays_at(metadata_field, value_field, 1).unwrap(),
1280 Variant::from("hello")
1281 );
1282
1283 assert!(!result.is_null(2));
1285 assert!(value_field.is_null(2));
1286 assert_eq!(typed_value_field.value(2), 100);
1287
1288 assert!(result.is_null(3));
1290
1291 assert!(!result.is_null(4));
1293 assert!(!value_field.is_null(4)); assert_eq!(
1295 variant_from_arrays_at(metadata_field, value_field, 4).unwrap(),
1296 Variant::Null
1297 );
1298 assert!(typed_value_field.is_null(4));
1299
1300 assert!(!result.is_null(5));
1302 assert!(value_field.is_null(5)); assert!(!typed_value_field.is_null(5));
1304 assert_eq!(typed_value_field.value(5), 3);
1305 }
1306
1307 #[test]
1308 fn test_primitive_different_target_types() {
1309 let input = VariantArray::from_iter(vec![
1310 Variant::from(42i32),
1311 Variant::from(3.15f64),
1312 Variant::from("not_a_number"),
1313 ]);
1314
1315 let result_int32 = shred_variant(&input, &DataType::Int32).unwrap();
1317 let typed_value_int32 = result_int32
1318 .typed_value_column()
1319 .unwrap()
1320 .as_any()
1321 .downcast_ref::<arrow::array::Int32Array>()
1322 .unwrap();
1323 assert_eq!(typed_value_int32.value(0), 42);
1324 assert_eq!(typed_value_int32.value(1), 3);
1325 assert!(typed_value_int32.is_null(2)); let result_float64 = shred_variant(&input, &DataType::Float64).unwrap();
1329 let typed_value_float64 = result_float64
1330 .typed_value_column()
1331 .unwrap()
1332 .as_any()
1333 .downcast_ref::<Float64Array>()
1334 .unwrap();
1335 assert_eq!(typed_value_float64.value(0), 42.0); assert_eq!(typed_value_float64.value(1), 3.15);
1337 assert!(typed_value_float64.is_null(2)); }
1339
1340 #[test]
1341 fn test_largeutf8_shredding() {
1342 let input = VariantArray::from_iter(vec![
1343 Some(Variant::from("hello")),
1344 Some(Variant::from(42i64)),
1345 None,
1346 Some(Variant::Null),
1347 Some(Variant::from("world")),
1348 ]);
1349
1350 let result = shred_variant(&input, &DataType::LargeUtf8).unwrap();
1351 let metadata = result.metadata_column();
1352 let value = result.value_column();
1353 let typed_value = result
1354 .typed_value_column()
1355 .unwrap()
1356 .as_any()
1357 .downcast_ref::<LargeStringArray>()
1358 .unwrap();
1359
1360 assert_eq!(result.len(), 5);
1361
1362 assert!(result.is_valid(0));
1364 assert!(value.is_null(0));
1365 assert_eq!(typed_value.value(0), "hello");
1366
1367 assert!(result.is_valid(1));
1369 assert!(value.is_valid(1));
1370 assert!(typed_value.is_null(1));
1371 assert_eq!(
1372 variant_from_arrays_at(metadata, value, 1).unwrap(),
1373 Variant::from(42i64)
1374 );
1375
1376 assert!(result.is_null(2));
1378 assert!(value.is_null(2));
1379 assert!(typed_value.is_null(2));
1380
1381 assert!(result.is_valid(3));
1383 assert!(value.is_valid(3));
1384 assert!(typed_value.is_null(3));
1385 assert_eq!(
1386 variant_from_arrays_at(metadata, value, 3).unwrap(),
1387 Variant::Null
1388 );
1389
1390 assert!(result.is_valid(4));
1392 assert!(value.is_null(4));
1393 assert_eq!(typed_value.value(4), "world");
1394 }
1395
1396 #[test]
1397 fn test_largebinary_shredding() {
1398 let input = VariantArray::from_iter(vec![
1399 Some(Variant::from(&b"\x00\x01\x02"[..])),
1400 Some(Variant::from("not_binary")),
1401 None,
1402 Some(Variant::Null),
1403 Some(Variant::from(&b"\xff\xaa"[..])),
1404 ]);
1405
1406 let result = shred_variant(&input, &DataType::LargeBinary).unwrap();
1407 let metadata = result.metadata_column();
1408 let value = result.value_column();
1409 let typed_value = result
1410 .typed_value_column()
1411 .unwrap()
1412 .as_any()
1413 .downcast_ref::<LargeBinaryArray>()
1414 .unwrap();
1415
1416 assert_eq!(result.len(), 5);
1417
1418 assert!(result.is_valid(0));
1420 assert!(value.is_null(0));
1421 assert_eq!(typed_value.value(0), &[0x00, 0x01, 0x02]);
1422
1423 assert!(result.is_valid(1));
1425 assert!(value.is_valid(1));
1426 assert!(typed_value.is_null(1));
1427 assert_eq!(
1428 variant_from_arrays_at(metadata, value, 1).unwrap(),
1429 Variant::from("not_binary")
1430 );
1431
1432 assert!(result.is_null(2));
1434 assert!(value.is_null(2));
1435 assert!(typed_value.is_null(2));
1436
1437 assert!(result.is_valid(3));
1439 assert!(value.is_valid(3));
1440 assert!(typed_value.is_null(3));
1441 assert_eq!(
1442 variant_from_arrays_at(metadata, value, 3).unwrap(),
1443 Variant::Null
1444 );
1445
1446 assert!(result.is_valid(4));
1448 assert!(value.is_null(4));
1449 assert_eq!(typed_value.value(4), &[0xff, 0xaa]);
1450 }
1451
1452 #[test]
1453 fn test_invalid_shredded_types_rejected() {
1454 let input = VariantArray::from_iter([Variant::from(42)]);
1455
1456 let invalid_types = vec![
1457 DataType::UInt8,
1458 DataType::Float16,
1459 DataType::Decimal256(38, 10),
1460 DataType::Date64,
1461 DataType::Time32(TimeUnit::Second),
1462 DataType::Time64(TimeUnit::Nanosecond),
1463 DataType::Timestamp(TimeUnit::Millisecond, None),
1464 DataType::FixedSizeBinary(17),
1465 DataType::Union(
1466 UnionFields::from_fields(vec![
1467 Field::new("int_field", DataType::Int32, false),
1468 Field::new("str_field", DataType::Utf8, true),
1469 ]),
1470 UnionMode::Dense,
1471 ),
1472 DataType::Map(
1473 Arc::new(Field::new(
1474 Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
1475 DataType::Struct(Fields::from(vec![
1476 Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
1477 Field::new(Field::MAP_VALUE_FIELD_DEFAULT_NAME, DataType::Int32, true),
1478 ])),
1479 false,
1480 )),
1481 false,
1482 ),
1483 DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
1484 DataType::RunEndEncoded(
1485 Arc::new(Field::new("run_ends", DataType::Int32, false)),
1486 Arc::new(Field::new("values", DataType::Utf8, true)),
1487 ),
1488 ];
1489
1490 for data_type in invalid_types {
1491 let err = shred_variant(&input, &data_type).unwrap_err();
1492 assert!(
1493 matches!(err, ArrowError::InvalidArgumentError(_)),
1494 "expected InvalidArgumentError for {:?}, got {:?}",
1495 data_type,
1496 err
1497 );
1498 }
1499 }
1500
1501 #[test]
1502 fn test_array_shredding_as_list() {
1503 let input = build_variant_array(vec![
1504 VariantRow::List(vec![
1506 VariantValue::from(1i64),
1507 VariantValue::from(2i64),
1508 VariantValue::from(3i64),
1509 ]),
1510 VariantRow::List(vec![
1512 VariantValue::from(1i64),
1513 VariantValue::from("two"),
1514 VariantValue::from(Variant::Null),
1515 ]),
1516 VariantRow::Value(VariantValue::from("not a list")),
1518 VariantRow::Null,
1520 VariantRow::List(vec![]),
1522 ]);
1523 let list_schema = DataType::List(Arc::new(Field::new("item", DataType::Int64, true)));
1524 let result = shred_variant(&input, &list_schema).unwrap();
1525 assert_eq!(result.len(), 5);
1526
1527 assert_list_structure_and_elements::<Int64Type, i32>(
1528 &result,
1529 5,
1530 &[0, 3, 6, 6, 6, 6],
1531 &[Some(3), Some(3), None, None, Some(0)],
1532 &[None, None, Some(Variant::from("not a list")), None, None],
1533 (
1534 &[Some(1), Some(2), Some(3), Some(1), None, None],
1535 &[
1536 None,
1537 None,
1538 None,
1539 None,
1540 Some(Variant::from("two")),
1541 Some(Variant::Null),
1542 ],
1543 ),
1544 );
1545 }
1546
1547 #[test]
1548 fn test_array_shredding_as_large_list() {
1549 let input = build_variant_array(vec![
1550 VariantRow::List(vec![VariantValue::from(1i64), VariantValue::from(2i64)]),
1552 VariantRow::Value(VariantValue::from("not a list")),
1554 VariantRow::List(vec![]),
1556 ]);
1557 let list_schema = DataType::LargeList(Arc::new(Field::new("item", DataType::Int64, true)));
1558 let result = shred_variant(&input, &list_schema).unwrap();
1559 assert_eq!(result.len(), 3);
1560
1561 assert_list_structure_and_elements::<Int64Type, i64>(
1562 &result,
1563 3,
1564 &[0, 2, 2, 2],
1565 &[Some(2), None, Some(0)],
1566 &[None, Some(Variant::from("not a list")), None],
1567 (&[Some(1), Some(2)], &[None, None]),
1568 );
1569 }
1570
1571 #[test]
1572 fn test_array_shredding_as_list_view() {
1573 let input = build_variant_array(vec![
1574 VariantRow::List(vec![
1576 VariantValue::from(1i64),
1577 VariantValue::from(2i64),
1578 VariantValue::from(3i64),
1579 ]),
1580 VariantRow::List(vec![
1582 VariantValue::from(1i64),
1583 VariantValue::from("two"),
1584 VariantValue::from(Variant::Null),
1585 ]),
1586 VariantRow::Value(VariantValue::from("not a list")),
1588 VariantRow::Null,
1590 VariantRow::List(vec![]),
1592 ]);
1593 let list_schema = DataType::ListView(Arc::new(Field::new("item", DataType::Int64, true)));
1594 let result = shred_variant(&input, &list_schema).unwrap();
1595 assert_eq!(result.len(), 5);
1596
1597 assert_list_structure_and_elements::<Int64Type, i32>(
1598 &result,
1599 5,
1600 &[0, 3, 6, 6, 6],
1601 &[Some(3), Some(3), None, None, Some(0)],
1602 &[None, None, Some(Variant::from("not a list")), None, None],
1603 (
1604 &[Some(1), Some(2), Some(3), Some(1), None, None],
1605 &[
1606 None,
1607 None,
1608 None,
1609 None,
1610 Some(Variant::from("two")),
1611 Some(Variant::Null),
1612 ],
1613 ),
1614 );
1615 }
1616
1617 #[test]
1618 fn test_array_shredding_as_large_list_view() {
1619 let input = build_variant_array(vec![
1620 VariantRow::List(vec![VariantValue::from(1i64), VariantValue::from(2i64)]),
1622 VariantRow::Value(VariantValue::from("fallback")),
1624 VariantRow::List(vec![]),
1626 ]);
1627 let list_schema =
1628 DataType::LargeListView(Arc::new(Field::new("item", DataType::Int64, true)));
1629 let result = shred_variant(&input, &list_schema).unwrap();
1630 assert_eq!(result.len(), 3);
1631
1632 assert_list_structure_and_elements::<Int64Type, i64>(
1633 &result,
1634 3,
1635 &[0, 2, 2],
1636 &[Some(2), None, Some(0)],
1637 &[None, Some(Variant::from("fallback")), None],
1638 (&[Some(1), Some(2)], &[None, None]),
1639 );
1640 }
1641
1642 #[test]
1643 fn test_array_shredding_as_fixed_size_list() {
1644 let input = build_variant_array(vec![
1645 VariantRow::List(vec![VariantValue::from(1i64), VariantValue::from(2i64)]),
1646 VariantRow::Value(VariantValue::from("This should not be shredded")),
1647 VariantRow::List(vec![VariantValue::from(3i64), VariantValue::from(4i64)]),
1648 ]);
1649
1650 let list_schema =
1651 DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int64, true)), 2);
1652 let result = shred_variant(&input, &list_schema).unwrap();
1653 assert_eq!(result.len(), 3);
1654
1655 assert!(result.is_valid(0));
1658 assert!(result.value_column().is_null(0));
1659 assert!(result.typed_value_column().unwrap().is_valid(0));
1660
1661 assert!(result.is_valid(1));
1665 assert!(result.value_column().is_valid(1));
1666 assert!(result.typed_value_column().unwrap().is_null(1));
1667
1668 assert!(result.is_valid(2));
1671 assert!(result.value_column().is_null(2));
1672 assert!(result.typed_value_column().unwrap().is_valid(2));
1673
1674 let typed_value = result.typed_value_column().unwrap();
1675 let fixed_size_list = typed_value
1676 .as_any()
1677 .downcast_ref::<FixedSizeListArray>()
1678 .expect("Expected FixedSizeListArray");
1679
1680 assert_eq!(fixed_size_list.len(), 3);
1682 assert_eq!(fixed_size_list.value_length(), 2);
1683
1684 let val0 = fixed_size_list.value(0);
1686 let val0_struct = val0.as_any().downcast_ref::<StructArray>().unwrap();
1687 let val0_typed = val0_struct.column_by_name("typed_value").unwrap();
1688 let val0_ints = val0_typed.as_any().downcast_ref::<Int64Array>().unwrap();
1689 assert_eq!(val0_ints.values(), &[1i64, 2i64]);
1690
1691 assert!(fixed_size_list.is_null(1));
1694
1695 let val2 = fixed_size_list.value(2);
1697 let val2_struct = val2.as_any().downcast_ref::<StructArray>().unwrap();
1698 let val2_typed = val2_struct.column_by_name("typed_value").unwrap();
1699 let val2_ints = val2_typed.as_any().downcast_ref::<Int64Array>().unwrap();
1700 assert_eq!(val2_ints.values(), &[3i64, 4i64]);
1701 }
1702
1703 #[test]
1704 fn test_array_shredding_as_fixed_size_list_wrong_size() {
1705 let input = build_variant_array(vec![VariantRow::List(vec![
1706 VariantValue::from(1i64),
1707 VariantValue::from(2i64),
1708 VariantValue::from(3i64),
1709 ])]);
1710 let list_schema =
1711 DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int64, true)), 2);
1712
1713 let err = shred_variant(&input, &list_schema).unwrap_err();
1714 assert!(
1715 err.to_string()
1716 .contains("Expected fixed size list of size 2, got size 3"),
1717 "got: {err}",
1718 );
1719 }
1720
1721 #[test]
1722 fn test_array_shredding_with_array_elements() {
1723 let input = build_variant_array(vec![
1724 VariantRow::List(vec![
1726 VariantValue::List(vec![VariantValue::from(1i64), VariantValue::from(2i64)]),
1727 VariantValue::List(vec![VariantValue::from(3i64), VariantValue::from(4i64)]),
1728 VariantValue::List(vec![]),
1729 ]),
1730 VariantRow::List(vec![
1732 VariantValue::List(vec![
1733 VariantValue::from(5i64),
1734 VariantValue::from("bad"),
1735 VariantValue::from(Variant::Null),
1736 ]),
1737 VariantValue::from("not a list inner"),
1738 VariantValue::Null,
1739 ]),
1740 VariantRow::Value(VariantValue::from("not a list")),
1742 VariantRow::Null,
1744 ]);
1745 let inner_field = Arc::new(Field::new("item", DataType::Int64, true));
1746 let inner_list_schema = DataType::List(inner_field);
1747 let list_schema = DataType::List(Arc::new(Field::new(
1748 "item",
1749 inner_list_schema.clone(),
1750 true,
1751 )));
1752 let result = shred_variant(&input, &list_schema).unwrap();
1753 assert_eq!(result.len(), 4);
1754
1755 let typed_value = result
1756 .typed_value_column()
1757 .unwrap()
1758 .as_any()
1759 .downcast_ref::<ListArray>()
1760 .unwrap();
1761
1762 assert_list_structure::<i32>(
1763 &result,
1764 4,
1765 &[0, 3, 6, 6, 6],
1766 &[Some(3), Some(3), None, None],
1767 &[None, None, Some(Variant::from("not a list")), None],
1768 );
1769
1770 let outer_elements =
1771 ShreddedVariantFieldArray::try_new(typed_value.values().as_ref()).unwrap();
1772 assert_eq!(outer_elements.len(), 6);
1773 let outer_values = outer_elements
1774 .typed_value_column()
1775 .unwrap()
1776 .as_any()
1777 .downcast_ref::<ListArray>()
1778 .unwrap();
1779 let outer_fallbacks = outer_elements.value_column();
1780
1781 let outer_metadata = Arc::new(BinaryViewArray::from_iter_values(std::iter::repeat_n(
1782 EMPTY_VARIANT_METADATA_BYTES,
1783 outer_elements.len(),
1784 )));
1785 let outer_variant = VariantArray::from_parts(
1786 outer_metadata,
1787 outer_fallbacks.clone(),
1788 Some(Arc::new(outer_values.clone())),
1789 None,
1790 );
1791
1792 assert_list_structure_and_elements::<Int64Type, i32>(
1793 &outer_variant,
1794 outer_elements.len(),
1795 &[0, 2, 4, 4, 7, 7, 7],
1796 &[Some(2), Some(2), Some(0), Some(3), None, None],
1797 &[
1798 None,
1799 None,
1800 None,
1801 None,
1802 Some(Variant::from("not a list inner")),
1803 Some(Variant::Null),
1804 ],
1805 (
1806 &[Some(1), Some(2), Some(3), Some(4), Some(5), None, None],
1807 &[
1808 None,
1809 None,
1810 None,
1811 None,
1812 None,
1813 Some(Variant::from("bad")),
1814 Some(Variant::Null),
1815 ],
1816 ),
1817 );
1818 }
1819
1820 #[test]
1821 fn test_array_shredding_with_object_elements() {
1822 let input = build_variant_array(vec![
1823 VariantRow::List(vec![
1825 VariantValue::Object(vec![
1826 ("id", VariantValue::from(1i64)),
1827 ("name", VariantValue::from("Alice")),
1828 ]),
1829 VariantValue::Object(vec![("id", VariantValue::from(Variant::Null))]),
1830 ]),
1831 VariantRow::Value(VariantValue::from("not a list")),
1833 VariantRow::Null,
1835 ]);
1836
1837 let object_fields = Fields::from(vec![
1839 Field::new("id", DataType::Int64, true),
1840 Field::new("name", DataType::Utf8, true),
1841 ]);
1842 let list_schema = DataType::List(Arc::new(Field::new(
1843 "item",
1844 DataType::Struct(object_fields),
1845 true,
1846 )));
1847 let result = shred_variant(&input, &list_schema).unwrap();
1848 assert_eq!(result.len(), 3);
1849
1850 assert_list_structure::<i32>(
1851 &result,
1852 3,
1853 &[0, 2, 2, 2],
1854 &[Some(2), None, None],
1855 &[None, Some(Variant::from("not a list")), None],
1856 );
1857
1858 let typed_value = result
1860 .typed_value_column()
1861 .unwrap()
1862 .as_any()
1863 .downcast_ref::<ListArray>()
1864 .unwrap();
1865 let element_array =
1866 ShreddedVariantFieldArray::try_new(typed_value.values().as_ref()).unwrap();
1867 assert_eq!(element_array.len(), 2);
1868 let element_objects = element_array
1869 .typed_value_column()
1870 .unwrap()
1871 .as_any()
1872 .downcast_ref::<arrow::array::StructArray>()
1873 .unwrap();
1874
1875 let id_field =
1877 ShreddedVariantFieldArray::try_new(element_objects.column_by_name("id").unwrap())
1878 .unwrap();
1879 let id_values = id_field.value_column();
1880 let id_typed_values = id_field
1881 .typed_value_column()
1882 .unwrap()
1883 .as_any()
1884 .downcast_ref::<Int64Array>()
1885 .unwrap();
1886 assert!(id_values.is_null(0));
1887 assert_eq!(id_typed_values.value(0), 1);
1888 assert!(id_values.is_valid(1));
1890 assert_eq!(
1891 Variant::new(
1892 EMPTY_VARIANT_METADATA_BYTES,
1893 binary_array_value(id_values.as_ref(), 1).unwrap()
1894 ),
1895 Variant::Null
1896 );
1897 assert!(id_typed_values.is_null(1));
1898
1899 let name_field =
1901 ShreddedVariantFieldArray::try_new(element_objects.column_by_name("name").unwrap())
1902 .unwrap();
1903 let name_values = name_field.value_column();
1904 let name_typed_values = name_field
1905 .typed_value_column()
1906 .unwrap()
1907 .as_any()
1908 .downcast_ref::<StringArray>()
1909 .unwrap();
1910 assert!(name_values.is_null(0));
1911 assert_eq!(name_typed_values.value(0), "Alice");
1912 assert!(name_values.is_null(1));
1914 assert!(name_typed_values.is_null(1));
1915 }
1916
1917 #[test]
1918 fn test_object_shredding_comprehensive() -> Result<()> {
1919 let input = build_variant_array(vec![
1920 VariantRow::Object(vec![
1922 ("score", VariantValue::from(95.5f64)),
1923 ("age", VariantValue::from(30i64)),
1924 ]),
1925 VariantRow::Object(vec![
1927 ("score", VariantValue::from(87.2f64)),
1928 ("age", VariantValue::from(25i64)),
1929 ("email", VariantValue::from("bob@example.com")),
1930 ]),
1931 VariantRow::Object(vec![("age", VariantValue::from(35i64))]),
1933 VariantRow::Object(vec![
1935 ("score", VariantValue::from("ninety-five")),
1936 ("age", VariantValue::from("thirty")),
1937 ]),
1938 VariantRow::Value(VariantValue::from("not an object")),
1940 VariantRow::Object(vec![]),
1942 VariantRow::Null,
1944 VariantRow::Object(vec![("foo", VariantValue::from(10))]),
1946 VariantRow::Object(vec![
1948 ("score", VariantValue::from(66.67f64)),
1949 ("foo", VariantValue::from(10)),
1950 ]),
1951 ]);
1952
1953 let target_schema = ShreddedSchemaBuilder::default()
1956 .with_path("score", &DataType::Float64)?
1957 .with_path("age", &DataType::Int64)?
1958 .build();
1959
1960 let result = shred_variant(&input, &target_schema).unwrap();
1961
1962 assert!(result.typed_value_column().is_some());
1964 assert_eq!(result.len(), 9);
1965
1966 let metadata = result.metadata_column();
1967 let value = result.value_column();
1968 let typed_value = result
1969 .typed_value_column()
1970 .unwrap()
1971 .as_any()
1972 .downcast_ref::<arrow::array::StructArray>()
1973 .unwrap();
1974
1975 let score_field =
1977 ShreddedVariantFieldArray::try_new(typed_value.column_by_name("score").unwrap())
1978 .unwrap();
1979 let age_field =
1980 ShreddedVariantFieldArray::try_new(typed_value.column_by_name("age").unwrap()).unwrap();
1981
1982 let score_value = score_field.value_column();
1983 let score_typed_value = score_field
1984 .typed_value_column()
1985 .unwrap()
1986 .as_any()
1987 .downcast_ref::<Float64Array>()
1988 .unwrap();
1989 let age_value = age_field.value_column();
1990 let age_typed_value = age_field
1991 .typed_value_column()
1992 .unwrap()
1993 .as_any()
1994 .downcast_ref::<Int64Array>()
1995 .unwrap();
1996
1997 struct ShreddedValue<'m, 'v, T> {
1999 value: Option<Variant<'m, 'v>>,
2000 typed_value: Option<T>,
2001 }
2002 struct ShreddedStruct<'m, 'v> {
2003 score: ShreddedValue<'m, 'v, f64>,
2004 age: ShreddedValue<'m, 'v, i64>,
2005 }
2006 fn get_value<'m, 'v>(
2007 i: usize,
2008 metadata: &'m dyn Array,
2009 value: &'v dyn Array,
2010 ) -> Variant<'m, 'v> {
2011 variant_from_arrays_at(metadata, value, i).unwrap()
2012 }
2013 let expect = |i, expected_result: Option<ShreddedValue<ShreddedStruct>>| {
2014 match expected_result {
2015 Some(ShreddedValue {
2016 value: expected_value,
2017 typed_value: expected_typed_value,
2018 }) => {
2019 assert!(result.is_valid(i));
2020 match expected_value {
2021 Some(expected_value) => {
2022 assert!(value.is_valid(i));
2023 assert_eq!(
2024 expected_value,
2025 get_value(i, metadata.as_ref(), value.as_ref())
2026 );
2027 }
2028 None => {
2029 assert!(value.is_null(i));
2030 }
2031 }
2032 match expected_typed_value {
2033 Some(ShreddedStruct {
2034 score: expected_score,
2035 age: expected_age,
2036 }) => {
2037 assert!(typed_value.is_valid(i));
2038 assert!(score_field.is_valid(i)); assert!(age_field.is_valid(i)); match expected_score.value {
2041 Some(expected_score_value) => {
2042 assert!(score_value.is_valid(i));
2043 assert_eq!(
2044 expected_score_value,
2045 get_value(i, metadata.as_ref(), score_value.as_ref())
2046 );
2047 }
2048 None => {
2049 assert!(score_value.is_null(i));
2050 }
2051 }
2052 match expected_score.typed_value {
2053 Some(expected_score) => {
2054 assert!(score_typed_value.is_valid(i));
2055 assert_eq!(expected_score, score_typed_value.value(i));
2056 }
2057 None => {
2058 assert!(score_typed_value.is_null(i));
2059 }
2060 }
2061 match expected_age.value {
2062 Some(expected_age_value) => {
2063 assert!(age_value.is_valid(i));
2064 assert_eq!(
2065 expected_age_value,
2066 get_value(i, metadata.as_ref(), age_value.as_ref())
2067 );
2068 }
2069 None => {
2070 assert!(age_value.is_null(i));
2071 }
2072 }
2073 match expected_age.typed_value {
2074 Some(expected_age) => {
2075 assert!(age_typed_value.is_valid(i));
2076 assert_eq!(expected_age, age_typed_value.value(i));
2077 }
2078 None => {
2079 assert!(age_typed_value.is_null(i));
2080 }
2081 }
2082 }
2083 None => {
2084 assert!(typed_value.is_null(i));
2085 }
2086 }
2087 }
2088 None => {
2089 assert!(result.is_null(i));
2090 }
2091 };
2092 };
2093
2094 expect(
2096 0,
2097 Some(ShreddedValue {
2098 value: None,
2099 typed_value: Some(ShreddedStruct {
2100 score: ShreddedValue {
2101 value: None,
2102 typed_value: Some(95.5),
2103 },
2104 age: ShreddedValue {
2105 value: None,
2106 typed_value: Some(30),
2107 },
2108 }),
2109 }),
2110 );
2111
2112 let mut builder = VariantBuilder::new();
2114 builder
2115 .new_object()
2116 .with_field("email", "bob@example.com")
2117 .finish();
2118 let (m, v) = builder.finish();
2119 let expected_value = Variant::new(&m, &v);
2120
2121 expect(
2122 1,
2123 Some(ShreddedValue {
2124 value: Some(expected_value),
2125 typed_value: Some(ShreddedStruct {
2126 score: ShreddedValue {
2127 value: None,
2128 typed_value: Some(87.2),
2129 },
2130 age: ShreddedValue {
2131 value: None,
2132 typed_value: Some(25),
2133 },
2134 }),
2135 }),
2136 );
2137
2138 expect(
2140 2,
2141 Some(ShreddedValue {
2142 value: None,
2143 typed_value: Some(ShreddedStruct {
2144 score: ShreddedValue {
2145 value: None,
2146 typed_value: None,
2147 },
2148 age: ShreddedValue {
2149 value: None,
2150 typed_value: Some(35),
2151 },
2152 }),
2153 }),
2154 );
2155
2156 expect(
2158 3,
2159 Some(ShreddedValue {
2160 value: None,
2161 typed_value: Some(ShreddedStruct {
2162 score: ShreddedValue {
2163 value: Some(Variant::from("ninety-five")),
2164 typed_value: None,
2165 },
2166 age: ShreddedValue {
2167 value: Some(Variant::from("thirty")),
2168 typed_value: None,
2169 },
2170 }),
2171 }),
2172 );
2173
2174 expect(
2176 4,
2177 Some(ShreddedValue {
2178 value: Some(Variant::from("not an object")),
2179 typed_value: None,
2180 }),
2181 );
2182
2183 expect(
2185 5,
2186 Some(ShreddedValue {
2187 value: None,
2188 typed_value: Some(ShreddedStruct {
2189 score: ShreddedValue {
2190 value: None,
2191 typed_value: None,
2192 },
2193 age: ShreddedValue {
2194 value: None,
2195 typed_value: None,
2196 },
2197 }),
2198 }),
2199 );
2200
2201 expect(6, None);
2203
2204 let object_with_foo_field = |i| {
2206 use parquet_variant::{ParentState, ValueBuilder, VariantMetadata};
2207 let metadata = VariantMetadata::new(binary_array_value(metadata.as_ref(), i).unwrap());
2208 let mut metadata_builder = ReadOnlyMetadataBuilder::new(&metadata);
2209 let mut value_builder = ValueBuilder::new();
2210 let state = ParentState::variant(&mut value_builder, &mut metadata_builder);
2211 ObjectBuilder::new(state, false)
2212 .with_field("foo", 10)
2213 .finish();
2214 (metadata, value_builder.into_inner())
2215 };
2216
2217 let (m, v) = object_with_foo_field(7);
2219 expect(
2220 7,
2221 Some(ShreddedValue {
2222 value: Some(Variant::new_with_metadata(m, &v)),
2223 typed_value: Some(ShreddedStruct {
2224 score: ShreddedValue {
2225 value: None,
2226 typed_value: None,
2227 },
2228 age: ShreddedValue {
2229 value: None,
2230 typed_value: None,
2231 },
2232 }),
2233 }),
2234 );
2235
2236 let (m, v) = object_with_foo_field(8);
2238 expect(
2239 8,
2240 Some(ShreddedValue {
2241 value: Some(Variant::new_with_metadata(m, &v)),
2242 typed_value: Some(ShreddedStruct {
2243 score: ShreddedValue {
2244 value: None,
2245 typed_value: Some(66.67),
2246 },
2247 age: ShreddedValue {
2248 value: None,
2249 typed_value: None,
2250 },
2251 }),
2252 }),
2253 );
2254 Ok(())
2255 }
2256
2257 #[test]
2258 fn test_object_shredding_with_array_field() {
2259 let input = build_variant_array(vec![
2260 VariantRow::Object(vec![(
2262 "scores",
2263 VariantValue::List(vec![VariantValue::from(10i64), VariantValue::from(20i64)]),
2264 )]),
2265 VariantRow::Object(vec![(
2267 "scores",
2268 VariantValue::List(vec![
2269 VariantValue::from("oops"),
2270 VariantValue::from(Variant::Null),
2271 ]),
2272 )]),
2273 VariantRow::Object(vec![]),
2275 VariantRow::Value(VariantValue::from("not an object")),
2277 VariantRow::Null,
2279 ]);
2280 let list_field = Arc::new(Field::new("item", DataType::Int64, true));
2281 let inner_list_schema = DataType::List(list_field);
2282 let schema = DataType::Struct(Fields::from(vec![Field::new(
2283 "scores",
2284 inner_list_schema.clone(),
2285 true,
2286 )]));
2287
2288 let result = shred_variant(&input, &schema).unwrap();
2289 assert_eq!(result.len(), 5);
2290
2291 let value_field = result.value_column();
2293 let typed_struct = result
2294 .typed_value_column()
2295 .unwrap()
2296 .as_any()
2297 .downcast_ref::<arrow::array::StructArray>()
2298 .unwrap();
2299
2300 assert!(value_field.is_null(0));
2302 assert!(value_field.is_null(1));
2303 assert!(value_field.is_null(2));
2304 assert!(value_field.is_valid(3));
2305 assert_eq!(
2306 variant_from_arrays_at(result.metadata_column(), value_field, 3).unwrap(),
2307 Variant::from("not an object")
2308 );
2309 assert!(value_field.is_null(4));
2310
2311 assert!(typed_struct.is_valid(0));
2313 assert!(typed_struct.is_valid(1));
2314 assert!(typed_struct.is_valid(2));
2315 assert!(typed_struct.is_null(3));
2316 assert!(typed_struct.is_null(4));
2317
2318 let scores_field =
2320 ShreddedVariantFieldArray::try_new(typed_struct.column_by_name("scores").unwrap())
2321 .unwrap();
2322 assert_list_structure_and_elements::<Int64Type, i32>(
2323 &VariantArray::from_parts(
2324 Arc::new(BinaryViewArray::from_iter_values(std::iter::repeat_n(
2325 EMPTY_VARIANT_METADATA_BYTES,
2326 scores_field.len(),
2327 ))),
2328 scores_field.value_column().clone(),
2329 Some(scores_field.typed_value_column().unwrap().clone()),
2330 None,
2331 ),
2332 scores_field.len(),
2333 &[0i32, 2, 4, 4, 4, 4],
2334 &[Some(2), Some(2), None, None, None],
2335 &[None, None, None, None, None],
2336 (
2337 &[Some(10), Some(20), None, None],
2338 &[None, None, Some(Variant::from("oops")), Some(Variant::Null)],
2339 ),
2340 );
2341 }
2342
2343 #[test]
2344 fn test_object_different_schemas() -> Result<()> {
2345 let input = build_variant_array(vec![VariantRow::Object(vec![
2347 ("id", VariantValue::from(123i32)),
2348 ("age", VariantValue::from(25i64)),
2349 ("score", VariantValue::from(95.5f64)),
2350 ])]);
2351
2352 let schema1 = ShreddedSchemaBuilder::default()
2354 .with_path("id", &DataType::Int32)?
2355 .build();
2356 let result1 = shred_variant(&input, &schema1).unwrap();
2357 let value_field1 = result1.value_column();
2358 assert!(!value_field1.is_null(0)); let schema2 = ShreddedSchemaBuilder::default()
2362 .with_path("id", &DataType::Int32)?
2363 .with_path("age", &DataType::Int64)?
2364 .build();
2365 let result2 = shred_variant(&input, &schema2).unwrap();
2366 let value_field2 = result2.value_column();
2367 assert!(!value_field2.is_null(0)); let schema3 = ShreddedSchemaBuilder::default()
2371 .with_path("id", &DataType::Int32)?
2372 .with_path("age", &DataType::Int64)?
2373 .with_path("score", &DataType::Float64)?
2374 .build();
2375 let result3 = shred_variant(&input, &schema3).unwrap();
2376 let value_field3 = result3.value_column();
2377 assert!(value_field3.is_null(0)); Ok(())
2380 }
2381
2382 #[test]
2383 fn test_uuid_shredding_in_objects() -> Result<()> {
2384 let mock_uuid_1 = Uuid::new_v4();
2385 let mock_uuid_2 = Uuid::new_v4();
2386 let mock_uuid_3 = Uuid::new_v4();
2387
2388 let input = build_variant_array(vec![
2389 VariantRow::Object(vec![
2391 ("id", VariantValue::from(mock_uuid_1)),
2392 ("session_id", VariantValue::from(mock_uuid_2)),
2393 ]),
2394 VariantRow::Object(vec![
2396 ("id", VariantValue::from(mock_uuid_2)),
2397 ("session_id", VariantValue::from(mock_uuid_3)),
2398 ("name", VariantValue::from("test_user")),
2399 ]),
2400 VariantRow::Object(vec![("id", VariantValue::from(mock_uuid_1))]),
2402 VariantRow::Object(vec![
2404 ("id", VariantValue::from(mock_uuid_3)),
2405 ("session_id", VariantValue::from("not-a-uuid")),
2406 ]),
2407 VariantRow::Object(vec![
2409 ("id", VariantValue::from(12345i64)),
2410 ("session_id", VariantValue::from(mock_uuid_1)),
2411 ]),
2412 VariantRow::Null,
2414 ]);
2415
2416 let target_schema = ShreddedSchemaBuilder::default()
2417 .with_path("id", DataType::FixedSizeBinary(16))?
2418 .with_path("session_id", DataType::FixedSizeBinary(16))?
2419 .build();
2420
2421 let result = shred_variant(&input, &target_schema).unwrap();
2422
2423 assert!(result.typed_value_column().is_some());
2424 assert_eq!(result.len(), 6);
2425
2426 let metadata = result.metadata_column();
2427 let value = result.value_column();
2428 let typed_value = result
2429 .typed_value_column()
2430 .unwrap()
2431 .as_any()
2432 .downcast_ref::<arrow::array::StructArray>()
2433 .unwrap();
2434
2435 let id_field =
2437 ShreddedVariantFieldArray::try_new(typed_value.column_by_name("id").unwrap()).unwrap();
2438 let session_id_field =
2439 ShreddedVariantFieldArray::try_new(typed_value.column_by_name("session_id").unwrap())
2440 .unwrap();
2441
2442 let id_value = id_field.value_column();
2443 let id_typed_value = id_field
2444 .typed_value_column()
2445 .unwrap()
2446 .as_any()
2447 .downcast_ref::<FixedSizeBinaryArray>()
2448 .unwrap();
2449 let session_id_value = session_id_field.value_column();
2450 let session_id_typed_value = session_id_field
2451 .typed_value_column()
2452 .unwrap()
2453 .as_any()
2454 .downcast_ref::<FixedSizeBinaryArray>()
2455 .unwrap();
2456
2457 assert!(result.is_valid(0));
2459
2460 assert!(value.is_null(0)); assert!(id_value.is_null(0));
2462 assert!(session_id_value.is_null(0));
2463
2464 assert!(typed_value.is_valid(0));
2465 assert!(id_typed_value.is_valid(0));
2466 assert!(session_id_typed_value.is_valid(0));
2467
2468 assert_eq!(id_typed_value.value(0), mock_uuid_1.as_bytes());
2469 assert_eq!(session_id_typed_value.value(0), mock_uuid_2.as_bytes());
2470
2471 assert!(result.is_valid(1));
2473
2474 assert!(value.is_valid(1)); assert!(typed_value.is_valid(1));
2476
2477 assert!(id_value.is_null(1));
2478 assert!(id_typed_value.is_valid(1));
2479 assert_eq!(id_typed_value.value(1), mock_uuid_2.as_bytes());
2480
2481 assert!(session_id_value.is_null(1));
2482 assert!(session_id_typed_value.is_valid(1));
2483 assert_eq!(session_id_typed_value.value(1), mock_uuid_3.as_bytes());
2484
2485 let row_1_variant = variant_from_arrays_at(metadata, value, 1).unwrap();
2487 let Variant::Object(obj) = row_1_variant else {
2488 panic!("Expected object");
2489 };
2490
2491 assert_eq!(obj.get("name"), Some(Variant::from("test_user")));
2492
2493 assert!(result.is_valid(2));
2495
2496 assert!(value.is_null(2)); assert!(typed_value.is_valid(2));
2498
2499 assert!(id_value.is_null(2));
2500 assert!(id_typed_value.is_valid(2));
2501 assert_eq!(id_typed_value.value(2), mock_uuid_1.as_bytes());
2502
2503 assert!(session_id_value.is_null(2));
2504 assert!(session_id_typed_value.is_null(2)); assert!(result.is_valid(3));
2508
2509 assert!(value.is_null(3)); assert!(typed_value.is_valid(3));
2511
2512 assert!(id_value.is_null(3));
2513 assert!(id_typed_value.is_valid(3));
2514 assert_eq!(id_typed_value.value(3), mock_uuid_3.as_bytes());
2515
2516 assert!(session_id_value.is_valid(3)); assert!(session_id_typed_value.is_null(3));
2518 let session_id_variant = variant_from_arrays_at(metadata, session_id_value, 3).unwrap();
2519 assert_eq!(session_id_variant, Variant::from("not-a-uuid"));
2520
2521 assert!(result.is_valid(4));
2523
2524 assert!(value.is_null(4)); assert!(typed_value.is_valid(4));
2526
2527 assert!(id_value.is_valid(4)); assert!(id_typed_value.is_null(4));
2529 let id_variant = variant_from_arrays_at(metadata, id_value, 4).unwrap();
2530 assert_eq!(id_variant, Variant::from(12345i64));
2531
2532 assert!(session_id_value.is_null(4));
2533 assert!(session_id_typed_value.is_valid(4));
2534 assert_eq!(session_id_typed_value.value(4), mock_uuid_1.as_bytes());
2535
2536 assert!(result.is_null(5));
2538
2539 Ok(())
2540 }
2541
2542 #[test]
2543 fn test_spec_compliance() {
2544 let input = VariantArray::from_iter(vec![Variant::from(42i64), Variant::from("hello")]);
2545
2546 let result = shred_variant(&input, &DataType::Int64).unwrap();
2547
2548 let inner_struct = result.inner();
2550 assert!(inner_struct.column_by_name("metadata").is_some());
2551 assert!(inner_struct.column_by_name("value").is_some());
2552 assert!(inner_struct.column_by_name("typed_value").is_some());
2553
2554 assert_eq!(
2556 result.metadata_column().len(),
2557 input.metadata_column().len()
2558 );
2559 assert_eq!(
2562 result.metadata_column().len(),
2563 input.metadata_column().len()
2564 );
2565
2566 assert_eq!(result.len(), input.len());
2568 assert!(result.typed_value_column().is_some());
2569
2570 let value_field = result.value_column();
2573 let typed_value_field = result
2574 .typed_value_column()
2575 .unwrap()
2576 .as_any()
2577 .downcast_ref::<Int64Array>()
2578 .unwrap();
2579
2580 for i in 0..result.len() {
2581 if !result.is_null(i) {
2582 let value_is_null = value_field.is_null(i);
2583 let typed_value_is_null = typed_value_field.is_null(i);
2584 assert!(
2586 value_is_null || typed_value_is_null,
2587 "Row {}: both value and typed_value are non-null for primitive shredding",
2588 i
2589 );
2590 }
2591 }
2592 }
2593
2594 #[test]
2595 fn test_variant_schema_builder_simple() -> Result<()> {
2596 let shredding_type = ShreddedSchemaBuilder::default()
2597 .with_path("a", &DataType::Int64)?
2598 .with_path("b", &DataType::Float64)?
2599 .build();
2600
2601 assert_eq!(
2602 shredding_type,
2603 DataType::Struct(Fields::from(vec![
2604 Field::new("a", DataType::Int64, true),
2605 Field::new("b", DataType::Float64, true),
2606 ]))
2607 );
2608
2609 Ok(())
2610 }
2611
2612 #[test]
2613 fn test_variant_schema_builder_nested() -> Result<()> {
2614 let shredding_type = ShreddedSchemaBuilder::default()
2615 .with_path("a", &DataType::Int64)?
2616 .with_path("b.c", &DataType::Utf8)?
2617 .with_path("b.d", &DataType::Float64)?
2618 .build();
2619
2620 assert_eq!(
2621 shredding_type,
2622 DataType::Struct(Fields::from(vec![
2623 Field::new("a", DataType::Int64, true),
2624 Field::new(
2625 "b",
2626 DataType::Struct(Fields::from(vec![
2627 Field::new("c", DataType::Utf8, true),
2628 Field::new("d", DataType::Float64, true),
2629 ])),
2630 true
2631 ),
2632 ]))
2633 );
2634
2635 Ok(())
2636 }
2637
2638 #[test]
2639 fn test_variant_schema_builder_with_path_variant_path_arg() -> Result<()> {
2640 let path = VariantPath::from_iter([VariantPathElement::from("a.b")]);
2641 let shredding_type = ShreddedSchemaBuilder::default()
2642 .with_path(path, &DataType::Int64)?
2643 .build();
2644
2645 match shredding_type {
2646 DataType::Struct(fields) => {
2647 assert_eq!(fields.len(), 1);
2648 assert_eq!(fields[0].name(), "a.b");
2649 assert_eq!(fields[0].data_type(), &DataType::Int64);
2650 }
2651 _ => panic!("expected struct data type"),
2652 }
2653
2654 Ok(())
2655 }
2656
2657 #[test]
2658 fn test_variant_schema_builder_custom_nullability() -> Result<()> {
2659 let shredding_type = ShreddedSchemaBuilder::default()
2660 .with_path(
2661 "foo",
2662 Arc::new(Field::new("should_be_renamed", DataType::Utf8, false)),
2663 )?
2664 .with_path("bar", (&DataType::Int64, false))?
2665 .build();
2666
2667 let DataType::Struct(fields) = shredding_type else {
2668 panic!("expected struct data type");
2669 };
2670
2671 let foo = fields.iter().find(|f| f.name() == "foo").unwrap();
2672 assert_eq!(foo.data_type(), &DataType::Utf8);
2673 assert!(!foo.is_nullable());
2674
2675 let bar = fields.iter().find(|f| f.name() == "bar").unwrap();
2676 assert_eq!(bar.data_type(), &DataType::Int64);
2677 assert!(!bar.is_nullable());
2678
2679 Ok(())
2680 }
2681
2682 #[test]
2683 fn test_variant_schema_builder_with_shred_variant() -> Result<()> {
2684 let input = build_variant_array(vec![
2685 VariantRow::Object(vec![
2686 ("time", VariantValue::from(1234567890i64)),
2687 ("hostname", VariantValue::from("server1")),
2688 ("extra", VariantValue::from(42)),
2689 ]),
2690 VariantRow::Object(vec![
2691 ("time", VariantValue::from(9876543210i64)),
2692 ("hostname", VariantValue::from("server2")),
2693 ]),
2694 VariantRow::Null,
2695 ]);
2696
2697 let shredding_type = ShreddedSchemaBuilder::default()
2698 .with_path("time", &DataType::Int64)?
2699 .with_path("hostname", &DataType::Utf8)?
2700 .build();
2701
2702 let result = shred_variant(&input, &shredding_type).unwrap();
2703
2704 assert_eq!(
2705 result.data_type(),
2706 &DataType::Struct(Fields::from(vec![
2707 Field::new("metadata", DataType::BinaryView, false),
2708 Field::new("value", DataType::BinaryView, true),
2709 Field::new(
2710 "typed_value",
2711 DataType::Struct(Fields::from(vec![
2712 Field::new(
2713 "hostname",
2714 DataType::Struct(Fields::from(vec![
2715 Field::new("value", DataType::BinaryView, true),
2716 Field::new("typed_value", DataType::Utf8, true),
2717 ])),
2718 false,
2719 ),
2720 Field::new(
2721 "time",
2722 DataType::Struct(Fields::from(vec![
2723 Field::new("value", DataType::BinaryView, true),
2724 Field::new("typed_value", DataType::Int64, true),
2725 ])),
2726 false,
2727 ),
2728 ])),
2729 true,
2730 ),
2731 ]))
2732 );
2733
2734 assert_eq!(result.len(), 3);
2735 assert!(result.typed_value_column().is_some());
2736
2737 let typed_value = result
2738 .typed_value_column()
2739 .unwrap()
2740 .as_any()
2741 .downcast_ref::<arrow::array::StructArray>()
2742 .unwrap();
2743
2744 let time_field =
2745 ShreddedVariantFieldArray::try_new(typed_value.column_by_name("time").unwrap())
2746 .unwrap();
2747 let hostname_field =
2748 ShreddedVariantFieldArray::try_new(typed_value.column_by_name("hostname").unwrap())
2749 .unwrap();
2750
2751 let time_typed = time_field
2752 .typed_value_column()
2753 .unwrap()
2754 .as_any()
2755 .downcast_ref::<Int64Array>()
2756 .unwrap();
2757 let hostname_typed = hostname_field
2758 .typed_value_column()
2759 .unwrap()
2760 .as_any()
2761 .downcast_ref::<arrow::array::StringArray>()
2762 .unwrap();
2763
2764 assert!(!result.is_null(0));
2766 assert_eq!(time_typed.value(0), 1234567890);
2767 assert_eq!(hostname_typed.value(0), "server1");
2768
2769 assert!(!result.is_null(1));
2771 assert_eq!(time_typed.value(1), 9876543210);
2772 assert_eq!(hostname_typed.value(1), "server2");
2773
2774 assert!(result.is_null(2));
2776
2777 Ok(())
2778 }
2779
2780 #[test]
2781 fn test_variant_schema_builder_conflicting_path() -> Result<()> {
2782 let shredding_type = ShreddedSchemaBuilder::default()
2783 .with_path("a", &DataType::Int64)?
2784 .with_path("a", &DataType::Float64)?
2785 .build();
2786
2787 assert_eq!(
2788 shredding_type,
2789 DataType::Struct(Fields::from(
2790 vec![Field::new("a", DataType::Float64, true),]
2791 ))
2792 );
2793
2794 Ok(())
2795 }
2796
2797 #[test]
2798 fn test_variant_schema_builder_root_path() -> Result<()> {
2799 let path = VariantPath::new(vec![]);
2800 let shredding_type = ShreddedSchemaBuilder::default()
2801 .with_path(path, &DataType::Int64)?
2802 .build();
2803
2804 assert_eq!(shredding_type, DataType::Int64);
2805
2806 Ok(())
2807 }
2808
2809 #[test]
2810 fn test_variant_schema_builder_empty_path() -> Result<()> {
2811 let shredding_type = ShreddedSchemaBuilder::default()
2812 .with_path("", &DataType::Int64)?
2813 .build();
2814
2815 assert_eq!(shredding_type, DataType::Int64);
2816 Ok(())
2817 }
2818
2819 #[test]
2820 fn test_variant_schema_builder_default() {
2821 let shredding_type = ShreddedSchemaBuilder::default().build();
2822 assert_eq!(shredding_type, DataType::Null);
2823 }
2824}