From 9182a202a18907e4f6290777d5af1c8df7655280 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 1 Oct 2026 10:06:22 -0600 Subject: [PATCH 1/5] fix: defer Parquet conversion errors until a row group is decoded, as Spark does Spark checks a Parquet type conversion in `getUpdater`, which it calls only while decoding a row group, so a file with nothing to decode (empty, or with every row group pruned) reads even when a column's type can't be converted. The native scan rejected BINARY read as a non-string type, decimal narrowing, an integer read as a too-narrow decimal, and every scalar/complex mismatch when it opened the file. Since #5681 applies those rules to nested fields, such a file failed the query where Spark and 1.0.0 return rows. Defer all of them through `RejectOnNonEmpty`, except the shape mismatches that Spark also fails when it opens the file because it can't clip the file's type to the requested one. Closes #6506. --- native/core/src/parquet/schema_adapter.rs | 264 +++++++++++++++--- .../comet/parquet/ParquetReadSuite.scala | 44 +++ 2 files changed, 275 insertions(+), 33 deletions(-) diff --git a/native/core/src/parquet/schema_adapter.rs b/native/core/src/parquet/schema_adapter.rs index c7ad0418578..93f28df9f47 100644 --- a/native/core/src/parquet/schema_adapter.rs +++ b/native/core/src/parquet/schema_adapter.rs @@ -466,10 +466,12 @@ fn reject_on_non_empty_expr( enum ConversionCheck { /// Spark has an updater for the pair (for a same-shape complex pair: for every leaf). Accept, - /// Spark rejects the pair; raised at plan time. + /// Spark rejects the pair when it opens the file, before it decodes anything; raised at + /// plan time. Reject(DataFusionError), /// Spark rejects the pair, but only while decoding a row group, so the rejection is - /// deferred to runtime via [`RejectOnNonEmpty`] (SPARK-26709). Carries the offending + /// deferred to runtime via [`RejectOnNonEmpty`] (SPARK-26709). A file with nothing to + /// decode, empty or with every row group pruned, reads. Carries the offending /// leaf's column path and physical / requested types for the error message. RejectOnNonEmpty { column: String, @@ -483,6 +485,10 @@ enum ConversionCheck { /// for a top-level column, `s, x` for a nested leaf, mirroring /// `Arrays.toString(descriptor.getPath())`). The rules and their order are exactly those the /// adapter applies to top-level columns; [`check_conversion`] applies them to nested leaves. +/// +/// Spark calls `getUpdater` only while it decodes a row group, so every rejection here is +/// deferred to runtime except a shape mismatch that Spark already fails when it opens the +/// file (#6506). fn check_leaf_conversion( physical_type: &DataType, target_type: &DataType, @@ -505,6 +511,55 @@ fn check_leaf_conversion( target_type: target_type.clone(), }; + // Scalar/complex mismatch (e.g. TIMESTAMP read as ARRAY): + // Spark's vectorized reader rejects with + // SchemaColumnConvertNotSupportedException (SPARK-45604). Same-shape + // complex pairs never reach this leaf check (`check_conversion` walks their + // leaves instead), so two complex types here differ in shape (e.g. STRUCT + // read as ARRAY), which Spark rejects just the same. + // + // Checked first because the shape decides when Spark rejects. It fails when + // it opens the file if it can't clip the file's type to the requested one, + // which it can't for a group (`ParquetToSparkSchemaConverter`) or for a + // primitive read as a struct, or as an array or map with a complex element + // (`ParquetReadSupport.clipParquetType`). It doesn't clip a primitive read as + // an array or map of primitives, so only `getUpdater` rejects that, while + // decoding, as in SPARK-45604. + let is_complex = |t: &DataType| { + matches!( + t, + DataType::Struct(_) + | DataType::List(_) + | DataType::LargeList(_) + | DataType::FixedSizeList(_, _) + | DataType::ListView(_) + | DataType::LargeListView(_) + | DataType::Map(_, _) + ) + }; + if is_complex(physical_type) { + return reject(); + } + if is_complex(target_type) { + let elements_are_primitive = match target_type { + DataType::List(item) + | DataType::LargeList(item) + | DataType::FixedSizeList(item, _) + | DataType::ListView(item) + | DataType::LargeListView(item) => !is_complex(item.data_type()), + DataType::Map(entries, _) => matches!( + entries.data_type(), + DataType::Struct(kv) if kv.iter().all(|f| !is_complex(f.data_type())) + ), + _ => false, + }; + return if elements_are_primitive { + reject_on_non_empty() + } else { + reject() + }; + } + // Reject reading a string/binary Parquet column as anything else. Spark's // `ParquetVectorUpdaterFactory.getUpdater` BINARY case allows StringType / // BinaryType, or DecimalType only when the column carries a @@ -513,7 +568,7 @@ fn check_leaf_conversion( // nulls, parse strings, or surface as a generic Arrow type-mismatch error. // See #4088 and #4351. if is_string_or_binary(physical_type) && !is_string_or_binary(target_type) { - return reject(); + return reject_on_non_empty(); } // Reject reading a primitive numeric Parquet column as StringType / @@ -547,13 +602,13 @@ fn check_leaf_conversion( let src_int_precision = i32::from(*src_p) - i32::from(*src_s); let dst_int_precision = i32::from(*dst_p) - i32::from(*dst_s); if dst_s < src_s || dst_int_precision < src_int_precision { - return reject(); + return reject_on_non_empty(); } } // Integer-to-decimal narrowing. Spark's `canReadAsDecimal` requires // `precision - scale >= 10` for an INT32 source and `>= 20` for INT64. - // Unconditional in all Spark versions, so reject at plan time. See #4344. + // Unconditional in all Spark versions. See #4344. let int_decimal_min_int_precision = match physical_type { DataType::Int8 | DataType::Int16 | DataType::Int32 => Some(10i32), DataType::Int64 => Some(20i32), @@ -567,7 +622,7 @@ fn check_leaf_conversion( if let Some((dst_p, dst_s)) = dst_precision_scale { let dst_int_precision = i32::from(dst_p) - i32::from(dst_s); if dst_int_precision < min_int_precision { - return reject(); + return reject_on_non_empty(); } } } @@ -664,28 +719,6 @@ fn check_leaf_conversion( return reject_on_non_empty(); } - // Scalar/complex mismatch (e.g. TIMESTAMP read as ARRAY): - // Spark's vectorized reader rejects with - // SchemaColumnConvertNotSupportedException (SPARK-45604). Same-shape - // complex pairs never reach this leaf check (`check_conversion` walks their - // leaves instead), so two complex types here differ in shape (e.g. STRUCT - // read as ARRAY), which Spark rejects just the same. - let is_complex = |t: &DataType| { - matches!( - t, - DataType::Struct(_) - | DataType::List(_) - | DataType::LargeList(_) - | DataType::FixedSizeList(_, _) - | DataType::ListView(_) - | DataType::LargeListView(_) - | DataType::Map(_, _) - ) - }; - if is_complex(physical_type) || is_complex(target_type) { - return reject(); - } - ConversionCheck::Accept } @@ -1632,9 +1665,9 @@ pub(crate) mod test { use arrow::array::cast::AsArray; use arrow::array::UInt32Array; use arrow::array::{ - Array, ArrayRef, BinaryArray, Date32Array, Decimal128Array, DictionaryArray, - FixedSizeListArray, Float32Array, Float64Array, Int32Array, Int64Array, LargeListArray, - ListArray, MapArray, StringArray, StructArray, TimestampMicrosecondArray, + new_null_array, Array, ArrayRef, BinaryArray, Date32Array, Decimal128Array, + DictionaryArray, FixedSizeListArray, Float32Array, Float64Array, Int32Array, Int64Array, + LargeListArray, ListArray, MapArray, StringArray, StructArray, TimestampMicrosecondArray, TimestampMillisecondArray, }; use arrow::buffer::OffsetBuffer; @@ -1649,7 +1682,8 @@ pub(crate) mod test { use datafusion::datasource::source::DataSourceExec; use datafusion::execution::object_store::ObjectStoreUrl; use datafusion::execution::TaskContext; - use datafusion::physical_expr::expressions::Column; + use datafusion::logical_expr::Operator; + use datafusion::physical_expr::expressions::{BinaryExpr, Column, Literal}; use datafusion::physical_expr::PhysicalExpr; use datafusion::physical_plan::{ExecutionPlan, SendableRecordBatchStream}; use datafusion_comet_spark_expr::test_common::file_util::get_temp_filename; @@ -2228,6 +2262,16 @@ pub(crate) mod test { batch: &RecordBatch, required_schema: SchemaRef, options: SparkParquetOptions, + ) -> Result { + scan_parquet_with_predicate(batch, required_schema, options, None) + } + + /// [`scan_parquet`], with `predicate` (if any) given to the scan for row-group pruning. + fn scan_parquet_with_predicate( + batch: &RecordBatch, + required_schema: SchemaRef, + options: SparkParquetOptions, + predicate: Option>, ) -> Result { let filename = get_temp_filename(); let filename = filename.as_path().as_os_str().to_str().unwrap().to_string(); @@ -2242,7 +2286,10 @@ pub(crate) mod test { let expr_adapter_factory: Arc = Arc::new(SparkPhysicalExprAdapterFactory::new(options, None)); - let parquet_source = ParquetSource::new(required_schema); + let mut parquet_source = ParquetSource::new(required_schema); + if let Some(predicate) = predicate { + parquet_source = parquet_source.with_predicate(predicate); + } let files = FileGroup::new(vec![PartitionedFile::from_path(filename)?]); let file_scan_config = @@ -2994,6 +3041,157 @@ pub(crate) mod test { Ok(()) } + /// Spark checks a conversion only while it decodes a row group, so a file with nothing to + /// decode reads even when a column, top-level or nested, has a type Spark rejects. A file + /// with a row still fails (#6506). + #[tokio::test] + async fn rejected_conversions_pass_for_empty_file() -> Result<(), DataFusionError> { + for (physical_type, target_type) in [ + (DataType::Utf8, DataType::Int32), + (DataType::Binary, DataType::Decimal128(37, 1)), + (DataType::Decimal128(10, 4), DataType::Decimal128(10, 2)), + (DataType::Int32, DataType::Decimal128(9, 0)), + (DataType::Int64, DataType::Int32), + (DataType::Int32, list_type(DataType::Int32)), + (DataType::Int32, map_type(DataType::Int32)), + ] { + for nested in [false, true] { + let shape = |t: &DataType| { + if nested { + struct_type(vec![("x", t.clone())]) + } else { + t.clone() + } + }; + let label = format!("{physical_type} read as {target_type}, nested: {nested}"); + let file_schema = Arc::new(Schema::new(vec![Field::new( + "s", + shape(&physical_type), + true, + )])); + let required_schema = Arc::new(Schema::new(vec![Field::new( + "s", + shape(&target_type), + true, + )])); + + let empty = RecordBatch::new_empty(Arc::clone(&file_schema)); + let mut stream = + scan_parquet(&empty, Arc::clone(&required_schema), default_options())?; + while let Some(batch) = stream.next().await { + let batch = batch.unwrap_or_else(|e| panic!("{label}: {e}")); + assert_eq!(batch.num_rows(), 0, "{label}"); + } + + let one_row = RecordBatch::try_new( + Arc::clone(&file_schema), + vec![new_null_array(&shape(&physical_type), 1)], + )?; + let mut stream = scan_parquet(&one_row, required_schema, default_options())?; + let err = stream.next().await.unwrap().expect_err(&label); + let column = if nested { "[[s, x]]" } else { "[[s]]" }; + assert!(err.to_string().contains(column), "{label}: {err}"); + } + } + Ok(()) + } + + /// The same holds for a file whose only row group a filter prunes (#6506). + #[tokio::test] + async fn rejected_conversion_passes_for_pruned_row_group() -> Result<(), DataFusionError> { + for nested in [false, true] { + let decimals: ArrayRef = Arc::new( + Decimal128Array::from(vec![12_345i128]) + .with_precision_and_scale(10, 4) + .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?, + ); + let (s, read_type) = if nested { + let fields = Fields::from(vec![Field::new("d", DataType::Decimal128(10, 4), true)]); + ( + Arc::new(StructArray::try_new(fields, vec![decimals], None)?) as ArrayRef, + struct_type(vec![("d", DataType::Decimal128(10, 2))]), + ) + } else { + (decimals, DataType::Decimal128(10, 2)) + }; + let batch = RecordBatch::try_new( + Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int32, false), + Field::new("s", s.data_type().clone(), true), + ])), + vec![Arc::new(Int32Array::from(vec![1])), s], + )?; + let required_schema = Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int32, false), + Field::new("s", read_type, true), + ])); + let id_is_100: Arc = Arc::new(BinaryExpr::new( + Arc::new(Column::new("id", 0)), + Operator::Eq, + Arc::new(Literal::new(ScalarValue::Int32(Some(100)))), + )); + + let mut stream = scan_parquet_with_predicate( + &batch, + Arc::clone(&required_schema), + default_options(), + Some(id_is_100), + )?; + while let Some(read) = stream.next().await { + assert_eq!(read?.num_rows(), 0, "nested: {nested}"); + } + + let mut stream = scan_parquet(&batch, required_schema, default_options())?; + let err = stream + .next() + .await + .unwrap() + .expect_err("the row group fails when it is decoded"); + let column = if nested { "[[s, d]]" } else { "[[s]]" }; + assert!(err.to_string().contains(column), "nested: {nested}: {err}"); + } + Ok(()) + } + + /// A shape mismatch fails when the file is opened only where Spark fails then too, because + /// it can't clip the file's type to the requested one. Spark doesn't clip a primitive read + /// as an array or map of primitives, so only decoding a row group rejects that (#6506). + #[test] + fn shape_mismatch_rejects_at_open_only_where_spark_cannot_clip() -> Result<(), DataFusionError> + { + let options = default_options(); + let int_list = list_type(DataType::Int32); + for (physical, target) in [ + (DataType::Int32, int_list.clone()), + (DataType::Int32, map_type(DataType::Int32)), + ] { + assert!( + matches!( + check_conversion(&physical, &target, "a", &options)?, + ConversionCheck::RejectOnNonEmpty { .. } + ), + "{physical} read as {target}" + ); + } + for (physical, target) in [ + (int_list.clone(), DataType::Int32), + (DataType::Int32, struct_type(vec![("y", DataType::Int32)])), + (DataType::Utf8, struct_type(vec![("y", DataType::Int32)])), + (DataType::Int32, list_type(int_list.clone())), + (DataType::Int32, map_type(int_list.clone())), + (struct_type(vec![("y", DataType::Int32)]), int_list), + ] { + assert!( + matches!( + check_conversion(&physical, &target, "a", &options)?, + ConversionCheck::Reject(_) + ), + "{physical} read as {target}" + ); + } + Ok(()) + } + /// Positive: `INT32 -> bigint` inside a struct still converts when type promotion is /// allowed (Spark 4.x behaviour). #[tokio::test] diff --git a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala index 9b6bceeef25..e165b6961de 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala @@ -1649,6 +1649,50 @@ abstract class ParquetReadSuite extends CometTestBase { } } + test("native scan reads files with nothing to decode whose types Spark rejects") { + // Regression guard for #6506. Spark checks a conversion only while it decodes a row group, + // so an empty file, or one whose only row group a filter prunes, reads even when a column + // has a type the read schema can't convert. The native scan used to reject such a file + // when it opened it. + withSQLConf(SQLConf.USE_V1_SOURCE_LIST.key -> "parquet") { + val cases = Seq( + // (written expression of the read type, written expression of another type, read type) + ("named_struct('x', 1)", "named_struct('x', 'a')", "struct"), + ( + "named_struct('d', cast(1.23 as decimal(10,2)))", + "named_struct('d', cast(1.2345 as decimal(10,4)))", + "struct"), + ("array(1)", "array('a')", "array"), + ("map('k', 1)", "map('k', 'v')", "map"), + ("named_struct('x', array(1))", "named_struct('x', 1)", "struct>"), + ("1", "'a'", "int")) + cases.foreach { case (matching, mismatched, readType) => + withClue(s"$mismatched read as $readType: ") { + withTempPath { dir => + val path = dir.getCanonicalPath + spark.sql(s"select $matching as s").write.parquet(path) + spark.sql(s"select $mismatched as s where false").write.mode("append").parquet(path) + checkSparkAnswerAndOperator(spark.read.schema(s"s $readType").parquet(path)) + } + withTempPath { dir => + val path = dir.getCanonicalPath + spark.sql(s"select 100 as id, $matching as s").repartition(1).write.parquet(path) + spark + .sql(s"select 1 as id, $mismatched as s") + .repartition(1) + .write + .mode("append") + .parquet(path) + val df = spark.read.schema(s"id int, s $readType").parquet(path) + checkSparkAnswerAndOperator(df.where("id = 100")) + // The pruned row group still fails when it is decoded. + intercept[SparkException](df.collect()) + } + } + } + } + } + test("nested schema evolution follows Spark's per-version widening rules") { // Companion to "schema evolution": `INT32 -> bigint` inside a struct is gated by the same // per-Spark-version constant as the top level (see ShimCometConf), and accepted nested From 7576a0ad13fe47844842e1c8fdf1d06b0d5a2fb4 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 1 Oct 2026 10:42:46 -0600 Subject: [PATCH 2/5] fix: let an open-time Parquet rejection win over a deferred sibling `check_conversion` returned the first non-Accept verdict in leaf order, so a struct or map whose deferred leaf mismatch came before a shape mismatch Spark fails when it opens the file read an empty file without error. A `Reject` now wins wherever it is, as it already does across top-level columns. Also document that a pushed row filter can skip a deferred mismatch error, and fix two doc comments that described the old verdict rules. --- .../user-guide/latest/compatibility/scans.md | 6 + native/core/src/parquet/schema_adapter.rs | 134 ++++++++++++------ .../comet/parquet/ParquetReadSuite.scala | 14 ++ 3 files changed, 112 insertions(+), 42 deletions(-) diff --git a/docs/source/user-guide/latest/compatibility/scans.md b/docs/source/user-guide/latest/compatibility/scans.md index d4d45b100a9..9e5aa73fa37 100644 --- a/docs/source/user-guide/latest/compatibility/scans.md +++ b/docs/source/user-guide/latest/compatibility/scans.md @@ -125,6 +125,12 @@ and `INT32 → DOUBLE` widening that Spark 4.0+ accepts unconditionally; `Timest is rejected by Spark 3.x but accepted by Spark 4.0+). Comet aims to follow the per-version Spark behavior. +- **A pushed row filter can skip the mismatch error**. Like Spark, Comet raises most mismatches only + while it decodes a file's rows, so an empty file, or one whose row groups are all pruned, reads + without error. With `spark.comet.parquet.rowFilterPushdown.enabled=true` (off by default), the scan + decodes a column that the filter doesn't read only for the rows that pass the filter, so a file in + which no row passes also reads without error. Spark decodes every row of a row group it doesn't + prune, so it raises the error for that file. - **List conversion error paths assume Spark's standard encoding**. Comet inserts `list` before the element name when reporting a rejected array element conversion. Arrow's schema omits the repeated group name, so paths for legacy LIST encodings or custom group names may diff --git a/native/core/src/parquet/schema_adapter.rs b/native/core/src/parquet/schema_adapter.rs index 93f28df9f47..1a79c6dab8c 100644 --- a/native/core/src/parquet/schema_adapter.rs +++ b/native/core/src/parquet/schema_adapter.rs @@ -722,6 +722,25 @@ fn check_leaf_conversion( ConversionCheck::Accept } +/// Combine the verdicts of sibling fields. A `Reject` wins wherever it is, because Spark +/// raises it when it opens the file, before it decodes anything. Otherwise the first deferred +/// verdict wins, like Spark, which raises for the first offending column it decodes. +fn combine_checks( + checks: impl Iterator>, +) -> DataFusionResult { + let mut deferred = None; + for check in checks { + match check? { + ConversionCheck::Accept => {} + reject @ ConversionCheck::Reject(_) => return Ok(reject), + check @ ConversionCheck::RejectOnNonEmpty { .. } => { + deferred.get_or_insert(check); + } + } + } + Ok(deferred.unwrap_or(ConversionCheck::Accept)) +} + /// Check a physical/logical type pair the way Spark's vectorized reader does. Spark runs /// `getUpdater` on every *leaf* column regardless of nesting, so same-shape complex pairs /// (struct / list / map, at any depth) are walked and [`check_leaf_conversion`] is applied to @@ -731,8 +750,7 @@ fn check_leaf_conversion( /// encoding used by Spark. Legacy or custom group names cannot be recovered from this schema. /// Requested struct fields resolve to file fields with the same field-id / case-fold rules the runtime /// convert uses ([`match_struct_fields`]); requested fields missing from the file are skipped -/// (they read as null / default, as before). The first non-`Accept` verdict in leaf order -/// wins, like Spark, which raises for the first offending column it initializes. +/// (they read as null / default, as before). Sibling verdicts combine as in [`combine_checks`]. fn check_conversion( physical_type: &DataType, target_type: &DataType, @@ -746,22 +764,17 @@ fn check_conversion( match (physical_type, target_type) { (DataType::Struct(physical_fields), DataType::Struct(target_fields)) => { let physical_indices = match_struct_fields(physical_fields, target_fields, options)?; - for (target_field, physical_index) in target_fields.iter().zip(physical_indices) { - let Some(physical_index) = physical_index else { - continue; - }; - let physical_field = &physical_fields[physical_index]; - let check = check_conversion( - physical_field.data_type(), - target_field.data_type(), - &format!("{column}, {}", physical_field.name()), - options, - )?; - if !matches!(check, ConversionCheck::Accept) { - return Ok(check); - } - } - Ok(ConversionCheck::Accept) + combine_checks(target_fields.iter().zip(physical_indices).filter_map( + |(target_field, physical_index)| { + let physical_field = &physical_fields[physical_index?]; + Some(check_conversion( + physical_field.data_type(), + target_field.data_type(), + &format!("{column}, {}", physical_field.name()), + options, + )) + }, + )) } ( DataType::List(physical_item) @@ -796,22 +809,20 @@ fn check_conversion( target_type, ))); } - for (physical_field, target_field) in physical_kv.iter().zip(target_kv.iter()) { - let check = check_conversion( - physical_field.data_type(), - target_field.data_type(), - &format!( - "{column}, {}, {}", - physical_entries.name(), - physical_field.name() - ), - options, - )?; - if !matches!(check, ConversionCheck::Accept) { - return Ok(check); - } - } - return Ok(ConversionCheck::Accept); + return combine_checks(physical_kv.iter().zip(target_kv.iter()).map( + |(physical_field, target_field)| { + check_conversion( + physical_field.data_type(), + target_field.data_type(), + &format!( + "{column}, {}, {}", + physical_entries.name(), + physical_field.name() + ), + options, + ) + }, + )); } Ok(ConversionCheck::Reject(parquet_schema_convert_err( column, @@ -1554,14 +1565,14 @@ impl SparkPhysicalExprAdapter { } } -/// Defers a Parquet type-promotion rejection to runtime: returns an empty array +/// Defers a Parquet type conversion rejection to runtime: returns an empty array /// when the input batch has no rows, and raises `ParquetSchemaConvert` otherwise. /// /// Mirrors Spark's vectorized reader, which only invokes /// `ParquetVectorUpdaterFactory.getUpdater` while decoding a row group. A /// Parquet file with no row groups (e.g. one written from an empty DataFrame) /// never triggers the per-row-group check, so a partition mixing such a file -/// with another whose schema would otherwise fail the type-promotion check +/// with another whose schema would otherwise fail the conversion check /// (SPARK-26709) is still readable. #[derive(Debug, Eq)] struct RejectOnNonEmpty { @@ -3192,6 +3203,47 @@ pub(crate) mod test { Ok(()) } + /// Spark fails an unclippable shape when it opens the file, before it decodes a sibling it + /// would also reject, so the shape mismatch wins over a deferred sibling in either order. + #[test] + fn rejection_at_open_wins_over_a_deferred_sibling() -> Result<(), DataFusionError> { + let options = default_options(); + let int_list = list_type(DataType::Int32); + let map = |key: DataType, value: DataType| { + let entries = Fields::from(vec![ + Field::new("key", key, false), + Field::new("value", value, true), + ]); + DataType::Map( + Arc::new(Field::new("key_value", DataType::Struct(entries), false)), + false, + ) + }; + for (physical, target) in [ + ( + struct_type(vec![("a", DataType::Utf8), ("b", int_list.clone())]), + struct_type(vec![("a", DataType::Int32), ("b", DataType::Int32)]), + ), + ( + struct_type(vec![("b", int_list.clone()), ("a", DataType::Utf8)]), + struct_type(vec![("b", DataType::Int32), ("a", DataType::Int32)]), + ), + ( + map(DataType::Int64, int_list.clone()), + map(DataType::Int32, DataType::Int32), + ), + ] { + assert!( + matches!( + check_conversion(&physical, &target, "s", &options)?, + ConversionCheck::Reject(_) + ), + "{physical} read as {target}" + ); + } + Ok(()) + } + /// Positive: `INT32 -> bigint` inside a struct still converts when type promotion is /// allowed (Spark 4.x behaviour). #[tokio::test] @@ -3339,18 +3391,16 @@ pub(crate) mod test { #[test] fn issue_5783_fallback_deferred_rejection_checks_decoded_subtree() { + // An int read as a map of primitives is deferred, like `other`, and DataFusion's + // default adapter can't cast it, which forces the fallback. let logical = struct_schema(vec![ Field::new("other", DataType::Int32, true), - Field::new("force_fallback", DataType::Int32, true), + Field::new("force_fallback", map_type(DataType::Int32), true), ]); for duplicate in [false, true] { let mut fields = vec![ Field::new("other", DataType::Int64, true), - Field::new( - "force_fallback", - DataType::List(Arc::new(Field::new("element", DataType::Int32, true))), - true, - ), + Field::new("force_fallback", DataType::Int32, true), Field::new("dup", DataType::Int64, true), ]; if duplicate { diff --git a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala index e165b6961de..b4edfade018 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala @@ -1690,6 +1690,20 @@ abstract class ParquetReadSuite extends CometTestBase { } } } + // Spark fails a shape it can't clip when it opens the file, so an empty file with one + // still fails, even after a field whose type Spark only rejects while decoding. + withTempPath { dir => + val path = dir.getCanonicalPath + spark.sql("select named_struct('a', 1, 'b', 1) as s").write.parquet(path) + spark + .sql("select named_struct('a', 'x', 'b', array(1)) as s where false") + .write + .mode("append") + .parquet(path) + val (sparkError, cometError) = + checkSparkAnswerMaybeThrows(spark.read.schema("s struct").parquet(path)) + assert(sparkError.isDefined && cometError.isDefined, s"$sparkError, $cometError") + } } } From 54bd7cfec33c7bfa14169fa6040c51f277c8384a Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 1 Oct 2026 10:59:09 -0600 Subject: [PATCH 3/5] refactor: decide open-time versus decode-time Parquet rejection in one place The `getUpdater` rules are now a true/false predicate, `spark_has_updater`, and `check_leaf_conversion` alone decides when a rejection fires: an unclippable shape mismatch is an error when the file opens, and every other rejection is deferred. No rule can pick open-time rejection on its own any more, which is how #6506 came about. The `Reject` variant is gone because every caller turned it into an error, so `?` now makes an open-time rejection win over a deferred one. Test cleanups: a `map_type_with_key` helper replaces a copy of `map_type`, the pruned-row-group checks reuse `nested_rejection_message` and assert that Spark fails too, and the Scala test drops two `repartition(1)` calls a one-row select doesn't need. --- native/core/src/parquet/schema_adapter.rs | 279 +++++++----------- .../comet/parquet/ParquetReadSuite.scala | 14 +- 2 files changed, 118 insertions(+), 175 deletions(-) diff --git a/native/core/src/parquet/schema_adapter.rs b/native/core/src/parquet/schema_adapter.rs index 1a79c6dab8c..224612424f4 100644 --- a/native/core/src/parquet/schema_adapter.rs +++ b/native/core/src/parquet/schema_adapter.rs @@ -462,13 +462,11 @@ fn reject_on_non_empty_expr( } /// Outcome of checking one Parquet (physical) -> Spark (logical) type pair against the -/// conversion rules of Spark's vectorized Parquet reader. +/// conversion rules of Spark's vectorized Parquet reader. A pair Spark rejects when it opens +/// the file, before it decodes anything, is an error from [`check_conversion`] instead. enum ConversionCheck { /// Spark has an updater for the pair (for a same-shape complex pair: for every leaf). Accept, - /// Spark rejects the pair when it opens the file, before it decodes anything; raised at - /// plan time. - Reject(DataFusionError), /// Spark rejects the pair, but only while decoding a row group, so the rejection is /// deferred to runtime via [`RejectOnNonEmpty`] (SPARK-26709). A file with nothing to /// decode, empty or with every row group pruned, reads. Carries the offending @@ -480,86 +478,79 @@ enum ConversionCheck { }, } -/// Apply the rejection matrix of Spark's `ParquetVectorUpdaterFactory.getUpdater` to a single -/// physical/logical leaf pair. `column` is the Spark-style column path used in the error (`a` -/// for a top-level column, `s, x` for a nested leaf, mirroring -/// `Arrays.toString(descriptor.getPath())`). The rules and their order are exactly those the -/// adapter applies to top-level columns; [`check_conversion`] applies them to nested leaves. +/// Whether Parquet stores `data_type` as a group: a struct, a list or a map. +fn is_complex(data_type: &DataType) -> bool { + matches!( + data_type, + DataType::Struct(_) + | DataType::List(_) + | DataType::LargeList(_) + | DataType::FixedSizeList(_, _) + | DataType::ListView(_) + | DataType::LargeListView(_) + | DataType::Map(_, _) + ) +} + +/// Check a pair that [`check_conversion`] doesn't walk: two primitives, or two types of +/// different shape. `column` is the Spark-style column path used in the error (`a` for a +/// top-level column, `s, x` for a nested leaf, mirroring `Arrays.toString(descriptor.getPath())`). /// -/// Spark calls `getUpdater` only while it decodes a row group, so every rejection here is -/// deferred to runtime except a shape mismatch that Spark already fails when it opens the -/// file (#6506). +/// A shape mismatch (e.g. TIMESTAMP read as ARRAY, or STRUCT read as ARRAY) fails +/// when Spark opens the file if Spark can't clip the file's type to the requested one: a group +/// read as another type (`ParquetToSparkSchemaConverter`), or a primitive read as a struct, or +/// as an array or map with a complex element (`ParquetReadSupport.clipParquetType`). Every +/// other pair Spark rejects, including a primitive read as an array or map of primitives +/// (SPARK-45604), is rejected only by `getUpdater`, which Spark calls while decoding a row +/// group, so the rejection is deferred to runtime (#6506). fn check_leaf_conversion( physical_type: &DataType, target_type: &DataType, column: &str, options: &SparkParquetOptions, -) -> ConversionCheck { +) -> DataFusionResult { if physical_type == target_type { - return ConversionCheck::Accept; - } - let reject = || { - ConversionCheck::Reject(parquet_schema_convert_err( - column, - physical_type, - target_type, - )) - }; - let reject_on_non_empty = || ConversionCheck::RejectOnNonEmpty { + return Ok(ConversionCheck::Accept); + } + if is_complex(physical_type) || is_complex(target_type) { + let is_unclipped = !is_complex(physical_type) + && match target_type { + DataType::List(item) + | DataType::LargeList(item) + | DataType::FixedSizeList(item, _) + | DataType::ListView(item) + | DataType::LargeListView(item) => !is_complex(item.data_type()), + DataType::Map(entries, _) => matches!( + entries.data_type(), + DataType::Struct(kv) if kv.iter().all(|f| !is_complex(f.data_type())) + ), + _ => false, + }; + if !is_unclipped { + return Err(parquet_schema_convert_err( + column, + physical_type, + target_type, + )); + } + } else if spark_has_updater(physical_type, target_type, options) { + return Ok(ConversionCheck::Accept); + } + Ok(ConversionCheck::RejectOnNonEmpty { column: column.to_string(), physical_type: physical_type.clone(), target_type: target_type.clone(), - }; - - // Scalar/complex mismatch (e.g. TIMESTAMP read as ARRAY): - // Spark's vectorized reader rejects with - // SchemaColumnConvertNotSupportedException (SPARK-45604). Same-shape - // complex pairs never reach this leaf check (`check_conversion` walks their - // leaves instead), so two complex types here differ in shape (e.g. STRUCT - // read as ARRAY), which Spark rejects just the same. - // - // Checked first because the shape decides when Spark rejects. It fails when - // it opens the file if it can't clip the file's type to the requested one, - // which it can't for a group (`ParquetToSparkSchemaConverter`) or for a - // primitive read as a struct, or as an array or map with a complex element - // (`ParquetReadSupport.clipParquetType`). It doesn't clip a primitive read as - // an array or map of primitives, so only `getUpdater` rejects that, while - // decoding, as in SPARK-45604. - let is_complex = |t: &DataType| { - matches!( - t, - DataType::Struct(_) - | DataType::List(_) - | DataType::LargeList(_) - | DataType::FixedSizeList(_, _) - | DataType::ListView(_) - | DataType::LargeListView(_) - | DataType::Map(_, _) - ) - }; - if is_complex(physical_type) { - return reject(); - } - if is_complex(target_type) { - let elements_are_primitive = match target_type { - DataType::List(item) - | DataType::LargeList(item) - | DataType::FixedSizeList(item, _) - | DataType::ListView(item) - | DataType::LargeListView(item) => !is_complex(item.data_type()), - DataType::Map(entries, _) => matches!( - entries.data_type(), - DataType::Struct(kv) if kv.iter().all(|f| !is_complex(f.data_type())) - ), - _ => false, - }; - return if elements_are_primitive { - reject_on_non_empty() - } else { - reject() - }; - } + }) +} +/// Whether Spark's `ParquetVectorUpdaterFactory.getUpdater` has an updater that reads the +/// primitive Parquet (physical) type as the requested (logical) type. The same rules apply to +/// top-level columns and, through [`check_conversion`], to nested leaves. +fn spark_has_updater( + physical_type: &DataType, + target_type: &DataType, + options: &SparkParquetOptions, +) -> bool { // Reject reading a string/binary Parquet column as anything else. Spark's // `ParquetVectorUpdaterFactory.getUpdater` BINARY case allows StringType / // BinaryType, or DecimalType only when the column carries a @@ -568,14 +559,11 @@ fn check_leaf_conversion( // nulls, parse strings, or surface as a generic Arrow type-mismatch error. // See #4088 and #4351. if is_string_or_binary(physical_type) && !is_string_or_binary(target_type) { - return reject_on_non_empty(); + return false; } // Reject reading a primitive numeric Parquet column as StringType / - // BinaryType. Spark has no `int -> string` etc. updater. Defer to - // runtime via `RejectOnNonEmpty` so empty Parquet files (SPARK-26709) - // pass and the JVM shim translates to - // `SchemaColumnConvertNotSupportedException`. + // BinaryType. Spark has no `int -> string` etc. updater. let physical_is_primitive_numeric = matches!( physical_type, DataType::Boolean @@ -587,7 +575,7 @@ fn check_leaf_conversion( | DataType::Float64 ); if physical_is_primitive_numeric && is_string_or_binary(target_type) { - return reject_on_non_empty(); + return false; } // Decimal-to-decimal narrowing. Spark's `isDecimalTypeMatched` (the @@ -602,7 +590,7 @@ fn check_leaf_conversion( let src_int_precision = i32::from(*src_p) - i32::from(*src_s); let dst_int_precision = i32::from(*dst_p) - i32::from(*dst_s); if dst_s < src_s || dst_int_precision < src_int_precision { - return reject_on_non_empty(); + return false; } } @@ -622,7 +610,7 @@ fn check_leaf_conversion( if let Some((dst_p, dst_s)) = dst_precision_scale { let dst_int_precision = i32::from(dst_p) - i32::from(dst_s); if dst_int_precision < min_int_precision { - return reject_on_non_empty(); + return false; } } } @@ -630,8 +618,7 @@ fn check_leaf_conversion( // Type promotion (widening). When `allow_type_promotion` is false, // reject the three widenings (INT32→INT64, FLOAT→DOUBLE, INT32→DOUBLE) // that Spark 3.x's vectorized reader rejects. The flag tracks Comet's - // per-Spark-version constant in ShimCometConf. Deferred to runtime so - // empty files (SPARK-26709) pass. + // per-Spark-version constant in ShimCometConf. if !options.allow_type_promotion { let is_disallowed_promotion = matches!( (physical_type, target_type), @@ -640,7 +627,7 @@ fn check_leaf_conversion( | (DataType::Int32, DataType::Float64) ); if is_disallowed_promotion { - return reject_on_non_empty(); + return false; } } @@ -657,7 +644,7 @@ fn check_leaf_conversion( // - `Date32 -> Timestamp(LTZ)`: Spark only allows Date -> TimestampNTZ. // - `Timestamp -> Date32`: no Timestamp updater branches into Date. // - // Deferred to runtime (SPARK-26709). See #4297. + // See #4297. let is_spark_rejected_conversion = matches!( (physical_type, target_type), // Long -> narrower int. @@ -693,14 +680,13 @@ fn check_leaf_conversion( | (DataType::Timestamp(_, _), DataType::Date32) ); if is_spark_rejected_conversion { - return reject_on_non_empty(); + return false; } // Spark 3.x refuses to read a Parquet TimestampLTZ column as // TimestampNTZ (SPARK-36182); Spark 4.0 (SPARK-47447) lifted that. // The flag tracks Comet's per-Spark-version constant in - // ShimCometConf. Deferred to runtime so empty files (SPARK-26709) - // still pass. See #4219. + // ShimCometConf. See #4219. // // This catches all LTZ physical encodings: TIMESTAMP_MICROS / // TIMESTAMP_MILLIS arrive as `Timestamp(_, Some(_))` directly, and @@ -716,29 +702,26 @@ fn check_leaf_conversion( ) ) { - return reject_on_non_empty(); + return false; } - ConversionCheck::Accept + true } -/// Combine the verdicts of sibling fields. A `Reject` wins wherever it is, because Spark -/// raises it when it opens the file, before it decodes anything. Otherwise the first deferred +/// Combine the verdicts of sibling fields. A rejection at open is an error, so it wins wherever +/// it is, as in Spark, which raises it before it decodes anything. Otherwise the first deferred /// verdict wins, like Spark, which raises for the first offending column it decodes. fn combine_checks( checks: impl Iterator>, ) -> DataFusionResult { - let mut deferred = None; + let mut combined = ConversionCheck::Accept; for check in checks { - match check? { - ConversionCheck::Accept => {} - reject @ ConversionCheck::Reject(_) => return Ok(reject), - check @ ConversionCheck::RejectOnNonEmpty { .. } => { - deferred.get_or_insert(check); - } + let check = check?; + if matches!(combined, ConversionCheck::Accept) { + combined = check; } } - Ok(deferred.unwrap_or(ConversionCheck::Accept)) + Ok(combined) } /// Check a physical/logical type pair the way Spark's vectorized reader does. Spark runs @@ -750,7 +733,8 @@ fn combine_checks( /// encoding used by Spark. Legacy or custom group names cannot be recovered from this schema. /// Requested struct fields resolve to file fields with the same field-id / case-fold rules the runtime /// convert uses ([`match_struct_fields`]); requested fields missing from the file are skipped -/// (they read as null / default, as before). Sibling verdicts combine as in [`combine_checks`]. +/// (they read as null / default, as before). A pair Spark rejects when it opens the file is an +/// error, and sibling verdicts combine as in [`combine_checks`]. fn check_conversion( physical_type: &DataType, target_type: &DataType, @@ -803,11 +787,11 @@ fn check_conversion( (physical_entries.data_type(), target_entries.data_type()) { if physical_kv.len() != 2 || target_kv.len() != 2 { - return Ok(ConversionCheck::Reject(parquet_schema_convert_err( + return Err(parquet_schema_convert_err( column, physical_type, target_type, - ))); + )); } return combine_checks(physical_kv.iter().zip(target_kv.iter()).map( |(physical_field, target_field)| { @@ -824,18 +808,13 @@ fn check_conversion( }, )); } - Ok(ConversionCheck::Reject(parquet_schema_convert_err( + Err(parquet_schema_convert_err( column, physical_type, target_type, - ))) + )) } - _ => Ok(check_leaf_conversion( - physical_type, - target_type, - column, - options, - )), + _ => check_leaf_conversion(physical_type, target_type, column, options), } } @@ -1281,7 +1260,6 @@ impl SparkPhysicalExprAdapter { &self.parquet_options, )? { ConversionCheck::Accept => {} - ConversionCheck::Reject(err) => return Err(err), ConversionCheck::RejectOnNonEmpty { column, physical_type: leaf_physical_type, @@ -1390,7 +1368,6 @@ impl SparkPhysicalExprAdapter { &self.parquet_options, )? { ConversionCheck::Accept => {} - ConversionCheck::Reject(err) => return Err(err), ConversionCheck::RejectOnNonEmpty { column, physical_type: leaf_physical_type, @@ -2890,10 +2867,7 @@ pub(crate) mod test { ), ] { for (physical, target) in [(&valid, &invalid), (&invalid, &valid)] { - assert!(matches!( - check_conversion(physical, target, "m", &options)?, - ConversionCheck::Reject(_) - )); + assert!(check_conversion(physical, target, "m", &options).is_err()); } } Ok(()) @@ -2902,8 +2876,13 @@ pub(crate) mod test { /// Build the Arrow `Map` type with Parquet's `key_value` / `key` / /// `value` names, as Spark-written files surface it. fn map_type(value_type: DataType) -> DataType { + map_type_with_key(DataType::Utf8, value_type) + } + + /// [`map_type`] with keys of `key_type`. + fn map_type_with_key(key_type: DataType, value_type: DataType) -> DataType { let entries = Fields::from(vec![ - Field::new("key", DataType::Utf8, false), + Field::new("key", key_type, false), Field::new("value", value_type, true), ]); DataType::Map( @@ -3067,24 +3046,17 @@ pub(crate) mod test { (DataType::Int32, map_type(DataType::Int32)), ] { for nested in [false, true] { - let shape = |t: &DataType| { - if nested { + let schema = |t: &DataType| { + let t = if nested { struct_type(vec![("x", t.clone())]) } else { t.clone() - } + }; + Arc::new(Schema::new(vec![Field::new("s", t, true)])) }; let label = format!("{physical_type} read as {target_type}, nested: {nested}"); - let file_schema = Arc::new(Schema::new(vec![Field::new( - "s", - shape(&physical_type), - true, - )])); - let required_schema = Arc::new(Schema::new(vec![Field::new( - "s", - shape(&target_type), - true, - )])); + let file_schema = schema(&physical_type); + let required_schema = schema(&target_type); let empty = RecordBatch::new_empty(Arc::clone(&file_schema)); let mut stream = @@ -3096,7 +3068,7 @@ pub(crate) mod test { let one_row = RecordBatch::try_new( Arc::clone(&file_schema), - vec![new_null_array(&shape(&physical_type), 1)], + vec![new_null_array(file_schema.field(0).data_type(), 1)], )?; let mut stream = scan_parquet(&one_row, required_schema, default_options())?; let err = stream.next().await.unwrap().expect_err(&label); @@ -3111,11 +3083,8 @@ pub(crate) mod test { #[tokio::test] async fn rejected_conversion_passes_for_pruned_row_group() -> Result<(), DataFusionError> { for nested in [false, true] { - let decimals: ArrayRef = Arc::new( - Decimal128Array::from(vec![12_345i128]) - .with_precision_and_scale(10, 4) - .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?, - ); + let decimals: ArrayRef = + Arc::new(Decimal128Array::from(vec![12_345i128]).with_precision_and_scale(10, 4)?); let (s, read_type) = if nested { let fields = Fields::from(vec![Field::new("d", DataType::Decimal128(10, 4), true)]); ( @@ -3152,14 +3121,9 @@ pub(crate) mod test { assert_eq!(read?.num_rows(), 0, "nested: {nested}"); } - let mut stream = scan_parquet(&batch, required_schema, default_options())?; - let err = stream - .next() - .await - .unwrap() - .expect_err("the row group fails when it is decoded"); + let err = nested_rejection_message(&batch, required_schema).await; let column = if nested { "[[s, d]]" } else { "[[s]]" }; - assert!(err.to_string().contains(column), "nested: {nested}: {err}"); + assert!(err.contains(column), "nested: {nested}: {err}"); } Ok(()) } @@ -3193,10 +3157,7 @@ pub(crate) mod test { (struct_type(vec![("y", DataType::Int32)]), int_list), ] { assert!( - matches!( - check_conversion(&physical, &target, "a", &options)?, - ConversionCheck::Reject(_) - ), + check_conversion(&physical, &target, "a", &options).is_err(), "{physical} read as {target}" ); } @@ -3206,19 +3167,9 @@ pub(crate) mod test { /// Spark fails an unclippable shape when it opens the file, before it decodes a sibling it /// would also reject, so the shape mismatch wins over a deferred sibling in either order. #[test] - fn rejection_at_open_wins_over_a_deferred_sibling() -> Result<(), DataFusionError> { + fn rejection_at_open_wins_over_a_deferred_sibling() { let options = default_options(); let int_list = list_type(DataType::Int32); - let map = |key: DataType, value: DataType| { - let entries = Fields::from(vec![ - Field::new("key", key, false), - Field::new("value", value, true), - ]); - DataType::Map( - Arc::new(Field::new("key_value", DataType::Struct(entries), false)), - false, - ) - }; for (physical, target) in [ ( struct_type(vec![("a", DataType::Utf8), ("b", int_list.clone())]), @@ -3229,19 +3180,15 @@ pub(crate) mod test { struct_type(vec![("b", DataType::Int32), ("a", DataType::Int32)]), ), ( - map(DataType::Int64, int_list.clone()), - map(DataType::Int32, DataType::Int32), + map_type_with_key(DataType::Int64, int_list.clone()), + map_type_with_key(DataType::Int32, DataType::Int32), ), ] { assert!( - matches!( - check_conversion(&physical, &target, "s", &options)?, - ConversionCheck::Reject(_) - ), + check_conversion(&physical, &target, "s", &options).is_err(), "{physical} read as {target}" ); } - Ok(()) } /// Positive: `INT32 -> bigint` inside a struct still converts when type promotion is diff --git a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala index b4edfade018..b5f51022810 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala @@ -1676,17 +1676,13 @@ abstract class ParquetReadSuite extends CometTestBase { } withTempPath { dir => val path = dir.getCanonicalPath - spark.sql(s"select 100 as id, $matching as s").repartition(1).write.parquet(path) - spark - .sql(s"select 1 as id, $mismatched as s") - .repartition(1) - .write - .mode("append") - .parquet(path) + spark.sql(s"select 100 as id, $matching as s").write.parquet(path) + spark.sql(s"select 1 as id, $mismatched as s").write.mode("append").parquet(path) val df = spark.read.schema(s"id int, s $readType").parquet(path) checkSparkAnswerAndOperator(df.where("id = 100")) - // The pruned row group still fails when it is decoded. - intercept[SparkException](df.collect()) + // The pruned row group still fails in both engines when it is decoded. + val (sparkError, cometError) = checkSparkAnswerMaybeThrows(df) + assert(sparkError.isDefined && cometError.isDefined, s"$sparkError, $cometError") } } } From 3c252caccfa954300d404b495762d667377fa6e1 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 1 Oct 2026 19:34:19 -0600 Subject: [PATCH 4/5] fix: preserve Parquet conversion failures before row filtering --- .../user-guide/latest/compatibility/scans.md | 11 +- .../eager_page_index_reader_factory.rs | 141 +++++++- native/core/src/parquet/parquet_exec.rs | 10 +- native/core/src/parquet/schema_adapter.rs | 306 +++++++++++++++++- .../comet/parquet/ParquetReadSuite.scala | 62 ++++ 5 files changed, 505 insertions(+), 25 deletions(-) diff --git a/docs/source/user-guide/latest/compatibility/scans.md b/docs/source/user-guide/latest/compatibility/scans.md index 9e5aa73fa37..b6208bef832 100644 --- a/docs/source/user-guide/latest/compatibility/scans.md +++ b/docs/source/user-guide/latest/compatibility/scans.md @@ -125,12 +125,11 @@ and `INT32 → DOUBLE` widening that Spark 4.0+ accepts unconditionally; `Timest is rejected by Spark 3.x but accepted by Spark 4.0+). Comet aims to follow the per-version Spark behavior. -- **A pushed row filter can skip the mismatch error**. Like Spark, Comet raises most mismatches only - while it decodes a file's rows, so an empty file, or one whose row groups are all pruned, reads - without error. With `spark.comet.parquet.rowFilterPushdown.enabled=true` (off by default), the scan - decodes a column that the filter doesn't read only for the rows that pass the filter, so a file in - which no row passes also reads without error. Spark decodes every row of a row group it doesn't - prune, so it raises the error for that file. +- **Conversion errors precede pushed row filters**. Empty files and files whose row groups or + pages are all pruned read without decode-time conversion errors. When row-filter pushdown is + enabled, Comet checks rejected conversions before reading surviving data pages, so a row filter + cannot suppress the error by discarding every row. Legacy LIST shape mismatches that Spark cannot + clip are rejected when the file is opened, including for empty and fully pruned files. - **List conversion error paths assume Spark's standard encoding**. Comet inserts `list` before the element name when reporting a rejected array element conversion. Arrow's schema omits the repeated group name, so paths for legacy LIST encodings or custom group names may 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..5749dcd138a 100644 --- a/native/core/src/parquet/eager_page_index_reader_factory.rs +++ b/native/core/src/parquet/eager_page_index_reader_factory.rs @@ -45,8 +45,7 @@ //! //! Filed upstream as apache/datafusion#23978. Once the opener merges its deferred page-index //! load back into `FileMetadataCache` instead of bypassing it, the eager policy can go, but the -//! factory cannot: `get_metadata` is the one per-file hook that sees the raw footer, and two -//! other things hang off it. +//! factory cannot: `get_metadata` is the per-file hook that sees the raw footer. //! //! The first is Spark's missing field id check. `ParquetReadSupport` refuses to open a file //! whose Parquet schema carries no field id when the requested schema carries one, unless @@ -59,8 +58,17 @@ //! //! 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}; +//! +//! Conversion checks also use the raw footer to preserve legacy LIST clipping failures. With +//! row filters enabled, a deferred conversion failure is raised on the first data-page request. +//! Footer, Bloom-filter and page-index reads still proceed, so format pruning can discard a file +//! without decoding it, but row selection cannot hide the failure. + +use crate::parquet::parquet_support::SparkParquetOptions; +use crate::parquet::schema_adapter::{ + check_file_conversions, RejectOnNonEmpty, REPEATED_PRIMITIVE_KEY, +}; +use arrow::datatypes::{DataType, FieldRef, Schema, SchemaRef}; use async_trait::async_trait; use bytes::Bytes; use datafusion::common::Result as DFResult; @@ -93,10 +101,11 @@ use parquet::file::metadata::{ FooterTail, PageIndexPolicy, ParquetMetaData, ParquetMetaDataReader, }; use parquet::schema::types::{ColumnDescPtr, SchemaDescriptor, Type as ParquetType}; +use std::collections::HashMap; use std::fmt::{Debug, Display, Formatter}; use std::ops::Range; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; -use std::sync::{Arc, OnceLock}; +use std::sync::{Arc, Mutex, OnceLock}; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) enum ScanIoSource { @@ -181,6 +190,8 @@ pub struct EagerPageIndexReaderFactory { // Refuse a file whose Parquet schema carries no field id, as Spark's `ParquetReadSupport` // does when the requested schema carries one and `ignoreMissing` is not set. require_field_ids: bool, + conversion_check: Option<(SchemaRef, SparkParquetOptions, bool)>, + deferred_rejections: Arc>>, } impl EagerPageIndexReaderFactory { @@ -210,9 +221,23 @@ impl EagerPageIndexReaderFactory { scan_io_metrics, spark_variant_schema: false, require_field_ids: false, + conversion_check: None, + deferred_rejections: Arc::new(Mutex::new(HashMap::new())), } } + /// Retain raw LIST encoding for open-time validation. When row filters are enabled, + /// enforce deferred failures at the first data-page read, after format pruning. + pub(crate) fn with_conversion_check( + mut self, + required: SchemaRef, + options: SparkParquetOptions, + row_filters: bool, + ) -> Self { + self.conversion_check = Some((required, options, row_filters)); + self + } + pub fn with_spark_variant_schema(mut self, enabled: bool) -> Self { self.spark_variant_schema = enabled; self @@ -254,6 +279,15 @@ impl ParquetFileReaderFactory for EagerPageIndexReaderFactory { metrics, ); + // DataFusion replaces the metadata reader after Bloom-filter pruning. The replacement + // receives already-loaded metadata, so retain deferred failures in this scan's factory. + let deferred_rejection = self + .deferred_rejections + .lock() + .unwrap() + .get(&partitioned_file.object_meta.location) + .filter(|(meta, _)| meta == &partitioned_file.object_meta) + .map(|(_, rejection)| rejection.clone()); Ok(Box::new(EagerPageIndexReader { file_metrics, scan_io_metrics: Arc::clone(&self.scan_io_metrics), @@ -263,6 +297,9 @@ impl ParquetFileReaderFactory for EagerPageIndexReaderFactory { metadata_size_hint, spark_variant_schema: self.spark_variant_schema, require_field_ids: self.require_field_ids, + conversion_check: self.conversion_check.clone(), + deferred_rejections: Arc::clone(&self.deferred_rejections), + deferred_rejection, })) } } @@ -279,6 +316,9 @@ struct EagerPageIndexReader { metadata_size_hint: Option, spark_variant_schema: bool, require_field_ids: bool, + conversion_check: Option<(SchemaRef, SparkParquetOptions, bool)>, + deferred_rejections: Arc>>, + deferred_rejection: Option, } // Arrow infers ENUM as Binary, losing the distinction from raw binary that Spark needs. @@ -411,6 +451,53 @@ fn with_spark_arrow_schema(metadata: Arc) -> ParquetResult Schema { + fn mark(field: &FieldRef, columns: &[ColumnDescPtr], index: &mut usize) -> FieldRef { + let data_type = match field.data_type() { + DataType::Struct(fields) => DataType::Struct( + fields + .iter() + .map(|f| mark(f, columns, index)) + .collect::>() + .into(), + ), + DataType::List(item) => DataType::List(mark(item, columns, index)), + DataType::LargeList(item) => DataType::LargeList(mark(item, columns, index)), + DataType::FixedSizeList(item, size) => { + DataType::FixedSizeList(mark(item, columns, index), *size) + } + DataType::Map(entries, sorted) => DataType::Map(mark(entries, columns, index), *sorted), + _ => { + let repeated = columns.get(*index).is_some_and(|column| { + column.self_type().get_basic_info().repetition() + == parquet::basic::Repetition::REPEATED + }); + *index += 1; + let mut field = field.as_ref().clone(); + if repeated { + let mut metadata = field.metadata().clone(); + metadata.insert(REPEATED_PRIMITIVE_KEY.to_string(), "true".to_string()); + field = field.with_metadata(metadata); + } + return Arc::new(field); + } + }; + Arc::new(field.as_ref().clone().with_data_type(data_type)) + } + let mut index = 0; + Schema::new_with_metadata( + schema + .fields() + .iter() + .map(|f| mark(f, parquet.columns(), &mut index)) + .collect::>(), + schema.metadata().clone(), + ) +} + impl AsyncFileReader for EagerPageIndexReader { /// Reads a metadata range, counting its requested size before I/O and its returned /// bytes only on success. The returned future borrows this reader; store errors retain @@ -449,6 +536,12 @@ impl AsyncFileReader for EagerPageIndexReader { where Self: Send, { + if ranges.iter().any(|range| !range.is_empty()) { + if let Some(rejection) = &self.deferred_rejection { + let error = ParquetError::External(Box::new(rejection.error())); + return async move { Err(error) }.boxed(); + } + } let requested = ranges_bytes(&ranges); self.file_metrics.bytes_scanned.add(requested); let scan_io_metrics = Arc::clone(&self.scan_io_metrics); @@ -550,11 +643,43 @@ impl AsyncFileReader for EagerPageIndexReader { }, ))); } - if spark_variant_schema { - with_spark_arrow_schema(metadata) + let metadata = if spark_variant_schema { + with_spark_arrow_schema(metadata)? } else { - Ok(metadata) + metadata + }; + if let Some((required, conversion_options, row_filters)) = + self.conversion_check.as_ref() + { + let file = metadata.file_metadata(); + let needs_check = *row_filters + || file.schema_descr().columns().iter().any(|column| { + column.self_type().get_basic_info().repetition() + == parquet::basic::Repetition::REPEATED + }); + if !needs_check { + return Ok(metadata); + } + let schema = + parquet_to_arrow_schema(file.schema_descr(), file.key_value_metadata())?; + let schema = mark_repeated_primitives(&schema, file.schema_descr()); + let rejection = check_file_conversions( + Arc::clone(required), + Arc::new(schema), + conversion_options.clone(), + ) + .map_err(|e| ParquetError::External(Box::new(e)))?; + if *row_filters { + if let Some(rejection) = &rejection { + self.deferred_rejections.lock().unwrap().insert( + object_meta.location.clone(), + (object_meta.clone(), rejection.clone()), + ); + } + self.deferred_rejection = rejection; + } } + Ok(metadata) } .boxed() } diff --git a/native/core/src/parquet/parquet_exec.rs b/native/core/src/parquet/parquet_exec.rs index 5c5d26d3922..8149c3656b7 100644 --- a/native/core/src/parquet/parquet_exec.rs +++ b/native/core/src/parquet/parquet_exec.rs @@ -196,7 +196,15 @@ pub(crate) fn init_datasource_exec( parquet_source.metrics(), ) .with_spark_variant_schema(projects_variant) - .with_require_field_ids(require_field_ids), + .with_require_field_ids(require_field_ids) + .with_conversion_check( + Arc::clone(&required_schema), + spark_parquet_options.clone(), + session_config.options().execution.parquet.pushdown_filters + && data_filters + .as_ref() + .is_some_and(|filters| !filters.is_empty()), + ), ); parquet_source = parquet_source.with_parquet_file_reader_factory(reader_factory); diff --git a/native/core/src/parquet/schema_adapter.rs b/native/core/src/parquet/schema_adapter.rs index 224612424f4..6bcc8d3518b 100644 --- a/native/core/src/parquet/schema_adapter.rs +++ b/native/core/src/parquet/schema_adapter.rs @@ -771,12 +771,25 @@ fn check_conversion( | DataType::FixedSizeList(target_item, _) | DataType::ListView(target_item) | DataType::LargeListView(target_item), - ) => check_conversion( - physical_item.data_type(), - target_item.data_type(), - &format!("{column}, list, {}", physical_item.name()), - options, - ), + ) => { + if physical_item + .metadata() + .contains_key(REPEATED_PRIMITIVE_KEY) + && is_complex(target_item.data_type()) + { + return Err(parquet_schema_convert_err( + column, + physical_type, + target_type, + )); + } + check_conversion( + physical_item.data_type(), + target_item.data_type(), + &format!("{column}, list, {}", physical_item.name()), + options, + ) + } ( DataType::Map(physical_entries, physical_sorted), DataType::Map(target_entries, target_sorted), @@ -1542,6 +1555,33 @@ impl SparkPhysicalExprAdapter { } } +/// Private, in-memory hint used only while checking the raw footer's LIST encoding. +pub(crate) const REPEATED_PRIMITIVE_KEY: &str = "comet.parquet.repeated_primitive"; + +/// Check only requested columns, using the same resolver and conversion rules as projection +/// adaptation. Return the first deferred failure, but continue to let open-time errors win. +pub(crate) fn check_file_conversions( + required: SchemaRef, + physical: SchemaRef, + options: SparkParquetOptions, +) -> DataFusionResult> { + let adapter = SparkPhysicalExprAdapterFactory::new(options, None) + .create(Arc::clone(&required), physical)?; + let mut deferred = None; + for (index, field) in required.fields().iter().enumerate() { + let expr = adapter.rewrite(Arc::new(Column::new(field.name(), index)))?; + expr.apply(|node| { + if deferred.is_none() { + if let Some(rejection) = node.downcast_ref::() { + deferred = Some(rejection.clone()); + } + } + Ok(datafusion::common::tree_node::TreeNodeRecursion::Continue) + })?; + } + Ok(deferred) +} + /// Defers a Parquet type conversion rejection to runtime: returns an empty array /// when the input batch has no rows, and raises `ParquetSchemaConvert` otherwise. /// @@ -1551,8 +1591,8 @@ impl SparkPhysicalExprAdapter { /// never triggers the per-row-group check, so a partition mixing such a file /// with another whose schema would otherwise fail the conversion check /// (SPARK-26709) is still readable. -#[derive(Debug, Eq)] -struct RejectOnNonEmpty { +#[derive(Debug, Clone, Eq)] +pub(crate) struct RejectOnNonEmpty { child: Arc, target_field: FieldRef, column: String, @@ -1590,6 +1630,17 @@ impl Display for RejectOnNonEmpty { } } +impl RejectOnNonEmpty { + pub(crate) fn error(&self) -> SparkError { + SparkError::ParquetSchemaConvert { + file_path: String::new(), + column: self.column.clone(), + physical_type: self.physical_type.clone(), + spark_type: self.spark_type.clone(), + } + } +} + impl PhysicalExpr for RejectOnNonEmpty { fn data_type(&self, _input_schema: &Schema) -> DataFusionResult { Ok(self.target_field.data_type().clone()) @@ -2268,13 +2319,41 @@ pub(crate) mod test { writer.write(batch)?; writer.close()?; + scan_parquet_file(filename, required_schema, options, predicate, false) + } + + fn scan_parquet_file( + filename: String, + required_schema: SchemaRef, + options: SparkParquetOptions, + predicate: Option>, + row_filters: bool, + ) -> Result { + use crate::parquet::eager_page_index_reader_factory::{ + EagerPageIndexReaderFactory, ScanIoSource, + }; + use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet; + let context = Arc::new(TaskContext::default()); + let reader_factory = EagerPageIndexReaderFactory::new( + Arc::new(object_store::local::LocalFileSystem::new()), + context + .runtime_env() + .cache_manager + .get_file_metadata_cache(), + ScanIoSource::Local, + &ExecutionPlanMetricsSet::new(), + ) + .with_conversion_check(Arc::clone(&required_schema), options.clone(), row_filters); let object_store_url = ObjectStoreUrl::local_filesystem(); // Create expression adapter factory for Spark-compatible schema adaptation let expr_adapter_factory: Arc = Arc::new(SparkPhysicalExprAdapterFactory::new(options, None)); - let mut parquet_source = ParquetSource::new(required_schema); + let mut parquet_source = ParquetSource::new(required_schema) + .with_pushdown_filters(row_filters) + .with_enable_page_index(true) + .with_parquet_file_reader_factory(Arc::new(reader_factory)); if let Some(predicate) = predicate { parquet_source = parquet_source.with_predicate(predicate); } @@ -2287,7 +2366,7 @@ pub(crate) mod test { .build(); let parquet_exec = DataSourceExec::new(Arc::new(file_scan_config)); - parquet_exec.execute(0, Arc::new(TaskContext::default())) + parquet_exec.execute(0, context) } /// Create a Parquet file containing a single batch and then read the batch back using @@ -3128,6 +3207,213 @@ pub(crate) mod test { Ok(()) } + #[tokio::test] + async fn rejected_conversion_precedes_row_filter() -> Result<(), DataFusionError> { + use parquet::file::properties::WriterProperties; + for nested in [false, true] { + let strings: ArrayRef = Arc::new(StringArray::from(vec!["bad", "bad"])); + let (values, read_type) = if nested { + let fields = Fields::from(vec![Field::new("x", DataType::Utf8, true)]); + ( + Arc::new(StructArray::try_new(fields, vec![strings], None)?) as ArrayRef, + struct_type(vec![("x", DataType::Int32)]), + ) + } else { + (strings, DataType::Int32) + }; + let schema = Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int32, false), + Field::new("s", values.data_type().clone(), true), + ])); + let batch = RecordBatch::try_new( + Arc::clone(&schema), + vec![Arc::new(Int32Array::from(vec![1, 3])), values], + )?; + let required = Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int32, false), + Field::new("s", read_type, true), + ])); + let filename = get_temp_filename().to_str().unwrap().to_string(); + let props = WriterProperties::builder() + .set_dictionary_enabled(false) + .build(); + let mut writer = ArrowWriter::try_new(File::create(&filename)?, schema, Some(props))?; + writer.write(&batch)?; + writer.close()?; + for row_filters in [false, true] { + for id in [2, 100] { + let predicate: Arc = Arc::new(BinaryExpr::new( + Arc::new(Column::new("id", 0)), + Operator::Eq, + Arc::new(Literal::new(ScalarValue::Int32(Some(id)))), + )); + let mut stream = scan_parquet_file( + filename.clone(), + Arc::clone(&required), + default_options(), + Some(predicate), + row_filters, + )?; + if id == 100 { + while let Some(batch) = stream.next().await { + assert_eq!(batch?.num_rows(), 0); + } + } else { + let error = stream + .next() + .await + .unwrap() + .expect_err("conversion must precede row selection") + .to_string(); + let column = if nested { "[[s, x]]" } else { "[[s]]" }; + assert!( + error.contains(column) + && error.contains("Expected: int") + && error.contains("Found: BINARY"), + "{error}" + ); + } + } + } + } + Ok(()) + } + + /// Page-index pruning, like row-group pruning, must finish before the data-read guard. + #[tokio::test] + async fn deferred_conversion_respects_page_pruning() -> Result<(), DataFusionError> { + use parquet::file::properties::{EnabledStatistics, WriterProperties}; + let schema = Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int32, false), + Field::new("s", DataType::Utf8, false), + ])); + let batch = RecordBatch::try_new( + Arc::clone(&schema), + vec![ + Arc::new(Int32Array::from(vec![1, 3])), + Arc::new(StringArray::from(vec!["bad", "bad"])), + ], + )?; + let required = Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int32, false), + Field::new("s", DataType::Int32, false), + ])); + let filename = get_temp_filename().to_str().unwrap().to_string(); + let props = WriterProperties::builder() + .set_dictionary_enabled(false) + .set_statistics_enabled(EnabledStatistics::Page) + .set_data_page_row_count_limit(1) + .set_write_batch_size(1) + .build(); + let mut writer = ArrowWriter::try_new(File::create(&filename)?, schema, Some(props))?; + writer.write(&batch)?; + writer.close()?; + let predicate: Arc = Arc::new(BinaryExpr::new( + Arc::new(Column::new("id", 0)), + Operator::Eq, + Arc::new(Literal::new(ScalarValue::Int32(Some(2)))), + )); + for row_filters in [false, true] { + let mut stream = scan_parquet_file( + filename.clone(), + Arc::clone(&required), + default_options(), + Some(Arc::clone(&predicate)), + row_filters, + )?; + while let Some(batch) = stream.next().await { + assert_eq!(batch?.num_rows(), 0); + } + } + Ok(()) + } + + #[tokio::test] + async fn legacy_list_shape_rejects_empty_and_pruned_files() -> Result<(), DataFusionError> { + use parquet::data_type::Int32Type as ParquetInt32; + use parquet::file::writer::SerializedFileWriter; + use parquet::schema::parser::parse_message_type; + for nested in [false, true] { + for legacy in [false, true] { + for empty in [false, true] { + let list = if legacy { + "REPEATED INT32 element;" + } else { + "REPEATED GROUP list { REQUIRED INT32 element; }" + }; + let field = format!("REQUIRED GROUP a (LIST) {{ {list} }}"); + let field = if nested { + format!("REQUIRED GROUP s {{ {field} }}") + } else { + field + }; + let schema = Arc::new(parse_message_type(&format!( + "message schema {{ REQUIRED INT32 id; {field} }}" + ))?); + let filename = get_temp_filename().to_str().unwrap().to_string(); + let mut writer = SerializedFileWriter::new( + File::create(&filename)?, + schema, + Default::default(), + )?; + if !empty { + let mut row_group = writer.next_row_group()?; + let mut id = row_group.next_column()?.unwrap(); + id.typed::().write_batch(&[1], None, None)?; + id.close()?; + let mut a = row_group.next_column()?.unwrap(); + a.typed::() + .write_batch(&[1], Some(&[1]), Some(&[0]))?; + a.close()?; + row_group.close()?; + } + writer.close()?; + for element in [list_type(DataType::Int32), map_type(DataType::Int32)] { + let required = Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int32, false), + if nested { + Field::new("s", struct_type(vec![("a", list_type(element))]), false) + } else { + Field::new("a", list_type(element), false) + }, + ])); + let predicate: Arc = Arc::new(BinaryExpr::new( + Arc::new(Column::new("id", 0)), + Operator::Eq, + Arc::new(Literal::new(ScalarValue::Int32(Some(100)))), + )); + let mut stream = scan_parquet_file( + filename.clone(), + required, + default_options(), + Some(predicate), + true, + )?; + if legacy { + let error = stream + .next() + .await + .unwrap() + .expect_err("legacy repeated primitive cannot be clipped") + .to_string(); + let column = if nested { + "Column: [[s, a]]" + } else { + "Column: [[a]]" + }; + assert!(error.contains(column), "{error}"); + } else { + while let Some(batch) = stream.next().await { + assert_eq!(batch?.num_rows(), 0); + } + } + } + } + } + } + Ok(()) + } + /// A shape mismatch fails when the file is opened only where Spark fails then too, because /// it can't clip the file's type to the requested one. Spark doesn't clip a primitive read /// as an array or map of primitives, so only decoding a row group rejects that (#6506). diff --git a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala index b5f51022810..b4f0d2683b9 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala @@ -1703,6 +1703,68 @@ abstract class ParquetReadSuite extends CometTestBase { } } + test("native scan rejects legacy LIST shape mismatches before decoding") { + for (legacy <- Seq(false, true); empty <- Seq(false, true)) { + withSQLConf( + SQLConf.USE_V1_SOURCE_LIST.key -> "parquet", + SQLConf.PARQUET_WRITE_LEGACY_FORMAT.key -> legacy.toString) { + withTempPath { dir => + val path = dir.getCanonicalPath + val rows = spark.sql("select 1 as id, array(1) as a") + (if (empty) rows.where("false") else rows).write.parquet(path) + for (element <- Seq("array", "map")) { + val df = spark.read + .schema(s"id int, a array<$element>") + .parquet(path) + .where("id = 100") + if (legacy) { + checkCometOperators(stripAQEPlan(df.queryExecution.executedPlan)) + val (sparkError, cometError) = checkSparkAnswerMaybeThrows(df) + assert(sparkError.isDefined && cometError.isDefined, s"$sparkError, $cometError") + } else { + checkSparkAnswerAndOperator(df) + } + } + } + } + } + } + + test("native scan preserves conversion errors with row filter pushdown") { + for (pushdown <- Seq(false, true); nested <- Seq(false, true)) { + withSQLConf( + SQLConf.USE_V1_SOURCE_LIST.key -> "parquet", + CometConf.COMET_PARQUET_ROW_FILTER_PUSHDOWN_ENABLED.key -> pushdown.toString) { + withTempPath { dir => + val path = dir.getCanonicalPath + val value = if (nested) "named_struct('x', 'bad')" else "'bad'" + val readType = if (nested) "struct" else "int" + spark + .sql(s"select 1 as id, $value as s union all select 3 as id, $value as s") + .coalesce(1) + .write + .option("parquet.enable.dictionary", "false") + .parquet(path) + val df = spark.read.schema(s"id int, s $readType").parquet(path) + // Statistics retain [1, 3], but no row passes id = 2. + checkCometOperators(stripAQEPlan(df.queryExecution.executedPlan)) + val (sparkError, cometError) = checkSparkAnswerMaybeThrows(df.where("id = 2")) + Seq("Spark" -> sparkError, "Comet" -> cometError).foreach { case (engine, error) => + val chain = error.toSeq.flatMap(causeChain) + assert( + chain.exists( + _.isInstanceOf[org.apache.spark.sql.execution.datasources.SchemaColumnConvertNotSupportedException]), + s"$engine: ${chain.mkString("\n")}") + } + // Format pruning must still suppress the decode-time mismatch (#6506). + checkSparkAnswerAndOperator(df.where("id = 100")) + // A mismatch in an unrequested column must not affect the scan. + checkSparkAnswerAndOperator(df.select("id").where("id = 2")) + } + } + } + } + test("nested schema evolution follows Spark's per-version widening rules") { // Companion to "schema evolution": `INT32 -> bigint` inside a struct is gated by the same // per-Spark-version constant as the top level (see ShimCometConf), and accepted nested From d9b754d7031f286b98c143baba4cc46bb2e0c369 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 1 Oct 2026 20:54:42 -0600 Subject: [PATCH 5/5] fix: preserve structured Spark errors in footer validation --- .../eager_page_index_reader_factory.rs | 67 ++++++++++++++++++- .../comet/parquet/ParquetReadSuite.scala | 27 ++++++++ 2 files changed, 92 insertions(+), 2 deletions(-) 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 5749dcd138a..dbc7d0f3b80 100644 --- a/native/core/src/parquet/eager_page_index_reader_factory.rs +++ b/native/core/src/parquet/eager_page_index_reader_factory.rs @@ -71,7 +71,7 @@ use crate::parquet::schema_adapter::{ use arrow::datatypes::{DataType, FieldRef, Schema, SchemaRef}; use async_trait::async_trait; use bytes::Bytes; -use datafusion::common::Result as DFResult; +use datafusion::common::{DataFusionError, Result as DFResult}; use datafusion::datasource::physical_plan::parquet::metadata::DFParquetMetadata; use datafusion::datasource::physical_plan::parquet::{ ParquetFileMetrics, ParquetFileReaderFactory, @@ -668,7 +668,12 @@ impl AsyncFileReader for EagerPageIndexReader { Arc::new(schema), conversion_options.clone(), ) - .map_err(|e| ParquetError::External(Box::new(e)))?; + .map_err(|e| match e { + // Keep structured Spark errors directly inside the Parquet wrapper so + // the JNI error converter preserves Spark's exception type and class. + DataFusionError::External(source) => ParquetError::External(source), + other => ParquetError::External(Box::new(other)), + })?; if *row_filters { if let Some(rejection) = &rejection { self.deferred_rejections.lock().unwrap().insert( @@ -1308,6 +1313,64 @@ mod tests { assert!(metadata_for(false, without_ids).await.is_ok()); } + #[tokio::test] + async fn conversion_check_preserves_structured_spark_errors() { + let store = Arc::new(InMemory::new()); + let schema = Arc::new(Schema::new(vec![ + arrow::datatypes::Field::new("B", DataType::Int32, false), + arrow::datatypes::Field::new("b", DataType::Int32, false), + ])); + let values: arrow::array::ArrayRef = Arc::new(Int32Array::from(vec![1, 2])); + let batch = + RecordBatch::try_new(Arc::clone(&schema), vec![Arc::clone(&values), values]).unwrap(); + let mut bytes = Vec::new(); + let mut writer = ArrowWriter::try_new(&mut bytes, schema, None).unwrap(); + writer.write(&batch).unwrap(); + writer.close().unwrap(); + let size = bytes.len() as u64; + let location = "duplicates.parquet"; + store + .put(&Path::from(location), Bytes::from(bytes).into()) + .await + .unwrap(); + let runtime = datafusion::execution::runtime_env::RuntimeEnv::default(); + let metrics = ExecutionPlanMetricsSet::new(); + let required = Arc::new(Schema::new(vec![arrow::datatypes::Field::new( + "B", + DataType::Int32, + false, + )])); + let mut options = + SparkParquetOptions::new(datafusion_comet_spark_expr::EvalMode::Legacy, "UTC", false); + options.case_sensitive = false; + let factory = EagerPageIndexReaderFactory::new( + store, + runtime.cache_manager.get_file_metadata_cache(), + ScanIoSource::Local, + &metrics, + ) + .with_conversion_check(required, options, true); + let mut reader = factory + .create_reader( + 0, + PartitionedFile::new(location.to_string(), size), + None, + &metrics, + ) + .unwrap(); + let error = reader.get_metadata(None).await.unwrap_err(); + let ParquetError::External(source) = error else { + panic!("expected structured Spark error, got {error}"); + }; + assert!( + matches!( + source.downcast_ref::(), + Some(SparkError::DuplicateFieldCaseInsensitive { .. }) + ), + "unexpected error: {source}" + ); + } + #[test] fn variant_policy_preserves_footer_metadata_and_indexes() { let schema = Arc::new(Schema::new_with_metadata( diff --git a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala index b4f0d2683b9..fa2bedb832f 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala @@ -1730,6 +1730,33 @@ abstract class ParquetReadSuite extends CometTestBase { } } + test("native scan preserves duplicate field error types with row filter pushdown") { + withTempPath { dir => + val path = dir.getCanonicalPath + withSQLConf(SQLConf.CASE_SENSITIVE.key -> "true") { + spark.range(5).selectExpr("id as A", "id as B", "id as b").write.parquet(path) + } + for (pushdown <- Seq(false, true)) { + withSQLConf( + SQLConf.USE_V1_SOURCE_LIST.key -> "parquet", + SQLConf.CASE_SENSITIVE.key -> "false", + CometConf.COMET_PARQUET_ROW_FILTER_PUSHDOWN_ENABLED.key -> pushdown.toString) { + val df = spark.read.schema("A long, B long").parquet(path).where("A = 1") + checkCometOperators(stripAQEPlan(df.queryExecution.executedPlan)) + val (sparkError, cometError) = checkSparkAnswerMaybeThrows(df) + Seq("Spark" -> sparkError, "Comet" -> cometError).foreach { case (engine, error) => + val chain = error.toSeq.flatMap(causeChain) + assert( + chain.exists(e => + e.getClass.getName == "org.apache.spark.SparkRuntimeException" && + e.getMessage.contains("Found duplicate field")), + s"$engine: ${chain.mkString("\n")}") + } + } + } + } + } + test("native scan preserves conversion errors with row filter pushdown") { for (pushdown <- Seq(false, true); nested <- Seq(false, true)) { withSQLConf(