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]

Reply via email to