From 1bc743162e824c665bb65efc4b34b1f5b8d3e4f5 Mon Sep 17 00:00:00 2001 From: sdf-jkl Date: Thu, 9 Apr 2026 10:59:04 -0400 Subject: [PATCH 1/8] Add `variant_get` access as `Variant` --- parquet-variant-compute/src/variant_get.rs | 96 +++++++++++++++++++++- 1 file changed, 93 insertions(+), 3 deletions(-) diff --git a/parquet-variant-compute/src/variant_get.rs b/parquet-variant-compute/src/variant_get.rs index 29e28c850be6..599561146d4e 100644 --- a/parquet-variant-compute/src/variant_get.rs +++ b/parquet-variant-compute/src/variant_get.rs @@ -20,12 +20,13 @@ use arrow::{ datatypes::Field, error::Result, }; +use arrow_schema::extension::ExtensionType; use arrow_schema::{ArrowError, DataType, FieldRef}; use parquet_variant::{VariantPath, VariantPathElement}; -use crate::VariantArray; use crate::variant_array::BorrowedShreddingState; use crate::variant_to_arrow::make_variant_to_arrow_row_builder; +use crate::{VariantArray, VariantType, unshred_variant}; use arrow::array::AsArray; use std::sync::Arc; @@ -109,6 +110,11 @@ pub(crate) fn follow_shredded_path_element<'a>( } } +fn is_variant_extension(field: &Field) -> bool { + field.extension_type_name() == Some(VariantType::NAME) + && field.try_extension_type::().is_ok() +} + /// Follows the given path as far as possible through shredded variant fields. If the path ends on a /// shredded field, return it directly. Otherwise, use a row shredder to follow the rest of the path /// and extract the requested value on a per-row basis. @@ -131,7 +137,22 @@ fn shredded_get_path( // Helper that shreds a VariantArray to a specific type. let shred_basic_variant = |target: VariantArray, path: VariantPath<'_>, as_field: Option<&Field>| { - let as_type = as_field.map(|f| f.data_type()); + let requested_variant = as_field.is_some_and(is_variant_extension); + let target = if requested_variant { + unshred_variant(&target)? + } else { + target + }; + + if requested_variant && path.is_empty() { + return Ok(ArrayRef::from(target)); + } + + let as_type = if requested_variant { + None + } else { + as_field.map(|f| f.data_type()) + }; let mut builder = make_variant_to_arrow_row_builder( target.metadata_field(), path, @@ -179,6 +200,16 @@ fn shredded_get_path( } ShreddedPathStep::Missing => { let num_rows = input.len(); + if as_field.is_some_and(is_variant_extension) { + let all_nulls = Some(arrow::buffer::NullBuffer::from(vec![false; num_rows])); + let arr = VariantArray::from_parts( + input.metadata_field().clone(), + None, + None, + all_nulls, + ); + return Ok(ArrayRef::from(arr)); + } let arr = match as_field.map(|f| f.data_type()) { Some(data_type) => array::new_null_array(data_type, num_rows), None => Arc::new(array::NullArray::new(num_rows)) as _, @@ -222,7 +253,9 @@ fn shredded_get_path( // // For shredded/partially-shredded targets (`typed_value` present), recurse into each field // separately to take advantage of deeper shredding in child fields. - if let DataType::Struct(fields) = as_field.data_type() { + if !is_variant_extension(as_field) + && let DataType::Struct(fields) = as_field.data_type() + { if target.typed_value_field().is_none() { return shred_basic_variant(target, VariantPath::default(), Some(as_field)); } @@ -2038,6 +2071,63 @@ mod test { println!("Nested path 'a.x' result: {:?}", result); } + #[test] + fn test_variant_get_as_variant_from_unshredded_input() { + let (unshredded, _) = create_variant_get_as_variant_test_data(); + assert_variant_field_extraction_returns_unshredded_variant(&unshredded); + } + + #[test] + fn test_variant_get_as_variant_from_shredded_input() { + let (_, shredded) = create_variant_get_as_variant_test_data(); + assert_variant_field_extraction_returns_unshredded_variant(&shredded); + } + + fn create_variant_get_as_variant_test_data() -> (ArrayRef, ArrayRef) { + let input_json: ArrayRef = Arc::new(StringArray::from(vec![ + Some(r#"{"field_name": {"k": 100000}}"#), + Some(r#"{"field_name": {"k": "s"}}"#), + ])); + + let unshredded = ArrayRef::from(json_to_variant(&input_json).unwrap()); + let unshredded_variant = VariantArray::try_new(&unshredded).unwrap(); + + let as_type = DataType::Struct(Fields::from(vec![Field::new( + "field_name", + DataType::Struct(Fields::from(vec![Field::new("k", DataType::Int32, true)])), + true, + )])); + let shredded = ArrayRef::from(shred_variant(&unshredded_variant, &as_type).unwrap()); + + (unshredded, shredded) + } + + fn assert_variant_field_extraction_returns_unshredded_variant(input: &ArrayRef) { + let variant_output = VariantArray::try_new(input).unwrap().field("result"); + let options = GetOptions::new_with_path(VariantPath::try_from("field_name").unwrap()) + .with_as_type(Some(FieldRef::from(variant_output))); + + let result = variant_get(input, options).unwrap(); + let result_variant = VariantArray::try_new(&result).unwrap(); + + assert!(result_variant.typed_value_field().is_none()); + assert!(result_variant.value_field().is_some()); + + let expected_json: ArrayRef = Arc::new(StringArray::from(vec![ + Some(r#"{"k":100000}"#), + Some(r#"{"k":"s"}"#), + ])); + let expected = json_to_variant(&expected_json).unwrap(); + + assert_eq!(result_variant.len(), expected.len()); + for i in 0..result_variant.len() { + assert_eq!(result_variant.is_null(i), expected.is_null(i)); + if !result_variant.is_null(i) { + assert_eq!(result_variant.value(i), expected.value(i)); + } + } + } + /// Create test data for depth 0 (direct field access) /// [{"x": 42}, {"x": "foo"}, {"y": 10}] fn create_depth_0_test_data() -> ArrayRef { From 5bf820684cad2ae6747ed91839abb86a62570191 Mon Sep 17 00:00:00 2001 From: sdf-jkl Date: Fri, 24 Apr 2026 22:05:24 -0400 Subject: [PATCH 2/8] use `has_valid_extension_type` --- parquet-variant-compute/src/variant_get.rs | 13 ++++--------- 1 file changed, 4 insertions(+), 9 deletions(-) diff --git a/parquet-variant-compute/src/variant_get.rs b/parquet-variant-compute/src/variant_get.rs index eb82ea645b2a..313d1c11240a 100644 --- a/parquet-variant-compute/src/variant_get.rs +++ b/parquet-variant-compute/src/variant_get.rs @@ -21,7 +21,6 @@ use arrow::{ datatypes::Field, error::Result, }; -use arrow_schema::extension::ExtensionType; use arrow_schema::{ArrowError, DataType, FieldRef}; use parquet_variant::{VariantPath, VariantPathElement}; @@ -104,11 +103,6 @@ pub(crate) fn follow_shredded_path_element<'a>( } } -fn is_variant_extension(field: &Field) -> bool { - field.extension_type_name() == Some(VariantType::NAME) - && field.try_extension_type::().is_ok() -} - /// Follows the given path as far as possible through shredded variant fields. If the path ends on a /// shredded field, return it directly. Otherwise, use a row shredder to follow the rest of the path /// and extract the requested value on a per-row basis. @@ -131,7 +125,8 @@ fn shredded_get_path( // Helper that shreds a VariantArray to a specific type. let shred_basic_variant = |target: VariantArray, path: VariantPath<'_>, as_field: Option<&Field>| { - let requested_variant = as_field.is_some_and(is_variant_extension); + let requested_variant = + as_field.is_some_and(Field::has_valid_extension_type::); let target = if requested_variant { unshred_variant(&target)? } else { @@ -192,7 +187,7 @@ fn shredded_get_path( } ShreddedPathStep::Missing => { let num_rows = input.len(); - if as_field.is_some_and(is_variant_extension) { + if as_field.is_some_and(Field::has_valid_extension_type::) { let all_nulls = Some(arrow::buffer::NullBuffer::from(vec![false; num_rows])); let arr = VariantArray::from_parts( input.metadata_field().clone(), @@ -245,7 +240,7 @@ fn shredded_get_path( // // For shredded/partially-shredded targets (`typed_value` present), recurse into each field // separately to take advantage of deeper shredding in child fields. - if !is_variant_extension(as_field) + if !as_field.has_valid_extension_type::() && let DataType::Struct(fields) = as_field.data_type() { if target.typed_value_field().is_none() { From 2da7f8cb5129958467c617f672b6a9154f60aec4 Mon Sep 17 00:00:00 2001 From: sdf-jkl Date: Fri, 24 Apr 2026 23:01:32 -0400 Subject: [PATCH 3/8] fix MSRV --- parquet-variant-compute/src/variant_get.rs | 50 +++++++++++----------- 1 file changed, 25 insertions(+), 25 deletions(-) diff --git a/parquet-variant-compute/src/variant_get.rs b/parquet-variant-compute/src/variant_get.rs index 313d1c11240a..c3c4e2a2290d 100644 --- a/parquet-variant-compute/src/variant_get.rs +++ b/parquet-variant-compute/src/variant_get.rs @@ -240,32 +240,32 @@ fn shredded_get_path( // // For shredded/partially-shredded targets (`typed_value` present), recurse into each field // separately to take advantage of deeper shredding in child fields. - if !as_field.has_valid_extension_type::() - && let DataType::Struct(fields) = as_field.data_type() - { - if target.typed_value_field().is_none() { - return shred_basic_variant(target, VariantPath::default(), Some(as_field)); - } - - let children = fields - .iter() - .map(|field| { - shredded_get_path( - &target, - &[VariantPathElement::from(field.name().as_str())], - Some(field), - cast_options, - ) - }) - .collect::>>()?; - - let struct_nulls = target.nulls().cloned(); + if !as_field.has_valid_extension_type::() { + if let DataType::Struct(fields) = as_field.data_type() { + if target.typed_value_field().is_none() { + return shred_basic_variant(target, VariantPath::default(), Some(as_field)); + } - return Ok(Arc::new(StructArray::try_new( - fields.clone(), - children, - struct_nulls, - )?)); + let children = fields + .iter() + .map(|field| { + shredded_get_path( + &target, + &[VariantPathElement::from(field.name().as_str())], + Some(field), + cast_options, + ) + }) + .collect::>>()?; + + let struct_nulls = target.nulls().cloned(); + + return Ok(Arc::new(StructArray::try_new( + fields.clone(), + children, + struct_nulls, + )?)); + } } // Not a struct, so directly shred the variant as the requested type From 755b3fd2639e17ebc8ef27ed65ccbb0302fb4483 Mon Sep 17 00:00:00 2001 From: Konstantin Tarasov <33369833+sdf-jkl@users.noreply.github.com> Date: Tue, 26 May 2026 19:12:45 -0400 Subject: [PATCH 4/8] Apply suggestion from @scovich Co-authored-by: Ryan Johnson --- parquet-variant-compute/src/variant_get.rs | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/parquet-variant-compute/src/variant_get.rs b/parquet-variant-compute/src/variant_get.rs index da52a30fcbc7..c56f188ca790 100644 --- a/parquet-variant-compute/src/variant_get.rs +++ b/parquet-variant-compute/src/variant_get.rs @@ -258,12 +258,10 @@ fn shredded_get_path( }) .collect::>>()?; - let struct_nulls = target.nulls().cloned(); - return Ok(Arc::new(StructArray::try_new( fields.clone(), children, - struct_nulls, + target.nulls().cloned(), )?)); } } From 406d84cad64f2de165c75fea3d2e6b105cf14072 Mon Sep 17 00:00:00 2001 From: sdf-jkl Date: Thu, 28 May 2026 12:39:22 -0400 Subject: [PATCH 5/8] Add reasoning for propagating metadata column --- parquet-variant-compute/src/variant_get.rs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/parquet-variant-compute/src/variant_get.rs b/parquet-variant-compute/src/variant_get.rs index c56f188ca790..afe5f06b76d5 100644 --- a/parquet-variant-compute/src/variant_get.rs +++ b/parquet-variant-compute/src/variant_get.rs @@ -190,6 +190,8 @@ fn shredded_get_path( if as_field.is_some_and(Field::has_valid_extension_type::) { let all_nulls = Some(arrow::buffer::NullBuffer::from(vec![false; num_rows])); let arr = VariantArray::from_parts( + // Propagating metadata is not necessary for an all-NULL array, but is cheaper than constructing + // a new empty metadata array. (n * 3 bytes vs Arc bump) input.metadata_field().clone(), None, None, From fcdcd8f25a4c074d49499fdb9798240c7694eb78 Mon Sep 17 00:00:00 2001 From: Konstantin Tarasov <33369833+sdf-jkl@users.noreply.github.com> Date: Wed, 10 Jun 2026 15:43:33 -0400 Subject: [PATCH 6/8] Apply suggestions from code review Co-authored-by: Ryan Johnson --- parquet-variant-compute/src/variant_get.rs | 20 ++++++-------------- 1 file changed, 6 insertions(+), 14 deletions(-) diff --git a/parquet-variant-compute/src/variant_get.rs b/parquet-variant-compute/src/variant_get.rs index fdd1cb1c6b09..b2fcc3b88243 100644 --- a/parquet-variant-compute/src/variant_get.rs +++ b/parquet-variant-compute/src/variant_get.rs @@ -268,14 +268,10 @@ fn shredded_get_path( let num_rows = input.len(); if as_field.is_some_and(Field::has_valid_extension_type::) { let all_nulls = Some(arrow::buffer::NullBuffer::from(vec![false; num_rows])); - let arr = VariantArray::from_parts( - // Propagating metadata is not necessary for an all-NULL array, but is cheaper than constructing - // a new empty metadata array. (n * 3 bytes vs Arc bump) - input.metadata_field().clone(), - None, - None, - all_nulls, - ); + // Propagating metadata is not necessary for an all-NULL array, but is cheaper than constructing + // a new empty metadata array. (n * 3 bytes vs Arc bump) + let metadata = input.metadata_field().clone(); + Ok(ArrayRef::from(VariantArray::from_parts(metadata, None, None, all_nulls))) return Ok(ArrayRef::from(arr)); } let arr = match as_field.map(|f| f.data_type()) { @@ -330,12 +326,8 @@ fn shredded_get_path( let children = fields .iter() .map(|field| { - shredded_get_path( - &target, - &[VariantPathElement::from(field.name().as_str())], - Some(field), - cast_options, - ) + let path = &[VariantPathElement::from(field.name().as_str())]; + shredded_get_path(&target, path, Some(field), cast_options) }) .collect::>>()?; From 0277d912892f508a2f3e00d5d8a517ba2f865ecb Mon Sep 17 00:00:00 2001 From: sdf-jkl Date: Wed, 10 Jun 2026 15:59:38 -0400 Subject: [PATCH 7/8] fix --- parquet-variant-compute/src/variant_get.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/parquet-variant-compute/src/variant_get.rs b/parquet-variant-compute/src/variant_get.rs index b2fcc3b88243..7b7102ee15c2 100644 --- a/parquet-variant-compute/src/variant_get.rs +++ b/parquet-variant-compute/src/variant_get.rs @@ -271,8 +271,9 @@ fn shredded_get_path( // Propagating metadata is not necessary for an all-NULL array, but is cheaper than constructing // a new empty metadata array. (n * 3 bytes vs Arc bump) let metadata = input.metadata_field().clone(); - Ok(ArrayRef::from(VariantArray::from_parts(metadata, None, None, all_nulls))) - return Ok(ArrayRef::from(arr)); + return Ok(ArrayRef::from(VariantArray::from_parts( + metadata, None, None, all_nulls, + ))); } let arr = match as_field.map(|f| f.data_type()) { Some(data_type) => array::new_null_array(data_type, num_rows), From cf7e545f852623a6818bc520cf272711741a93a6 Mon Sep 17 00:00:00 2001 From: sdf-jkl Date: Thu, 11 Jun 2026 12:58:17 -0400 Subject: [PATCH 8/8] only unshred if as_type Variant Field is not shredded + nit --- parquet-variant-compute/src/variant_get.rs | 74 +++++++++++++++++++--- 1 file changed, 64 insertions(+), 10 deletions(-) diff --git a/parquet-variant-compute/src/variant_get.rs b/parquet-variant-compute/src/variant_get.rs index 7b7102ee15c2..04b026fac51b 100644 --- a/parquet-variant-compute/src/variant_get.rs +++ b/parquet-variant-compute/src/variant_get.rs @@ -201,17 +201,37 @@ fn shredded_get_path( VariantArray::from_parts(metadata, value, typed_value, accumulated_nulls) }; - // Helper that shreds a VariantArray to a specific type. + // Helper that extracts the value at `path` and casts it to the requested type, or returns it as + // an unshredded binary variant when `Variant` output is requested. let shred_basic_variant = |target: VariantArray, path: VariantPath<'_>, as_field: Option<&Field>| { + // A `VariantType` extension on `as_field` requests `Variant` output: return an + // unshredded binary variant instead of casting to a concrete Arrow type. let requested_variant = as_field.is_some_and(Field::has_valid_extension_type::); + + // A `typed_value` in that field requests shredded output -- a `VariantArray` with + // `typed_value` columns. We produce only unshredded variant output. Shredded output is + // tracked in https://github.com/apache/arrow-rs/issues/8153. Reject such a request + // instead of silently dropping the shredding it asked for. + if requested_variant && requested_field_is_shredded(as_field) { + return Err(ArrowError::NotYetImplemented( + "variant_get with shredded `Variant` output is not yet supported".to_string(), + )); + } + + // Collapse any shredding back to binary. Only the `NotShredded` step below passes a + // non-empty `path`, and there `target` is already a plain `value` column (no + // `typed_value`) -- so `unshred_variant` hits its clone fast-path, with nothing deeper + // to shred. The builder then walks any remaining path per-row, emitting variant output + // because `as_type` is `None`. let target = if requested_variant { unshred_variant(&target)? } else { target }; + // Path exhausted, variant requested: return the target directly. if requested_variant && path.is_empty() { return Ok(ArrayRef::from(target)); } @@ -271,9 +291,8 @@ fn shredded_get_path( // Propagating metadata is not necessary for an all-NULL array, but is cheaper than constructing // a new empty metadata array. (n * 3 bytes vs Arc bump) let metadata = input.metadata_field().clone(); - return Ok(ArrayRef::from(VariantArray::from_parts( - metadata, None, None, all_nulls, - ))); + let arr = VariantArray::from_parts(metadata, None, None, all_nulls); + return Ok(ArrayRef::from(arr)); } let arr = match as_field.map(|f| f.data_type()) { Some(data_type) => array::new_null_array(data_type, num_rows), @@ -344,6 +363,17 @@ fn shredded_get_path( shred_basic_variant(target, VariantPath::default(), Some(as_field)) } +/// Returns true if `as_field` requests *shredded* `Variant` output. +/// +/// Its struct carries a `typed_value` field naming the type to shred to. +/// A plain variant request has only `metadata` and `value`. +fn requested_field_is_shredded(as_field: Option<&Field>) -> bool { + as_field.is_some_and(|f| match f.data_type() { + DataType::Struct(fields) => fields.iter().any(|field| field.name() == "typed_value"), + _ => false, + }) +} + fn try_perfect_shredding(variant_array: &VariantArray, as_field: &Field) -> Option { // Try to return the typed value directly when we have a perfect shredding match. if matches!(as_field.data_type(), DataType::Struct(_)) { @@ -453,7 +483,7 @@ mod test { use std::str::FromStr; use std::sync::Arc; - use super::{GetOptions, variant_get}; + use super::{GetOptions, requested_field_is_shredded, variant_get}; use crate::variant_array::{ShreddedVariantFieldArray, StructArrayBuilder}; use crate::{ ShreddedSchemaBuilder, VariantArray, VariantArrayBuilder, cast_to_variant, json_to_variant, @@ -473,6 +503,7 @@ mod test { use arrow::datatypes::DataType::{Int16, Int32, Int64}; use arrow::datatypes::i256; use arrow::util::display::FormatOptions; + use arrow_schema::ArrowError; use arrow_schema::DataType::{Boolean, Float32, Float64, Int8}; use arrow_schema::{DataType, Field, FieldRef, Fields, IntervalUnit, TimeUnit}; use chrono::DateTime; @@ -2363,13 +2394,34 @@ mod test { #[test] fn test_variant_get_as_variant_from_unshredded_input() { let (unshredded, _) = create_variant_get_as_variant_test_data(); - assert_variant_field_extraction_returns_unshredded_variant(&unshredded); + let unshredded_field = VariantArray::try_new(&unshredded).unwrap().field("result"); + assert_variant_field_extraction_returns_unshredded_variant(&unshredded, &unshredded_field); } #[test] fn test_variant_get_as_variant_from_shredded_input() { + let (unshredded, shredded) = create_variant_get_as_variant_test_data(); + let unshredded_field = VariantArray::try_new(&unshredded).unwrap().field("result"); + assert_variant_field_extraction_returns_unshredded_variant(&shredded, &unshredded_field); + } + + #[test] + fn test_variant_get_as_shredded_variant_is_not_yet_supported() { let (_, shredded) = create_variant_get_as_variant_test_data(); - assert_variant_field_extraction_returns_unshredded_variant(&shredded); + // Deriving the request field from the shredded array yields a `VariantType` field whose + // struct carries a `typed_value` -- a request to shred the output. That is unsupported + // (https://github.com/apache/arrow-rs/issues/8153) and must error, not silently return a + // plain binary variant. + let shredded_field = VariantArray::try_new(&shredded).unwrap().field("result"); + assert!(requested_field_is_shredded(Some(&shredded_field))); + + let options = GetOptions::new_with_path(VariantPath::try_from("field_name").unwrap()) + .with_as_type(Some(FieldRef::from(shredded_field))); + let err = variant_get(&shredded, options).unwrap_err(); + assert!( + matches!(err, ArrowError::NotYetImplemented(_)), + "expected NotYetImplemented, got {err:?}" + ); } fn create_variant_get_as_variant_test_data() -> (ArrayRef, ArrayRef) { @@ -2391,10 +2443,12 @@ mod test { (unshredded, shredded) } - fn assert_variant_field_extraction_returns_unshredded_variant(input: &ArrayRef) { - let variant_output = VariantArray::try_new(input).unwrap().field("result"); + fn assert_variant_field_extraction_returns_unshredded_variant( + input: &ArrayRef, + variant_field: &Field, + ) { let options = GetOptions::new_with_path(VariantPath::try_from("field_name").unwrap()) - .with_as_type(Some(FieldRef::from(variant_output))); + .with_as_type(Some(FieldRef::from(variant_field.clone()))); let result = variant_get(input, options).unwrap(); let result_variant = VariantArray::try_new(&result).unwrap();