diff --git a/dev/diffs/4.1.3.diff b/dev/diffs/4.1.3.diff index a8d9383ced1..623d62d5bb3 100644 --- a/dev/diffs/4.1.3.diff +++ b/dev/diffs/4.1.3.diff @@ -3303,29 +3303,6 @@ index 09ed6955a51..52d998aab46 100644 ) } test(s"parquet widening conversion $fromType -> $toType") { -diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetVariantShreddingSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetVariantShreddingSuite.scala -index 1cc6d3afbee..1ca791bc0cc 100644 ---- a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetVariantShreddingSuite.scala -+++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetVariantShreddingSuite.scala -@@ -29,7 +29,7 @@ import org.apache.parquet.schema.{LogicalTypeAnnotation, PrimitiveType, Type} - import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName - - import org.apache.spark.SparkException --import org.apache.spark.sql.{AnalysisException, QueryTest, Row} -+import org.apache.spark.sql.{AnalysisException, IgnoreComet, QueryTest, Row} - import org.apache.spark.sql.internal.SQLConf - import org.apache.spark.sql.internal.SQLConf.ParquetOutputTimestampType - import org.apache.spark.sql.test.SharedSparkSession -@@ -230,7 +230,8 @@ class ParquetVariantShreddingSuite extends QueryTest with ParquetTest with Share - } - } - -- test("variant logical type annotation - ignore variant annotation") { -+ test("variant logical type annotation - ignore variant annotation", -+ IgnoreComet("https://github.com/apache/datafusion-comet/issues/5741")) { - Seq(true, false).foreach { ignoreVariantAnnotation => - withSQLConf(SQLConf.PARQUET_ANNOTATE_VARIANT_LOGICAL_TYPE.key -> "true", - SQLConf.PARQUET_IGNORE_VARIANT_ANNOTATION.key -> ignoreVariantAnnotation.toString, diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/debug/DebuggingSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/debug/DebuggingSuite.scala index b8f3ea3c6f3..bbd44221288 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/debug/DebuggingSuite.scala diff --git a/dev/diffs/4.2.0.diff b/dev/diffs/4.2.0.diff index 60ad822d6d2..27f636bcfd5 100644 --- a/dev/diffs/4.2.0.diff +++ b/dev/diffs/4.2.0.diff @@ -3384,29 +3384,6 @@ index 7ccd664f6c7..0f16c2f2fad 100644 ) } test(s"parquet widening conversion $fromType -> $toType") { -diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetVariantShreddingSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetVariantShreddingSuite.scala -index 5c0739b6276..ab12b78b7db 100644 ---- a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetVariantShreddingSuite.scala -+++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetVariantShreddingSuite.scala -@@ -29,7 +29,7 @@ import org.apache.parquet.schema.{LogicalTypeAnnotation, PrimitiveType, Type} - import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName - - import org.apache.spark.SparkException --import org.apache.spark.sql.{AnalysisException, Row} -+import org.apache.spark.sql.{AnalysisException, IgnoreComet, Row} - import org.apache.spark.sql.internal.SQLConf - import org.apache.spark.sql.internal.SQLConf.ParquetOutputTimestampType - import org.apache.spark.sql.test.SharedSparkSession -@@ -230,7 +230,8 @@ class ParquetVariantShreddingSuite extends ParquetTest with SharedSparkSession { - } - } - -- test("variant logical type annotation - ignore variant annotation") { -+ test("variant logical type annotation - ignore variant annotation", -+ IgnoreComet("https://github.com/apache/datafusion-comet/issues/5741")) { - Seq(true, false).foreach { ignoreVariantAnnotation => - withSQLConf(SQLConf.PARQUET_ANNOTATE_VARIANT_LOGICAL_TYPE.key -> "true", - SQLConf.PARQUET_IGNORE_VARIANT_ANNOTATION.key -> ignoreVariantAnnotation.toString, diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/debug/DebuggingSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/debug/DebuggingSuite.scala index b8f3ea3c6f3..bbd44221288 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/debug/DebuggingSuite.scala diff --git a/native/common/src/error.rs b/native/common/src/error.rs index 8e6385065a3..25274d8cdd9 100644 --- a/native/common/src/error.rs +++ b/native/common/src/error.rs @@ -240,6 +240,18 @@ pub enum SparkError { #[error("Spark read schema expects field Ids, but Parquet file schema doesn't contain any field Ids. Please remove the field ids from Spark schema or ignore missing ids by setting `spark.sql.parquet.fieldId.read.ignoreMissing = true`")] ParquetMissingFieldIds { file_path: String }, + /// A Parquet field carries the VARIANT logical type annotation but the requested Spark read + /// type is not `VariantType`, while `spark.sql.parquet.ignoreVariantAnnotation` is false. + /// Mirrors the `checkConversionRequirement` in Spark's + /// `ParquetToSparkSchemaConverter.convertGroupField`, which raises `_LEGACY_ERROR_TEMP_3071` + /// during schema conversion -- before any row is read, so an empty file is rejected too. + #[error("[_LEGACY_ERROR_TEMP_3071] Invalid Spark read type: expected {column} to be variant type but found {spark_type}")] + ParquetVariantAnnotationMismatch { + file_path: String, + column: String, + spark_type: String, + }, + /// Schema mismatch when reading a Parquet column under a requested schema /// that's incompatible with the physical column type. Translated by the JVM /// shim into Spark's `SchemaColumnConvertNotSupportedException`. The @@ -355,6 +367,9 @@ impl SparkError { SparkError::DuplicateFieldCaseInsensitive { .. } => "DuplicateFieldCaseInsensitive", SparkError::DuplicateFieldByFieldId { .. } => "DuplicateFieldByFieldId", SparkError::ParquetMissingFieldIds { .. } => "ParquetMissingFieldIds", + SparkError::ParquetVariantAnnotationMismatch { .. } => { + "ParquetVariantAnnotationMismatch" + } SparkError::ParquetSchemaConvert { .. } => "ParquetSchemaConvert", SparkError::CannotReadFile { .. } => "CannotReadFile", SparkError::Arrow(_) => "Arrow", @@ -605,6 +620,17 @@ impl SparkError { "matchedFields": matched_fields, }) } + SparkError::ParquetVariantAnnotationMismatch { + file_path, + column, + spark_type, + } => { + serde_json::json!({ + "filePath": file_path, + "column": column, + "sparkType": spark_type, + }) + } SparkError::ParquetMissingFieldIds { file_path } => { serde_json::json!({ "filePath": file_path, @@ -723,6 +749,13 @@ impl SparkError { // the FAILED_READ_FILE SparkException Spark 4 raises at the task boundary. SparkError::ParquetMissingFieldIds { .. } => "java/lang/RuntimeException", + // ParquetVariantAnnotationMismatch - the shim rebuilds Spark's AnalysisException and + // wraps it in a FAILED_READ_FILE SparkException, matching what Spark's own file scan + // produces when its schema converter rejects the read type. + SparkError::ParquetVariantAnnotationMismatch { .. } => { + "org/apache/spark/sql/AnalysisException" + } + // ParquetSchemaConvert - converted to SchemaColumnConvertNotSupportedException by the shim SparkError::ParquetSchemaConvert { .. } => { "org/apache/spark/sql/execution/datasources/SchemaColumnConvertNotSupportedException" @@ -828,6 +861,10 @@ impl SparkError { // Parquet schema mismatch — translated to SchemaColumnConvertNotSupportedException // by the JVM shim. The shim wraps it in the version-appropriate // SparkException error class, so no error class is exposed here. + // ParquetVariantAnnotationMismatch — the shim rebuilds Spark's AnalysisException with + // its own error class and wraps it via cannotReadFilesError, so none is exposed here. + SparkError::ParquetVariantAnnotationMismatch { .. } => None, + SparkError::ParquetSchemaConvert { .. } => None, // CannotReadFile — the JVM shim wraps it via cannotReadFilesError, which supplies the diff --git a/native/core/src/execution/operators/dynamic_filter/join/tests.rs b/native/core/src/execution/operators/dynamic_filter/join/tests.rs index 5e2a4f39235..ae470e3a3e8 100644 --- a/native/core/src/execution/operators/dynamic_filter/join/tests.rs +++ b/native/core/src/execution/operators/dynamic_filter/join/tests.rs @@ -670,6 +670,7 @@ fn parquet_probe( false, false, false, + false, ) .unwrap(); (file, scan) @@ -775,6 +776,7 @@ async fn reader_filter_crosses_null_check_conjunction_and_retains_residual() { false, false, false, + false, ) .unwrap(); let checks = [("key", 0), ("payload", 1), ("other", 2)].map(|(name, index)| { diff --git a/native/core/src/execution/operators/dynamic_filter/join/tests/schema_errors.rs b/native/core/src/execution/operators/dynamic_filter/join/tests/schema_errors.rs index 8fa236f23d8..2b587e58e2d 100644 --- a/native/core/src/execution/operators/dynamic_filter/join/tests/schema_errors.rs +++ b/native/core/src/execution/operators/dynamic_filter/join/tests/schema_errors.rs @@ -86,6 +86,7 @@ fn scan( false, false, false, + false, ) .unwrap() } diff --git a/native/core/src/execution/operators/dynamic_filter/join/tests/schema_errors/partition_columns.rs b/native/core/src/execution/operators/dynamic_filter/join/tests/schema_errors/partition_columns.rs index 3a4d55ae06e..6f97666d7d9 100644 --- a/native/core/src/execution/operators/dynamic_filter/join/tests/schema_errors/partition_columns.rs +++ b/native/core/src/execution/operators/dynamic_filter/join/tests/schema_errors/partition_columns.rs @@ -71,6 +71,7 @@ fn partitioned_scan( false, false, false, + false, ) .unwrap() } diff --git a/native/core/src/execution/operators/dynamic_filter/join/tests/timestamp_errors.rs b/native/core/src/execution/operators/dynamic_filter/join/tests/timestamp_errors.rs index 2da9af30cb5..242664331fe 100644 --- a/native/core/src/execution/operators/dynamic_filter/join/tests/timestamp_errors.rs +++ b/native/core/src/execution/operators/dynamic_filter/join/tests/timestamp_errors.rs @@ -105,6 +105,7 @@ async fn assert_timestamp_overflow_preserved(nested: bool) { false, false, false, + false, ) .unwrap(); let join = single_key_join_plans( diff --git a/native/core/src/execution/operators/dynamic_filter/topk/tests.rs b/native/core/src/execution/operators/dynamic_filter/topk/tests.rs index 46e72071386..50aed9c818b 100644 --- a/native/core/src/execution/operators/dynamic_filter/topk/tests.rs +++ b/native/core/src/execution/operators/dynamic_filter/topk/tests.rs @@ -196,6 +196,7 @@ fn parquet_scan( false, false, false, + false, ) .unwrap() } diff --git a/native/core/src/execution/operators/dynamic_filter/topk/tests/timestamp.rs b/native/core/src/execution/operators/dynamic_filter/topk/tests/timestamp.rs index cdefa621c43..909dd50583c 100644 --- a/native/core/src/execution/operators/dynamic_filter/topk/tests/timestamp.rs +++ b/native/core/src/execution/operators/dynamic_filter/topk/tests/timestamp.rs @@ -95,6 +95,7 @@ fn timestamp_input( false, false, false, + false, ) .unwrap(); (file, scan) diff --git a/native/core/src/execution/operators/iceberg_scan.rs b/native/core/src/execution/operators/iceberg_scan.rs index 9f465329ca0..e57925eff51 100644 --- a/native/core/src/execution/operators/iceberg_scan.rs +++ b/native/core/src/execution/operators/iceberg_scan.rs @@ -230,6 +230,14 @@ impl IcebergScanExec { let scan_metrics = scan_result.metrics().clone(); let stream = scan_result.stream(); + // `ignore_variant_annotation` stays at its default of false, so the Parquet VARIANT + // annotation check in `SparkPhysicalExprAdapterFactory::create` is always enforced here. + // `spark.sql.parquet.ignoreVariantAnnotation` is a Parquet-datasource conf that the + // Iceberg read path never consults, so there is deliberately no escape hatch: an Iceberg + // table whose file annotates a field as VARIANT while the table schema asks for an + // ordinary struct is the same defect apache/datafusion-comet#5741 describes. Reaching it + // needs a table schema that disagrees with its own data files, since `CometScanRule` falls + // back for any Iceberg table carrying a Variant column. let spark_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); let adapter_factory = SparkPhysicalExprAdapterFactory::new(spark_options, None); diff --git a/native/core/src/execution/planner.rs b/native/core/src/execution/planner.rs index 2f992c4843f..7017cc7cde3 100644 --- a/native/core/src/execution/planner.rs +++ b/native/core/src/execution/planner.rs @@ -1803,6 +1803,7 @@ impl PhysicalPlanner { common.encryption_enabled, common.use_field_id, common.require_field_ids, + common.ignore_variant_annotation, )?; Ok(( vec![], diff --git a/native/core/src/parquet/eager_page_index_reader_factory.rs b/native/core/src/parquet/eager_page_index_reader_factory.rs index cadcf4c9f4e..0bcf5df2199 100644 --- a/native/core/src/parquet/eager_page_index_reader_factory.rs +++ b/native/core/src/parquet/eager_page_index_reader_factory.rs @@ -60,7 +60,7 @@ //! The second is the Variant footer rewrite, `with_spark_arrow_schema`, which replaces the //! Arrow schema hint in the footer for scans that project Variant. -use arrow::datatypes::{DataType, FieldRef, Schema}; +use arrow::datatypes::{DataType, Field, FieldRef, Fields, Schema}; use async_trait::async_trait; use bytes::Bytes; use datafusion::common::Result as DFResult; @@ -93,6 +93,7 @@ use parquet::file::metadata::{ FooterTail, PageIndexPolicy, ParquetMetaData, ParquetMetaDataReader, }; use parquet::schema::types::{ColumnDescPtr, SchemaDescriptor, Type as ParquetType}; +use parquet::variant::VariantType; use std::fmt::{Debug, Display, Formatter}; use std::ops::Range; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; @@ -364,22 +365,158 @@ fn spark_enum_schema(schema: &SchemaDescriptor) -> ParquetResult> ))) } -/// Arrow restores advisory `ARROW:schema` types that can differ from Spark's physical Parquet -/// interpretation. Replace that hint with physical inference and the ENUM string mapping. -/// Rebuild only the returned metadata; the shared cache retains the original footer. -/// https://github.com/apache/datafusion-comet/issues/5477 -fn with_spark_arrow_schema(metadata: Arc) -> ParquetResult> { +/// `arrow_schema::extension`'s keys for an extension type's name and metadata. The `arrow` facade +/// does not re-export that module, and `Field` offers no way to drop an extension type, so +/// clearing one means removing these two keys, as elsewhere in this crate. +const EXTENSION_TYPE_NAME_KEY: &str = "ARROW:extension:name"; +const EXTENSION_TYPE_METADATA_KEY: &str = "ARROW:extension:metadata"; + +/// Pair `hinted` with `physical` and take the Variant marker from `physical` wherever the two +/// disagree, returning `None` when every marker already agrees. `physical` is what the Parquet +/// annotations alone say, so this makes the annotation the marker's only source of truth. +/// +/// The two schemas describe the same Parquet leaves, so fields pair up by position, which is how +/// `complex.rs` pairs them when it applies the hint. A hint may still name a different container +/// (`LargeList` for a `List`) or a different field count, and where the shapes disagree the walk +/// leaves that subtree's markers alone rather than guessing: the reader honors the hint's shape, +/// and a hint that far from the file is not one this rewrite can speak for. +fn align_variant_markers(hinted: &Schema, physical: &Schema) -> Option { + if hinted.fields().len() != physical.fields().len() { + return None; + } + let mut changed = false; + let fields = hinted + .fields() + .iter() + .zip(physical.fields()) + .map(|(hinted, physical)| align_variant_field(hinted, physical, &mut changed)) + .collect::(); + changed.then(|| Schema::new_with_metadata(fields, hinted.metadata().clone())) +} + +/// Puts a realigned child back into the container it came from, leaving the container kind as the +/// hint declared it. +type RebuildContainer = fn(FieldRef, &DataType) -> DataType; + +/// The single child of a container field, paired with the rebuild that puts a new child back in +/// the same container. `None` for a field that holds no nested field. +fn variant_marker_child(data_type: &DataType) -> Option<(&FieldRef, RebuildContainer)> { + let rebuild: RebuildContainer = |child, original| match original { + DataType::List(_) => DataType::List(child), + DataType::LargeList(_) => DataType::LargeList(child), + DataType::ListView(_) => DataType::ListView(child), + DataType::LargeListView(_) => DataType::LargeListView(child), + DataType::FixedSizeList(_, len) => DataType::FixedSizeList(child, *len), + DataType::Map(_, sorted) => DataType::Map(child, *sorted), + other => other.clone(), + }; + match data_type { + DataType::List(child) + | DataType::LargeList(child) + | DataType::ListView(child) + | DataType::LargeListView(child) + | DataType::FixedSizeList(child, _) + | DataType::Map(child, _) => Some((child, rebuild)), + _ => None, + } +} + +/// One field of [`align_variant_markers`]. Sets `changed` when it moves a marker, so the caller +/// can leave the footer alone when the hint and the annotations already agree. +fn align_variant_field(hinted: &FieldRef, physical: &FieldRef, changed: &mut bool) -> FieldRef { + let data_type = match (hinted.data_type(), physical.data_type()) { + (DataType::Struct(hinted_fields), DataType::Struct(physical_fields)) + if hinted_fields.len() == physical_fields.len() => + { + DataType::Struct( + hinted_fields + .iter() + .zip(physical_fields) + .map(|(hinted, physical)| align_variant_field(hinted, physical, changed)) + .collect(), + ) + } + (hinted_type, physical_type) => { + match ( + variant_marker_child(hinted_type), + variant_marker_child(physical_type), + ) { + (Some((hinted_child, rebuild)), Some((physical_child, _))) => rebuild( + align_variant_field(hinted_child, physical_child, changed), + hinted_type, + ), + _ => hinted_type.clone(), + } + } + }; + + let annotated = physical.has_valid_extension_type::(); + if annotated == hinted.has_valid_extension_type::() { + return Arc::new(hinted.as_ref().clone().with_data_type(data_type)); + } + *changed = true; + let field = Field::new(hinted.name(), data_type, hinted.is_nullable()); + if annotated { + Arc::new( + field + .with_metadata(hinted.metadata().clone()) + .with_extension_type(VariantType), + ) + } else { + let mut metadata = hinted.metadata().clone(); + metadata.remove(EXTENSION_TYPE_NAME_KEY); + metadata.remove(EXTENSION_TYPE_METADATA_KEY); + Arc::new(field.with_metadata(metadata)) + } +} + +/// Make the file's VARIANT annotations, not its `ARROW:schema` hint, decide which fields the +/// reader hands back as Variant. +/// +/// `complex.rs`'s `convert_field` copies a hinted field's metadata verbatim and never derives the +/// extension type from the Parquet logical type, so with a hint present the marker the reader +/// produces is the hint's, whatever the file is annotated as. `schema_adapter`'s +/// `check_variant_annotation` reads that marker as the file's annotation, and Spark's converter +/// reads the Parquet schema alone, so a hint that disagrees makes Comet reject a column Spark +/// reads (hint marked, group not annotated) or read one Spark rejects (the reverse). arrow-rs's +/// own `ArrowWriter` writes the first shape when built without the `variant_experimental` +/// feature, where `logical_type_for_struct` returns `None` but the hint still records the +/// extension. +/// +/// A file with no hint already takes its markers from the annotations, and Spark's writer emits +/// no hint, so the common path costs one scan of the key-value metadata. Rebuild only the +/// returned metadata; the shared cache retains the original footer. +fn with_reconciled_variant_markers( + metadata: Arc, +) -> ParquetResult> { let file = metadata.file_metadata(); - let has_arrow_schema = file.key_value_metadata().is_some_and(|key_values| { + let key_values = file.key_value_metadata(); + if !key_values.is_some_and(|key_values| { key_values .iter() .any(|key_value| key_value.key == ARROW_SCHEMA_META_KEY) - }); - let enum_schema = spark_enum_schema(file.schema_descr())?; - if !has_arrow_schema && enum_schema.is_none() { + }) { return Ok(metadata); } + // The hint applied, which is the schema the reader itself computes, against the annotations + // alone. Re-encoding the first with only its markers moved keeps every other way the hint + // shapes the read, dictionary-encoded columns among them. + let hinted = parquet_to_arrow_schema(file.schema_descr(), key_values)?; + let physical = parquet_to_arrow_schema(file.schema_descr(), None)?; + let Some(reconciled) = align_variant_markers(&hinted, &physical) else { + return Ok(metadata); + }; + Ok(with_arrow_schema_hint(&metadata, Some(reconciled))) +} + +/// Replace the footer's `ARROW:schema` hint with `schema`, or drop the hint when `None`. +/// Rebuilds only the returned metadata, leaving the cached footer as it was read. +fn with_arrow_schema_hint( + metadata: &Arc, + schema: Option, +) -> Arc { + let file = metadata.file_metadata(); let mut key_values = file .key_value_metadata() .into_iter() @@ -387,7 +524,7 @@ fn with_spark_arrow_schema(metadata: Arc) -> ParquetResult>(); - if let Some(schema) = enum_schema { + if let Some(schema) = schema { key_values.push(KeyValue { key: ARROW_SCHEMA_META_KEY.to_string(), value: Some(encode_arrow_schema(&schema)), @@ -402,13 +539,32 @@ fn with_spark_arrow_schema(metadata: Arc) -> ParquetResult) -> ParquetResult> { + let file = metadata.file_metadata(); + let has_arrow_schema = file.key_value_metadata().is_some_and(|key_values| { + key_values + .iter() + .any(|key_value| key_value.key == ARROW_SCHEMA_META_KEY) + }); + let enum_schema = spark_enum_schema(file.schema_descr())?; + if !has_arrow_schema && enum_schema.is_none() { + return Ok(metadata); + } + + Ok(with_arrow_schema_hint(&metadata, enum_schema)) } impl AsyncFileReader for EagerPageIndexReader { @@ -553,7 +709,7 @@ impl AsyncFileReader for EagerPageIndexReader { if spark_variant_schema { with_spark_arrow_schema(metadata) } else { - Ok(metadata) + with_reconciled_variant_markers(metadata) } } .boxed() diff --git a/native/core/src/parquet/parquet_exec.rs b/native/core/src/parquet/parquet_exec.rs index 5c5d26d3922..6992cca1323 100644 --- a/native/core/src/parquet/parquet_exec.rs +++ b/native/core/src/parquet/parquet_exec.rs @@ -86,6 +86,7 @@ pub(crate) fn init_datasource_exec( encryption_enabled: bool, use_field_id: bool, require_field_ids: bool, + ignore_variant_annotation: bool, ) -> Result, ExecutionError> { // Computed once and reused below for `try_pushdown_filters`. `copied_config()` clones only // `SessionConfig` (an `Arc` plus a small extensions map); `SessionContext:: @@ -103,6 +104,7 @@ pub(crate) fn init_datasource_exec( &session_config.options().execution.parquet, ); spark_parquet_options.use_field_id = use_field_id; + spark_parquet_options.ignore_variant_annotation = ignore_variant_annotation; // Spark can discard filtered-out values before timestamp conversion using statistics, // dictionary, and row-level filters. Comet cannot mirror every pruning path, so applying // checked conversion in a filtered scan can fail on values Spark never reads. Preserve the @@ -232,7 +234,8 @@ pub(crate) fn init_datasource_exec( }; let expr_adapter_factory: Arc = Arc::new( - SparkPhysicalExprAdapterFactory::new(spark_parquet_options, default_values), + SparkPhysicalExprAdapterFactory::new(spark_parquet_options, default_values) + .with_required_schema(Arc::clone(&required_schema)), ); let file_groups = file_groups @@ -489,6 +492,7 @@ mod tests { false, false, false, + false, ) .unwrap() } @@ -669,6 +673,7 @@ mod tests { false, false, false, + false, ) .unwrap(); diff --git a/native/core/src/parquet/parquet_exec/variant_tests.rs b/native/core/src/parquet/parquet_exec/variant_tests.rs index ec857019656..c2fd7f38e2a 100644 --- a/native/core/src/parquet/parquet_exec/variant_tests.rs +++ b/native/core/src/parquet/parquet_exec/variant_tests.rs @@ -163,6 +163,7 @@ async fn scan_variant_file(filename: PathBuf) -> VariantArray { false, false, false, + false, ) .unwrap(); let mut stream = scan.execute(0, session_ctx.task_ctx()).unwrap(); @@ -228,6 +229,7 @@ async fn unread_variant_does_not_override_arrow_schema_hint() { false, false, false, + false, ) .unwrap(); let mut stream = scan.execute(0, session.task_ctx()).unwrap(); @@ -258,6 +260,7 @@ fn encrypted_projected_variant_is_rejected_before_reader_creation() { true, false, false, + false, ); assert!(result .unwrap_err() @@ -349,3 +352,167 @@ async fn variant_scan_reads_wide_physical_decimal_as_decimal128() { } } } + +/// The two-field storage Spark writes for a Variant, declared without the Variant marker, which is +/// the shape of a hand-written `struct` read schema. +fn plain_variant_storage() -> DataType { + DataType::Struct(Fields::from(vec![ + Field::new("value", DataType::Binary, false), + Field::new("metadata", DataType::Binary, false), + ])) +} + +/// The relation data schema `CometNativeScan` sends: every root, with `v` as a plain struct. +fn id_and_plain_variant_schema() -> SchemaRef { + Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int32, true), + Field::new("v", plain_variant_storage(), true), + ])) +} + +/// Writes `id INT` next to a VARIANT-annotated `v`. With `rows == 0` the file has no row group. +fn write_id_and_annotated_variant(rows: usize) -> PathBuf { + let file_schema = Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int32, true), + Field::new("v", plain_variant_storage(), true).with_extension_type(VariantType), + ])); + let filename = get_temp_filename(); + let mut writer = ArrowWriter::try_new_with_options( + File::create(&filename).unwrap(), + Arc::clone(&file_schema), + ArrowWriterOptions::new().with_skip_arrow_metadata(true), + ) + .unwrap(); + if rows > 0 { + let DataType::Struct(storage_fields) = plain_variant_storage() else { + unreachable!() + }; + let variant = StructArray::new( + storage_fields, + vec![ + Arc::new(BinaryArray::from(vec![Some(&[12u8, 1u8][..])])) as ArrayRef, + Arc::new(BinaryArray::from(vec![Some(&[1u8, 0u8, 0u8][..])])) as ArrayRef, + ], + None, + ); + let batch = RecordBatch::try_new( + file_schema, + vec![ + Arc::new(arrow::array::Int32Array::from(vec![1])) as ArrayRef, + Arc::new(variant) as ArrayRef, + ], + ) + .unwrap(); + writer.write(&batch).unwrap(); + } + writer.close().unwrap(); + filename +} + +/// Scans through `init_datasource_exec` the way `CometNativeScan` drives it: the full data +/// schema, the Spark read schema, and a projection into the data schema. Returns the row count or +/// the error the scan raised. +async fn scan_with_read_schema( + filename: PathBuf, + required_schema: SchemaRef, + projection: Vec, +) -> Result { + let partitioned_file = + PartitionedFile::from_path(filename.to_string_lossy().into_owned()).unwrap(); + let session_ctx = Arc::new(SessionContext::new()); + let scan = init_datasource_exec( + required_schema, + Some(id_and_plain_variant_schema()), + None, + ObjectStoreUrl::local_filesystem(), + ObjectStoreBackend::Local, + vec![vec![partitioned_file]], + Some(projection), + None, + None, + "UTC", + false, + false, + false, + false, + &session_ctx, + false, + false, + false, + false, + ) + .unwrap(); + let mut stream = scan.execute(0, session_ctx.task_ctx())?; + let mut rows = 0; + while let Some(batch) = stream.next().await { + rows += batch?.num_rows(); + } + Ok(rows) +} + +fn read_schema_id_only() -> SchemaRef { + Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, true)])) +} + +fn assert_variant_annotation_rejected(result: Result) { + let err = result.expect_err("reading the annotated `v` as a plain struct must be rejected"); + assert!( + err.to_string().contains("_LEGACY_ERROR_TEMP_3071"), + "unexpected error: {err}" + ); +} + +/// Spark's schema pruning keeps unrequested roots in the relation data schema, so `v` reaches the +/// native scan even when only `id` is read. The check must follow the read schema, not the data +/// schema, or selecting `id` alone fails on a column the query never touches. +#[tokio::test] +async fn annotated_root_outside_the_read_schema_is_not_rejected() { + let rows = scan_with_read_schema( + write_id_and_annotated_variant(1), + read_schema_id_only(), + vec![0], + ) + .await + .expect("`v` is not in the read schema, so its annotation must not be checked"); + assert_eq!(rows, 1); +} + +/// Once `v` is part of the read schema, the plain struct request is rejected as Spark rejects it. +#[tokio::test] +async fn annotated_root_inside_the_read_schema_is_rejected() { + assert_variant_annotation_rejected( + scan_with_read_schema( + write_id_and_annotated_variant(1), + id_and_plain_variant_schema(), + vec![0, 1], + ) + .await, + ); +} + +/// Scoping the check to the read schema must not move it to row groups: a requested annotated root +/// in a file with no row group is still rejected, matching Spark's schema-conversion-time check. +#[tokio::test] +async fn requested_annotated_root_is_rejected_on_an_empty_file() { + assert_variant_annotation_rejected( + scan_with_read_schema( + write_id_and_annotated_variant(0), + id_and_plain_variant_schema(), + vec![0, 1], + ) + .await, + ); +} + +/// The empty-file counterpart of `annotated_root_outside_the_read_schema_is_not_rejected`. +#[tokio::test] +async fn unrequested_annotated_root_is_not_rejected_on_an_empty_file() { + let rows = scan_with_read_schema( + write_id_and_annotated_variant(0), + read_schema_id_only(), + vec![0], + ) + .await + .expect("`v` is not in the read schema, so its annotation must not be checked"); + assert_eq!(rows, 0); +} diff --git a/native/core/src/parquet/parquet_support.rs b/native/core/src/parquet/parquet_support.rs index e8987a79da6..927be4e490e 100644 --- a/native/core/src/parquet/parquet_support.rs +++ b/native/core/src/parquet/parquet_support.rs @@ -118,6 +118,11 @@ pub struct SparkParquetOptions { /// (overflow -> NULL), because Spark may discard values through pruning paths that /// DataFusion cannot fully mirror before conversion. pub checked_timestamp_overflow: bool, + /// When true (`spark.sql.parquet.ignoreVariantAnnotation`, Spark 4.1+), a Parquet field + /// carrying the VARIANT logical type annotation is read as its plain underlying struct. + /// When false (Spark's default), requesting anything but Spark's VariantType for such a + /// field is rejected, mirroring `ParquetToSparkSchemaConverter.convertGroupField`. + pub ignore_variant_annotation: bool, } impl SparkParquetOptions { @@ -134,6 +139,7 @@ impl SparkParquetOptions { allow_type_promotion: false, allow_timestamp_ltz_to_ntz: false, checked_timestamp_overflow: true, + ignore_variant_annotation: false, } } @@ -150,6 +156,7 @@ impl SparkParquetOptions { allow_type_promotion: false, allow_timestamp_ltz_to_ntz: false, checked_timestamp_overflow: true, + ignore_variant_annotation: false, } } } diff --git a/native/core/src/parquet/schema_adapter.rs b/native/core/src/parquet/schema_adapter.rs index c7ad0418578..2ffeceee179 100644 --- a/native/core/src/parquet/schema_adapter.rs +++ b/native/core/src/parquet/schema_adapter.rs @@ -54,6 +54,12 @@ pub struct SparkPhysicalExprAdapterFactory { /// Default values for columns that may be missing from the physical schema. /// The key is the Column (containing name and index). default_values: Option>, + /// The Spark read schema (`requiredSchema`), which is what Spark clips the Parquet schema to + /// before converting it. The logical file schema can be wider: Spark's schema pruning keeps + /// unrequested top-level fields in the relation's data schema, and Comet passes that whole + /// schema with a separate projection. Roots outside this schema are never converted by + /// Spark, so the VARIANT annotation check skips them. `None` checks every root. + required_schema: Option, } impl SparkPhysicalExprAdapterFactory { @@ -65,8 +71,15 @@ impl SparkPhysicalExprAdapterFactory { Self { parquet_options, default_values, + required_schema: None, } } + + /// Restrict the VARIANT annotation check to the roots Spark actually reads. + pub fn with_required_schema(mut self, required_schema: SchemaRef) -> Self { + self.required_schema = Some(required_schema); + self + } } /// True when a root field of `schema` carries a field id. Root only on purpose: it gates the @@ -425,6 +438,172 @@ fn is_string_or_binary(dt: &DataType) -> bool { ) } +/// Approximate Spark's `DataType.sql` for the variant rejection message. `spark_catalog_name` +/// bottoms out at "unknown" for nested types, which is exactly the shape a hand-written Variant +/// read schema has, so the nested cases are spelled out here. +/// +/// This is close to `DataType.sql` but not identical. Spark's `StructField.sql` wraps names in +/// `QuotingUtils.quoteIfNeeded` and appends the nullability and comment DDL, none of which is +/// reproduced here, so a field named `has space` renders bare. The whole message already differs +/// from Spark's by design (see `variant_annotation_err`), and the rendering exists to tell a user +/// which type they asked for, so it is kept simple rather than made byte-identical. +fn spark_read_type_name(dt: &DataType) -> String { + match dt { + DataType::Struct(fields) => { + let rendered = fields + .iter() + .map(|field| { + format!( + "{}: {}", + field.name(), + spark_read_type_name(field.data_type()) + ) + }) + .collect::>() + .join(", "); + format!("STRUCT<{rendered}>") + } + DataType::List(item) | DataType::LargeList(item) => { + format!("ARRAY<{}>", spark_read_type_name(item.data_type())) + } + DataType::Map(entries, _) => match entries.data_type() { + DataType::Struct(fields) if fields.len() == 2 => format!( + "MAP<{}, {}>", + spark_read_type_name(fields[0].data_type()), + spark_read_type_name(fields[1].data_type()) + ), + other => format!("MAP<{}>", spark_read_type_name(other)), + }, + other => spark_catalog_name(other).to_uppercase(), + } +} + +/// Build the carrier for a Parquet field whose VARIANT logical type annotation is incompatible +/// with the requested Spark read type. The JVM shim turns these into Spark's +/// `_LEGACY_ERROR_TEMP_3071` `AnalysisException`. +/// +/// `column` is the dotted path to the offending field. Spark's own message interpolates the +/// Parquet `Type` itself, roughly `optional group v (VARIANT(1)) { ... }`, so the two messages +/// deliberately differ. `ParquetVariantShreddingSuite` matches on `Invalid Spark read type` and +/// accepts either, and a path is the more useful of the two once the field is nested. +fn variant_annotation_err(column: &str, target_type: &DataType) -> DataFusionError { + DataFusionError::External(Box::new(SparkError::ParquetVariantAnnotationMismatch { + file_path: String::new(), + column: column.to_string(), + spark_type: spark_read_type_name(target_type), + })) +} + +/// Whether `field` is marked as a Variant. On a physical field the marker is the Parquet VARIANT +/// logical type annotation, which arrow-rs surfaces as the `arrow.parquet.variant` Arrow +/// extension type for struct fields, list elements and map values alike. On a requested field it +/// is the same marker Comet's serde applies for `VariantType`. `check_variant_annotation` relies +/// on both sides using it, so it compares them through this one predicate. +/// +/// The marker does not carry the annotation's spec version, and Spark keys its handling off that +/// version. `ParquetToSparkSchemaConverter.convertGroupField` matches +/// `case v: VariantLogicalTypeAnnotation if v.getSpecVersion == 1`, and an annotation failing +/// that guard falls through to `case _ => throw unrecognizedParquetTypeError(...)`, which is +/// `PARQUET_TYPE_NOT_RECOGNIZED`. arrow-rs matches `LogicalType::Variant(_)` for any version, so +/// this predicate cannot tell the two apart. Spark's writer emits version 1. An Arrow-based +/// writer emits no version at all, which parquet-java reads back as 0. +/// +/// Both engines therefore reject a non-v1 annotation, with different error classes: Spark raises +/// `PARQUET_TYPE_NOT_RECOGNIZED` and Comet raises `_LEGACY_ERROR_TEMP_3071`. +/// `variant_annotation_without_a_spec_version_is_also_rejected` pins that. +/// +/// The one behavioral gap is `ignore_variant_annotation` on a non-v1 annotation. Spark's guard +/// fails before it reaches its own ignore branch, so Spark still raises +/// `PARQUET_TYPE_NOT_RECOGNIZED`, while Comet skips the check entirely and reads the plain +/// struct. Comet is the more permissive of the two there. Closing it would mean pairing the +/// requested schema against the Parquet `SchemaDescriptor` rather than the Arrow schema, which +/// duplicates the field-id and case-folding rules `check_variant_annotation` reuses. Closing it +/// means either reading the version from the Parquet schema or arrow-rs carrying it on the +/// extension type. +fn is_variant_marked(field: &Field) -> bool { + field.has_valid_extension_type::() +} + +/// Reject reading a Parquet field that carries the VARIANT logical type annotation as anything +/// other than Spark's `VariantType`, mirroring the `checkConversionRequirement` in Spark's +/// `ParquetToSparkSchemaConverter.convertGroupField`. +/// +/// Comet's scan never runs Spark's schema converter, so this is the only place the file's +/// annotation is ever compared against the requested type. The two sides are compared +/// symmetrically through `is_variant_marked`: a marked request is a legitimate Variant read (see +/// `parquet_exec::init_datasource_exec`'s `projects_variant`) and must not be rejected. +/// Identifying a Variant by its `value`/`metadata` child names instead would misclassify ordinary +/// structs, which is what apache/datafusion-comet#5741 rules out. +/// +/// Runs once per file when the reader opens it, before any row group is inspected, so a file with +/// no row groups is rejected too. That matches Spark, which rejects while converting the schema, +/// and is the opposite of the type-promotion checks in this module, which `RejectOnNonEmpty` +/// defers to execution to match Spark's per-row-group behavior. +fn check_variant_annotation( + logical: &FieldRef, + physical: &FieldRef, + parquet_options: &SparkParquetOptions, + path: &mut Vec, +) -> DataFusionResult<()> { + path.push(logical.name().clone()); + if is_variant_marked(physical) && !is_variant_marked(logical) { + return Err(variant_annotation_err(&path.join("."), logical.data_type())); + } + match (logical.data_type(), physical.data_type()) { + (DataType::Struct(logical_fields), DataType::Struct(physical_fields)) => { + // Resolve nested fields exactly as `spark_parquet_convert` does when it reads them, + // including field-id precedence and case folding, so the field inspected here is the + // one the read pulls from the file. A logical field with no counterpart in the file is + // null-filled and carries no annotation to check. + // + // A resolution failure means duplicate or case-ambiguous children, which is the read + // `check_decoded_field_names` and the nested convert already reject, each with the + // error Spark raises for its own case (#5786). No pairing exists to inspect, so leave + // the rejection to them rather than pre-empting it with this error at plan time. + if let Ok(physical_indices) = + match_struct_fields(physical_fields, logical_fields, parquet_options) + { + for (logical_child, physical_index) in logical_fields.iter().zip(physical_indices) { + if let Some(i) = physical_index { + check_variant_annotation( + logical_child, + &physical_fields[i], + parquet_options, + path, + )?; + } + } + } + } + (DataType::List(logical_item), DataType::List(physical_item)) + | (DataType::LargeList(logical_item), DataType::LargeList(physical_item)) + | (DataType::List(logical_item), DataType::LargeList(physical_item)) + | (DataType::LargeList(logical_item), DataType::List(physical_item)) => { + check_variant_annotation(logical_item, physical_item, parquet_options, path)?; + } + (DataType::Map(logical_entries, _), DataType::Map(physical_entries, _)) => { + // Map children are read by position, not by name: `MapArray::keys` and `values` take + // columns 0 and 1, and Spark's converter takes `getChild(0)` and `getChild(1)`. The + // file may name them anything, so pair them positionally rather than through + // `match_struct_fields`. + if let (DataType::Struct(logical_children), DataType::Struct(physical_children)) = + (logical_entries.data_type(), physical_entries.data_type()) + { + path.push(logical_entries.name().clone()); + for (logical_child, physical_child) in + logical_children.iter().zip(physical_children.iter()) + { + check_variant_annotation(logical_child, physical_child, parquet_options, path)?; + } + path.pop(); + } + } + _ => {} + } + path.pop(); + Ok(()) +} + /// Build a Spark-shaped `SchemaColumnConvertNotSupportedException` carrier for a /// rejected Parquet -> Spark conversion. `column` is the Spark-style column path (`a`, or /// `s, x` for a nested leaf); the bracketed wrapping mirrors @@ -973,6 +1152,54 @@ impl PhysicalExprAdapterFactory for SparkPhysicalExprAdapterFactory { HashMap::new() }; + // Compare the file's VARIANT annotations against the requested types before handing the + // schemas to the default adapter. `adapted_physical_schema` is used so that field-id and + // case-insensitive resolution has already aligned the two sides' top-level names. + if !self.parquet_options.ignore_variant_annotation { + let mut physical_by_folded: HashMap<&str, usize> = HashMap::new(); + for (i, name) in physical_folded.iter().enumerate() { + physical_by_folded.entry(name.as_str()).or_insert(i); + } + // When the read schema is known, check each requested root against its requested + // type and skip roots Spark would not read. Requested roots share their Spark names + // with the logical file schema, so the same fold pairs them. + let required_by_folded: Option> = self + .required_schema + .as_ref() + .map(|required| -> DataFusionResult<_> { + let mut map = HashMap::new(); + for (field, folded) in required + .fields() + .iter() + .zip(fold_schema_names(required, case_sensitive)?) + { + map.entry(folded).or_insert(field); + } + Ok(map) + }) + .transpose()?; + let mut path = Vec::new(); + for (logical_field, folded) in logical_file_schema.fields().iter().zip(&logical_folded) + { + let requested_field = match &required_by_folded { + Some(required) => match required.get(folded) { + Some(field) => *field, + None => continue, + }, + None => logical_field, + }; + if let Some(&i) = physical_by_folded.get(folded.as_str()) { + path.clear(); + check_variant_annotation( + requested_field, + &adapted_physical_schema.fields()[i], + &self.parquet_options, + &mut path, + )?; + } + } + } + let default_factory = DefaultPhysicalExprAdapterFactory; let default_adapter = default_factory.create( Arc::clone(&logical_file_schema), @@ -1624,6 +1851,9 @@ impl PhysicalExpr for RejectOnNonEmpty { #[cfg(test)] pub(crate) mod test { use crate::parquet::cast_column::CometCastColumnExpr; + use crate::parquet::eager_page_index_reader_factory::{ + EagerPageIndexReaderFactory, ScanIoSource, + }; use crate::parquet::parquet_support::SparkParquetOptions; use crate::parquet::schema_adapter::{ check_conversion, is_pure_structural_narrowing, ConversionCheck, @@ -1640,7 +1870,8 @@ pub(crate) mod test { use arrow::buffer::OffsetBuffer; use arrow::datatypes::SchemaRef; use arrow::datatypes::{ - DataType, Field, Fields, Int32Type, Int64Type, Schema, TimeUnit, TimestampMicrosecondType, + DataType, Field, FieldRef, Fields, Int32Type, Int64Type, Schema, TimeUnit, + TimestampMicrosecondType, }; use arrow::record_batch::RecordBatch; use datafusion::common::{DataFusionError, ScalarValue}; @@ -1648,20 +1879,30 @@ pub(crate) mod test { use datafusion::datasource::physical_plan::{FileGroup, FileScanConfigBuilder, ParquetSource}; use datafusion::datasource::source::DataSourceExec; use datafusion::execution::object_store::ObjectStoreUrl; + use datafusion::execution::runtime_env::RuntimeEnv; use datafusion::execution::TaskContext; use datafusion::physical_expr::expressions::Column; use datafusion::physical_expr::PhysicalExpr; + use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet; use datafusion::physical_plan::{ExecutionPlan, SendableRecordBatchStream}; use datafusion_comet_spark_expr::test_common::file_util::get_temp_filename; use datafusion_comet_spark_expr::EvalMode; use datafusion_physical_expr_adapter::PhysicalExprAdapterFactory; use futures::StreamExt; - use parquet::arrow::ArrowWriter; + use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder; + use parquet::arrow::arrow_writer::ArrowWriterOptions; use parquet::arrow::PARQUET_FIELD_ID_META_KEY; + use parquet::arrow::{ArrowWriter, ARROW_SCHEMA_META_KEY}; + use parquet::basic::{LogicalType, Repetition, Type as PhysicalType}; + use parquet::data_type::{ByteArray, ByteArrayType}; + use parquet::file::metadata::KeyValue; + use parquet::file::writer::SerializedFileWriter; + use parquet::schema::printer::print_schema; + use parquet::schema::types::Type as ParquetType; use parquet::variant::VariantType; use std::collections::HashMap; use std::fs::File; - use std::sync::Arc; + use std::sync::{Arc, Mutex}; /// Build field metadata carrying a Parquet field id, for the field-id remap tests. fn id_meta(id: &str) -> HashMap { @@ -4162,6 +4403,930 @@ pub(crate) mod test { assert!(!is_pure_structural_narrowing(&physical, &target, &opts).unwrap()); } + /// Verification probe for apache/datafusion-comet#5741. + /// + /// Captures the `physical_file_schema` DataFusion hands to + /// `PhysicalExprAdapterFactory::create`, so the tests below can assert both that `create` + /// was reached at all and what the file schema looked like when it was. + #[derive(Debug)] + struct ProbeFactory { + inner: SparkPhysicalExprAdapterFactory, + seen: Arc>>, + } + + impl PhysicalExprAdapterFactory for ProbeFactory { + fn create( + &self, + logical_file_schema: SchemaRef, + physical_file_schema: SchemaRef, + ) -> datafusion::common::Result< + Arc, + > { + *self.seen.lock().unwrap() = Some(Arc::clone(&physical_file_schema)); + self.inner.create(logical_file_schema, physical_file_schema) + } + } + + /// Recursively drop every Arrow extension marker from `field`, producing the schema a Spark + /// user's hand-written DDL yields: the same structure, no variant identity. Mirrors the + /// `.schema("v struct")` in Spark's + /// `ParquetVariantShreddingSuite`, test + /// `variant logical type annotation - ignore variant annotation`. + fn strip_extension_markers(field: &FieldRef) -> FieldRef { + let data_type = match field.data_type() { + DataType::Struct(fields) => { + DataType::Struct(fields.iter().map(strip_extension_markers).collect()) + } + DataType::List(item) => DataType::List(strip_extension_markers(item)), + DataType::Map(entries, sorted) => { + DataType::Map(strip_extension_markers(entries), *sorted) + } + other => other.clone(), + }; + let mut metadata = field.metadata().clone(); + metadata.remove("ARROW:extension:name"); + metadata.remove("ARROW:extension:metadata"); + Arc::new(Field::new(field.name(), data_type, field.is_nullable()).with_metadata(metadata)) + } + + /// Write `batch` to a temp Parquet file, then scan it back through `DataSourceExec` asking for + /// the same schema with every extension marker stripped. Reports the Parquet schema the writer + /// actually produced alongside the physical file schema `create` observed (`None` when `create` + /// never ran). + /// + /// Returning the Parquet schema separates the two ways a shape can fail: the *writer* never + /// emitting the VARIANT annotation, versus the *reader* not surfacing it as an Arrow extension + /// type. Only the second one is a problem for #5741; the first is a limitation of the test. + /// + /// `skip_arrow_metadata` controls whether the file embeds an `ARROW:schema` key-value entry. + /// Spark's writer embeds none, so the Spark-shaped cases pass `true`: a marker on the way back + /// can then only have come from the Parquet annotation, and a file carrying the hint would let + /// every assertion below pass vacuously. + struct VariantProbe { + /// The Parquet schema the writer actually produced. + parquet_schema: String, + /// The physical file schema `create` observed, or `None` when it never ran. + observed_physical: Option, + /// Total rows read, or the error the scan raised. + scan: Result, + } + + struct ProbeOptions { + /// Whether the file omits the embedded `ARROW:schema` hint, as Spark's writer does. + skip_arrow_metadata: bool, + /// `spark.sql.parquet.ignoreVariantAnnotation`. + ignore_variant_annotation: bool, + /// Keep the Variant markers on the requested schema, modelling a genuine native Variant + /// projection (`projects_variant` in `parquet_exec.rs`) instead of a hand-written struct. + keep_request_markers: bool, + /// Request this schema instead of one derived from the file schema, for cases where + /// Spark's requested names differ from the names the file was written with. + requested_schema: Option, + } + + impl Default for ProbeOptions { + fn default() -> Self { + Self { + skip_arrow_metadata: true, + ignore_variant_annotation: false, + keep_request_markers: false, + requested_schema: None, + } + } + } + + async fn probe_variant_annotation( + batch: RecordBatch, + options: ProbeOptions, + ) -> Result { + let file_schema = batch.schema(); + let filename = get_temp_filename(); + let filename = filename.as_path().as_os_str().to_str().unwrap().to_string(); + let file = File::create(&filename)?; + let mut writer = ArrowWriter::try_new_with_options( + file, + Arc::clone(&file_schema), + ArrowWriterOptions::new().with_skip_arrow_metadata(options.skip_arrow_metadata), + )?; + writer.write(&batch)?; + writer.close()?; + + probe_variant_annotation_file(&filename, &file_schema, options).await + } + + /// The reading half of `probe_variant_annotation`: reads back the Parquet schema `filename` + /// actually holds, derives the requested schema from `file_schema` per `options`, and scans + /// the file through a `ProbeFactory`. Split out so a test can hand it a file the `ArrowWriter` + /// alone cannot produce, such as one whose `ARROW:schema` hint disagrees with its annotations. + async fn probe_variant_annotation_file( + filename: &str, + file_schema: &SchemaRef, + options: ProbeOptions, + ) -> Result { + let filename = filename.to_string(); + let file_schema = Arc::clone(file_schema); + + let mut printed = Vec::new(); + let reader = ParquetRecordBatchReaderBuilder::try_new(File::open(&filename)?)?; + print_schema(&mut printed, reader.parquet_schema().root_schema()); + let printed = String::from_utf8(printed).unwrap(); + + let requested = if let Some(requested) = options.requested_schema { + requested + } else if options.keep_request_markers { + Arc::clone(&file_schema) + } else { + Arc::new(Schema::new( + file_schema + .fields() + .iter() + .map(strip_extension_markers) + .collect::(), + )) + }; + + let mut spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + spark_parquet_options.ignore_variant_annotation = options.ignore_variant_annotation; + let seen = Arc::new(Mutex::new(None)); + let expr_adapter_factory: Arc = Arc::new(ProbeFactory { + inner: SparkPhysicalExprAdapterFactory::new(spark_parquet_options, None), + seen: Arc::clone(&seen), + }); + + // Install the reader factory `init_datasource_exec` installs, gated the way it gates it: + // the footer rewrites that decide which fields the reader marks as Variant live there, so + // a probe without it would read a hint the real scan never sees. + let projects_variant = requested + .fields() + .iter() + .any(|field| field.has_valid_extension_type::()); + let runtime = Arc::new(RuntimeEnv::default()); + let store = runtime.object_store(ObjectStoreUrl::local_filesystem())?; + let reader_factory = Arc::new( + EagerPageIndexReaderFactory::new( + store, + runtime.cache_manager.get_file_metadata_cache(), + ScanIoSource::Local, + &ExecutionPlanMetricsSet::new(), + ) + .with_spark_variant_schema(projects_variant), + ); + let parquet_source = + ParquetSource::new(requested).with_parquet_file_reader_factory(reader_factory); + let files = FileGroup::new(vec![PartitionedFile::from_path(filename)?]); + let file_scan_config = FileScanConfigBuilder::new( + ObjectStoreUrl::local_filesystem(), + Arc::new(parquet_source), + ) + .with_file_groups(vec![files]) + .with_expr_adapter(Some(expr_adapter_factory)) + .build(); + let parquet_exec = DataSourceExec::new(Arc::new(file_scan_config)); + let scan = async { + let mut stream = parquet_exec.execute(0, Arc::new(TaskContext::default()))?; + let mut rows = 0; + while let Some(batch) = stream.next().await { + rows += batch?.num_rows(); + } + Ok(rows) + } + .await; + + let observed_physical = seen.lock().unwrap().clone(); + Ok(VariantProbe { + parquet_schema: printed, + observed_physical, + scan, + }) + } + + /// The two-field Variant storage group Spark writes when shredding is off. + fn variant_storage_fields() -> Fields { + Fields::from(vec![ + Field::new("value", DataType::Binary, false), + Field::new("metadata", DataType::Binary, false), + ]) + } + + /// A `StructArray` holding one Variant value: the encoded integer `1`, matching the + /// `Row(Array[Byte](12, 1), Array[Byte](1, 0, 0))` the Spark suite expects. + fn variant_storage_array() -> StructArray { + let value = Arc::new(BinaryArray::from(vec![Some(&[12u8, 1u8][..])])) as ArrayRef; + let meta = Arc::new(BinaryArray::from(vec![Some(&[1u8, 0u8, 0u8][..])])) as ArrayRef; + StructArray::new(variant_storage_fields(), vec![value, meta], None) + } + + /// An Arrow field named `name` carrying the Variant storage type and the + /// `arrow.parquet.variant` extension marker, which parquet-rs's writer turns into the Parquet + /// VARIANT logical type annotation. + fn variant_field(name: &str) -> FieldRef { + Arc::new( + Field::new(name, DataType::Struct(variant_storage_fields()), true) + .with_extension_type(VariantType), + ) + } + + /// Assert that `probe_variant_annotation` saw a live `create` call on a file the writer really + /// annotated, then hand the observed physical schema to `locate` to check the marker survived + /// at the shape-specific position. + fn assert_annotation_survived( + shape: &str, + probe: &VariantProbe, + locate: impl Fn(&SchemaRef) -> Option, + ) { + let printed = &probe.parquet_schema; + assert!( + printed.contains("VARIANT"), + "{shape}: parquet-rs did not write the VARIANT annotation, so this file cannot test \ + the reader at all. Parquet schema:\n{printed}" + ); + let observed = probe.observed_physical.clone().unwrap_or_else(|| { + panic!( + "{shape}: PhysicalExprAdapterFactory::create was never called: the logical and \ + physical file schemas compared equal, so the variant annotation did not survive \ + into the physical schema" + ) + }); + let field = locate(&observed).unwrap_or_else(|| { + panic!("{shape}: could not locate the variant field in {observed:?}") + }); + assert!( + field.has_valid_extension_type::(), + "{shape}: physical file schema lost the variant annotation: {field:?}" + ); + } + + /// One row of `map`, the `mv` column of the Spark suite, with the map value + /// child named `value_name` in the file. + fn map_of_variant_batch(value_name: &str) -> Result { + let value_field = variant_field(value_name); + let entry_fields = Fields::from(vec![ + Arc::new(Field::new("key", DataType::Utf8, false)), + Arc::clone(&value_field), + ]); + let entries = StructArray::new( + entry_fields.clone(), + vec![ + Arc::new(StringArray::from(vec!["v2"])) as ArrayRef, + Arc::new(variant_storage_array()) as ArrayRef, + ], + None, + ); + let entries_field = Arc::new(Field::new("entries", DataType::Struct(entry_fields), false)); + let map = MapArray::new( + Arc::clone(&entries_field), + OffsetBuffer::new(vec![0, 1].into()), + entries, + None, + false, + ); + let schema = Arc::new(Schema::new(vec![Field::new( + "mv", + DataType::Map(entries_field, false), + true, + )])); + Ok(RecordBatch::try_new(schema, vec![Arc::new(map)])?) + } + + /// Assert the scan was rejected with Spark's `_LEGACY_ERROR_TEMP_3071` shape, naming `column` + /// (a dotted path) and the requested type. + fn assert_variant_annotation_rejected(probe: &VariantProbe, column: &str, spark_type: &str) { + let err = probe + .scan + .as_ref() + .expect_err("reading a VARIANT-annotated field as a plain struct must be rejected"); + let msg = err.to_string(); + assert!( + msg.contains("_LEGACY_ERROR_TEMP_3071") + && msg.contains("Invalid Spark read type") + && msg.contains(column) + && msg.contains(spark_type), + "unexpected error for {column}: {msg}" + ); + } + + /// The plain struct a hand-written Spark read schema produces for a Variant column. + fn plain_variant_storage_sql() -> &'static str { + "STRUCT" + } + + /// #5741: reading a VARIANT-annotated top-level field as + /// `struct` must raise Spark's `_LEGACY_ERROR_TEMP_3071` + /// instead of silently returning the storage + /// struct. This is the `v` column of the Spark suite's `false` arm. + #[tokio::test] + async fn variant_annotation_read_as_plain_struct_is_rejected() -> Result<(), DataFusionError> { + let schema = Arc::new(Schema::new(vec![variant_field("v")])); + let batch = RecordBatch::try_new(schema, vec![Arc::new(variant_storage_array())])?; + let probe = probe_variant_annotation(batch, ProbeOptions::default()).await?; + assert_variant_annotation_rejected(&probe, "v", plain_variant_storage_sql()); + Ok(()) + } + + /// The `ns struct` case: the annotation is on a nested field, and the rejection + /// must name the dotted path rather than just the root column. + #[tokio::test] + async fn variant_annotation_inside_struct_is_rejected() -> Result<(), DataFusionError> { + let outer_fields = Fields::from(vec![variant_field("nv")]); + let outer = StructArray::new( + outer_fields.clone(), + vec![Arc::new(variant_storage_array()) as ArrayRef], + None, + ); + let schema = Arc::new(Schema::new(vec![Field::new( + "ns", + DataType::Struct(outer_fields), + true, + )])); + let batch = RecordBatch::try_new(schema, vec![Arc::new(outer)])?; + let probe = probe_variant_annotation(batch, ProbeOptions::default()).await?; + assert_variant_annotation_rejected(&probe, "ns.nv", plain_variant_storage_sql()); + Ok(()) + } + + /// The `av array` case. + #[tokio::test] + async fn variant_annotation_as_list_element_is_rejected() -> Result<(), DataFusionError> { + let element = variant_field("element"); + let list = ListArray::new( + Arc::clone(&element), + OffsetBuffer::new(vec![0, 1].into()), + Arc::new(variant_storage_array()), + None, + ); + let schema = Arc::new(Schema::new(vec![Field::new( + "av", + DataType::List(element), + true, + )])); + let batch = RecordBatch::try_new(schema, vec![Arc::new(list)])?; + let probe = probe_variant_annotation(batch, ProbeOptions::default()).await?; + assert_variant_annotation_rejected(&probe, "av.element", plain_variant_storage_sql()); + Ok(()) + } + + /// The `mv map` case. + #[tokio::test] + async fn variant_annotation_as_map_value_is_rejected() -> Result<(), DataFusionError> { + let batch = map_of_variant_batch("value")?; + let probe = probe_variant_annotation(batch, ProbeOptions::default()).await?; + // The path walks through the Parquet map's `entries` group, as Spark's own error does. + assert_variant_annotation_rejected(&probe, "mv.entries.value", plain_variant_storage_sql()); + Ok(()) + } + + /// In field-id read mode a nested logical field resolves to the file field with the matching + /// ID, not the matching name -- that is what `spark_parquet_convert` does when it builds the + /// read. Pairing by name here instead would inspect the wrong file field (or none) and let the + /// annotation through unchecked, so the rejection must follow the ID. + #[test] + fn variant_annotation_is_found_through_a_nested_field_id_match() { + let storage = DataType::Struct(variant_storage_fields()); + let logical = Arc::new(Schema::new(vec![Field::new( + "ns", + DataType::Struct(Fields::from(vec![Field::new( + "renamed_since_write", + storage.clone(), + true, + ) + .with_metadata(id_meta("7"))])), + true, + )])); + // The file holds the Variant under a different name, plus a decoy that would win a + // name-only match and carries no annotation. + let physical = Arc::new(Schema::new(vec![Field::new( + "ns", + DataType::Struct(Fields::from(vec![ + Field::new("renamed_since_write", DataType::Int32, true) + .with_metadata(id_meta("9")), + Field::new("nv", storage, true) + .with_metadata(id_meta("7")) + .with_extension_type(VariantType), + ])), + true, + )])); + + let mut options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + options.use_field_id = true; + let err = SparkPhysicalExprAdapterFactory::new(options, None) + .create(logical, physical) + .expect_err("an ID-resolved Variant field must still be rejected"); + let msg = err.to_string(); + assert!( + msg.contains("_LEGACY_ERROR_TEMP_3071") && msg.contains("ns.renamed_since_write"), + "unexpected error: {msg}" + ); + } + + /// Write a one-row Parquet file holding a single top-level group annotated as VARIANT with + /// `spec_version`, using the low-level writer so the annotation is exactly what Spark's writer + /// emits. `ArrowWriter` cannot do this: `logical_type_for_struct` hardcodes + /// `LogicalType::variant(None)`, which is a different annotation from Spark's `VARIANT(1)`. + fn write_annotated_variant_file(spec_version: Option) -> String { + let value = Arc::new( + ParquetType::primitive_type_builder("value", PhysicalType::BYTE_ARRAY) + .with_repetition(Repetition::REQUIRED) + .build() + .unwrap(), + ); + let metadata = Arc::new( + ParquetType::primitive_type_builder("metadata", PhysicalType::BYTE_ARRAY) + .with_repetition(Repetition::REQUIRED) + .build() + .unwrap(), + ); + let variant = Arc::new( + ParquetType::group_type_builder("v") + .with_repetition(Repetition::REQUIRED) + .with_logical_type(Some(LogicalType::variant(spec_version))) + .with_fields(vec![value, metadata]) + .build() + .unwrap(), + ); + let schema = Arc::new( + ParquetType::group_type_builder("schema") + .with_fields(vec![variant]) + .build() + .unwrap(), + ); + + let filename = get_temp_filename(); + let filename = filename.as_path().as_os_str().to_str().unwrap().to_string(); + let file = File::create(&filename).unwrap(); + let mut writer = SerializedFileWriter::new(file, schema, Default::default()).unwrap(); + let mut row_group = writer.next_row_group().unwrap(); + for bytes in [vec![12u8, 1u8], vec![1u8, 0u8, 0u8]] { + let mut column = row_group.next_column().unwrap().unwrap(); + column + .typed::() + .write_batch(&[ByteArray::from(bytes)], None, None) + .unwrap(); + column.close().unwrap(); + } + row_group.close().unwrap(); + writer.close().unwrap(); + filename + } + + /// Scan `filename` asking for the plain `struct` a hand-written + /// Spark read schema produces, and report the scan outcome. + async fn scan_as_plain_variant_storage(filename: String) -> Result { + let requested = Arc::new(Schema::new(vec![Field::new( + "v", + DataType::Struct(variant_storage_fields()), + false, + )])); + let expr_adapter_factory: Arc = + Arc::new(SparkPhysicalExprAdapterFactory::new( + SparkParquetOptions::new(EvalMode::Legacy, "UTC", false), + None, + )); + let file_scan_config = FileScanConfigBuilder::new( + ObjectStoreUrl::local_filesystem(), + Arc::new(ParquetSource::new(requested)), + ) + .with_file_groups(vec![FileGroup::new(vec![PartitionedFile::from_path( + filename, + )?])]) + .with_expr_adapter(Some(expr_adapter_factory)) + .build(); + let exec = DataSourceExec::new(Arc::new(file_scan_config)); + let mut stream = exec.execute(0, Arc::new(TaskContext::default()))?; + let mut rows = 0; + while let Some(batch) = stream.next().await { + rows += batch?.num_rows(); + } + Ok(rows) + } + + /// The annotation Spark itself writes is `VARIANT(1)`, and every other rejection test in this + /// module goes through `ArrowWriter`, which emits `VARIANT(None)` instead. This one writes the + /// real Spark shape so the rejection is pinned against the annotation users actually have on + /// disk, not only against the one the test writer happens to produce. + #[tokio::test] + async fn variant_annotation_with_spec_version_1_is_rejected() -> Result<(), DataFusionError> { + let err = scan_as_plain_variant_storage(write_annotated_variant_file(Some(1))) + .await + .expect_err("a VARIANT(1) field read as a plain struct must be rejected"); + let msg = err.to_string(); + assert!( + msg.contains("_LEGACY_ERROR_TEMP_3071") && msg.contains("Invalid Spark read type"), + "unexpected error: {msg}" + ); + Ok(()) + } + + /// An annotation with no spec version, which is what an Arrow-based writer emits and + /// parquet-java reads back as 0. Spark rejects this too, but through its + /// `unrecognizedParquetTypeError` catch-all rather than the variant branch, so the error class + /// differs: `PARQUET_TYPE_NOT_RECOGNIZED` there against `_LEGACY_ERROR_TEMP_3071` here. The + /// `arrow.parquet.variant` extension type carries no version, so this check cannot tell the + /// file apart from `VARIANT(1)`. See `check_variant_annotation` for the full comparison. If a + /// future arrow-rs exposes the version, this test is the one to revisit. + #[tokio::test] + async fn variant_annotation_without_a_spec_version_is_also_rejected( + ) -> Result<(), DataFusionError> { + let err = scan_as_plain_variant_storage(write_annotated_variant_file(None)) + .await + .expect_err("Comet rejects an unversioned VARIANT annotation that Spark would accept"); + assert!( + err.to_string().contains("_LEGACY_ERROR_TEMP_3071"), + "unexpected error: {err}" + ); + Ok(()) + } + + /// `spark.sql.parquet.ignoreVariantAnnotation=true` is the explicit opt-in to read the + /// annotated group as its plain underlying struct. The `true` arm of the Spark suite: the read + /// must succeed and return the Variant storage bytes unchanged. + #[tokio::test] + async fn variant_annotation_is_ignored_when_conf_is_set() -> Result<(), DataFusionError> { + let top_level = Arc::new(Schema::new(vec![variant_field("v")])); + let cases = vec![ + ( + "top level", + RecordBatch::try_new(top_level, vec![Arc::new(variant_storage_array())])?, + ), + // The conf short-circuits the whole walk, so a nested annotation must be ignored too. + ("map value", map_of_variant_batch("value")?), + ]; + for (shape, batch) in cases { + let probe = probe_variant_annotation( + batch, + ProbeOptions { + ignore_variant_annotation: true, + ..Default::default() + }, + ) + .await?; + let rows = probe.scan.unwrap_or_else(|err| { + panic!( + "{shape}: ignoreVariantAnnotation=true must allow the plain struct read: {err}" + ) + }); + assert_eq!(rows, 1, "{shape}: unexpected row count"); + } + Ok(()) + } + + /// The other half of the symmetric check: when the *requested* schema is itself marked as a + /// Variant -- a genuine native Variant projection, which #5794 made reachable -- the annotation + /// agrees with the request and the read must not be rejected. Without this test, "simplifying" + /// the check to fire on the physical marker alone would break Variant projection with nothing + /// turning red. + #[tokio::test] + async fn marked_variant_request_is_not_rejected() -> Result<(), DataFusionError> { + let schema = Arc::new(Schema::new(vec![variant_field("v")])); + let batch = RecordBatch::try_new(schema, vec![Arc::new(variant_storage_array())])?; + let probe = probe_variant_annotation( + batch, + ProbeOptions { + keep_request_markers: true, + ..Default::default() + }, + ) + .await?; + assert_eq!( + probe + .scan + .expect("a marked Variant request must not trip the annotation check"), + 1 + ); + Ok(()) + } + + /// Spark rejects this read inside `ParquetToSparkSchemaConverter`, while converting the schema, + /// so a file with no rows fails just the same. That is the opposite of the type-promotion + /// checks in this module, which `RejectOnNonEmpty` defers to match Spark's per-row-group + /// behavior (see `parquet_empty_file_disallowed_widening`, SPARK-26709). Pinning it here keeps + /// a later reader from "fixing" this check to match its neighbours. + #[tokio::test] + async fn variant_annotation_is_rejected_on_an_empty_file() -> Result<(), DataFusionError> { + let schema = Arc::new(Schema::new(vec![variant_field("v")])); + let empty = StructArray::new_null(variant_storage_fields(), 0); + let batch = RecordBatch::try_new(schema, vec![Arc::new(empty)])?; + let probe = probe_variant_annotation(batch, ProbeOptions::default()).await?; + assert_variant_annotation_rejected(&probe, "v", plain_variant_storage_sql()); + Ok(()) + } + + /// Map children are read by position, so a VARIANT annotation on the second child must be + /// rejected even when the file names that child something other than the `value` Spark + /// requests. Pairing by name would find no `value` in the file and skip the annotation while + /// the read still consumes the child positionally. The error path uses the requested names. + #[tokio::test] + async fn variant_annotation_on_a_renamed_map_value_is_rejected() -> Result<(), DataFusionError> + { + let requested_entries = Fields::from(vec![ + Field::new("key", DataType::Utf8, false), + Field::new("value", DataType::Struct(variant_storage_fields()), true), + ]); + let requested = Arc::new(Schema::new(vec![Field::new( + "mv", + DataType::Map( + Arc::new(Field::new( + "entries", + DataType::Struct(requested_entries), + false, + )), + false, + ), + true, + )])); + let probe = probe_variant_annotation( + map_of_variant_batch("payload")?, + ProbeOptions { + requested_schema: Some(requested), + ..Default::default() + }, + ) + .await?; + assert!( + probe.parquet_schema.contains("payload"), + "the file must name the map value `payload` for this case to mean anything" + ); + assert_variant_annotation_rejected(&probe, "mv.entries.value", plain_variant_storage_sql()); + Ok(()) + } + + /// A requested root is walked with its requested type, not its data-schema type, so an + /// annotated nested field that the read schema prunes away is not checked. Spark clips the + /// Parquet schema to the read schema before converting it and never looks at `ns.nv` here. + #[test] + fn annotated_field_pruned_from_the_read_schema_is_not_rejected() { + let storage = DataType::Struct(variant_storage_fields()); + let nested = |nv: Field| { + Arc::new(Schema::new(vec![Field::new( + "ns", + DataType::Struct(Fields::from(vec![ + Field::new("a", DataType::Int32, true), + nv, + ])), + true, + )])) + }; + let logical = nested(Field::new("nv", storage.clone(), true)); + let physical = nested(Field::new("nv", storage, true).with_extension_type(VariantType)); + let required = Arc::new(Schema::new(vec![Field::new( + "ns", + DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])), + true, + )])); + let result = SparkPhysicalExprAdapterFactory::new( + SparkParquetOptions::new(EvalMode::Legacy, "UTC", false), + None, + ) + .with_required_schema(required) + .create(logical, physical); + assert!( + result.is_ok(), + "`ns.nv` is outside the read schema: {:?}", + result.err() + ); + } + + /// A top-level Variant column: the `v variant` case of the Spark suite. + #[tokio::test] + async fn variant_annotation_survives_at_top_level() -> Result<(), DataFusionError> { + let schema = Arc::new(Schema::new(vec![variant_field("v")])); + let batch = RecordBatch::try_new(schema, vec![Arc::new(variant_storage_array())])?; + let probe = probe_variant_annotation(batch, ProbeOptions::default()).await?; + assert_annotation_survived("top level", &probe, |s| { + s.field_with_name("v").ok().map(|f| Arc::new(f.clone())) + }); + Ok(()) + } + + /// A file that *does* embed an `ARROW:schema` hint, which Spark never writes but Arrow-based + /// writers do. #5794 strips that hint before the reader sees it, but only when the *requested* + /// schema projects a Variant (`projects_variant` in `parquet_exec.rs`); in the #5741 shape the + /// request is a plain struct, so the hint survives and the reader may rebuild the field from + /// it instead of from the Parquet annotation. This pins down whether the marker still reaches + /// the adapter on that path. + #[tokio::test] + async fn variant_annotation_survives_with_embedded_arrow_schema() -> Result<(), DataFusionError> + { + let schema = Arc::new(Schema::new(vec![variant_field("v")])); + let batch = RecordBatch::try_new(schema, vec![Arc::new(variant_storage_array())])?; + let probe = probe_variant_annotation( + batch, + ProbeOptions { + skip_arrow_metadata: false, + ..Default::default() + }, + ) + .await?; + assert_annotation_survived("embedded arrow schema", &probe, |s| { + s.field_with_name("v").ok().map(|f| Arc::new(f.clone())) + }); + Ok(()) + } + + /// The test above writes a file whose hint and annotation agree. They need not: the physical + /// side of `check_variant_annotation` is the reader's Arrow field, and `convert_field` in + /// parquet-rs's `arrow/schema/complex.rs` copies each field's metadata straight from the + /// `ARROW:schema` hint, never deriving the extension type from the Parquet logical type when a + /// hint is present. A plain-struct request leaves `projects_variant` false, so + /// `init_datasource_exec` keeps that hint, and the marker this check reads then comes from the + /// hint rather than from the file. + /// + /// arrow-rs's own `ArrowWriter` produces such a file when built without the + /// `variant_experimental` feature: `logical_type_for_struct` returns `None`, so the group + /// carries no annotation, while the hint still records `arrow.parquet.variant`. Spark reads + /// only the Parquet schema, infers `struct` and reads the + /// column, and so does Comet on main. The check must not reject it. + #[tokio::test] + async fn hint_marking_an_unannotated_group_is_read_as_a_struct() -> Result<(), DataFusionError> + { + let filename = + write_with_foreign_hint(&unmarked_storage_schema(), &marked_storage_schema())?; + let probe = probe_variant_annotation_file( + &filename, + &unmarked_storage_schema(), + ProbeOptions::default(), + ) + .await?; + assert!( + !probe.parquet_schema.contains("VARIANT"), + "this file must carry no annotation for the test to mean anything. Parquet schema:\n{}", + probe.parquet_schema + ); + assert_eq!( + probe.scan?, 1, + "a group the writer never annotated must read as a plain struct, as Spark reads it" + ); + Ok(()) + } + + /// The other half of the same disagreement: the group is annotated but the hint does not mark + /// it, so the hint would hide the annotation from `check_variant_annotation` and hand back the + /// storage bytes Spark refuses to return. The rejection must follow the file, not the hint. + #[tokio::test] + async fn annotated_group_the_hint_leaves_unmarked_is_still_rejected( + ) -> Result<(), DataFusionError> { + let filename = + write_with_foreign_hint(&marked_storage_schema(), &unmarked_storage_schema())?; + let probe = probe_variant_annotation_file( + &filename, + &unmarked_storage_schema(), + ProbeOptions::default(), + ) + .await?; + assert!( + probe.parquet_schema.contains("VARIANT"), + "this file must carry the annotation for the test to mean anything. Parquet schema:\n{}", + probe.parquet_schema + ); + assert_variant_annotation_rejected(&probe, "v", plain_variant_storage_sql()); + Ok(()) + } + + /// The Variant storage group with the `arrow.parquet.variant` marker, which parquet-rs's + /// writer turns into the VARIANT annotation. + fn marked_storage_schema() -> SchemaRef { + Arc::new(Schema::new(vec![variant_field("v")])) + } + + /// The same group with no marker, so the writer annotates nothing. + fn unmarked_storage_schema() -> SchemaRef { + Arc::new(Schema::new(vec![Field::new( + "v", + DataType::Struct(variant_storage_fields()), + true, + )])) + } + + /// Write one row of `data_schema` to a temp file that carries no hint of its own, then attach + /// the `ARROW:schema` hint arrow-rs writes for `hint_schema`. The annotation follows + /// `data_schema` and the hint follows `hint_schema`, so the two disagree on their Variant + /// markers — a file no single `ArrowWriter` call produces, and the shape `convert_field` + /// resolves in the hint's favor. The hint is harvested from a throwaway file rather than + /// hand-built, so its bytes are the reader's own encoding. + fn write_with_foreign_hint( + data_schema: &SchemaRef, + hint_schema: &SchemaRef, + ) -> Result { + let write = |schema: &SchemaRef, skip_hint: bool| -> Result { + let filename = get_temp_filename(); + let filename = filename.as_path().as_os_str().to_str().unwrap().to_string(); + let mut writer = ArrowWriter::try_new_with_options( + File::create(&filename)?, + Arc::clone(schema), + ArrowWriterOptions::new().with_skip_arrow_metadata(skip_hint), + )?; + writer.write(&RecordBatch::try_new( + Arc::clone(schema), + vec![Arc::new(variant_storage_array())], + )?)?; + writer.close()?; + Ok(filename) + }; + + let donor = write(hint_schema, false)?; + let hint = ParquetRecordBatchReaderBuilder::try_new(File::open(&donor)?)? + .metadata() + .file_metadata() + .key_value_metadata() + .into_iter() + .flatten() + .find(|key_value| key_value.key == ARROW_SCHEMA_META_KEY) + .and_then(|key_value| key_value.value.clone()) + .expect("arrow-rs wrote no ARROW:schema hint to harvest"); + + let filename = get_temp_filename(); + let filename = filename.as_path().as_os_str().to_str().unwrap().to_string(); + let mut writer = ArrowWriter::try_new_with_options( + File::create(&filename)?, + Arc::clone(data_schema), + ArrowWriterOptions::new().with_skip_arrow_metadata(true), + )?; + writer.write(&RecordBatch::try_new( + Arc::clone(data_schema), + vec![Arc::new(variant_storage_array())], + )?)?; + writer.append_key_value_metadata(KeyValue::new( + ARROW_SCHEMA_META_KEY.to_string(), + Some(hint), + )); + writer.close()?; + Ok(filename) + } + + /// A Variant nested inside a struct: the `ns struct` case. + #[tokio::test] + async fn variant_annotation_survives_inside_struct() -> Result<(), DataFusionError> { + let inner = variant_field("nv"); + let outer_fields = Fields::from(vec![Arc::clone(&inner)]); + let outer = StructArray::new( + outer_fields.clone(), + vec![Arc::new(variant_storage_array()) as ArrayRef], + None, + ); + let schema = Arc::new(Schema::new(vec![Field::new( + "ns", + DataType::Struct(outer_fields), + true, + )])); + let batch = RecordBatch::try_new(schema, vec![Arc::new(outer)])?; + let probe = probe_variant_annotation(batch, ProbeOptions::default()).await?; + assert_annotation_survived("struct field", &probe, |s| { + match s.field_with_name("ns").ok()?.data_type() { + DataType::Struct(fields) => fields.iter().find(|f| f.name() == "nv").cloned(), + _ => None, + } + }); + Ok(()) + } + + /// A Variant as a list element: the `av array` case. parquet-rs converts list + /// elements through a different branch of `complex.rs` than struct fields, so a struct-field + /// result does not carry over to here. + #[tokio::test] + async fn variant_annotation_survives_as_list_element() -> Result<(), DataFusionError> { + let element = variant_field("element"); + let values = variant_storage_array(); + let list = ListArray::new( + Arc::clone(&element), + OffsetBuffer::new(vec![0, 1].into()), + Arc::new(values), + None, + ); + let schema = Arc::new(Schema::new(vec![Field::new( + "av", + DataType::List(element), + true, + )])); + let batch = RecordBatch::try_new(schema, vec![Arc::new(list)])?; + let probe = probe_variant_annotation(batch, ProbeOptions::default()).await?; + assert_annotation_survived("list element", &probe, |s| { + match s.field_with_name("av").ok()?.data_type() { + DataType::List(item) => Some(Arc::clone(item)), + _ => None, + } + }); + Ok(()) + } + + /// A Variant as a map value: the `mv map` case. Map key/value fields take yet + /// another conversion branch, independent of both struct fields and list elements. + #[tokio::test] + async fn variant_annotation_survives_as_map_value() -> Result<(), DataFusionError> { + let batch = map_of_variant_batch("value")?; + let probe = probe_variant_annotation(batch, ProbeOptions::default()).await?; + assert_annotation_survived("map value", &probe, |s| { + match s.field_with_name("mv").ok()?.data_type() { + DataType::Map(entries, _) => match entries.data_type() { + DataType::Struct(fields) => { + fields.iter().find(|f| f.name() == "value").cloned() + } + _ => None, + }, + _ => None, + } + }); + Ok(()) + } + /// A requested schema that repeats an id is declined at planning time and never reaches /// the native scan, so the duplicate can only sit in the file. The file holds `s` with /// `x` and `y` both carrying id 1 beside `z` with id 2, and the read asks for `x` (id 1), diff --git a/native/proto/src/proto/operator.proto b/native/proto/src/proto/operator.proto index d78ee8df39c..68dcc12e3ee 100644 --- a/native/proto/src/proto/operator.proto +++ b/native/proto/src/proto/operator.proto @@ -182,6 +182,11 @@ message NativeScanCommon { // True when Spark supplied a Parquet data filter, including when Comet could // not serialize any of those filters into data_filters. bool has_data_filters = 19; + // True when spark.sql.parquet.ignoreVariantAnnotation is set (Spark 4.1+). When + // false, reading a Parquet field carrying the VARIANT logical type annotation as + // anything other than Spark's VariantType is rejected, matching the check + // Spark's ParquetToSparkSchemaConverter performs during schema conversion. + bool ignore_variant_annotation = 20; } message NativeScan { diff --git a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala index 241405c0ea5..1f91fa7af37 100644 --- a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala +++ b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala @@ -24,10 +24,15 @@ import org.apache.parquet.crypto.DecryptionPropertiesFactory import org.apache.parquet.crypto.keytools.{KeyToolkit, PropertiesDrivenCryptoFactory} import org.apache.spark.sql.internal.SQLConf +import org.apache.comet.CometSparkSessionExtensions.isSpark41Plus + object CometParquetUtils { private val PARQUET_FIELD_ID_WRITE_ENABLED = "spark.sql.parquet.fieldId.write.enabled" private val PARQUET_FIELD_ID_READ_ENABLED = "spark.sql.parquet.fieldId.read.enabled" private val IGNORE_MISSING_PARQUET_FIELD_ID = "spark.sql.parquet.fieldId.read.ignoreMissing" + // `SQLConf.PARQUET_IGNORE_VARIANT_ANNOTATION`, read by key because the conf only exists from + // Spark 4.1 while this file compiles against 3.4 through 4.1. + private val IGNORE_VARIANT_ANNOTATION = "spark.sql.parquet.ignoreVariantAnnotation" // Field-metadata key arrow-rs writes when it lifts Parquet field IDs into the Arrow schema // (`parquet::arrow::PARQUET_FIELD_ID_META_KEY`). Spark's local key for the same concept is @@ -57,6 +62,18 @@ object CometParquetUtils { def ignoreMissingIds(conf: SQLConf): Boolean = conf.getConfString(IGNORE_MISSING_PARQUET_FIELD_ID, "false").toBoolean + /** + * Whether the native reader should read a VARIANT-annotated Parquet field as its plain + * underlying struct instead of rejecting the read. + * + * Spark only validates the annotation from 4.1: `ParquetToSparkSchemaConverter` gained the + * `VariantLogicalTypeAnnotation` branch and `spark.sql.parquet.ignoreVariantAnnotation` in that + * release, and neither exists in 3.4, 3.5 or 4.0. Comet's check mirrors that rule, so on + * earlier versions it is switched off rather than inventing a failure Spark does not have. + */ + def ignoreVariantAnnotation(conf: SQLConf): Boolean = + !isSpark41Plus || conf.getConfString(IGNORE_VARIANT_ANNOTATION, "false").toBoolean + /** * Checks if the given Hadoop configuration contains any unsupported encryption settings. * diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometNativeScan.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometNativeScan.scala index 286937635b7..b10c990d3e7 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometNativeScan.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometNativeScan.scala @@ -353,6 +353,12 @@ object CometNativeScan extends CometOperatorSerde[CometScanExec] with CometTypeS commonBuilder.setRequireFieldIds( hasFieldIds && !scan.conf.getConf(SQLConf.IGNORE_MISSING_PARQUET_FIELD_ID)) + // Spark rejects reading a VARIANT-annotated Parquet field as anything but VariantType + // unless this conf is on; the native reader enforces the same rule from the file's + // annotation, which is the only place that annotation is visible. + commonBuilder.setIgnoreVariantAnnotation( + CometParquetUtils.ignoreVariantAnnotation(scan.conf)) + commonBuilder.setAllowTypePromotion(CometConf.COMET_SCHEMA_EVOLUTION_ENABLED) commonBuilder.setAllowTimestampLtzToNtz(CometConf.COMET_ALLOW_TIMESTAMP_LTZ_AS_NTZ) diff --git a/spark/src/main/spark-4.x/org/apache/spark/sql/comet/shims/ShimSparkErrorConverter.scala b/spark/src/main/spark-4.x/org/apache/spark/sql/comet/shims/ShimSparkErrorConverter.scala index 30e25423357..0cf658417ef 100644 --- a/spark/src/main/spark-4.x/org/apache/spark/sql/comet/shims/ShimSparkErrorConverter.scala +++ b/spark/src/main/spark-4.x/org/apache/spark/sql/comet/shims/ShimSparkErrorConverter.scala @@ -24,6 +24,7 @@ import java.io.FileNotFoundException import scala.util.matching.Regex import org.apache.spark.{QueryContext, SparkException, SparkIllegalArgumentException} +import org.apache.spark.sql.AnalysisException import org.apache.spark.sql.errors.QueryExecutionErrors import org.apache.spark.sql.execution.datasources.SchemaColumnConvertNotSupportedException import org.apache.spark.sql.types._ @@ -363,6 +364,24 @@ trait ShimSparkErrorConverter { val missingPath = params.get("filePath").map(_.toString).getOrElse("") Some(QueryExecutionErrors.cannotReadFilesError(missingCause, missingPath)) + case "ParquetVariantAnnotationMismatch" => + // Mirror Spark's `ParquetToSparkSchemaConverter.checkConversionRequirement`, which throws + // an AnalysisException while converting the schema. Spark's file scan surfaces that as a + // FAILED_READ_FILE SparkException with the AnalysisException as its cause, which is the + // shape `ParquetVariantShreddingSuite` asserts. The rejection happens while adapting the + // schema, so the native side has no file path, and the message reads "reading file " with + // nothing after it. `SparkErrorConverter`'s fallback to the per-task file list does not + // help: it is populated only for RDDs built by `CometNativeScanExec`, and this error + // surfaces through the projection above the scan. `ParquetSchemaConvert` below has the + // same empty path for the same reason. + val variantCause = new AnalysisException( + errorClass = "_LEGACY_ERROR_TEMP_3071", + messageParameters = Map( + "msg" -> ("Invalid Spark read type: expected " + params("column") + + " to be variant type but found " + params("sparkType")))) + val variantPath = params.get("filePath").map(_.toString).getOrElse("") + Some(QueryExecutionErrors.cannotReadFilesError(variantCause, variantPath)) + case "ParquetSchemaConvert" => // Mirror Spark 4.0's FileDataSourceV2: wrap the // SchemaColumnConvertNotSupportedException in a FAILED_READ_FILE