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 5e07d4d0353e test(variant): pin top-level sibling widening (#20126)
5e07d4d0353e is described below

commit 5e07d4d0353edba703397df4752dfad1c4e62044
Author: voonhous <[email protected]>
AuthorDate: Tue Sep 29 15:14:41 2026 +0800

    test(variant): pin top-level sibling widening (#20126)
    
    The nested leg from #19783 needs the addMissingFields struct walk: a
    type change on one member of s makes the whole struct an implicit type
    change, and every member of s is reconciled. A top-level variant beside
    a widened top-level column takes another path: buildImplicitSchemaChange
    Info reconciles per column, v requested as the PushVariantIntoScan
    projection struct is declared equal to the file's VariantType and never
    walked, and only n gets a type-change Cast. Expected to pass; nothing
    pinned it.
    
    TestVariantShreddingMixedLayouts gets the top-level twin over an int
    base file, a still-int SQL update and a bigint DataFrame upsert that
    widens n and opens a second file group. On COW the upsert carries new
    keys only, so the int file group stays; on MOR it also carries id 2, so
    the int base file merges with an int log block and a bigint log block.
    Both record types on COW; native parquet and avro data blocks on MOR;
    pushVariantIntoScan on and off with the plan pinned per arm.
    
    Green on master as-is; no production change.
    
    Part of #18285.
---
 .../schema/TestVariantShreddingMixedLayouts.scala  | 145 +++++++++++++++++++++
 1 file changed, 145 insertions(+)

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 210eeca3476f..984c97b38989 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
@@ -1393,6 +1393,151 @@ class TestVariantShreddingMixedLayouts extends 
HoodieSparkSqlTestBase with Varia
     }
   }
 
+  test("Implicit widening of a top-level sibling keeps the variant 
projection") {
+    assume(HoodieSparkUtils.gteqSpark4_1, SPARK_4_1_GATE)
+
+    // The nested leg below needs the struct walk in 
SparkSchemaTransformUtils.addMissingFields: a
+    // type change on one member of `s` makes the whole struct an implicit 
type change, so every
+    // member is reconciled. A TOP-LEVEL variant beside a widened top-level 
column takes a
+    // different path: buildImplicitSchemaChangeInfo reconciles per top-level 
column, `v` requested
+    // as the PushVariantIntoScan projection struct is declared equal to the 
file's VariantType and
+    // never walked, and only `n` gets a type-change Cast. Expected to pass, 
but nothing pinned it.
+    // This leg pins it on COW (base files only) and on MOR, where the widened 
rows also arrive
+    // through log blocks and the file-group reader merges rows from an int 
base file, an int log
+    // block and a bigint log block.
+    def rowsSql(lo: Int, hi: Int, nType: String, kPrefix: String, ts: Long): 
String =
+      s"""select cast(id as int) as id,
+         | parse_json(concat('{"k":"$kPrefix', id, '"}')) as v,
+         | cast(id as $nType) as n,
+         | ${ts}L as ts from range($lo, $hi, 1, 1)""".stripMargin
+
+    def runLeg(label: String, tableType: String, recordType: HoodieRecordType, 
tableProps: Seq[String],
+               writeOptions: Seq[(String, String)], widenFrom: Int,
+               assertLogs: (String, String) => Unit): Unit = {
+      Seq("true", "false").foreach { pushIntoScan =>
+        withSQLConf("spark.sql.variant.pushVariantIntoScan" -> pushIntoScan) {
+          withVariantTable(s"$label pushVariantIntoScan=$pushIntoScan", 
tableType, props = tableProps,
+            extraCols = "n int", recordTypes = Seq(recordType)) { (tableName, 
tablePath, leg) =>
+            // Commit 1: the int base file.
+            withWriteLayout(Forced("k string")) {
+              spark.sql(s"insert into $tableName ${rowsSql(0, 3, "int", "x", 
1000L)}")
+            }
+            assertVariantLayout(tablePath, shredded = true, leg)
+
+            // Commit 2, still int (a SQL update coerces to the table schema): 
the int log block on
+            // MOR, a rewrite of the base file on COW.
+            withWriteLayout(Forced("k string")) {
+              spark.sql(s"update $tableName set " +
+                """v = parse_json(concat('{"k":"y', id, '"}')), n = cast(id as 
int), ts = 1001 """ +
+                "where id >= 1")
+            }
+
+            // Commit 3, the widening, through the DataFrame API like the 
nested leg: a SQL write
+            // coerces to the table schema, a DataFrame whose n is bigint 
evolves it in place. On
+            // COW the upsert carries only NEW keys (widenFrom = 3): an upsert 
of an existing key
+            // would rewrite the int file group as bigint and lose the int 
arm. On MOR it also
+            // carries id 2 (widenFrom = 2), which lands as a bigint log block 
on the int file
+            // group. The new keys open a second file group with a bigint base 
file because the
+            // small-file limit is 0. Every knob is an explicit option: 
df.write collects
+            // spark.hoodie.* only and drops bare hoodie.* session confs, so 
neither
+            // withWriteLayout, withRecordType's merger and block-format confs 
nor the
+            // tblproperties reach it.
+            var writer = spark.sql(rowsSql(widenFrom, 6, "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",
+                if (tableType == "cow") "COPY_ON_WRITE" else "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)
+
+            // Every parquet file, base versions and parquet logs alike, is 
shredded under the
+            // forced layout. COW keeps older base-file versions of the first 
group, so the file
+            // groups are counted by file id, not by file.
+            assertVariantLayout(tablePath, shredded = true, leg)
+            val fileGroups = listDataParquetFiles(tablePath).map(new 
HadoopPath(_).getName)
+              .filter(FSUtils.isBaseFile(_)).map(FSUtils.getFileId(_)).distinct
+            assert(fileGroups.size == 2, s"[$leg] expected the int and the 
bigint file group, got $fileGroups")
+            assertLogs(tablePath, leg)
+
+            // Same reason as the nested leg: the SQL catalog keeps reporting 
n as int, so the
+            // reads go through a path-based view that sees the widened table 
schema.
+            val widenedView = s"${tableName}_widened"
+            
spark.read.format("hudi").load(tablePath).createOrReplaceTempView(widenedView)
+            val widenedN = spark.table(widenedView).schema("n").dataType
+            assert(widenedN == LongType,
+              s"[$leg] the DataFrame write should have widened n to bigint, 
got $widenedN")
+
+            // id 0 from the int base file, ids 1 to widenFrom-1 from the 
update, the rest from the
+            // bigint write.
+            def expectedK(id: Int): String = if (id == 0) "x0" else if (id < 
widenFrom) s"y$id" else s"z$id"
+            checkAnswer(s"select id, variant_get(v, '$$.k', 'string'), n from 
$widenedView order by id")(
+              (0 until 6).map(id => Seq(id, expectedK(id), id.toLong)): _*)
+            // One filter per slot: the int base, the update, id 2 (the update 
on COW, the bigint
+            // log block on MOR) and the bigint base file.
+            checkAnswer(s"select id from $widenedView where variant_get(v, 
'$$.k', 'string') = 'x0'")(Seq(0))
+            checkAnswer(s"select id from $widenedView where variant_get(v, 
'$$.k', 'string') = 'y1'")(Seq(1))
+            checkAnswer(s"select id from $widenedView " +
+              s"where variant_get(v, '$$.k', 'string') = 
'${expectedK(2)}'")(Seq(2))
+            checkAnswer(s"select id from $widenedView where variant_get(v, 
'$$.k', 'string') = 'z4'")(Seq(4))
+            checkAnswer(s"select id from $widenedView where n >= 3 order by 
id")(Seq(3), Seq(4), Seq(5))
+            checkAnswer(s"select id, cast(v as string) from $widenedView order 
by id")(
+              (0 until 6).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(v, '$$.k', 'string'), n from 
$widenedView") == pushed,
+              s"[$leg] PushVariantIntoScan $verdict rewritten v into a 
projection struct")
+            spark.catalog.dropTempView(widenedView)
+          }
+        }
+      }
+    }
+
+    // COW, SPARK records: base files only, the int file group and the bigint 
one.
+    runLeg("cow top-level widening", "cow", HoodieRecordType.SPARK, Seq.empty,
+      Seq("hoodie.write.record.merge.impl.classes" -> 
classOf[DefaultSparkRecordMerger].getName),
+      widenFrom = 3, assertLogs = (_, _) => ())
+    // COW, AVRO records: the same two file groups written and merged through 
the avro record path.
+    runLeg("cow top-level widening, avro records", "cow", 
HoodieRecordType.AVRO, Seq.empty,
+      Seq("hoodie.write.record.merge.impl.classes" -> 
classOf[HoodieAvroRecordMerger].getName),
+      widenFrom = 3, assertLogs = (_, _) => ())
+    // MOR, SPARK records on the current table version: native parquet log 
files. The two data
+    // blocks are the update (int, coerced to the table schema) and the upsert 
(bigint).
+    runLeg("mor top-level widening, parquet log blocks", "mor", 
HoodieRecordType.SPARK,
+      Seq("hoodie.compact.inline = 'false'"),
+      Seq("hoodie.write.record.merge.impl.classes" -> 
classOf[DefaultSparkRecordMerger].getName),
+      widenFrom = 2, assertLogs = { (tablePath, leg) =>
+        val blockTypes = listLogBlockTypes(tablePath)
+        assert(blockTypes.count(_ == HoodieLogBlockType.PARQUET_DATA_BLOCK) == 
2,
+          s"[$leg] expected the update's int block and the upsert's bigint 
block as parquet data blocks, " +
+            s"found: $blockTypes")
+      })
+    // MOR, avro data blocks: mirrors the nested MOR avro leg (table version 9 
is the only way to
+    // get avro blocks). The DataFrame write pins hoodie.write.table.version 
too, because
+    // 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 top-level widening, avro data blocks", "mor", 
HoodieRecordType.AVRO,
+      Seq("hoodie.compact.inline = 'false'",
+        "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"),
+      widenFrom = 2, assertLogs = { (tablePath, leg) =>
+        val blockTypes = listLogBlockTypes(tablePath)
+        assert(blockTypes.count(_ == HoodieLogBlockType.AVRO_DATA_BLOCK) == 2,
+          s"[$leg] expected two avro data blocks, found: $blockTypes")
+        assert(!blockTypes.contains(HoodieLogBlockType.PARQUET_DATA_BLOCK),
+          s"[$leg] this leg must not write parquet data blocks, found: 
$blockTypes")
+      })
+  }
+
   test("Implicit widening of a sibling keeps the nested variant projection") {
     assume(HoodieSparkUtils.gteqSpark4_1, SPARK_4_1_GATE)
 

Reply via email to