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)