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 0d40ce4b22df test(variant): sibling widening on MOR log blocks (#20125)
0d40ce4b22df is described below

commit 0d40ce4b22dffb5198a95abe35049eeef56849b5
Author: voonhous <[email protected]>
AuthorDate: Tue Sep 29 13:48:17 2026 +0800

    test(variant): sibling widening on MOR log blocks (#20125)
    
    The COW leg from #19783 pins that the PushVariantIntoScan projection
    struct at s.inner survives SparkSchemaTransformUtils.addMissingFields
    when a sibling member of s changed type implicitly (file int, table
    long). On MOR the file-group reader merges rows from the base file and
    from log blocks, and each source reconciles its own file schema against
    the widened requested schema before the merge, so all of them have to
    agree on the row shape. Nothing pinned that arm.
    
    TestVariantShreddingMixedLayouts gets the MOR twin of that leg: an int
    base file, an int log block written before the widening, then one
    DataFrame upsert that widens s.n to bigint and lands as a bigint log
    block on the first file group plus a bigint base file in a second one.
    Every slot is read back through variant_get, a filter and a whole-value
    cast, with pushVariantIntoScan on and off and the plan pinned per arm.
    Two legs: native parquet log blocks (SPARK record type) and avro data
    blocks (AVRO record type on table version 9, the DataFrame write pinning
    hoodie.write.table.version so auto-upgrade does not flip the format).
    
    Green on master as-is; no production change.
    
    Part of #18285.
---
 .../schema/TestVariantShreddingMixedLayouts.scala  | 138 ++++++++++++++++++++-
 .../dml/schema/VariantShreddingTestSupport.scala   |   9 ++
 2 files changed, 145 insertions(+), 2 deletions(-)

diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantShreddingMixedLayouts.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantShreddingMixedLayouts.scala
index a9921d8a8d0c..210eeca3476f 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantShreddingMixedLayouts.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantShreddingMixedLayouts.scala
@@ -19,10 +19,10 @@
 
 package org.apache.spark.sql.hudi.dml.schema
 
-import org.apache.hudi.{DataSourceReadOptions, DataSourceWriteOptions, 
HoodieSchemaConversionUtils, HoodieSparkUtils, HoodieTableSchema, 
SparkAdapterSupport}
+import org.apache.hudi.{DataSourceReadOptions, DataSourceWriteOptions, 
DefaultSparkRecordMerger, HoodieSchemaConversionUtils, HoodieSparkUtils, 
HoodieTableSchema, SparkAdapterSupport}
 import org.apache.hudi.common.config.HoodieMetadataConfig
 import org.apache.hudi.common.fs.FSUtils
-import org.apache.hudi.common.model.HoodieFileFormat
+import org.apache.hudi.common.model.{HoodieAvroRecordMerger, HoodieFileFormat}
 import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType
 import org.apache.hudi.common.schema.{HoodieSchema, HoodieSchemaField, 
HoodieSchemaType}
 import org.apache.hudi.common.table.log.block.HoodieLogBlock.HoodieLogBlockType
@@ -1452,6 +1452,140 @@ class TestVariantShreddingMixedLayouts extends 
HoodieSparkSqlTestBase with Varia
     }
   }
 
+  test("Implicit widening of a sibling keeps the nested variant projection 
through MOR log blocks") {
+    assume(HoodieSparkUtils.gteqSpark4_1, SPARK_4_1_GATE)
+
+    // The MOR twin of the COW leg above, which pins the base-file arm of the 
addMissingFields fix.
+    // Here the file-group reader merges rows from base files and log blocks, 
and each source
+    // reconciles its own file schema against the widened requested schema: 
base parquet and
+    // parquet log blocks through the parquet reader's implicit type-change 
path, avro data blocks
+    // through Avro schema resolution followed by the avro-path variant 
projector. The merge needs
+    // one row shape out of all of them - with pushVariantIntoScan on, the 
projection struct at
+    // s.inner next to a long s.n - which is where #19775's crashes lived. So 
the slots below span
+    // an int base file, an int log block, a bigint log block and a bigint 
base file, and every
+    // read has to find the variant payload of each.
+    def rowsSql(lo: Int, hi: Int, nType: String, kPrefix: String, ts: Long): 
String =
+      s"""select cast(id as int) as id,
+         | named_struct('inner', parse_json(concat('{"k":"$kPrefix', id, 
'"}')),
+         |   'n', cast(id as $nType)) as s,
+         | ${ts}L as ts from range($lo, $hi, 1, 1)""".stripMargin
+    // id 0 is served by the int base file, id 1 by the int log block, id 2 by 
the bigint log
+    // block, and ids 3-4 by the bigint base file.
+    def expectedK(id: Int): String = if (id == 0) "x0" else if (id == 1) "y1" 
else s"z$id"
+    def nTypeOf(schema: HoodieSchema): HoodieSchemaType =
+      
schema.getField("s").get.schema().getNonNullType.getField("n").get.schema().getNonNullType.getType
+
+    def runLeg(label: String,
+               recordType: HoodieRecordType,
+               tableProps: Seq[String],
+               writeOptions: Seq[(String, String)],
+               assertLogs: (String, String) => Unit): Unit = {
+      Seq("true", "false").foreach { pushIntoScan =>
+        withSQLConf("spark.sql.variant.pushVariantIntoScan" -> pushIntoScan) {
+          withNestedOnlyVariantTable(s"$label 
pushVariantIntoScan=$pushIntoScan", "mor",
+            props = Seq("hoodie.compact.inline = 'false'") ++ tableProps,
+            recordTypes = Seq(recordType), structMembers = "inner: variant, n: 
int") {
+            (tableName, tablePath, leg) =>
+            withWriteLayout(Forced("k string")) {
+              spark.sql(s"insert into $tableName ${rowsSql(0, 3, "int", "x", 
1000L)}")
+            }
+            val intInstant = latestCompletedInstant(tablePath)
+            assertNestedBaseLayout(tablePath, intInstant, shredded = true, leg)
+
+            // A log block written BEFORE the widening: a SQL update coerces 
to the table schema,
+            // so s.n is still int in it.
+            withWriteLayout(Unshredded) {
+              spark.sql(s"update $tableName set " +
+                """s = named_struct('inner', parse_json(concat('{"k":"y', id, 
'"}')), """ +
+                "'n', cast(id as int)), ts = 1001 where id >= 1")
+            }
+
+            // The widening goes through the DataFrame API, as in the COW leg: 
a SQL insert would
+            // coerce to int. Every knob the write needs is an explicit option 
- the write path
+            // collects spark.hoodie.* only and drops bare hoodie.* session 
confs, so neither
+            // withWriteLayout, nor the merger and block-format confs 
withRecordType sets, nor the
+            // tblproperties reach a df.write. id 2 is an existing key and 
lands in a second, bigint
+            // log block of the first file group; ids 3-4 open a new file 
group with a bigint base
+            // file because the small-file limit is 0.
+            var writer = spark.sql(rowsSql(2, 5, "bigint", "z", 
1002L)).write.format("hudi")
+              .options(layoutConfs(Forced("k string")).toMap)
+              .option("hoodie.table.name", tableName)
+              .option("hoodie.datasource.write.recordkey.field", "id")
+              .option("hoodie.datasource.write.precombine.field", "ts")
+              .option("hoodie.datasource.write.operation", "upsert")
+              .option("hoodie.datasource.write.table.type", "MERGE_ON_READ")
+              .option("hoodie.parquet.small.file.limit", "0")
+              .option("hoodie.compact.inline", "false")
+            writeOptions.foreach { case (key, value) => writer = 
writer.option(key, value) }
+            writer.mode(SaveMode.Append).save(tablePath)
+            assertLogs(tablePath, leg)
+            // Two file groups: the int base file with its two log blocks, and 
the bigint base file
+            // the upsert opened for ids 3-4.
+            val baseFiles = listDataParquetFiles(tablePath).map(new 
HadoopPath(_).getName).filter(FSUtils.isBaseFile(_))
+            assert(baseFiles.size == 2, s"[$leg] expected an int and a bigint 
base file, got $baseFiles")
+            // The log blocks really disagree on s.n: the update's block is 
int, the upsert's bigint.
+            val logNTypes = logBlockSchemas(tablePath).map(nTypeOf).toSet
+            assert(logNTypes == Set(HoodieSchemaType.INT, 
HoodieSchemaType.LONG),
+              s"[$leg] expected one int and one bigint log block for s.n, got 
$logNTypes")
+
+            // Path-based view for the same reason as the COW leg: the SQL 
catalog keeps reporting
+            // s.n as int, while the path-based relation reads the widened 
commit schema.
+            val widenedView = s"${tableName}_widened"
+            
spark.read.format("hudi").load(tablePath).createOrReplaceTempView(widenedView)
+            val widenedN = 
spark.table(widenedView).schema("s").dataType.asInstanceOf[StructType]("n").dataType
+            assert(widenedN == LongType,
+              s"[$leg] the DataFrame write should have widened s.n to bigint, 
got $widenedN")
+
+            checkAnswer(s"select id, variant_get(s.inner, '$$.k', 'string'), 
s.n from $widenedView order by id")(
+              (0 until 5).map(id => Seq(id, expectedK(id), id.toLong)): _*)
+            // One filter per physical slot: int base, int log, bigint log, 
bigint base.
+            checkAnswer(s"select id from $widenedView where 
variant_get(s.inner, '$$.k', 'string') = 'x0'")(Seq(0))
+            checkAnswer(s"select id from $widenedView where 
variant_get(s.inner, '$$.k', 'string') = 'y1'")(Seq(1))
+            checkAnswer(s"select id from $widenedView where 
variant_get(s.inner, '$$.k', 'string') = 'z2'")(Seq(2))
+            checkAnswer(s"select id from $widenedView where 
variant_get(s.inner, '$$.k', 'string') = 'z4'")(Seq(4))
+            // The widened sibling itself, read back across every slot.
+            checkAnswer(s"select id from $widenedView where s.n >= 1 order by 
id")(
+              Seq(1), Seq(2), Seq(3), Seq(4))
+            checkAnswer(s"select id, cast(s.inner as string) from $widenedView 
order by id")(
+              (0 until 5).map(id => Seq(id, s"""{"k":"${expectedK(id)}"}""")): 
_*)
+            // Both arms expect the very same rows; only the plan tells them 
apart.
+            val pushed = pushIntoScan.toBoolean
+            val verdict = if (pushed) "should have" else "must not have"
+            assert(variantProjectionPushedIntoScan(
+              s"select id, variant_get(s.inner, '$$.k', 'string'), s.n from 
$widenedView") == pushed,
+              s"[$leg] PushVariantIntoScan $verdict rewritten s.inner into a 
projection struct")
+            spark.catalog.dropTempView(widenedView)
+          }
+        }
+      }
+    }
+
+    // Parquet log blocks: the SPARK record type on the current table version, 
where the append
+    // handle writes native parquet log files. They are read by the parquet 
reader with the
+    // projected struct threaded into the requested schema, like the base 
files.
+    runLeg("mor widening beside a nested variant, parquet log blocks", 
HoodieRecordType.SPARK, Seq.empty,
+      Seq("hoodie.write.record.merge.impl.classes" -> 
classOf[DefaultSparkRecordMerger].getName),
+      (tablePath, leg) => 
assert(listDataParquetFiles(tablePath).exists(_.endsWith(".log.parquet")),
+        s"[$leg] the update and the upsert should have written native parquet 
log files"))
+
+    // Avro data blocks: mirrors the nested MOR avro leg, since a pre-native 
table version 9 is the
+    // only way to get them. The DataFrame write pins 
hoodie.write.table.version as well:
+    // hoodie.write.auto.upgrade defaults to true, and a default write config 
would upgrade the
+    // table and switch it to native parquet logs.
+    runLeg("mor widening beside a nested variant, avro data blocks", 
HoodieRecordType.AVRO,
+      Seq("hoodie.write.table.version = '9'", 
"hoodie.logfile.data.block.format = 'avro'"),
+      Seq("hoodie.write.record.merge.impl.classes" -> 
classOf[HoodieAvroRecordMerger].getName,
+        "hoodie.logfile.data.block.format" -> "avro",
+        "hoodie.write.table.version" -> "9"),
+      (tablePath, leg) => {
+        val blockTypes = listLogBlockTypes(tablePath)
+        assert(blockTypes.contains(HoodieLogBlockType.AVRO_DATA_BLOCK),
+          s"[$leg] expected an avro data block in the log files, found: 
$blockTypes")
+        assert(!blockTypes.contains(HoodieLogBlockType.PARQUET_DATA_BLOCK),
+          s"[$leg] this leg must not write parquet data blocks, found: 
$blockTypes")
+      })
+  }
+
   test("An array element shredded through a declared write schema reads 
natively and survives compaction") {
     assume(HoodieSparkUtils.gteqSpark4_1, SPARK_4_1_GATE)
 
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/VariantShreddingTestSupport.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/VariantShreddingTestSupport.scala
index 8c0d5a980170..3e141478b488 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/VariantShreddingTestSupport.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/VariantShreddingTestSupport.scala
@@ -24,6 +24,7 @@ import org.apache.hudi.common.fs.FSUtils
 import org.apache.hudi.common.model.HoodieLogFile
 import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType
 import org.apache.hudi.common.model.WriteOperationType
+import org.apache.hudi.common.schema.HoodieSchema
 import org.apache.hudi.common.table.TableSchemaResolver
 import org.apache.hudi.common.table.log.HoodieLogFormat
 import org.apache.hudi.common.table.log.block.{HoodieDataBlock, HoodieLogBlock}
@@ -747,6 +748,14 @@ trait VariantShreddingTestSupport { self: 
HoodieSparkSqlTestBase =>
   protected def listLogBlockTypes(tablePath: String): Seq[HoodieLogBlockType] =
     mapLogBlocks(tablePath)(_.getBlockType)
 
+  /** The writer schema in the header of every data block in the table's log 
files, in log order. */
+  protected def logBlockSchemas(tablePath: String): Seq[HoodieSchema] =
+    mapLogBlocks(tablePath) {
+      case dataBlock: HoodieDataBlock =>
+        
Some(HoodieSchema.parse(dataBlock.getLogBlockHeader.get(HoodieLogBlock.HeaderMetadataType.SCHEMA)))
+      case _ => None
+    }.flatten
+
   /** Whether any data block in the table's log files carries a partial-update 
schema (IS_PARTIAL). */
   protected def hasPartialLogBlock(tablePath: String): Boolean =
     mapLogBlocks(tablePath) {

Reply via email to