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]