This is an automated email from the ASF dual-hosted git repository.

danny0405 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 f51c36fd6439 [HUDI-9683] Fix field renames in 
HoodieInternalRowUtils.genUnsafeRowWriter (#13672)
f51c36fd6439 is described below

commit f51c36fd64398718c039b26c15767e793b82904e
Author: Jon Vexler <[email protected]>
AuthorDate: Sun Aug 3 03:00:06 2025 -0400

    [HUDI-9683] Fix field renames in HoodieInternalRowUtils.genUnsafeRowWriter 
(#13672)
    
    * fix field renaming in spark projection
    * should not use full path so that we match the output of the merger
    * update comments that were wrong
    * avoid string handling when field is not renamed
    
    ---------
    
    Co-authored-by: Jonathan Vexler <=>
    Co-authored-by: danny0405 <[email protected]>
---
 .../apache/spark/sql/HoodieInternalRowUtils.scala  |  28 +++---
 .../hudi/common/TestHoodieInternalRowUtils.scala   | 102 +++++++++++++++++++++
 2 files changed, 113 insertions(+), 17 deletions(-)

diff --git 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/spark/sql/HoodieInternalRowUtils.scala
 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/spark/sql/HoodieInternalRowUtils.scala
index 4a91609e19ba..e5c5761c73b3 100644
--- 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/spark/sql/HoodieInternalRowUtils.scala
+++ 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/spark/sql/HoodieInternalRowUtils.scala
@@ -188,25 +188,14 @@ object HoodieInternalRowUtils {
       fieldNamesStack.push(newField.name)
 
       val (fieldWriter, prevFieldPos): (RowFieldUpdater, Int) =
-        prevStructType.getFieldIndex(newField.name) match {
+        prevStructType.getFieldIndex(lookupRenamedField(newField.name, 
createFullName(fieldNamesStack), renamedColumnsMap)) match {
           case Some(prevFieldPos) =>
             val prevField = prevStructType(prevFieldPos)
             (newWriterRenaming(prevField.dataType, newField.dataType, 
renamedColumnsMap, fieldNamesStack), prevFieldPos)
 
           case None =>
-            val newFieldQualifiedName = createFullName(fieldNamesStack)
-            val prevFieldName: String = 
lookupRenamedField(newFieldQualifiedName, renamedColumnsMap)
-
-            // Handle rename
-            prevStructType.getFieldIndex(prevFieldName) match {
-              case Some(prevFieldPos) =>
-                val prevField = prevStructType.fields(prevFieldPos)
-                (newWriterRenaming(prevField.dataType, newField.dataType, 
renamedColumnsMap, fieldNamesStack), prevFieldPos)
-
-              case None =>
-                val updater: RowFieldUpdater = (fieldUpdater, ordinal, _) => 
fieldUpdater.setNullAt(ordinal)
-                (updater, -1)
-            }
+            val updater: RowFieldUpdater = (fieldUpdater, ordinal, _) => 
fieldUpdater.setNullAt(ordinal)
+            (updater, -1)
         }
 
       fieldWriters += fieldWriter
@@ -415,9 +404,14 @@ object HoodieInternalRowUtils {
     }
   }
 
-  private def lookupRenamedField(newFieldQualifiedName: String, 
renamedColumnsMap: JMap[String, String]) = {
-    val prevFieldQualifiedName = 
renamedColumnsMap.getOrDefault(newFieldQualifiedName, "")
-    val prevFieldQualifiedNameParts = prevFieldQualifiedName.split("\\.")
+  private def lookupRenamedField(newFieldName: String,
+                                 newFieldQualifiedName: String,
+                                 renamedColumnsMap: JMap[String, String]): 
String = {
+    val renamed = renamedColumnsMap.get(newFieldQualifiedName)
+    if (renamed == null) {
+      return newFieldName
+    }
+    val prevFieldQualifiedNameParts = renamed.split("\\.")
     val prevFieldName = 
prevFieldQualifiedNameParts(prevFieldQualifiedNameParts.length - 1)
 
     prevFieldName
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/TestHoodieInternalRowUtils.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/TestHoodieInternalRowUtils.scala
index d0d045da5931..989b2b19c225 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/TestHoodieInternalRowUtils.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/TestHoodieInternalRowUtils.scala
@@ -124,6 +124,108 @@ class TestHoodieInternalRowUtils extends FunSuite with 
Matchers with BeforeAndAf
     assertEquals(serDe.deserializeRow(newRow2), Row("Andrew", 25, Row("Mission 
st", "SF")));
   }
 
+  test("Test rewrite row with renamed columns") {
+    // Original schema
+    val oldSchema = StructType(Seq(
+      StructField("first_name", StringType),
+      StructField("years_old", IntegerType)
+    ))
+
+    // Renamed schema
+    val newSchema = StructType(Seq(
+      StructField("name", StringType),
+      StructField("age", IntegerType)
+    ))
+
+    // Rename mapping: new -> old
+    val renameMap: java.util.Map[String, String] = new java.util.HashMap()
+    renameMap.put("name", "first_name")
+    renameMap.put("age", "years_old")
+
+    // Sample row
+    val oldRowData = sparkSession.sparkContext.parallelize(Seq(Row("Alice", 
30)))
+    val oldRow = sparkSession.createDataFrame(oldRowData, 
oldSchema).queryExecution.toRdd.first()
+
+    // Generate writer with rename map
+    val rowWriter = HoodieInternalRowUtils.genUnsafeRowWriter(oldSchema, 
newSchema, renameMap, JCollections.emptyMap())
+    val newRow = rowWriter(oldRow)
+
+    val serDe = sparkAdapter.createSparkRowSerDe(newSchema)
+    assertEquals(Row("Alice", 30), serDe.deserializeRow(newRow))
+  }
+
+  test("Test rewrite row with columns swap") {
+    // Original schema
+    val oldSchema = StructType(Seq(
+      StructField("first_name", StringType),
+      StructField("years_old", IntegerType)
+    ))
+
+    // Renamed schema
+    val newSchema = StructType(Seq(
+      StructField("years_old", StringType),
+      StructField("first_name", IntegerType)
+    ))
+
+    // Rename mapping: new -> old
+    val renameMap: java.util.Map[String, String] = new java.util.HashMap()
+    renameMap.put("years_old", "first_name")
+    renameMap.put("first_name", "years_old")
+
+    // Sample row
+    val oldRowData = sparkSession.sparkContext.parallelize(Seq(Row("Alice", 
30)))
+    val oldRow = sparkSession.createDataFrame(oldRowData, 
oldSchema).queryExecution.toRdd.first()
+
+    // Generate writer with rename map
+    val rowWriter = HoodieInternalRowUtils.genUnsafeRowWriter(oldSchema, 
newSchema, renameMap, JCollections.emptyMap())
+    val newRow = rowWriter(oldRow)
+
+    val serDe = sparkAdapter.createSparkRowSerDe(newSchema)
+    assertEquals(Row("Alice", 30), serDe.deserializeRow(newRow))
+  }
+
+  test("Test rewrite row with columns swap nested") {
+    // Original schema
+    val oldSchema = StructType(Seq(
+      StructField("first_name", StringType),
+      StructField("years_old", IntegerType),
+      StructField("address",
+        StructType(Seq(
+          StructField("city", StringType),
+          StructField("street", StringType)
+        )
+    ))))
+
+    // Renamed schema
+    val newSchema = StructType(Seq(
+      StructField("years_old", StringType),
+      StructField("first_name", IntegerType),
+      StructField("address",
+        StructType(Seq(
+          StructField("street", StringType),
+          StructField("city", StringType)
+        )
+        ))))
+
+    // Rename mapping: new -> old
+    val renameMap: java.util.Map[String, String] = new java.util.HashMap()
+    renameMap.put("years_old", "first_name")
+    renameMap.put("first_name", "years_old")
+    renameMap.put("address.city", "street")
+    renameMap.put("address.street", "city")
+
+    // Sample row
+    val oldRowData = sparkSession.sparkContext.parallelize(Seq(Row("Alice", 
30, Row("SF", "Mission st"))))
+    val oldRow = sparkSession.createDataFrame(oldRowData, 
oldSchema).queryExecution.toRdd.first()
+
+    // Generate writer with rename map
+    val rowWriter = HoodieInternalRowUtils.genUnsafeRowWriter(oldSchema, 
newSchema, renameMap, JCollections.emptyMap())
+    val newRow = rowWriter(oldRow)
+
+    val serDe = sparkAdapter.createSparkRowSerDe(newSchema)
+    assertEquals(Row("Alice", 30, Row("SF", "Mission st")), 
serDe.deserializeRow(newRow))
+  }
+
   /**
    * test record data type changes.
    * int => long/float/double/string

Reply via email to