From e41e60062e80a09af7ff859ba33087d770251f0a Mon Sep 17 00:00:00 2001 From: Liam Bao Date: Fri, 26 Dec 2025 18:05:15 +0800 Subject: [PATCH 1/3] [Variant] Support Shredded Lists/Array in `variant_get` --- parquet-variant-compute/src/shred_variant.rs | 203 +++--------------- .../src/variant_to_arrow.rs | 196 ++++++++++++++++- 2 files changed, 209 insertions(+), 190 deletions(-) diff --git a/parquet-variant-compute/src/shred_variant.rs b/parquet-variant-compute/src/shred_variant.rs index 45e7fc95c9f9..7f253d249dfb 100644 --- a/parquet-variant-compute/src/shred_variant.rs +++ b/parquet-variant-compute/src/shred_variant.rs @@ -19,19 +19,17 @@ use crate::variant_array::{ShreddedVariantFieldArray, StructArrayBuilder}; use crate::variant_to_arrow::{ - PrimitiveVariantToArrowRowBuilder, make_primitive_variant_to_arrow_row_builder, + ArrayVariantToArrowRowBuilder, PrimitiveVariantToArrowRowBuilder, + make_primitive_variant_to_arrow_row_builder, }; use crate::{VariantArray, VariantValueArrayBuilder}; -use arrow::array::{ - ArrayRef, BinaryViewArray, GenericListArray, GenericListViewArray, NullBufferBuilder, - OffsetSizeTrait, -}; -use arrow::buffer::{NullBuffer, OffsetBuffer, ScalarBuffer}; +use arrow::array::{ArrayRef, BinaryViewArray, NullBufferBuilder}; +use arrow::buffer::NullBuffer; use arrow::compute::CastOptions; -use arrow::datatypes::{ArrowNativeTypeOp, DataType, Field, FieldRef, Fields, TimeUnit}; +use arrow::datatypes::{DataType, Field, FieldRef, Fields, TimeUnit}; use arrow::error::{ArrowError, Result}; use indexmap::IndexMap; -use parquet_variant::{Variant, VariantBuilderExt, VariantList, VariantPath, VariantPathElement}; +use parquet_variant::{Variant, VariantBuilderExt, VariantPath, VariantPathElement}; use std::collections::BTreeMap; use std::sync::Arc; @@ -123,7 +121,8 @@ pub(crate) fn make_variant_to_shredded_variant_arrow_row_builder<'a>( DataType::List(_) | DataType::LargeList(_) | DataType::ListView(_) - | DataType::LargeListView(_) => { + | DataType::LargeListView(_) + | DataType::FixedSizeList(..) => { let typed_value_builder = VariantToShreddedArrayVariantRowBuilder::try_new( data_type, cast_options, @@ -131,11 +130,6 @@ pub(crate) fn make_variant_to_shredded_variant_arrow_row_builder<'a>( )?; VariantToShreddedVariantRowBuilder::Array(typed_value_builder) } - DataType::FixedSizeList(..) => { - return Err(ArrowError::NotYetImplemented( - "Shredding variant array values as fixed-size lists".to_string(), - )); - } // Supported shredded primitive types, see Variant shredding spec: // https://github.com/apache/parquet-format/blob/master/VariantShredding.md#shredded-value-types DataType::Boolean @@ -312,171 +306,6 @@ impl<'a> VariantToShreddedArrayVariantRowBuilder<'a> { } } -enum ArrayVariantToArrowRowBuilder<'a> { - List(VariantToListArrowRowBuilder<'a, i32, false>), - LargeList(VariantToListArrowRowBuilder<'a, i64, false>), - ListView(VariantToListArrowRowBuilder<'a, i32, true>), - LargeListView(VariantToListArrowRowBuilder<'a, i64, true>), -} - -impl<'a> ArrayVariantToArrowRowBuilder<'a> { - fn try_new( - data_type: &'a DataType, - cast_options: &'a CastOptions, - capacity: usize, - ) -> Result { - use ArrayVariantToArrowRowBuilder::*; - - // Make List/ListView builders without repeating the constructor boilerplate. - macro_rules! make_list_builder { - ($variant:ident, $offset:ty, $is_view:expr, $field:ident) => { - $variant(VariantToListArrowRowBuilder::<$offset, $is_view>::try_new( - $field.clone(), - $field.data_type(), - cast_options, - capacity, - )?) - }; - } - - let builder = match data_type { - DataType::List(field) => make_list_builder!(List, i32, false, field), - DataType::LargeList(field) => make_list_builder!(LargeList, i64, false, field), - DataType::ListView(field) => make_list_builder!(ListView, i32, true, field), - DataType::LargeListView(field) => make_list_builder!(LargeListView, i64, true, field), - other => { - return Err(ArrowError::InvalidArgumentError(format!( - "Casting to {other:?} is not applicable for array Variant types" - ))); - } - }; - Ok(builder) - } - - fn append_null(&mut self) { - match self { - Self::List(builder) => builder.append_null(), - Self::LargeList(builder) => builder.append_null(), - Self::ListView(builder) => builder.append_null(), - Self::LargeListView(builder) => builder.append_null(), - } - } - - fn append_value(&mut self, list: VariantList<'_, '_>) -> Result<()> { - match self { - Self::List(builder) => builder.append_value(list), - Self::LargeList(builder) => builder.append_value(list), - Self::ListView(builder) => builder.append_value(list), - Self::LargeListView(builder) => builder.append_value(list), - } - } - - fn finish(self) -> Result { - match self { - Self::List(builder) => builder.finish(), - Self::LargeList(builder) => builder.finish(), - Self::ListView(builder) => builder.finish(), - Self::LargeListView(builder) => builder.finish(), - } - } -} - -struct VariantToListArrowRowBuilder<'a, O, const IS_VIEW: bool> -where - O: OffsetSizeTrait + ArrowNativeTypeOp, -{ - field: FieldRef, - offsets: Vec, - element_builder: Box>, - nulls: NullBufferBuilder, - current_offset: O, -} - -impl<'a, O, const IS_VIEW: bool> VariantToListArrowRowBuilder<'a, O, IS_VIEW> -where - O: OffsetSizeTrait + ArrowNativeTypeOp, -{ - fn try_new( - field: FieldRef, - element_data_type: &'a DataType, - cast_options: &'a CastOptions, - capacity: usize, - ) -> Result { - if capacity >= isize::MAX as usize { - return Err(ArrowError::ComputeError( - "Capacity exceeds isize::MAX when reserving list offsets".to_string(), - )); - } - let mut offsets = Vec::with_capacity(capacity + 1); - offsets.push(O::ZERO); - let element_builder = make_variant_to_shredded_variant_arrow_row_builder( - element_data_type, - cast_options, - capacity, - false, - )?; - Ok(Self { - field, - offsets, - element_builder: Box::new(element_builder), - nulls: NullBufferBuilder::new(capacity), - current_offset: O::ZERO, - }) - } - - fn append_null(&mut self) { - self.offsets.push(self.current_offset); - self.nulls.append_null(); - } - - fn append_value(&mut self, list: VariantList<'_, '_>) -> Result<()> { - for element in list.iter() { - self.element_builder.append_value(element)?; - self.current_offset = self.current_offset.add_checked(O::ONE)?; - } - self.offsets.push(self.current_offset); - self.nulls.append_non_null(); - Ok(()) - } - - fn finish(mut self) -> Result { - let (value, typed_value, nulls) = self.element_builder.finish()?; - let element_array = - ShreddedVariantFieldArray::from_parts(Some(value), Some(typed_value), nulls); - let field = Arc::new( - self.field - .as_ref() - .clone() - .with_data_type(element_array.data_type().clone()), - ); - - if IS_VIEW { - // NOTE: `offsets` is never empty (constructor pushes an entry) - let mut sizes = Vec::with_capacity(self.offsets.len() - 1); - for i in 1..self.offsets.len() { - sizes.push(self.offsets[i] - self.offsets[i - 1]); - } - self.offsets.pop(); - let list_view_array = GenericListViewArray::::new( - field, - ScalarBuffer::from(self.offsets), - ScalarBuffer::from(sizes), - ArrayRef::from(element_array), - self.nulls.finish(), - ); - Ok(Arc::new(list_view_array)) - } else { - let list_array = GenericListArray::::new( - field, - OffsetBuffer::::new(ScalarBuffer::from(self.offsets)), - ArrayRef::from(element_array), - self.nulls.finish(), - ); - Ok(Arc::new(list_array)) - } - } -} - pub(crate) struct VariantToShreddedObjectVariantRowBuilder<'a> { value_builder: VariantValueArrayBuilder, typed_value_builders: IndexMap<&'a str, VariantToShreddedVariantRowBuilder<'a>>, @@ -1513,6 +1342,22 @@ mod tests { ); } + #[test] + fn test_array_shredding_as_fixed_size_list() { + let input = build_variant_array(vec![VariantRow::List(vec![ + VariantValue::from(1i64), + VariantValue::from(2i64), + VariantValue::from(3i64), + ])]); + let list_schema = + DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int64, true)), 2); + let err = shred_variant(&input, &list_schema).unwrap_err(); + assert_eq!( + err.to_string(), + "Not yet implemented: Converting unshredded variant arrays to arrow fixed-size lists" + ); + } + #[test] fn test_array_shredding_with_array_elements() { let input = build_variant_array(vec![ diff --git a/parquet-variant-compute/src/variant_to_arrow.rs b/parquet-variant-compute/src/variant_to_arrow.rs index 57d9944bb527..e3523c6fb09c 100644 --- a/parquet-variant-compute/src/variant_to_arrow.rs +++ b/parquet-variant-compute/src/variant_to_arrow.rs @@ -15,22 +15,26 @@ // specific language governing permissions and limitations // under the License. -use arrow::array::{ - ArrayRef, BinaryBuilder, BinaryLikeArrayBuilder, BinaryViewArray, BinaryViewBuilder, - BooleanBuilder, FixedSizeBinaryBuilder, LargeBinaryBuilder, LargeStringBuilder, NullArray, - NullBufferBuilder, PrimitiveBuilder, StringBuilder, StringLikeArrayBuilder, StringViewBuilder, +use crate::shred_variant::{ + VariantToShreddedVariantRowBuilder, make_variant_to_shredded_variant_arrow_row_builder, }; -use arrow::compute::{CastOptions, DecimalCast}; -use arrow::datatypes::{self, DataType, DecimalType}; -use arrow::error::{ArrowError, Result}; -use parquet_variant::{Variant, VariantPath}; - use crate::type_conversion::{ PrimitiveFromVariant, TimestampFromVariant, variant_to_unscaled_decimal, }; +use crate::variant_array::ShreddedVariantFieldArray; use crate::{VariantArray, VariantValueArrayBuilder}; - -use arrow_schema::TimeUnit; +use arrow::array::{ + ArrayRef, BinaryBuilder, BinaryLikeArrayBuilder, BinaryViewArray, BinaryViewBuilder, + BooleanBuilder, FixedSizeBinaryBuilder, GenericListArray, GenericListViewArray, + LargeBinaryBuilder, LargeStringBuilder, NullArray, NullBufferBuilder, OffsetSizeTrait, + PrimitiveBuilder, StringBuilder, StringLikeArrayBuilder, StringViewBuilder, +}; +use arrow::buffer::{OffsetBuffer, ScalarBuffer}; +use arrow::compute::{CastOptions, DecimalCast}; +use arrow::datatypes::{self, ArrowNativeTypeOp, DataType, DecimalType}; +use arrow::error::{ArrowError, Result}; +use arrow_schema::{FieldRef, TimeUnit}; +use parquet_variant::{Variant, VariantList, VariantPath}; use std::sync::Arc; /// Builder for converting primitive variant values to Arrow arrays. It is used by both @@ -427,6 +431,80 @@ pub(crate) fn make_primitive_variant_to_arrow_row_builder<'a>( Ok(builder) } +pub(crate) enum ArrayVariantToArrowRowBuilder<'a> { + List(VariantToListArrowRowBuilder<'a, i32, false>), + LargeList(VariantToListArrowRowBuilder<'a, i64, false>), + ListView(VariantToListArrowRowBuilder<'a, i32, true>), + LargeListView(VariantToListArrowRowBuilder<'a, i64, true>), +} + +impl<'a> ArrayVariantToArrowRowBuilder<'a> { + pub(crate) fn try_new( + data_type: &'a DataType, + cast_options: &'a CastOptions, + capacity: usize, + ) -> Result { + use ArrayVariantToArrowRowBuilder::*; + + // Make List/ListView builders without repeating the constructor boilerplate. + macro_rules! make_list_builder { + ($variant:ident, $offset:ty, $is_view:expr, $field:ident) => { + $variant(VariantToListArrowRowBuilder::<$offset, $is_view>::try_new( + $field.clone(), + $field.data_type(), + cast_options, + capacity, + )?) + }; + } + + let builder = match data_type { + DataType::List(field) => make_list_builder!(List, i32, false, field), + DataType::LargeList(field) => make_list_builder!(LargeList, i64, false, field), + DataType::ListView(field) => make_list_builder!(ListView, i32, true, field), + DataType::LargeListView(field) => make_list_builder!(LargeListView, i64, true, field), + DataType::FixedSizeList(..) => { + return Err(ArrowError::NotYetImplemented( + "Converting unshredded variant arrays to arrow fixed-size lists".to_string(), + )); + } + other => { + return Err(ArrowError::InvalidArgumentError(format!( + "Casting to {other:?} is not applicable for array Variant types" + ))); + } + }; + Ok(builder) + } + + pub(crate) fn append_null(&mut self) { + match self { + Self::List(builder) => builder.append_null(), + Self::LargeList(builder) => builder.append_null(), + Self::ListView(builder) => builder.append_null(), + Self::LargeListView(builder) => builder.append_null(), + } + } + + pub(crate) fn append_value(&mut self, list: VariantList<'_, '_>) -> Result<()> { + match self { + Self::List(builder) => builder.append_value(list), + Self::LargeList(builder) => builder.append_value(list), + Self::ListView(builder) => builder.append_value(list), + Self::LargeListView(builder) => builder.append_value(list), + } + } + + pub(crate) fn finish(self) -> Result { + match self { + Self::List(builder) => builder.finish(), + Self::LargeList(builder) => builder.finish(), + Self::ListView(builder) => builder.finish(), + Self::LargeListView(builder) => builder.finish(), + } + } +} + pub(crate) fn make_variant_to_arrow_row_builder<'a>( metadata: &BinaryViewArray, path: VariantPath<'a>, @@ -750,6 +828,102 @@ impl VariantToBinaryVariantArrowRowBuilder { } } +pub(crate) struct VariantToListArrowRowBuilder<'a, O, const IS_VIEW: bool> +where + O: OffsetSizeTrait + ArrowNativeTypeOp, +{ + field: FieldRef, + offsets: Vec, + element_builder: Box>, + nulls: NullBufferBuilder, + current_offset: O, +} + +impl<'a, O, const IS_VIEW: bool> VariantToListArrowRowBuilder<'a, O, IS_VIEW> +where + O: OffsetSizeTrait + ArrowNativeTypeOp, +{ + fn try_new( + field: FieldRef, + element_data_type: &'a DataType, + cast_options: &'a CastOptions, + capacity: usize, + ) -> Result { + if capacity >= isize::MAX as usize { + return Err(ArrowError::ComputeError( + "Capacity exceeds isize::MAX when reserving list offsets".to_string(), + )); + } + let mut offsets = Vec::with_capacity(capacity + 1); + offsets.push(O::ZERO); + let element_builder = make_variant_to_shredded_variant_arrow_row_builder( + element_data_type, + cast_options, + capacity, + false, + )?; + Ok(Self { + field, + offsets, + element_builder: Box::new(element_builder), + nulls: NullBufferBuilder::new(capacity), + current_offset: O::ZERO, + }) + } + + fn append_null(&mut self) { + self.offsets.push(self.current_offset); + self.nulls.append_null(); + } + + fn append_value(&mut self, list: VariantList<'_, '_>) -> Result<()> { + for element in list.iter() { + self.element_builder.append_value(element)?; + self.current_offset = self.current_offset.add_checked(O::ONE)?; + } + self.offsets.push(self.current_offset); + self.nulls.append_non_null(); + Ok(()) + } + + fn finish(mut self) -> Result { + let (value, typed_value, nulls) = self.element_builder.finish()?; + let element_array = + ShreddedVariantFieldArray::from_parts(Some(value), Some(typed_value), nulls); + let field = Arc::new( + self.field + .as_ref() + .clone() + .with_data_type(element_array.data_type().clone()), + ); + + if IS_VIEW { + // NOTE: `offsets` is never empty (constructor pushes an entry) + let mut sizes = Vec::with_capacity(self.offsets.len() - 1); + for i in 1..self.offsets.len() { + sizes.push(self.offsets[i] - self.offsets[i - 1]); + } + self.offsets.pop(); + let list_view_array = GenericListViewArray::::new( + field, + ScalarBuffer::from(self.offsets), + ScalarBuffer::from(sizes), + ArrayRef::from(element_array), + self.nulls.finish(), + ); + Ok(Arc::new(list_view_array)) + } else { + let list_array = GenericListArray::::new( + field, + OffsetBuffer::::new(ScalarBuffer::from(self.offsets)), + ArrayRef::from(element_array), + self.nulls.finish(), + ); + Ok(Arc::new(list_array)) + } + } +} + #[derive(Default)] struct FakeNullBuilder { item_count: usize, From 108824e4177a89020b5c8151f46696660b62c07a Mon Sep 17 00:00:00 2001 From: Liam Bao Date: Fri, 26 Dec 2025 20:14:27 +0800 Subject: [PATCH 2/3] Reorder varaint_to_arrow --- .../src/variant_to_arrow.rs | 180 +++++++++--------- 1 file changed, 90 insertions(+), 90 deletions(-) diff --git a/parquet-variant-compute/src/variant_to_arrow.rs b/parquet-variant-compute/src/variant_to_arrow.rs index e3523c6fb09c..dcde5bb5e9ce 100644 --- a/parquet-variant-compute/src/variant_to_arrow.rs +++ b/parquet-variant-compute/src/variant_to_arrow.rs @@ -37,6 +37,96 @@ use arrow_schema::{FieldRef, TimeUnit}; use parquet_variant::{Variant, VariantList, VariantPath}; use std::sync::Arc; +/// Builder for converting variant values into strongly typed Arrow arrays. +/// +/// Useful for variant_get kernels that need to extract specific paths from variant values, possibly +/// with casting of leaf values to specific types. +pub(crate) enum VariantToArrowRowBuilder<'a> { + Primitive(PrimitiveVariantToArrowRowBuilder<'a>), + BinaryVariant(VariantToBinaryVariantArrowRowBuilder), + + // Path extraction wrapper - contains a boxed enum for any of the above + WithPath(VariantPathRowBuilder<'a>), +} + +impl<'a> VariantToArrowRowBuilder<'a> { + pub fn append_null(&mut self) -> Result<()> { + use VariantToArrowRowBuilder::*; + match self { + Primitive(b) => b.append_null(), + BinaryVariant(b) => b.append_null(), + WithPath(path_builder) => path_builder.append_null(), + } + } + + pub fn append_value(&mut self, value: Variant<'_, '_>) -> Result { + use VariantToArrowRowBuilder::*; + match self { + Primitive(b) => b.append_value(&value), + BinaryVariant(b) => b.append_value(value), + WithPath(path_builder) => path_builder.append_value(value), + } + } + + pub fn finish(self) -> Result { + use VariantToArrowRowBuilder::*; + match self { + Primitive(b) => b.finish(), + BinaryVariant(b) => b.finish(), + WithPath(path_builder) => path_builder.finish(), + } + } +} + +pub(crate) fn make_variant_to_arrow_row_builder<'a>( + metadata: &BinaryViewArray, + path: VariantPath<'a>, + data_type: Option<&'a DataType>, + cast_options: &'a CastOptions, + capacity: usize, +) -> Result> { + use VariantToArrowRowBuilder::*; + + let mut builder = match data_type { + // If no data type was requested, build an unshredded VariantArray. + None => BinaryVariant(VariantToBinaryVariantArrowRowBuilder::new( + metadata.clone(), + capacity, + )), + Some(DataType::Struct(_)) => { + return Err(ArrowError::NotYetImplemented( + "Converting unshredded variant objects to arrow structs".to_string(), + )); + } + Some( + DataType::List(_) + | DataType::LargeList(_) + | DataType::ListView(_) + | DataType::LargeListView(_) + | DataType::FixedSizeList(..), + ) => { + return Err(ArrowError::NotYetImplemented( + "Converting unshredded variant arrays to arrow lists".to_string(), + )); + } + Some(data_type) => { + let builder = + make_primitive_variant_to_arrow_row_builder(data_type, cast_options, capacity)?; + Primitive(builder) + } + }; + + // Wrap with path extraction if needed + if !path.is_empty() { + builder = WithPath(VariantPathRowBuilder { + builder: Box::new(builder), + path, + }) + }; + + Ok(builder) +} + /// Builder for converting primitive variant values to Arrow arrays. It is used by both /// `VariantToArrowRowBuilder` (below) and `VariantToShreddedPrimitiveVariantRowBuilder` (in /// `shred_variant.rs`). @@ -85,18 +175,6 @@ pub(crate) enum PrimitiveVariantToArrowRowBuilder<'a> { BinaryView(VariantToBinaryArrowRowBuilder<'a, BinaryViewBuilder>), } -/// Builder for converting variant values into strongly typed Arrow arrays. -/// -/// Useful for variant_get kernels that need to extract specific paths from variant values, possibly -/// with casting of leaf values to specific types. -pub(crate) enum VariantToArrowRowBuilder<'a> { - Primitive(PrimitiveVariantToArrowRowBuilder<'a>), - BinaryVariant(VariantToBinaryVariantArrowRowBuilder), - - // Path extraction wrapper - contains a boxed enum for any of the above - WithPath(VariantPathRowBuilder<'a>), -} - impl<'a> PrimitiveVariantToArrowRowBuilder<'a> { pub fn append_null(&mut self) -> Result<()> { use PrimitiveVariantToArrowRowBuilder::*; @@ -231,35 +309,6 @@ impl<'a> PrimitiveVariantToArrowRowBuilder<'a> { } } -impl<'a> VariantToArrowRowBuilder<'a> { - pub fn append_null(&mut self) -> Result<()> { - use VariantToArrowRowBuilder::*; - match self { - Primitive(b) => b.append_null(), - BinaryVariant(b) => b.append_null(), - WithPath(path_builder) => path_builder.append_null(), - } - } - - pub fn append_value(&mut self, value: Variant<'_, '_>) -> Result { - use VariantToArrowRowBuilder::*; - match self { - Primitive(b) => b.append_value(&value), - BinaryVariant(b) => b.append_value(value), - WithPath(path_builder) => path_builder.append_value(value), - } - } - - pub fn finish(self) -> Result { - use VariantToArrowRowBuilder::*; - match self { - Primitive(b) => b.finish(), - BinaryVariant(b) => b.finish(), - WithPath(path_builder) => path_builder.finish(), - } - } -} - /// Creates a row builder that converts primitive `Variant` values into the requested Arrow data type. pub(crate) fn make_primitive_variant_to_arrow_row_builder<'a>( data_type: &'a DataType, @@ -505,55 +554,6 @@ impl<'a> ArrayVariantToArrowRowBuilder<'a> { } } -pub(crate) fn make_variant_to_arrow_row_builder<'a>( - metadata: &BinaryViewArray, - path: VariantPath<'a>, - data_type: Option<&'a DataType>, - cast_options: &'a CastOptions, - capacity: usize, -) -> Result> { - use VariantToArrowRowBuilder::*; - - let mut builder = match data_type { - // If no data type was requested, build an unshredded VariantArray. - None => BinaryVariant(VariantToBinaryVariantArrowRowBuilder::new( - metadata.clone(), - capacity, - )), - Some(DataType::Struct(_)) => { - return Err(ArrowError::NotYetImplemented( - "Converting unshredded variant objects to arrow structs".to_string(), - )); - } - Some( - DataType::List(_) - | DataType::LargeList(_) - | DataType::ListView(_) - | DataType::LargeListView(_) - | DataType::FixedSizeList(..), - ) => { - return Err(ArrowError::NotYetImplemented( - "Converting unshredded variant arrays to arrow lists".to_string(), - )); - } - Some(data_type) => { - let builder = - make_primitive_variant_to_arrow_row_builder(data_type, cast_options, capacity)?; - Primitive(builder) - } - }; - - // Wrap with path extraction if needed - if !path.is_empty() { - builder = WithPath(VariantPathRowBuilder { - builder: Box::new(builder), - path, - }) - }; - - Ok(builder) -} - /// A thin wrapper whose only job is to extract a specific path from a variant value and pass the /// result to a nested builder. pub(crate) struct VariantPathRowBuilder<'a> { From 1a741688ffe60944ba13394040da69092aa5ddcb Mon Sep 17 00:00:00 2001 From: Liam Bao Date: Fri, 26 Dec 2025 21:09:16 +0800 Subject: [PATCH 3/3] Support list in variant_get --- parquet-variant-compute/src/shred_variant.rs | 7 +- parquet-variant-compute/src/variant_get.rs | 187 +++++++++++++++++- .../src/variant_to_arrow.rs | 59 ++++-- 3 files changed, 226 insertions(+), 27 deletions(-) diff --git a/parquet-variant-compute/src/shred_variant.rs b/parquet-variant-compute/src/shred_variant.rs index 7f253d249dfb..d82eb15af598 100644 --- a/parquet-variant-compute/src/shred_variant.rs +++ b/parquet-variant-compute/src/shred_variant.rs @@ -274,7 +274,7 @@ impl<'a> VariantToShreddedArrayVariantRowBuilder<'a> { fn append_null(&mut self) -> Result<()> { self.value_builder.append_value(Variant::Null); - self.typed_value_builder.append_null(); + self.typed_value_builder.append_null()?; Ok(()) } @@ -284,12 +284,13 @@ impl<'a> VariantToShreddedArrayVariantRowBuilder<'a> { match variant { Variant::List(list) => { self.value_builder.append_null(); - self.typed_value_builder.append_value(list)?; + self.typed_value_builder + .append_value(&Variant::List(list))?; Ok(true) } other => { self.value_builder.append_value(other); - self.typed_value_builder.append_null(); + self.typed_value_builder.append_null()?; Ok(false) } } diff --git a/parquet-variant-compute/src/variant_get.rs b/parquet-variant-compute/src/variant_get.rs index 624c8ae128dc..0c3599b17d68 100644 --- a/parquet-variant-compute/src/variant_get.rs +++ b/parquet-variant-compute/src/variant_get.rs @@ -339,10 +339,11 @@ mod test { Array, ArrayRef, AsArray, BinaryArray, BinaryViewArray, BooleanArray, Date32Array, Date64Array, Decimal32Array, Decimal64Array, Decimal128Array, Decimal256Array, Float32Array, Float64Array, Int8Array, Int16Array, Int32Array, Int64Array, - LargeBinaryArray, LargeStringArray, NullBuilder, StringArray, StringViewArray, StructArray, + LargeBinaryArray, LargeListArray, LargeListViewArray, LargeStringArray, ListArray, + ListViewArray, NullBuilder, StringArray, StringViewArray, StructArray, Time32MillisecondArray, Time32SecondArray, Time64MicrosecondArray, Time64NanosecondArray, }; - use arrow::buffer::NullBuffer; + use arrow::buffer::{NullBuffer, OffsetBuffer, ScalarBuffer}; use arrow::compute::CastOptions; use arrow::datatypes::DataType::{Int16, Int32, Int64}; use arrow::datatypes::i256; @@ -351,8 +352,8 @@ mod test { use arrow_schema::{DataType, Field, FieldRef, Fields, IntervalUnit, TimeUnit}; use chrono::DateTime; use parquet_variant::{ - EMPTY_VARIANT_METADATA_BYTES, Variant, VariantDecimal4, VariantDecimal8, VariantDecimal16, - VariantDecimalType, VariantPath, + EMPTY_VARIANT_METADATA_BYTES, Variant, VariantBuilder, VariantDecimal4, VariantDecimal8, + VariantDecimal16, VariantDecimalType, VariantPath, }; fn single_variant_get_test(input_json: &str, path: VariantPath, expected_json: &str) { @@ -4158,4 +4159,182 @@ mod test { assert!(inner_values_result.is_null(1)); assert_eq!(inner_values_result.value(2), 333); } + + #[test] + fn test_variant_get_list_like_safe_cast() { + let string_array: ArrayRef = Arc::new(StringArray::from(vec![ + r#"[1, "two", 3]"#, + "\"not a list\"", + ])); + let variant_array = ArrayRef::from(json_to_variant(&string_array).unwrap()); + + let value_array: ArrayRef = { + let mut builder = VariantBuilder::new(); + builder.append_value("two"); + let (_, value_bytes) = builder.finish(); + Arc::new(BinaryViewArray::from(vec![ + None, + Some(value_bytes.as_slice()), + None, + ])) + }; + let typed_value_array: ArrayRef = Arc::new(Int64Array::from(vec![Some(1), None, Some(3)])); + let struct_fields = Fields::from(vec![ + Field::new("value", DataType::BinaryView, true), + Field::new("typed_value", DataType::Int64, true), + ]); + let struct_array: ArrayRef = Arc::new( + StructArray::try_new( + struct_fields.clone(), + vec![value_array.clone(), typed_value_array.clone()], + None, + ) + .unwrap(), + ); + + let request_field = Arc::new(Field::new("item", DataType::Int64, true)); + let result_field = Arc::new(Field::new("item", DataType::Struct(struct_fields), true)); + + let expectations = vec![ + ( + DataType::List(request_field.clone()), + Arc::new(ListArray::new( + result_field.clone(), + OffsetBuffer::new(ScalarBuffer::from(vec![0, 3, 3])), + struct_array.clone(), + Some(NullBuffer::from(vec![true, false])), + )) as ArrayRef, + ), + ( + DataType::LargeList(request_field.clone()), + Arc::new(LargeListArray::new( + result_field.clone(), + OffsetBuffer::new(ScalarBuffer::from(vec![0, 3, 3])), + struct_array.clone(), + Some(NullBuffer::from(vec![true, false])), + )) as ArrayRef, + ), + ( + DataType::ListView(request_field.clone()), + Arc::new(ListViewArray::new( + result_field.clone(), + ScalarBuffer::from(vec![0, 3]), + ScalarBuffer::from(vec![3, 0]), + struct_array.clone(), + Some(NullBuffer::from(vec![true, false])), + )) as ArrayRef, + ), + ( + DataType::LargeListView(request_field), + Arc::new(LargeListViewArray::new( + result_field, + ScalarBuffer::from(vec![0, 3]), + ScalarBuffer::from(vec![3, 0]), + struct_array, + Some(NullBuffer::from(vec![true, false])), + )) as ArrayRef, + ), + ]; + + for (request_type, expected) in expectations { + let options = GetOptions::new().with_as_type(Some(FieldRef::from(Field::new( + "result", + request_type.clone(), + true, + )))); + + let result = variant_get(&variant_array, options).unwrap(); + assert_eq!(result.data_type(), expected.data_type()); + assert_eq!(&result, &expected); + } + } + + #[test] + fn test_variant_get_list_like_unsafe_cast_errors_on_element_mismatch() { + let string_array: ArrayRef = + Arc::new(StringArray::from(vec![r#"[1, "two", 3]"#, "[4, 5]"])); + let variant_array = ArrayRef::from(json_to_variant(&string_array).unwrap()); + let cast_options = CastOptions { + safe: false, + ..Default::default() + }; + + let item_field = Arc::new(Field::new("item", DataType::Int64, true)); + let request_types = vec![ + DataType::List(item_field.clone()), + DataType::LargeList(item_field.clone()), + DataType::ListView(item_field.clone()), + DataType::LargeListView(item_field), + ]; + + for request_type in request_types { + let options = GetOptions::new() + .with_as_type(Some(FieldRef::from(Field::new( + "result", + request_type.clone(), + true, + )))) + .with_cast_options(cast_options.clone()); + + let err = variant_get(&variant_array, options).unwrap_err(); + assert!( + err.to_string() + .contains("Failed to extract primitive of type Int64") + ); + } + } + + #[test] + fn test_variant_get_list_like_unsafe_cast_errors_on_non_list() { + let string_array: ArrayRef = Arc::new(StringArray::from(vec!["[1, 2]", "\"not a list\""])); + let variant_array = ArrayRef::from(json_to_variant(&string_array).unwrap()); + let cast_options = CastOptions { + safe: false, + ..Default::default() + }; + let item_field = Arc::new(Field::new("item", Int64, true)); + let data_types = vec![ + DataType::List(item_field.clone()), + DataType::LargeList(item_field.clone()), + DataType::ListView(item_field.clone()), + DataType::LargeListView(item_field), + ]; + + for data_type in data_types { + let options = GetOptions::new() + .with_as_type(Some(FieldRef::from(Field::new("result", data_type, true)))) + .with_cast_options(cast_options.clone()); + + let err = variant_get(&variant_array, options).unwrap_err(); + assert!( + err.to_string() + .contains("Failed to extract list from variant"), + ); + } + } + + #[test] + fn test_variant_get_fixed_size_list_not_implemented() { + let string_array: ArrayRef = Arc::new(StringArray::from(vec!["[1, 2]", "\"not a list\""])); + let variant_array = ArrayRef::from(json_to_variant(&string_array).unwrap()); + let item_field = Arc::new(Field::new("item", Int64, true)); + for safe in [true, false] { + let options = GetOptions::new() + .with_as_type(Some(FieldRef::from(Field::new( + "result", + DataType::FixedSizeList(item_field.clone(), 2), + true, + )))) + .with_cast_options(CastOptions { + safe, + ..Default::default() + }); + + let err = variant_get(&variant_array, options).unwrap_err(); + assert!( + err.to_string() + .contains("Converting unshredded variant arrays to arrow fixed-size lists") + ); + } + } } diff --git a/parquet-variant-compute/src/variant_to_arrow.rs b/parquet-variant-compute/src/variant_to_arrow.rs index dcde5bb5e9ce..8cac1963ddea 100644 --- a/parquet-variant-compute/src/variant_to_arrow.rs +++ b/parquet-variant-compute/src/variant_to_arrow.rs @@ -34,7 +34,7 @@ use arrow::compute::{CastOptions, DecimalCast}; use arrow::datatypes::{self, ArrowNativeTypeOp, DataType, DecimalType}; use arrow::error::{ArrowError, Result}; use arrow_schema::{FieldRef, TimeUnit}; -use parquet_variant::{Variant, VariantList, VariantPath}; +use parquet_variant::{Variant, VariantPath}; use std::sync::Arc; /// Builder for converting variant values into strongly typed Arrow arrays. @@ -43,6 +43,7 @@ use std::sync::Arc; /// with casting of leaf values to specific types. pub(crate) enum VariantToArrowRowBuilder<'a> { Primitive(PrimitiveVariantToArrowRowBuilder<'a>), + Array(ArrayVariantToArrowRowBuilder<'a>), BinaryVariant(VariantToBinaryVariantArrowRowBuilder), // Path extraction wrapper - contains a boxed enum for any of the above @@ -54,6 +55,7 @@ impl<'a> VariantToArrowRowBuilder<'a> { use VariantToArrowRowBuilder::*; match self { Primitive(b) => b.append_null(), + Array(b) => b.append_null(), BinaryVariant(b) => b.append_null(), WithPath(path_builder) => path_builder.append_null(), } @@ -63,6 +65,7 @@ impl<'a> VariantToArrowRowBuilder<'a> { use VariantToArrowRowBuilder::*; match self { Primitive(b) => b.append_value(&value), + Array(b) => b.append_value(&value), BinaryVariant(b) => b.append_value(value), WithPath(path_builder) => path_builder.append_value(value), } @@ -72,6 +75,7 @@ impl<'a> VariantToArrowRowBuilder<'a> { use VariantToArrowRowBuilder::*; match self { Primitive(b) => b.finish(), + Array(b) => b.finish(), BinaryVariant(b) => b.finish(), WithPath(path_builder) => path_builder.finish(), } @@ -99,15 +103,15 @@ pub(crate) fn make_variant_to_arrow_row_builder<'a>( )); } Some( - DataType::List(_) + data_type @ (DataType::List(_) | DataType::LargeList(_) | DataType::ListView(_) | DataType::LargeListView(_) - | DataType::FixedSizeList(..), + | DataType::FixedSizeList(..)), ) => { - return Err(ArrowError::NotYetImplemented( - "Converting unshredded variant arrays to arrow lists".to_string(), - )); + let builder = + ArrayVariantToArrowRowBuilder::try_new(data_type, cast_options, capacity)?; + Array(builder) } Some(data_type) => { let builder = @@ -526,7 +530,7 @@ impl<'a> ArrayVariantToArrowRowBuilder<'a> { Ok(builder) } - pub(crate) fn append_null(&mut self) { + pub(crate) fn append_null(&mut self) -> Result<()> { match self { Self::List(builder) => builder.append_null(), Self::LargeList(builder) => builder.append_null(), @@ -535,12 +539,12 @@ impl<'a> ArrayVariantToArrowRowBuilder<'a> { } } - pub(crate) fn append_value(&mut self, list: VariantList<'_, '_>) -> Result<()> { + pub(crate) fn append_value(&mut self, value: &Variant<'_, '_>) -> Result { match self { - Self::List(builder) => builder.append_value(list), - Self::LargeList(builder) => builder.append_value(list), - Self::ListView(builder) => builder.append_value(list), - Self::LargeListView(builder) => builder.append_value(list), + Self::List(builder) => builder.append_value(value), + Self::LargeList(builder) => builder.append_value(value), + Self::ListView(builder) => builder.append_value(value), + Self::LargeListView(builder) => builder.append_value(value), } } @@ -837,6 +841,7 @@ where element_builder: Box>, nulls: NullBufferBuilder, current_offset: O, + cast_options: &'a CastOptions<'a>, } impl<'a, O, const IS_VIEW: bool> VariantToListArrowRowBuilder<'a, O, IS_VIEW> @@ -868,22 +873,36 @@ where element_builder: Box::new(element_builder), nulls: NullBufferBuilder::new(capacity), current_offset: O::ZERO, + cast_options, }) } - fn append_null(&mut self) { + fn append_null(&mut self) -> Result<()> { self.offsets.push(self.current_offset); self.nulls.append_null(); + Ok(()) } - fn append_value(&mut self, list: VariantList<'_, '_>) -> Result<()> { - for element in list.iter() { - self.element_builder.append_value(element)?; - self.current_offset = self.current_offset.add_checked(O::ONE)?; + fn append_value(&mut self, value: &Variant<'_, '_>) -> Result { + match value { + Variant::List(list) => { + for element in list.iter() { + self.element_builder.append_value(element)?; + self.current_offset = self.current_offset.add_checked(O::ONE)?; + } + self.offsets.push(self.current_offset); + self.nulls.append_non_null(); + Ok(true) + } + _ if self.cast_options.safe => { + self.append_null()?; + Ok(false) + } + _ => Err(ArrowError::CastError(format!( + "Failed to extract list from variant {:?}", + value + ))), } - self.offsets.push(self.current_offset); - self.nulls.append_non_null(); - Ok(()) } fn finish(mut self) -> Result {