This is an automated email from the ASF dual-hosted git repository.
voonhous pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git
The following commit(s) were added to refs/heads/master by this push:
new 5d002c001bf5 fix(spark): reject pushVariantIntoScan reads on Spark 4.0
(#20033)
5d002c001bf5 is described below
commit 5d002c001bf5290e70fafefe799b81683717d9a7
Author: voonhous <[email protected]>
AuthorDate: Thu Sep 24 11:52:54 2026 +0800
fix(spark): reject pushVariantIntoScan reads on Spark 4.0 (#20033)
Spark 4.0 ships spark.sql.variant.pushVariantIntoScan off (on from
4.1, SPARK-54454). Once set, 4.0 rewrites a VARIANT column into the
same projection struct 4.1 produces, but its readers cannot evaluate
it and Hudi does not align log records to it. Nothing on 4.0
recognised the struct, so the read fell through into schema-change
handling.
Spark 4.0 with the conf on is unsupported. Every 4.x adapter now
recognises the projection struct (isVariantProjectionStruct moves
from Spark4_1Adapter to BaseSpark4Adapter), and a new
SparkAdapter.validateVariantProjectionReadable hook, called first in
buildReaderWithPartitionValues, lets Spark4_0Adapter fail the scan on
the driver with the column, the conf and the fix in the message.
Closes #20032
---
.../org/apache/spark/sql/hudi/SparkAdapter.scala | 20 +++++++++---
.../HoodieFileGroupReaderBasedFileFormat.scala | 3 ++
.../TestBaseSpark4AdapterVariantMethods.scala | 37 ++++++++++++++++++++++
.../sql/hudi/dml/schema/TestVariantDataType.scala | 28 ++++++++++++++++
.../spark/sql/adapter/BaseSpark4Adapter.scala | 4 +++
.../apache/spark/sql/adapter/Spark4_0Adapter.scala | 17 +++++++++-
.../apache/spark/sql/adapter/Spark4_1Adapter.scala | 4 ---
.../apache/spark/sql/adapter/Spark4_2Adapter.scala | 4 ---
8 files changed, 104 insertions(+), 13 deletions(-)
diff --git
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/spark/sql/hudi/SparkAdapter.scala
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/spark/sql/hudi/SparkAdapter.scala
index be3966634165..2b9725b52049 100644
---
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/spark/sql/hudi/SparkAdapter.scala
+++
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/spark/sql/hudi/SparkAdapter.scala
@@ -506,10 +506,12 @@ trait SparkAdapter extends Serializable {
def isVariantShreddingStruct(structType: StructType): Boolean
/**
- * Checks if a StructType is the result of Spark 4.1's PushVariantIntoScan
rewriting — i.e.,
- * every child field carries `VariantMetadata` describing a pushed-down
variant extraction.
- *
- * Returns false on Spark versions earlier than 4.1 (the rewriting only
happens there).
+ * Checks if a StructType is the result of PushVariantIntoScan rewriting,
i.e. every child field
+ * carries `VariantMetadata` describing a pushed-down variant extraction.
The rewrite exists on
+ * every Spark 4.x and runs whenever spark.sql.variant.pushVariantIntoScan
is on (off by default
+ * on 4.0, on from 4.1), so every 4.x adapter recognises the shape; whether
the version can read
+ * it is [[validateVariantProjectionReadable]]'s question. Returns false on
Spark 3.x, which has
+ * no VariantType.
*/
def isVariantProjectionStruct(structType: StructType): Boolean = false
@@ -525,6 +527,16 @@ trait SparkAdapter extends Serializable {
case _ => false
}
+ /**
+ * Fails the read when `requiredSchema` carries a variant projection struct
this Spark version
+ * cannot evaluate. Spark 4.0 rewrites a variant column exactly as 4.1 does
once the conf is on,
+ * but its readers do not evaluate the projection and Hudi does not align
log records to it, so
+ * a read there would fall through to the schema-change path instead of
failing. The Spark 4.0
+ * adapter overrides this to throw; every other version reads the shape
(4.1+) or never sees it
+ * (3.x). See https://github.com/apache/hudi/issues/20032.
+ */
+ def validateVariantProjectionReadable(requiredSchema: StructType): Unit = ()
+
/**
* If `sparkRequiredSchema` contains any Spark 4.1 variant projection struct
(i.e., the
* same-named field in `sparkDataSchema` is `VariantType`), returns a row
transformer that
diff --git
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/HoodieFileGroupReaderBasedFileFormat.scala
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/HoodieFileGroupReaderBasedFileFormat.scala
index b6f5224105e2..823fda9001f6 100644
---
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/HoodieFileGroupReaderBasedFileFormat.scala
+++
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/HoodieFileGroupReaderBasedFileFormat.scala
@@ -291,6 +291,9 @@ class HoodieFileGroupReaderBasedFileFormat(tablePath:
String,
filters: Seq[Filter],
options: Map[String, String],
hadoopConf: Configuration):
PartitionedFile => Iterator[InternalRow] = {
+ // Driver side, once per scan: Spark 4.0 cannot read a PushVariantIntoScan
projection struct and
+ // has to fail here rather than in the schema-change path (#20032).
+ sparkAdapter.validateVariantProjectionReadable(requiredSchema)
val outputSchema = StructType(requiredSchema.fields ++
partitionSchema.fields)
val isCount = requiredSchema.isEmpty && !isMOR && !isIncremental
// Spark planner only adds the user-provided predicates (from `WHERE`
clause or `.filter()`)
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/adapter/TestBaseSpark4AdapterVariantMethods.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/adapter/TestBaseSpark4AdapterVariantMethods.scala
index b9a7dbe07500..e5d6b51aa65c 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/adapter/TestBaseSpark4AdapterVariantMethods.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/adapter/TestBaseSpark4AdapterVariantMethods.scala
@@ -20,6 +20,7 @@ package org.apache.spark.sql.adapter
import org.apache.hudi.HoodieSparkUtils
import org.apache.hudi.SparkAdapterSupport
import org.apache.hudi.common.schema.{HoodieSchema, HoodieSchemaType}
+import org.apache.hudi.exception.HoodieNotSupportedException
import org.apache.parquet.schema.PrimitiveType
import org.apache.parquet.schema.Type.Repetition
@@ -290,6 +291,42 @@ class TestBaseSpark4AdapterVariantMethods extends
SparkAdapterSupport {
method.invoke(module, args: _*)
}
+ @Test
+ def testValidateVariantProjectionReadableRejectsOnlySpark40(): Unit = {
+ assumeTrue(HoodieSparkUtils.isSpark4, "Only applies to Spark 4.x")
+ val projectionStruct = StructType(Seq(
+ StructField("0", StringType, metadata =
variantProjectionMetadata("$.k"))))
+ assertTrue(sparkAdapter.isVariantProjectionStruct(projectionStruct),
+ "every Spark 4.x adapter must recognise the projection struct shape")
+
+ // The projection sits at the root of the relation output or below a
struct member, the two
+ // places PushVariantIntoScan puts it, so the guard has to look in both.
+ val rootSchema = StructType(Seq(
+ StructField("id", IntegerType),
+ StructField("v", projectionStruct)))
+ val nestedSchema = StructType(Seq(
+ StructField("id", IntegerType),
+ StructField("s", StructType(Seq(
+ StructField("inner", projectionStruct),
+ StructField("other", IntegerType))))))
+ val plainSchema = StructType(Seq(
+ StructField("id", IntegerType),
+ StructField("v", sparkAdapter.getVariantDataType.get)))
+
+ // A native variant request is readable everywhere.
+ sparkAdapter.validateVariantProjectionReadable(plainSchema)
+ if (HoodieSparkUtils.isSpark4_0) {
+ Seq(rootSchema, nestedSchema).foreach { schema =>
+ val e = assertThrows(classOf[HoodieNotSupportedException],
+ () => sparkAdapter.validateVariantProjectionReadable(schema))
+
assertTrue(e.getMessage.contains("spark.sql.variant.pushVariantIntoScan"),
e.getMessage)
+ }
+ } else {
+ sparkAdapter.validateVariantProjectionReadable(rootSchema)
+ sparkAdapter.validateVariantProjectionReadable(nestedSchema)
+ }
+ }
+
@Test
def testBuildVariantProjectorRecursesIntoStructPaths(): Unit = {
assumeTrue(HoodieSparkUtils.gteqSpark4_1, "Variant projection structs only
exist on Spark 4.1+")
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantDataType.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantDataType.scala
index f89b51cfce65..ec8637383cbd 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantDataType.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantDataType.scala
@@ -198,6 +198,34 @@ class TestVariantDataType extends HoodieSparkSqlTestBase
with VariantShreddingTe
})
}
+ test("Test Spark 4.0 rejects reads rewritten by pushVariantIntoScan") {
+ // The conf is off by default on Spark 4.0 and on from 4.1 (SPARK-54454).
Spark 4.0 still runs
+ // the rewrite once it is set but cannot read the projection struct it
produces, so Hudi fails
+ // the scan up front instead of falling through to schema-change handling
(#20032).
+ assume(HoodieSparkUtils.isSpark4_0, "Guards the Spark 4.0 read path only")
+
+ Seq("cow", "mor").foreach { tableType =>
+ withVariantTable(s"pushVariantIntoScan $tableType", tableType) {
(tableName, _, _) =>
+ spark.sql(s"""insert into $tableName values (1, parse_json('{"key":
"v1"}'), 1000)""")
+ // A second write so the MOR leg reads a log file too.
+ spark.sql(s"""update $tableName set v = parse_json('{"key": "v2"}')
where id = 1""")
+
+ withSQLConf("spark.sql.variant.pushVariantIntoScan" -> "true") {
+ // An extraction and a whole-variant read are both rewritten into a
projection struct.
+ Seq(s"select id, variant_get(v, '$$.key', 'string') from $tableName",
+ s"select id, cast(v as string) from $tableName").foreach { sql =>
+ checkNestedExceptionContains(() => spark.sql(sql).collect())(
+ "spark.sql.variant.pushVariantIntoScan")
+ }
+ // A query that does not touch the variant column is not rewritten
and still reads.
+ checkAnswer(s"select id, ts from $tableName")(Seq(1, 1000))
+ }
+ // Back on the default the same reads work.
+ checkAnswer(s"select id, cast(v as string) from $tableName")(Seq(1,
"{\"key\":\"v2\"}"))
+ }
+ }
+ }
+
test("Test Query Log Only MOR Table With VARIANT column triggers
compaction") {
// Gated on Spark >= 4.1. Compaction writes the base file via the AVRO
shredding writer, which
// lays the variant group out as [metadata, value, typed_value]. Hudi's
Spark 4.0 reader does not
diff --git
a/hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/spark/sql/adapter/BaseSpark4Adapter.scala
b/hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/spark/sql/adapter/BaseSpark4Adapter.scala
index 75dc1e3b8e9f..2f049cbedce6 100644
---
a/hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/spark/sql/adapter/BaseSpark4Adapter.scala
+++
b/hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/spark/sql/adapter/BaseSpark4Adapter.scala
@@ -240,6 +240,10 @@ abstract class BaseSpark4Adapter extends SparkAdapter with
Logging {
dataType.isInstanceOf[VariantType]
}
+ override def isVariantProjectionStruct(structType: StructType): Boolean = {
+ VariantMetadata.isVariantStruct(structType)
+ }
+
override def createVariantValueWriter(
dataType: DataType,
writeValue: Consumer[Array[Byte]],
diff --git
a/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/spark/sql/adapter/Spark4_0Adapter.scala
b/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/spark/sql/adapter/Spark4_0Adapter.scala
index f622b266db9f..14031abf21d6 100644
---
a/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/spark/sql/adapter/Spark4_0Adapter.scala
+++
b/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/spark/sql/adapter/Spark4_0Adapter.scala
@@ -17,12 +17,13 @@
package org.apache.spark.sql.adapter
-import org.apache.hudi.{HoodiePartitionCDCFileGroupMapping,
HoodiePartitionFileSliceMapping, Spark40HoodiePartitionCDCFileGroupMapping,
Spark40HoodiePartitionFileSliceMapping}
+import org.apache.hudi.{HoodiePartitionCDCFileGroupMapping,
HoodiePartitionFileSliceMapping, HoodieSparkUtils,
Spark40HoodiePartitionCDCFileGroupMapping,
Spark40HoodiePartitionFileSliceMapping}
import org.apache.hudi.client.model.{HoodieInternalRow,
Spark40HoodieInternalRow}
import org.apache.hudi.common.model.FileSlice
import org.apache.hudi.common.schema.HoodieSchema
import org.apache.hudi.common.table.cdc.HoodieCDCFileSplit
import org.apache.hudi.common.util.{Option => HOption}
+import org.apache.hudi.exception.HoodieNotSupportedException
import org.apache.hadoop.conf.Configuration
import org.apache.parquet.schema.MessageType
@@ -118,6 +119,20 @@ class Spark4_0Adapter extends BaseSpark4Adapter {
Some(new Spark40LegacyHoodieParquetFileFormat(appendPartitionValues))
}
+ // Spark 4.0 rewrites a variant column into a projection struct once
+ // spark.sql.variant.pushVariantIntoScan is on (off by default there), but
its readers do not
+ // evaluate the projection and Hudi does not align log records to it, so
fail here instead of
+ // falling through to the schema-change path (#20032).
+ override def validateVariantProjectionReadable(requiredSchema: StructType):
Unit = {
+ requiredSchema.fields.find(f =>
containsVariantProjection(f.dataType)).foreach { field =>
+ throw new HoodieNotSupportedException(
+ s"Column '${field.name}' was rewritten by
spark.sql.variant.pushVariantIntoScan into a variant " +
+ s"projection struct, which Hudi does not support on Spark
${HoodieSparkUtils.getSparkVersion}. " +
+ "Set spark.sql.variant.pushVariantIntoScan=false (the Spark 4.0
default) or read the table " +
+ "with Spark 4.1+.")
+ }
+ }
+
override def createInternalRow(metaFields: Array[UTF8String],
sourceRow: InternalRow,
sourceContainsMetaFields: Boolean):
HoodieInternalRow = {
diff --git
a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/spark/sql/adapter/Spark4_1Adapter.scala
b/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/spark/sql/adapter/Spark4_1Adapter.scala
index ddaf0f581476..9bedeb03b078 100644
---
a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/spark/sql/adapter/Spark4_1Adapter.scala
+++
b/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/spark/sql/adapter/Spark4_1Adapter.scala
@@ -230,10 +230,6 @@ class Spark4_1Adapter extends BaseSpark4Adapter {
RebaseDateTime.RebaseSpec(LegacyBehaviorPolicy.withName(policy))
}
- override def isVariantProjectionStruct(structType: StructType): Boolean = {
- VariantMetadata.isVariantStruct(structType)
- }
-
// Spark 4.1 reconstructs shredded variants on read (SPARK-54410), so opt in
to the
// shared rewrite; Spark 4.0 stays on the default None.
override def buildFullVariantReadSchema(schema: StructType):
Option[StructType] =
diff --git
a/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/spark/sql/adapter/Spark4_2Adapter.scala
b/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/spark/sql/adapter/Spark4_2Adapter.scala
index b6425b03e5f6..bac1e8075e0a 100644
---
a/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/spark/sql/adapter/Spark4_2Adapter.scala
+++
b/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/spark/sql/adapter/Spark4_2Adapter.scala
@@ -230,10 +230,6 @@ class Spark4_2Adapter extends BaseSpark4Adapter {
RebaseDateTime.RebaseSpec(LegacyBehaviorPolicy.withName(policy))
}
- override def isVariantProjectionStruct(structType: StructType): Boolean = {
- VariantMetadata.isVariantStruct(structType)
- }
-
// Spark 4.2 reconstructs shredded variants on read (SPARK-54410), so opt in
to the
// shared rewrite; Spark 4.0 stays on the default None.
override def buildFullVariantReadSchema(schema: StructType):
Option[StructType] =