diff --git a/native/common/src/error.rs b/native/common/src/error.rs index cd912a41994..acd7dac124d 100644 --- a/native/common/src/error.rs +++ b/native/common/src/error.rs @@ -227,7 +227,7 @@ pub enum SparkError { /// Multiple Parquet fields share the same field id when the read schema requested an /// id-based lookup. Mirrors Spark's `_LEGACY_ERROR_TEMP_2094` /// (`foundDuplicateFieldInFieldIdLookupModeError`). - #[error("[_LEGACY_ERROR_TEMP_2094] Found duplicate field(s) by id: id={required_id} matches [{matched_fields}] in id-lookup mode")] + #[error("[_LEGACY_ERROR_TEMP_2094] Found duplicate field(s) by id: id={required_id} matches {matched_fields} in id-lookup mode")] DuplicateFieldByFieldId { required_id: i32, matched_fields: String, diff --git a/native/core/src/parquet/cast_column.rs b/native/core/src/parquet/cast_column.rs index 8c72455f4db..ad7b909768b 100644 --- a/native/core/src/parquet/cast_column.rs +++ b/native/core/src/parquet/cast_column.rs @@ -24,7 +24,7 @@ use arrow::{ record_batch::RecordBatch, }; -use crate::parquet::parquet_support::{spark_parquet_convert, SparkParquetOptions}; +use crate::parquet::parquet_support::{field_id, spark_parquet_convert, SparkParquetOptions}; use datafusion::common::format::DEFAULT_CAST_OPTIONS; use datafusion::common::{DataFusionError, Result as DataFusionResult}; use datafusion::logical_expr::ColumnarValue; @@ -37,36 +37,61 @@ use std::{ }; /// Returns true if two DataTypes are structurally equivalent (same data layout) -/// but may differ in field names within nested types. -fn types_differ_only_in_field_names(physical: &DataType, logical: &DataType) -> bool { +/// but may differ in field names within nested types. With `use_field_id`, a struct +/// field that carries a Parquet field id must also find that id on the file field at +/// its position, since Spark's `clipParquetGroupFields` resolves such a field by id. +fn types_differ_only_in_field_names( + physical: &DataType, + logical: &DataType, + use_field_id: bool, +) -> bool { match (physical, logical) { (DataType::List(pf), DataType::List(lf)) => { pf.is_nullable() == lf.is_nullable() && (pf.data_type() == lf.data_type() - || types_differ_only_in_field_names(pf.data_type(), lf.data_type())) + || types_differ_only_in_field_names( + pf.data_type(), + lf.data_type(), + use_field_id, + )) } (DataType::LargeList(pf), DataType::LargeList(lf)) => { pf.is_nullable() == lf.is_nullable() && (pf.data_type() == lf.data_type() - || types_differ_only_in_field_names(pf.data_type(), lf.data_type())) + || types_differ_only_in_field_names( + pf.data_type(), + lf.data_type(), + use_field_id, + )) } (DataType::Map(pf, p_sorted), DataType::Map(lf, l_sorted)) => { p_sorted == l_sorted && pf.is_nullable() == lf.is_nullable() && (pf.data_type() == lf.data_type() - || types_differ_only_in_field_names(pf.data_type(), lf.data_type())) + || types_differ_only_in_field_names( + pf.data_type(), + lf.data_type(), + use_field_id, + )) } (DataType::Struct(pfields), DataType::Struct(lfields)) => { // For Struct types, field names are semantically meaningful (they // identify different columns), so we require name equality here. // This distinguishes from List/Map wrapper field names ("item" vs - // "element") which are purely cosmetic. + // "element") which are purely cosmetic. Under field-id matching a + // requested id must sit on the file field at the same position, or + // the relabel would read the wrong column (#6192). pfields.len() == lfields.len() && pfields.iter().zip(lfields.iter()).all(|(pf, lf)| { pf.name() == lf.name() && pf.is_nullable() == lf.is_nullable() + && (!use_field_id || field_id(lf).is_none() || field_id(lf) == field_id(pf)) && (pf.data_type() == lf.data_type() - || types_differ_only_in_field_names(pf.data_type(), lf.data_type())) + || types_differ_only_in_field_names( + pf.data_type(), + lf.data_type(), + use_field_id, + )) }) } _ => false, @@ -155,6 +180,10 @@ pub struct CometCastColumnExpr { /// Spark parquet options for complex nested type conversions. /// When present, enables `spark_parquet_convert` as a fallback. parquet_options: Option, + /// True when the physical and target types differ only in nested field names, so a + /// metadata-only relabel is the whole conversion. Derived from the fields above once + /// at construction rather than by walking the type tree on every batch. + relabel_only: bool, } // Manually derive `PartialEq`/`Hash` as `Arc` does not @@ -208,17 +237,24 @@ impl CometCastColumnExpr { ))); } + let relabel_only = physical_type != target_type + && types_differ_only_in_field_names(physical_type, target_type, false); Ok(Self { expr, input_physical_field: physical_field, target_field, cast_options: cast_options.unwrap_or(DEFAULT_CAST_OPTIONS), parquet_options: None, + relabel_only, }) } /// Set Spark parquet options to enable complex nested type conversions. pub fn with_parquet_options(mut self, options: SparkParquetOptions) -> Self { + let physical_type = self.input_physical_field.data_type(); + let target_type = self.target_field.data_type(); + self.relabel_only = physical_type != target_type + && types_differ_only_in_field_names(physical_type, target_type, options.use_field_id); self.parquet_options = Some(options); self } @@ -267,34 +303,25 @@ impl PhysicalExpr for CometCastColumnExpr { return Ok(value); } - let input_physical_field = self.input_physical_field.data_type(); let target_field = self.target_field.data_type(); - match (input_physical_field, target_field) { - // Nested types that differ only in field names (e.g., List element named - // "item" vs "element", or Map entries named "key_value" vs "entries"). - // Re-label the array so the DataType metadata matches the logical schema. - (physical, logical) - if physical != logical && types_differ_only_in_field_names(physical, logical) => - { - match value { - ColumnarValue::Array(array) => { - let relabeled = relabel_array(array, logical); - Ok(ColumnarValue::Array(relabeled)) - } - other => Ok(other), + // Nested types that differ only in field names (e.g., List element named + // "item" vs "element", or Map entries named "key_value" vs "entries"). + // Re-label the array so the DataType metadata matches the logical schema. + if self.relabel_only { + return Ok(match value { + ColumnarValue::Array(array) => { + ColumnarValue::Array(relabel_array(array, target_field)) } - } - // Fallback: use spark_parquet_convert for complex nested type conversions - // (e.g., List → List, Map field selection, etc.) - _ => { - if let Some(parquet_options) = &self.parquet_options { - let converted = spark_parquet_convert(value, target_field, parquet_options)?; - Ok(converted) - } else { - Ok(value) - } - } + other => other, + }); + } + // Fallback: use spark_parquet_convert for complex nested type conversions + // (e.g., List → List, Map field selection, etc.) + if let Some(parquet_options) = &self.parquet_options { + spark_parquet_convert(value, target_field, parquet_options) + } else { + Ok(value) } } @@ -339,6 +366,60 @@ mod tests { use datafusion::physical_expr::expressions::Column; use datafusion_comet_spark_expr::EvalMode; + /// File struct `x` (id 1) = 42, `y` (id 2) = 43; the requested struct names them the same + /// but swaps the ids. Names and types match at every position, so only the field id check + /// in the relabel shortcut keeps it from firing: the read resolves by id and the result + /// must be `x` = 43, `y` = 42 (#6192). + #[test] + fn test_swapped_field_ids_bypass_relabel_shortcut() { + use crate::parquet::schema_adapter::test::struct_type_with_field_id; + + let physical_type = + struct_type_with_field_id(vec![("x", DataType::Int32, 1), ("y", DataType::Int32, 2)]); + let logical_type = + struct_type_with_field_id(vec![("x", DataType::Int32, 2), ("y", DataType::Int32, 1)]); + let DataType::Struct(physical_fields) = &physical_type else { + unreachable!() + }; + + let input_field = Arc::new(Field::new("s", physical_type.clone(), true)); + let target_field = Arc::new(Field::new("s", logical_type.clone(), true)); + + let columns: Vec = vec![ + Arc::new(Int32Array::from(vec![42])), + Arc::new(Int32Array::from(vec![43])), + ]; + let struct_arr = StructArray::new(physical_fields.clone(), columns, None); + let schema = Schema::new(vec![Arc::clone(&input_field)]); + let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(struct_arr)]).unwrap(); + + let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + opts.use_field_id = true; + + let col_expr: Arc = Arc::new(Column::new("s", 0)); + let cast_expr = CometCastColumnExpr::try_new(col_expr, input_field, target_field, None) + .unwrap() + .with_parquet_options(opts); + + let ColumnarValue::Array(arr) = cast_expr.evaluate(&batch).unwrap() else { + panic!("expected array result"); + }; + assert_eq!(arr.data_type(), &logical_type); + let result = arr.as_any().downcast_ref::().unwrap(); + let x = result + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let y = result + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(x.value(0), 43); + assert_eq!(y.value(0), 42); + } + #[test] fn test_rejects_millisecond_logical_timestamp() { for timezone in [None, Some(Arc::from("UTC"))] { diff --git a/native/core/src/parquet/parquet_support.rs b/native/core/src/parquet/parquet_support.rs index 521f95d4590..38f7df74a00 100644 --- a/native/core/src/parquet/parquet_support.rs +++ b/native/core/src/parquet/parquet_support.rs @@ -411,13 +411,26 @@ fn list_view_visibility( } /// Read the Parquet field id stored under arrow-rs's `PARQUET_FIELD_ID_META_KEY`. -fn field_id(field: &arrow::datatypes::Field) -> Option { +pub(crate) fn field_id(field: &arrow::datatypes::Field) -> Option { field .metadata() .get(PARQUET_FIELD_ID_META_KEY) .and_then(|v| v.parse::().ok()) } +/// Names of the fields carrying `id`, for the duplicate-id error message. Bracketed and +/// comma-joined the way Spark's `matchIdField` renders the list, so the message reads +/// `Found duplicate field(s) "1": [x, y] in id mapping mode` on both sides. +pub(crate) fn field_names_with_id(fields: &[FieldRef], id: i32) -> String { + let names = fields + .iter() + .filter(|f| field_id(f) == Some(id)) + .map(|f| f.name().as_str()) + .collect::>() + .join(", "); + format!("[{names}]") +} + /// Resolve each requested (`to`) struct field to the index of the file (`from`) field it reads /// from, or `None` when the file holds no such field. Mirrors Spark's `clipParquetGroupFields`: /// when the requested struct carries Parquet field IDs anywhere (and `use_field_id` is set), @@ -425,7 +438,8 @@ fn field_id(field: &arrow::datatypes::Field) -> Option { /// fallback); other fields match by name, folded with the same `toLowerCase(Locale.ROOT)` fold /// the top-level schema adapter uses when `case_sensitive` is false. A requested field whose /// folded name matches more than one file field is rejected. Case-insensitive matching retains -/// Spark's `foundDuplicateFieldInCaseInsensitiveModeError`. +/// Spark's `foundDuplicateFieldInCaseInsensitiveModeError`, and a requested ID that more than +/// one file field carries raises Spark's `_LEGACY_ERROR_TEMP_2094` (`matchIdField`). /// /// Shared by the runtime convert (`parquet_convert_struct_to_struct`) and the plan-time /// conversion check in `schema_adapter`, so both resolve nested fields identically. @@ -437,11 +451,12 @@ pub(crate) fn match_struct_fields( let should_match_by_id = parquet_options.use_field_id && to_fields.iter().any(|f| field_id(f).is_some()); - let from_id_to_index: HashMap = if should_match_by_id { + // `None` marks an id that more than one file field carries. + let from_id_to_index: HashMap> = if should_match_by_id { let mut map = HashMap::new(); for (i, field) in from_fields.iter().enumerate() { if let Some(id) = field_id(field) { - map.entry(id).or_insert(i); + map.entry(id).and_modify(|m| *m = None).or_insert(Some(i)); } } map @@ -475,7 +490,14 @@ pub(crate) fn match_struct_fields( |(to_pos, to_field)| match (should_match_by_id, field_id(to_field)) { // Spark treats a missing ID match as a missing column rather than // falling back to name match. - (true, Some(id)) => Ok(from_id_to_index.get(&id).copied()), + (true, Some(id)) => match from_id_to_index.get(&id) { + Some(None) => Err(SparkError::DuplicateFieldByFieldId { + required_id: id, + matched_fields: field_names_with_id(from_fields, id), + } + .into()), + index => Ok(index.copied().flatten()), + }, _ => match folded_to_indices.get(to_folded[to_pos].as_str()) { // Reject selected ambiguity before a decoder can multiply rows. Some(indices) if indices.len() > 1 => { @@ -2076,4 +2098,84 @@ mod tests { "unexpected error: {err}" ); } + + /// Two file struct fields share field id 1 and the requested struct asks for that id: + /// Spark's `matchIdField` raises `foundDuplicateFieldInFieldIdLookupModeError` + /// (`_LEGACY_ERROR_TEMP_2094`) rather than silently reading the first match. + #[test] + fn requested_duplicate_field_id_errors() { + use crate::parquet::parquet_support::{parquet_convert_array, SparkParquetOptions}; + use crate::parquet::schema_adapter::test::struct_type_with_field_id; + use arrow::array::{ArrayRef, Int32Array, StructArray}; + use arrow::datatypes::DataType; + use datafusion_comet_spark_expr::EvalMode; + use std::sync::Arc; + + let DataType::Struct(from_fields) = + struct_type_with_field_id(vec![("x", DataType::Int32, 1), ("y", DataType::Int32, 1)]) + else { + unreachable!() + }; + let from: ArrayRef = Arc::new(StructArray::new( + from_fields, + vec![ + Arc::new(Int32Array::from(vec![42])), + Arc::new(Int32Array::from(vec![43])), + ], + None, + )); + let to_type = struct_type_with_field_id(vec![("f", DataType::Int32, 1)]); + + let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + opts.use_field_id = true; + + let err = parquet_convert_array(from, &to_type, &opts).unwrap_err(); + let msg = err.to_string(); + assert!( + msg.contains("_LEGACY_ERROR_TEMP_2094") && msg.contains("id=1 matches [x, y]"), + "unexpected error: {msg}" + ); + } + + /// Companion to `requested_duplicate_field_id_errors`: a duplicated file id that no + /// requested field looks up stays harmless, as Spark only raises inside `matchIdField`. + #[test] + fn unrequested_duplicate_field_id_reads_fine() { + use crate::parquet::parquet_support::{parquet_convert_array, SparkParquetOptions}; + use crate::parquet::schema_adapter::test::struct_type_with_field_id; + use arrow::array::{Array, ArrayRef, Int32Array, StructArray}; + use arrow::datatypes::DataType; + use datafusion_comet_spark_expr::EvalMode; + use std::sync::Arc; + + let DataType::Struct(from_fields) = struct_type_with_field_id(vec![ + ("x", DataType::Int32, 1), + ("y", DataType::Int32, 1), + ("z", DataType::Int32, 2), + ]) else { + unreachable!() + }; + let from: ArrayRef = Arc::new(StructArray::new( + from_fields, + vec![ + Arc::new(Int32Array::from(vec![42])), + Arc::new(Int32Array::from(vec![43])), + Arc::new(Int32Array::from(vec![44])), + ], + None, + )); + let to_type = struct_type_with_field_id(vec![("f", DataType::Int32, 2)]); + + let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + opts.use_field_id = true; + + let result = parquet_convert_array(from, &to_type, &opts).unwrap(); + let result_struct = result.as_any().downcast_ref::().unwrap(); + let col = result_struct + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(col.value(0), 44); + } } diff --git a/native/core/src/parquet/schema_adapter.rs b/native/core/src/parquet/schema_adapter.rs index dc6a7799588..46d2ad7000d 100644 --- a/native/core/src/parquet/schema_adapter.rs +++ b/native/core/src/parquet/schema_adapter.rs @@ -18,7 +18,8 @@ use crate::parquet::cast_column::CometCastColumnExpr; use crate::parquet::name_fold::{fold_name, fold_names, fold_schema_names}; use crate::parquet::parquet_support::{ - duplicate_parquet_field_error, match_struct_fields, spark_parquet_convert, SparkParquetOptions, + duplicate_parquet_field_error, field_id, field_names_with_id, match_struct_fields, + spark_parquet_convert, SparkParquetOptions, }; use arrow::array::new_empty_array; use arrow::compute::can_cast_types; @@ -36,7 +37,7 @@ use datafusion_physical_expr_adapter::{ replace_columns_with_literals, DefaultPhysicalExprAdapterFactory, PhysicalExprAdapter, PhysicalExprAdapterFactory, }; -use parquet::{arrow::PARQUET_FIELD_ID_META_KEY, variant::VariantType}; +use parquet::variant::VariantType; use std::collections::{HashMap, HashSet}; use std::fmt::{self, Display}; use std::hash::{Hash, Hasher}; @@ -68,16 +69,8 @@ impl SparkPhysicalExprAdapterFactory { } } -/// Read the Parquet field id stored under arrow-rs's `PARQUET_FIELD_ID_META_KEY`. -fn parse_field_id(field: &Field) -> Option { - field - .metadata() - .get(PARQUET_FIELD_ID_META_KEY) - .and_then(|v| v.parse::().ok()) -} - fn schema_has_field_ids(schema: &SchemaRef) -> bool { - schema.fields().iter().any(|f| parse_field_id(f).is_some()) + schema.fields().iter().any(|f| field_id(f).is_some()) } /// Returns true when casting `physical_type` to `target_type` is a *pure* structural @@ -119,9 +112,7 @@ fn is_pure_structural_narrowing( // Comet matches by Parquet field id first when the target carries one; // DataFusion's generic cast has no field-id concept, so any field-id-bearing // target field is a potential divergence. - if parquet_options.use_field_id - && target_fields.iter().any(|f| parse_field_id(f).is_some()) - { + if parquet_options.use_field_id && target_fields.iter().any(|f| field_id(f).is_some()) { return Ok(false); } // Fold the source field names once (O(sources), not O(targets x sources)), matching @@ -221,7 +212,7 @@ fn remap_physical_schema( let mut id_to_phys_names: HashMap> = HashMap::new(); if should_match_by_id { for pf in physical_schema.fields() { - if let Some(id) = parse_field_id(pf) { + if let Some(id) = field_id(pf) { id_to_phys_names .entry(id) .or_default() @@ -229,15 +220,14 @@ fn remap_physical_schema( } } for lf in logical_schema.fields() { - if let Some(id) = parse_field_id(lf) { + if let Some(id) = field_id(lf) { if let Some(matches) = id_to_phys_names.get(&id) { if matches.len() > 1 { - return Err(DataFusionError::External(Box::new( - SparkError::DuplicateFieldByFieldId { - required_id: id, - matched_fields: matches.join(", "), - }, - ))); + return Err(SparkError::DuplicateFieldByFieldId { + required_id: id, + matched_fields: field_names_with_id(physical_schema.fields(), id), + } + .into()); } } } @@ -248,7 +238,7 @@ fn remap_physical_schema( let id_to_logical: HashMap = if should_match_by_id { let mut map = HashMap::new(); for lf in logical_schema.fields() { - if let Some(id) = parse_field_id(lf) { + if let Some(id) = field_id(lf) { map.entry(id).or_insert(lf); } } @@ -270,7 +260,7 @@ fn remap_physical_schema( .fields() .iter() .zip(&logical_folded) - .filter(|(field, _)| should_match_by_id && parse_field_id(field).is_some()) + .filter(|(field, _)| should_match_by_id && field_id(field).is_some()) .map(|(_, name)| name) .collect(); let mut occupied_names = HashSet::new(); @@ -284,7 +274,7 @@ fn remap_physical_schema( .map(|(phys_idx, field)| { // ID match first when the logical schema is ID-bearing. if should_match_by_id { - if let Some(phys_id) = parse_field_id(field) { + if let Some(phys_id) = field_id(field) { if let Some(logical_field) = id_to_logical.get(&phys_id) { if logical_field.name() != field.name() { name_map.insert(logical_field.name().clone(), field.name().clone()); @@ -329,7 +319,7 @@ fn remap_physical_schema( .iter() .enumerate() .find(|(j, lf)| { - let lf_has_id = should_match_by_id && parse_field_id(lf).is_some(); + let lf_has_id = should_match_by_id && field_id(lf).is_some(); !lf_has_id && logical_folded[*j] == physical_folded[phys_idx] }) .map(|(_, lf)| lf); @@ -955,7 +945,7 @@ impl PhysicalExprAdapterFactory for SparkPhysicalExprAdapterFactory { .fields() .iter() .zip(&logical_folded) - .filter(|(lf, _)| parse_field_id(lf).is_some()) + .filter(|(lf, _)| field_id(lf).is_some()) .map(|(_, folded)| folded.clone()) .collect::>(), ) @@ -975,14 +965,14 @@ impl PhysicalExprAdapterFactory for SparkPhysicalExprAdapterFactory { .fields() .iter() .filter(|f| duplicate_names.contains(f.name())) - .filter_map(|f| parse_field_id(f).map(|id| (id, f.name().clone()))) + .filter_map(|f| field_id(f).map(|id| (id, f.name().clone()))) .collect(); logical_file_schema .fields() .iter() .zip(&logical_folded) .filter_map(|(field, folded)| { - parse_field_id(field) + field_id(field) .and_then(|id| duplicated_ids.get(&id)) .map(|name| (folded.clone(), name.clone())) }) @@ -1050,6 +1040,15 @@ struct SparkPhysicalExprAdapter { /// Folded logical name -> byte-identical duplicate physical name. Populated only when /// matching by field ID, then checked before the ID-resolved name skip so decoded duplicate /// roots still fail. + /// + /// This rejects a read even when the requested id is unambiguous, say `a (id 1)` from a + /// file holding `a (id 1)`, `a (id 2)` and `b (id 3)`, where Spark's `matchIdField` finds + /// exactly one field. The rejection protects the decoder rather than diverging from a + /// Spark behaviour Comet could match: with this guard removed Comet returns two rows from + /// that one-row file for either id, Spark's vectorized reader returns id 1's value for + /// both ids, and parquet-mr fails outright, all from the name-based leaf lookup #5964 + /// describes. `CometNativeReaderSuite` "duplicate Parquet field names - root group and + /// unprojected root duplicates" pins the rejection end to end. id_duplicate_roots: HashMap, /// `logical_file_schema` field names pre-folded once (see `fold_names`), parallel to /// `logical_file_schema.fields()`. Lets the per-column rewrite fallbacks match by folded name @@ -1631,7 +1630,7 @@ impl PhysicalExpr for RejectOnNonEmpty { } #[cfg(test)] -mod test { +pub(crate) mod test { use crate::parquet::cast_column::CometCastColumnExpr; use crate::parquet::parquet_support::SparkParquetOptions; use crate::parquet::schema_adapter::{ @@ -3830,7 +3829,7 @@ mod test { ) } - fn struct_type_with_field_id(fields: Vec<(&str, DataType, i32)>) -> DataType { + pub(crate) fn struct_type_with_field_id(fields: Vec<(&str, DataType, i32)>) -> DataType { DataType::Struct( fields .into_iter() @@ -4170,4 +4169,62 @@ mod test { let target = struct_type(vec![("id", DataType::Int64)]); assert!(!is_pure_structural_narrowing(&physical, &target, &opts).unwrap()); } + + /// 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), + /// `y` (id 3) and `z` (id 2). Read positionally the names line up and all three values + /// come back, but Spark's `clipParquetSchema` raises because requested id 1 resolves to + /// two file fields. The check runs when the file schema is mapped, before any value is + /// handed back. + #[tokio::test] + async fn parquet_duplicate_file_field_id_rejected_when_requested() { + let file_type = struct_type_with_field_id(vec![ + ("x", DataType::Int64, 1), + ("y", DataType::Int64, 1), + ("z", DataType::Int64, 2), + ]); + let requested_type = struct_type_with_field_id(vec![ + ("x", DataType::Int64, 1), + ("y", DataType::Int64, 3), + ("z", DataType::Int64, 2), + ]); + let DataType::Struct(file_fields) = &file_type else { + unreachable!() + }; + let file_schema = Arc::new(Schema::new(vec![ + Field::new("s", file_type.clone(), true).with_metadata(id_meta("10")) + ])); + let required_schema = Arc::new(Schema::new(vec![ + Field::new("s", requested_type, true).with_metadata(id_meta("10")) + ])); + let children: Vec> = vec![ + Arc::new(Int64Array::from(vec![42])), + Arc::new(Int64Array::from(vec![43])), + Arc::new(Int64Array::from(vec![44])), + ]; + let col = Arc::new(arrow::array::StructArray::new( + file_fields.clone(), + children, + None, + )) as Arc; + let batch = RecordBatch::try_new(file_schema, vec![col]).unwrap(); + + let mut opts = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false); + opts.use_field_id = true; + + let err = match scan_parquet(&batch, required_schema, opts) { + Ok(mut stream) => stream + .next() + .await + .unwrap() + .expect_err("requested id 1 matches two file fields and must error"), + Err(err) => err, + }; + let msg = err.to_string(); + assert!( + msg.contains("_LEGACY_ERROR_TEMP_2094") && msg.contains("id=1 matches [x, y]"), + "expected duplicate field id error, got: {msg}" + ); + } } diff --git a/spark/src/test/scala/org/apache/comet/exec/CometNativeReaderSuite.scala b/spark/src/test/scala/org/apache/comet/exec/CometNativeReaderSuite.scala index 61fdd35bca5..ec94a708fcb 100644 --- a/spark/src/test/scala/org/apache/comet/exec/CometNativeReaderSuite.scala +++ b/spark/src/test/scala/org/apache/comet/exec/CometNativeReaderSuite.scala @@ -36,7 +36,7 @@ import org.apache.spark.sql.functions.{array, col} import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types._ -import org.apache.comet.CometConf +import org.apache.comet.{CometConf, CometNativeException} import org.apache.comet.CometSparkSessionExtensions.isSpark41Plus class CometNativeReaderSuite extends CometTestBase with AdaptiveSparkPlanHelper { @@ -117,6 +117,11 @@ class CometNativeReaderSuite extends CometTestBase with AdaptiveSparkPlanHelper val messages = causeMessages(error) assert(messages.contains("duplicate Parquet field name 'dup'"), messages) assert(!messages.toLowerCase.contains("case-insensitive"), messages) + // The refusal is Comet's own execution error, the class the root duplicate check + // raises too, and never Spark's INTERNAL_ERROR, which is reserved for engine bugs. + val chain = causeChain(error) + assert(chain.exists(_.isInstanceOf[CometNativeException]), messages) + assert(!messages.contains("INTERNAL_ERROR"), messages) } } } 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 d412918b572..287d4ecb14a 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala @@ -49,7 +49,7 @@ import org.apache.spark.sql.types._ import com.google.common.primitives.UnsignedLong import org.apache.comet.CometConf -import org.apache.comet.CometSparkSessionExtensions.isSpark40Plus +import org.apache.comet.CometSparkSessionExtensions.{isSpark40Plus, isSpark41Plus} import org.apache.comet.vector.CometVector abstract class ParquetReadSuite extends CometTestBase { @@ -2079,10 +2079,102 @@ abstract class ParquetReadSuite extends CometTestBase { } } - // Verbatim port of Spark `ParquetFieldIdIOSuite.test("multiple id matches")` so the shim - // error path is exercised on both 3.x and 4.x. The stock suite is the CI signal but it - // requires the Spark test jars and `withAllParquetReaders`; keeping a copy here lets us - // iterate locally. + // The two shapes of #6192. A nested column dropped and added back under its old name gets a + // fresh field id, so the file holds `struct` while the table reads + // `struct`, and a swapped pair reads `struct`. The + // names line up at every position, which is exactly what the metadata-only relabel shortcut + // in the native cast looks for, so without the field id check in that shortcut the scan + // hands back the file's values by position. Spark returns null for the re-added `x` and + // keeps `y`, and swaps the two values for the swapped ids. Each nesting is read on its own, + // so a wrong answer names the level that produced it. Before Spark 4.1 the vectorized reader + // raises on both reads below a list or map, since its column vector rejects the clipped + // struct, which carries a placeholder field for the unmatched id and the file's field order + // for the swapped ids, so that comparison with Spark runs from 4.1 on. The pinned rows hold + // everywhere. + test("nested field ids resolve by id below struct, list and map, not by position") { + def struct(xId: Int, yId: Int): StructType = new StructType() + .add("x", LongType, true, withId(xId)) + .add("y", LongType, true, withId(yId)) + def schema(inner: StructType): StructType = new StructType() + .add("id", LongType, true, withId(10)) + .add("s", inner, true, withId(11)) + .add("l", ArrayType(inner), true, withId(12)) + .add("m", MapType(StringType, inner), true, withId(13)) + val writeData = Seq( + Row(1L, Row(1L, 10L), Seq(Row(2L, 20L), Row(3L, 30L)), Map("k" -> Row(4L, 40L))), + Row(2L, Row(5L, 50L), Seq(Row(6L, 60L), null), Map("a" -> Row(7L, 70L), "b" -> null)), + Row(3L, Row(null, 80L), Seq(), Map()), + Row(4L, null, null, null)) + // Each read schema with what Spark answers for one file struct `(x, y)` under it. + val cases = Seq( + ("dropped and re-added", struct(3, 2), (r: Row) => Row(null, r.get(1))), + ("swapped", struct(2, 1), (r: Row) => Row(r.get(1), r.get(0)))) + val columns = Seq(("s", 1), ("l", 2), ("m", 3)) + + withSQLConf( + SQLConf.PARQUET_FIELD_ID_WRITE_ENABLED.key -> "true", + SQLConf.PARQUET_FIELD_ID_READ_ENABLED.key -> "true") { + withTempPath { dir => + spark + .createDataFrame(spark.sparkContext.parallelize(writeData), schema(struct(1, 2))) + .write + .mode("overwrite") + .parquet(dir.getCanonicalPath) + for ((label, readStruct, remap) <- cases; (column, index) <- columns) { + withClue(s"$label, column $column: ") { + def read(): DataFrame = spark.read + .schema(schema(readStruct)) + .parquet(dir.getCanonicalPath) + .select("id", column) + .sort("id") + def remapNested(value: Any): Any = value match { + case null => null + case r: Row => remap(r) + case seq: Seq[_] => seq.map(remapNested) + case map: Map[_, _] => map.map { case (k, v) => k -> remapNested(v) } + } + val expected = writeData.map(r => Row(r.get(0), remapNested(r.get(index)))) + + if (column == "s" || isSpark41Plus) { + checkSparkAnswerAndOperator(read()) + } + val plan = stripAQEPlan(read().queryExecution.executedPlan) + assert( + collect(plan) { case scan: CometNativeScanExec => scan }.nonEmpty, + s"expected CometNativeScanExec in the plan:\n$plan") + // Pin the values too, so the test states what Spark answers rather than only + // that Comet agrees with it. + checkAnswer(read(), expected) + } + } + } + } + } + + // The duplicate field id error Comet raises for `df` must carry the message Spark raises + // for the same read, with the matching file fields listed the way `matchIdField` lists + // them. `matched` is that list, for instance `"1": [x, y]`. + private def checkDuplicateFieldIdMessage(df: => DataFrame, matched: String): Unit = { + val (sparkError, cometError) = checkSparkAnswerMaybeThrows(df) + // The deepest cause carrying the text: the wrappers above it quote it with a stack trace. + def message(error: Option[Throwable]): String = + error.toSeq + .flatMap(causeChain) + .flatMap(e => Option(e.getMessage)) + .reverse + .find(_.contains("Found duplicate field(s)")) + .getOrElse(fail(s"expected a duplicate field id error, got: $error")) + val expected = message(sparkError) + assert( + expected.contains(s"Found duplicate field(s) $matched in id mapping mode"), + s"Spark did not raise the expected duplicate field id error: $expected") + assert(message(cometError) == expected) + } + + // Port of Spark `ParquetFieldIdIOSuite.test("multiple id matches")` so the shim error path + // is exercised on both 3.x and 4.x. The stock suite is the CI signal but it requires the + // Spark test jars and `withAllParquetReaders`. Keeping a copy here lets us iterate locally. + // On top of the port, the message is compared with Spark's in full. test("multiple id matches") { withSQLConf(SQLConf.PARQUET_FIELD_ID_READ_ENABLED.key -> "true") { withTempPath { dir => @@ -2109,6 +2201,52 @@ abstract class ParquetReadSuite extends CometTestBase { assert( cause.isInstanceOf[RuntimeException] && cause.getMessage.contains("Found duplicate field(s)")) + checkDuplicateFieldIdMessage( + spark.read.schema(readSchema).parquet(dir.getCanonicalPath), + """"1": [a, rand2]""") + } + } + } + + test("duplicate field id inside a struct is rejected when a requested id matches two fields") { + // The requested struct names one id that two file fields carry. A requested schema that + // repeats an id itself is declined at planning time, so this is the shape the native scan + // still has to refuse. The schema adapter refuses it while resolving the file's fields, + // with the message Spark raises for the same read. + withSQLConf(SQLConf.PARQUET_FIELD_ID_READ_ENABLED.key -> "true") { + withTempPath { dir => + val schema = + new StructType() + .add( + "s", + new StructType() + .add("x", LongType, true, withId(1)) + .add("y", LongType, true, withId(1)), + true, + withId(2)) + val readSchema = + new StructType() + .add("s", new StructType().add("x", LongType, true, withId(1)), true, withId(2)) + + val writeData = Seq(Row(Row(42L, 43L))) + spark + .createDataFrame(spark.sparkContext.parallelize(writeData), schema) + .write + .mode("overwrite") + .parquet(dir.getCanonicalPath) + + val df = spark.read.schema(readSchema).parquet(dir.getCanonicalPath) + val scans = stripAQEPlan(df.queryExecution.executedPlan).collect { + case scan: CometNativeScanExec => scan + } + assert(scans.nonEmpty, "expected CometNativeScanExec in the plan") + val cause = intercept[SparkException] { + df.collect() + }.getCause + assert( + cause.isInstanceOf[RuntimeException] && + cause.getMessage.contains("Found duplicate field(s)")) + checkDuplicateFieldIdMessage(df, """"1": [x, y]""") } } }