hudi-agent commented on code in PR #20036:
URL: https://github.com/apache/hudi/pull/20036#discussion_r4085150824
##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantShreddingMixedLayouts.scala:
##########
@@ -1346,6 +1348,182 @@ class TestVariantShreddingMixedLayouts extends
HoodieSparkSqlTestBase with Varia
})
}
+ test("Variants reached through array elements and map values read correctly
on every path") {
+ assume(HoodieSparkUtils.gteqSpark4_1, SPARK_4_1_GATE)
+
+ // PushVariantIntoScan rewrites struct paths only
(VariantInRelation.rewriteType), so a variant
+ // that is an array element or a map value reaches the reader as native
VariantType on both
+ // arms, and the reader context's overlay and row projector leave it alone
(#19783). "Left
+ // alone" has to mean read correctly, not merely unrewritten, so every
path is checked for the
+ // values themselves: parquet base files, parquet and avro log blocks, and
the merged MOR row,
+ // including a null element, a null map value and a null struct member
that arrive through the
+ // log, beside a top-level variant that IS rewritten on the conf-on arm in
the very same scan.
+ val collectionCols = "arr array<variant>, m map<string, variant>, items
array<struct<inner: variant>>"
+ def rowsSql(lo: Int, hi: Int): String =
+ s"""select cast(id as int) as id, parse_json(concat('{"k":"t', id,
'"}')) as v,
+ | array(parse_json(concat('{"k":"a', id, '"}')),
parse_json(concat('{"k":"b', id, '"}'))) as arr,
+ | map('p', parse_json(concat('{"k":"p', id, '"}')), 'q',
parse_json(concat('{"k":"q', id, '"}'))) as m,
+ | array(named_struct('inner', parse_json(concat('{"k":"i', id,
'"}')))) as items,
+ | 1000L as ts from range($lo, $hi, 1, 1)""".stripMargin
+
+ // ids 0-2 keep their inserted values, id 3 is updated, and id 4's update
nulls the map value
+ // under 'q', the struct member inside the array element and, on the SPARK
legs, the second
+ // array element. The AVRO legs keep that element non-null: the Avro write
path forces
+ // parquet-avro's old list structure, which cannot write a null array
element of ANY type
+ // ("Array contains a null element at 1"), a write-side limit unrelated to
variants.
+ def top(id: Int): String = if (id < 3) s"t$id" else s"T$id"
+ def first(id: Int): String = if (id < 3) s"a$id" else s"A$id"
+ def second(id: Int, nullElement: Boolean): String =
+ if (id < 3) s"b$id" else if (id == 3 || !nullElement) s"B$id" else null
+ def pVal(id: Int): String = if (id < 3) s"p$id" else s"P$id"
+ def qVal(id: Int): String = if (id < 3) s"q$id" else if (id == 3) "Q3"
else null
+ def item(id: Int): String = if (id < 3) s"i$id" else if (id == 3) "I3"
else null
+ def json(k: String): String = Option(k).map(k => s"""{"k":"$k"}""").orNull
+
+ def assertCollectionsNative(sql: String, pushed: Boolean, leg: String):
Unit = {
+ val scans = spark.sql(sql).queryExecution.sparkPlan.collect { case scan:
FileSourceScanExec => scan }
+ assert(scans.nonEmpty, s"[$leg] expected a file scan in the plan of:
$sql")
+ val required = scans.head.requiredSchema
+ def field(name: String): DataType = required.fields.find(_.name == name)
+ .getOrElse(fail(s"[$leg] scan does not read '$name':
${required.treeString}")).dataType
+ val adapter = SparkAdapterSupport.sparkAdapter
+ // The array element, the map value and the struct member inside the
array element are all
+ // native VariantType in the scan on both arms; only the top-level
column follows the conf.
+
assert(adapter.isVariantType(field("arr").asInstanceOf[ArrayType].elementType),
+ s"[$leg] arr's element should be native VariantType in the scan:
${required.treeString}")
+ assert(adapter.isVariantType(field("m").asInstanceOf[MapType].valueType),
+ s"[$leg] m's value should be native VariantType in the scan:
${required.treeString}")
+ val itemsElement =
field("items").asInstanceOf[ArrayType].elementType.asInstanceOf[StructType]
+ assert(adapter.isVariantType(itemsElement("inner").dataType),
+ s"[$leg] items' element member should be native VariantType in the
scan: ${required.treeString}")
+ assert(adapter.containsVariantProjection(field("v")) == pushed,
+ s"[$leg] the top-level v ${if (pushed) "should" else "must not"} be a
projection struct: ${required.treeString}")
+ }
+
+ def runLeg(tableName: String, tablePath: String, pushed: Boolean,
nullElement: Boolean,
+ expectedBlock: Option[HoodieLogBlockType], legLabel: String):
Unit = {
+ val leg = legLabel
Review Comment:
🤖 nit: `val leg = legLabel` is a pure alias, and `assertReads` then shadows
it with `val leg = s"$legLabel, $phase"` — could you just name the parameter
`leg` and call the inner one something like `legPhase` so the two don't collide?
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]