From f23c314799e8e606183c56780908379490b0b4a8 Mon Sep 17 00:00:00 2001 From: comphead Date: Mon, 28 Sep 2026 09:49:27 -0700 Subject: [PATCH] fix: fall back when struct field names collide case-insensitively Under spark.sql.caseSensitive=false Spark's Parquet reader matches each requested field to the file fields sharing its lowercased name and raises when more than one answers. Comet resolves nested names only while casting a column whose file type differs from the requested type. When the two types are equal, DataFusion's default expression adapter leaves the column bare, and the Parquet opener skips the adapter entirely when the schemas match and no predicate is pushed. A requested struct read from a file with the same shape therefore came back positionally, where Spark raises. Spark's analyzer rejects such a schema under the case-insensitive resolver, but a DataFrame analyzed under the case-sensitive one and planned under the case-insensitive one still reaches the scan. The names are in the plan, so CometNativeScan.isSupported now declines a required schema whose sibling field names collide under toLowerCase(Locale.ROOT), and Spark's reader reports the ambiguity. Closes #6136. --- .../user-guide/latest/compatibility/scans.md | 3 ++ .../org/apache/comet/DataTypeSupport.scala | 25 ++++++++++ .../serde/operator/CometNativeScan.scala | 15 +++++- .../comet/exec/CometNativeReaderSuite.scala | 48 ++++++++++++++++++- .../comet/rules/CometScanRuleSuite.scala | 23 ++++++++- 5 files changed, 111 insertions(+), 3 deletions(-) diff --git a/docs/source/user-guide/latest/compatibility/scans.md b/docs/source/user-guide/latest/compatibility/scans.md index 57363125331..d4d45b100a9 100644 --- a/docs/source/user-guide/latest/compatibility/scans.md +++ b/docs/source/user-guide/latest/compatibility/scans.md @@ -50,6 +50,9 @@ The following features are not supported and cause Comet to fall back to Spark: Comet Parquet scan regardless. - A read schema that repeats a Parquet field id, at the top level or within a struct, when `spark.sql.parquet.fieldId.read.enabled=true`. +- A read schema with sibling struct fields whose names collide case-insensitively, when + `spark.sql.caseSensitive=false`. Spark's analyzer normally rejects such a schema before the scan + is planned. The following limitation may produce incorrect results without falling back to Spark: diff --git a/spark/src/main/scala/org/apache/comet/DataTypeSupport.scala b/spark/src/main/scala/org/apache/comet/DataTypeSupport.scala index dca806a3347..4d482c156ea 100644 --- a/spark/src/main/scala/org/apache/comet/DataTypeSupport.scala +++ b/spark/src/main/scala/org/apache/comet/DataTypeSupport.scala @@ -19,6 +19,9 @@ package org.apache.comet +import java.util.Locale + +import scala.collection.mutable import scala.collection.mutable.ListBuffer import org.apache.spark.sql.execution.datasources.parquet.ParquetUtils @@ -90,6 +93,28 @@ object DataTypeSupport { def hasDuplicateFieldNames(fields: Array[StructField]): Boolean = fields.map(_.name).distinct.length != fields.length + /** + * True when two sibling struct fields anywhere in `dt` fold to one name under + * `toLowerCase(Locale.ROOT)`, the fold Spark's Parquet reader matches requested fields with + * when `spark.sql.caseSensitive` is off. + * + * Like [[hasDuplicateFieldIds]] this inspects only the requested schema and ignores field ids, + * so it can decline a read Spark would accept, costing native execution but not correctness. + */ + def hasCaseInsensitiveDuplicateFieldNames(dt: DataType): Boolean = dt match { + case StructType(fields) => + val folded = mutable.HashSet.empty[String] + fields.exists { f => + !folded.add(f.name.toLowerCase(Locale.ROOT)) || + hasCaseInsensitiveDuplicateFieldNames(f.dataType) + } + case ArrayType(elementType, _) => hasCaseInsensitiveDuplicateFieldNames(elementType) + case MapType(keyType, valueType, _) => + hasCaseInsensitiveDuplicateFieldNames(keyType) || + hasCaseInsensitiveDuplicateFieldNames(valueType) + case _ => false + } + /** * True when two of `fields` declare the same Parquet field id. * diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometNativeScan.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometNativeScan.scala index 080c86085b3..835e9db2f04 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometNativeScan.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometNativeScan.scala @@ -31,7 +31,7 @@ import org.apache.spark.sql.execution.datasources.parquet.ParquetUtils import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types.{StructField, StructType} -import org.apache.comet.{CometConf, ConfigEntry} +import org.apache.comet.{CometConf, ConfigEntry, DataTypeSupport} import org.apache.comet.CometConf.COMET_EXEC_ENABLED import org.apache.comet.CometSparkSessionExtensions.{hasFallbackReason, isSpark35Plus, isSpark41Plus, withFallbackReason} import org.apache.comet.objectstore.NativeConfig @@ -143,6 +143,19 @@ object CometNativeScan extends CometOperatorSerde[CometScanExec] with CometTypeS withFallbackReason(scanExec, unsupportedDefaultReason) } + // Spark's reader raises when a requested field folds to more than one file field. The native + // scan resolves nested names only while casting, so a column whose file type equals the + // requested one, and whose siblings are therefore the requested ones, is read positionally. + // A DataFrame analyzed case-sensitively gets such a schema past the analyzer (#6136). Not in + // CometScanTypeChecker, which also gates Iceberg scans that resolve fields by id. + if (!SQLConf.get.caseSensitiveAnalysis && + DataTypeSupport.hasCaseInsensitiveDuplicateFieldNames(scanExec.requiredSchema)) { + withFallbackReason( + scanExec, + "Native Parquet scan does not support a read schema whose field names collide " + + s"case-insensitively when ${SQLConf.CASE_SENSITIVE.key}=false") + } + if (scanExec.requiredSchema.exists(field => isVariantType(field.dataType))) { // Spark's strict legacy reader owns malformed-layout errors (SPARK-47546). // TODO: Remove this guard once the native reader implements Spark's strict Variant layout 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 ec94a708fcb..63a7c3f0420 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, CometNativeException} +import org.apache.comet.{CometConf, CometNativeException, ExtendedExplainInfo} import org.apache.comet.CometSparkSessionExtensions.isSpark41Plus class CometNativeReaderSuite extends CometTestBase with AdaptiveSparkPlanHelper { @@ -1368,6 +1368,52 @@ class CometNativeReaderSuite extends CometTestBase with AdaptiveSparkPlanHelper StructType(Seq(withFieldId("x", 1), withFieldId("y", 1)))) } + test("native scan declines a nested struct whose field names collide case-insensitively") { + // #6136: `s` in the file has the requested type, so the native scan read it positionally + // instead of raising like Spark. `writeDirect` reproduces the issue's metadata-free file. + withTempPath { dir => + writeDirect( + new Path(dir.getCanonicalPath, "case-colliding-names.parquet").toString, + """message spark_schema { + | optional group s { + | optional int64 x; + | optional int64 X; + | } + |} + """.stripMargin, + { rc: RecordConsumer => + rc.startMessage() + rc.startField("s", 0) + rc.startGroup() + rc.startField("x", 0) + rc.addLong(10L) + rc.endField("x", 0) + rc.startField("X", 1) + rc.addLong(20L) + rc.endField("X", 1) + rc.endGroup() + rc.endField("s", 0) + rc.endMessage() + }) + + // Spark's analyzer rejects this schema case-insensitively, so analyze case-sensitively and + // plan and run case-insensitively. + withSQLConf(SQLConf.CASE_SENSITIVE.key -> "true") { + val df = spark.read.schema("s struct").parquet(dir.getCanonicalPath) + withSQLConf(SQLConf.CASE_SENSITIVE.key -> "false") { + val plan = df.queryExecution.executedPlan + val reasons = new ExtendedExplainInfo().getFallbackReasons(plan) + assert(reasons.exists(_.contains("collide case-insensitively")), s"$reasons\n$plan") + val messages = causeMessages(intercept[Exception](df.collect())) + assert( + messages.contains( + """Found duplicate field(s) "x": [x, X] in case-insensitive mode"""), + messages) + } + } + } + } + /** Write a Parquet file using a raw RecordConsumer for full schema control. */ private def writeDirect( path: String, diff --git a/spark/src/test/scala/org/apache/comet/rules/CometScanRuleSuite.scala b/spark/src/test/scala/org/apache/comet/rules/CometScanRuleSuite.scala index 6bb4d6254a2..52c6159d11e 100644 --- a/spark/src/test/scala/org/apache/comet/rules/CometScanRuleSuite.scala +++ b/spark/src/test/scala/org/apache/comet/rules/CometScanRuleSuite.scala @@ -29,7 +29,7 @@ import org.apache.spark.sql.execution.adaptive.QueryStageExec import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types._ -import org.apache.comet.CometConf +import org.apache.comet.{CometConf, DataTypeSupport} import org.apache.comet.testing.{DataGenOptions, FuzzDataGenerator} /** @@ -239,4 +239,25 @@ class CometScanRuleSuite extends CometTestBase { } } + test("DataTypeSupport finds sibling field names that collide case-insensitively") { + // See https://github.com/apache/datafusion-comet/issues/6136. + def fields(names: String*): StructType = + StructType(names.map(StructField(_, LongType))) + + val colliding = fields("x", "X") + val detected = Seq[(String, DataType)]( + "root" -> colliding, + "nested struct" -> StructType(Seq(StructField("s", colliding))), + "array element" -> ArrayType(colliding), + "map key" -> MapType(colliding, LongType), + "map value" -> MapType(StringType, colliding)) + for ((label, dt) <- detected) { + assert(DataTypeSupport.hasCaseInsensitiveDuplicateFieldNames(dt), s"$label: missed") + } + + // Only siblings are compared, so neither a parent and child nor cousins collide. + val distinct = StructType(Seq(StructField("x", fields("X")), StructField("t", fields("x")))) + assert(!DataTypeSupport.hasCaseInsensitiveDuplicateFieldNames(distinct)) + } + }