Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions docs/source/user-guide/latest/compatibility/scans.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand Down
25 changes: 25 additions & 0 deletions spark/src/main/scala/org/apache/comet/DataTypeSupport.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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<x: bigint, X: bigint>").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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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}

/**
Expand Down Expand Up @@ -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))
}

}
Loading