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