sunchao commented on code in PR #5407:
URL: https://github.com/apache/datafusion-comet/pull/5407#discussion_r3836756839


##########
native/core/src/parquet/cast_column.rs:
##########
@@ -176,6 +186,275 @@ fn cast_timestamp_micros_to_millis_scalar(
     ScalarValue::TimestampMillisecond(new_val, target_tz)
 }
 
+fn normalize_variant_array(
+    array: &ArrayRef,
+    target_field: &FieldRef,
+) -> DataFusionResult<ArrayRef> {
+    let DataType::Struct(fields) = target_field.data_type() else {
+        return Err(DataFusionError::Execution(
+            "Variant extension field must use Struct storage".to_string(),
+        ));
+    };
+    if fields.len() != 2
+        || fields[0].name() != "value"
+        || fields[1].name() != "metadata"
+        || fields
+            .iter()
+            .any(|field| field.data_type() != &DataType::Binary)
+    {
+        return Err(DataFusionError::Execution(
+            "Variant output must contain Binary children [value, 
metadata]".to_string(),
+        ));
+    }
+
+    let array = decode_variant_metadata_dictionary(array)?;
+    let variant = 
prepare_variant_for_unshredding(&VariantArray::try_new(array.as_ref())?)?;
+    let unshredded = unshred_variant(&variant)?;
+    let value = unshredded.value_field().ok_or_else(|| {
+        DataFusionError::Execution("Unshredded Variant is missing its value 
field".to_string())
+    })?;
+    let value = cast(value.as_ref(), &DataType::Binary)?;
+    let metadata = cast(unshredded.metadata_field().as_ref(), 
&DataType::Binary)?;
+    let value = reorder_variant_values(
+        &value,
+        &metadata,
+        unshredded.inner().nulls(),
+        VariantObjectKeyOrder::SparkUtf16,
+        false,
+    )?;
+    let output = StructArray::try_new(
+        fields.clone(),
+        vec![value, metadata],
+        unshredded.inner().nulls().cloned(),
+    )?;
+    Ok(Arc::new(output))
+}
+
+/// Arrow's unshredder fully validates any residual `value` in a partially 
shredded object. Spark
+/// writes object keys in Java UTF-16 order, so put that residual value in 
Arrow UTF-8 order only
+/// while it passes through the upstream unshredder.
+fn prepare_variant_for_unshredding(variant: &VariantArray) -> 
DataFusionResult<VariantArray> {
+    let (Some(value), Some(_)) = (variant.value_field(), 
variant.typed_value_field()) else {
+        return Ok(variant.clone());
+    };
+
+    let value = cast(value.as_ref(), &DataType::Binary)?;
+    let metadata = cast(variant.metadata_field().as_ref(), &DataType::Binary)?;
+    let value = reorder_variant_values(
+        &value,
+        &metadata,
+        variant.inner().nulls(),
+        VariantObjectKeyOrder::ArrowUtf8,
+        true,
+    )?;
+
+    let value_index = variant
+        .inner()
+        .fields()
+        .iter()
+        .position(|field| field.name() == "value")
+        .unwrap();
+    let mut fields = 
variant.inner().fields().iter().cloned().collect::<Vec<_>>();
+    fields[value_index] = Arc::new(
+        fields[value_index]
+            .as_ref()
+            .clone()
+            .with_data_type(DataType::Binary),
+    );
+    let mut columns = variant.inner().columns().to_vec();
+    columns[value_index] = value;
+    let array = StructArray::try_new(fields.into(), columns, 
variant.inner().nulls().cloned())?;
+    Ok(VariantArray::try_new(&array)?)
+}
+
+/// Arrow-rs parquet-variant-compute allows dictionary-encoded metadata in its 
contract, but 58.4's
+/// `VariantArray::try_new` validates only Binary, LargeBinary, and 
BinaryView. Decode just that
+/// child and keep the physical struct otherwise unchanged.
+/// 
https://github.com/apache/arrow-rs/blob/0ff81c1215cc026a1de93ce3d2078df1ecba6f09/parquet-variant-compute/src/variant_array.rs#L276-L310
+fn decode_variant_metadata_dictionary(array: &ArrayRef) -> 
DataFusionResult<ArrayRef> {
+    let Some(struct_array) = array.as_any().downcast_ref::<StructArray>() else 
{
+        return Ok(Arc::clone(array));
+    };
+    let Some((metadata_index, metadata_field)) = struct_array
+        .fields()
+        .iter()
+        .enumerate()
+        .find(|(_, field)| field.name() == "metadata")
+    else {
+        return Ok(Arc::clone(array));
+    };
+    let DataType::Dictionary(_, value_type) = metadata_field.data_type() else {
+        return Ok(Arc::clone(array));
+    };
+
+    let decoded = cast(struct_array.column(metadata_index).as_ref(), 
value_type)?;
+    let mut fields = struct_array.fields().iter().cloned().collect::<Vec<_>>();
+    fields[metadata_index] = Arc::new(
+        metadata_field
+            .as_ref()
+            .clone()
+            .with_data_type(decoded.data_type().clone()),
+    );
+    let mut columns = struct_array.columns().to_vec();
+    columns[metadata_index] = decoded;
+    Ok(Arc::new(StructArray::try_new(
+        fields.into(),
+        columns,
+        struct_array.nulls().cloned(),
+    )?))
+}
+
+/// Supplies sort-only field names whose Rust ordering matches Java 
`String.compareTo` ordering.
+/// The original metadata dictionary still supplies the field IDs written to 
the Variant value.
+#[derive(Debug)]
+struct SparkMetadataBuilder<'a, 'm> {
+    metadata: &'a VariantMetadata<'m>,
+    sort_keys: Vec<String>,
+}
+
+impl<'a, 'm> SparkMetadataBuilder<'a, 'm> {
+    fn new(metadata: &'a VariantMetadata<'m>) -> Self {
+        let sort_keys = metadata
+            .iter()
+            .map(|field_name| {
+                field_name
+                    .encode_utf16()
+                    .map(|unit| char::from_u32(0x10000 + 
u32::from(unit)).unwrap())
+                    .collect()
+            })
+            .collect();
+        Self {
+            metadata,
+            sort_keys,
+        }
+    }
+}
+
+impl MetadataBuilder for SparkMetadataBuilder<'_, '_> {
+    fn try_upsert_field_name(&mut self, field_name: &str) -> Result<u32, 
ArrowError> {
+        self.metadata
+            .get_entry(field_name)
+            .map(|(field_id, _)| field_id)
+            .ok_or_else(|| {
+                ArrowError::InvalidArgumentError(format!(
+                    "Field name '{field_name}' not found in metadata 
dictionary"
+                ))
+            })
+    }
+
+    fn field_name(&self, field_id: usize) -> &str {
+        &self.sort_keys[field_id]
+    }
+
+    fn num_field_names(&self) -> usize {
+        self.metadata.len()
+    }
+
+    fn truncate_field_names(&mut self, new_size: usize) {
+        debug_assert_eq!(self.metadata.len(), new_size);
+    }
+
+    fn finish(&mut self) -> usize {
+        self.metadata.size()
+    }
+}
+
+#[derive(Clone, Copy)]
+enum VariantObjectKeyOrder {
+    ArrowUtf8,
+    SparkUtf16,
+}
+
+fn is_compatible_variant(variant: &Variant<'_, '_>, order: 
VariantObjectKeyOrder) -> bool {
+    match variant {
+        Variant::Object(object) => {
+            let mut previous = None;
+            object.iter().all(|(name, value)| {
+                let ordered = previous
+                    .map(|previous: &str| match order {
+                        VariantObjectKeyOrder::ArrowUtf8 => previous <= name,
+                        VariantObjectKeyOrder::SparkUtf16 => {
+                            previous.encode_utf16().cmp(name.encode_utf16())
+                                != std::cmp::Ordering::Greater
+                        }
+                    })
+                    .unwrap_or(true);
+                previous = Some(name);
+                ordered && is_compatible_variant(&value, order)
+            })
+        }
+        Variant::List(list) => list
+            .iter()
+            .all(|value| is_compatible_variant(&value, order)),
+        _ => true,
+    }
+}
+
+/// Reorder object keys for either Arrow's UTF-8 order or Spark's Java UTF-16 
order. Preserve
+/// already-compatible values byte-for-byte and retain the original metadata 
dictionary.
+fn reorder_variant_values(
+    value: &ArrayRef,
+    metadata: &ArrayRef,
+    parent_nulls: Option<&NullBuffer>,
+    order: VariantObjectKeyOrder,
+    allow_null_value: bool,
+) -> DataFusionResult<ArrayRef> {
+    let value = value.as_any().downcast_ref::<BinaryArray>().unwrap();
+    let metadata = metadata.as_any().downcast_ref::<BinaryArray>().unwrap();
+    let mut output = BinaryBuilder::new();
+
+    for index in 0..value.len() {
+        if parent_nulls.is_some_and(|nulls| nulls.is_null(index)) {
+            output.append_null();
+            continue;
+        }
+        if value.is_null(index) {
+            if allow_null_value {
+                output.append_null();
+                continue;
+            }
+            return Err(DataFusionError::Execution(format!(
+                "Variant value is null at row {index}"
+            )));
+        }
+        if metadata.is_null(index) {
+            return Err(DataFusionError::Execution(format!(
+                "Variant metadata is null at row {index}"
+            )));
+        }
+
+        let metadata = VariantMetadata::try_new(metadata.value(index))?;

Review Comment:
   [P2] Allow empty object keys in Variant metadata
   
   Could we preserve valid empty dictionary entries here? On Spark 4.0.4, 
writing `parse_json('{"":1}')` to ordinary unshredded Parquet and then reading 
`v` with `spark.sql.variant.pushVariantIntoScan=false` selects 
`CometNativeScan` but now fails with `offsets not monotonically increasing`. 
Spark reads the same file, and the previous normalizer accepts the identical 
bytes. Spark encodes this empty key with metadata `01 01 00 00`: one dictionary 
entry with two equal offsets. Arrow/Parquet 58.4.0's full metadata validator 
requires strictly increasing offsets when the sorted bit is unset, so this new 
unconditional call rejects the value even though no key reordering is needed. A 
nested empty key fails the same way; ordinary keys and empty string values 
pass. Please allow these valid empty keys and add a native-read regression.



##########
spark/src/test/resources/sql-tests/expressions/misc/variant.sql:
##########
@@ -55,6 +81,65 @@ SELECT id FROM test_variant WHERE variant_get(v, '$.a', 
'int') = 1
 query expect_fallback(type VariantType)
 SELECT COUNT(*) FROM test_variant WHERE v IS NOT NULL
 
+query expect_fallback(type VariantType)
+SELECT CAST(v AS STRING) FROM test_variant
+
+-- A Variant existence default is read from Spark's table schema and applied 
only when an old
+-- Parquet file does not contain the column. variant_get remains a Spark 
expression, while the
+-- ordinary Parquet scan and missing-column substitution stay native.
+statement
+CREATE TABLE test_variant_defaults_sql(id INT) USING parquet
+
+statement
+INSERT INTO test_variant_defaults_sql VALUES (1)
+
+statement
+ALTER TABLE test_variant_defaults_sql ADD COLUMNS(
+  v VARIANT DEFAULT parse_json('{"a":1}'), n INT DEFAULT 7)
+
+statement
+INSERT INTO test_variant_defaults_sql VALUES (2, parse_json('{"a":2}'), 8)
+
+statement
+SET spark.sql.parquet.enableVectorizedReader=false
+
+statement
+SET spark.comet.scan.allowDisabledParquetVectorizedReader=true
+
+query expect_fallback(type VariantType)
+SELECT id, variant_get(v, '$.a', 'int') AS a, n
+FROM test_variant_defaults_sql ORDER BY id
+
+statement
+SET spark.sql.parquet.enableVectorizedReader=true
+
+statement
+SET spark.comet.scan.allowDisabledParquetVectorizedReader=false
+
+-- Arrow and Spark order supplementary Unicode object keys differently. Force 
a shredded field so
+-- the native scan reconstructs the whole 32-field value before Spark's 
variant_get binary search.
+statement
+SET spark.sql.variant.writeShredding.enabled=true
+
+statement
+SET spark.sql.variant.forceShreddingSchemaForTest=k00 BIGINT

Review Comment:
   [P2] Scope the forced shredding schema to this SQL fixture
   
   Could we scope or restore `spark.sql.variant.forceShreddingSchemaForTest`? 
Because this key is absent from the fixture's header configs, the runner does 
not restore this `SET` when `variant.sql` finishes. On Spark 4.1/4.2, 
`writeShredding.enabled` is then restored to `true`, and later ordinary Parquet 
writes enter Spark's test-only forced-schema path. Both the [4.1 expression 
job](https://github.com/apache/datafusion-comet/actions/runs/32589534235/job/97071757612)
 and [4.2 expression 
job](https://github.com/apache/datafusion-comet/actions/runs/32589534235/job/97071757703)
 show `variant.sql` passing followed by 13 other fixture failures, starting 
with `lag_lead.sql`: `ParquetWriteSupport.writeFields` throws `Index 3 out of 
bounds for length 3`. Running the unchanged fixture through its actual runner 
reproduces a passing ordinary write before it, the same failing write 
afterward, and recovery after unsetting only this key. Please include this 
setting in the fixture's scoped configs 
 or restore it so subsequent tests retain their original configuration.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to