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] =

Reply via email to