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 b84f61823610 fix(reader): Remove redundant partition value conversion
in RecordContext (#19201)
b84f61823610 is described below
commit b84f61823610647fbf56dba53d260741a3e38d23
Author: Shuo Cheng <[email protected]>
AuthorDate: Tue Jul 7 19:24:17 2026 +0800
fix(reader): Remove redundant partition value conversion in RecordContext
(#19201)
---
.../hudi/BaseSparkInternalRecordContext.java | 11 ---------
.../apache/hudi/common/engine/RecordContext.java | 4 ----
.../common/table/read/HoodieFileGroupReader.java | 2 +-
...stSparkFileFormatInternalRowReaderContext.scala | 26 ++++++++--------------
4 files changed, 10 insertions(+), 33 deletions(-)
diff --git
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java
index e8bacf6d3b78..bdda6deb1d76 100644
---
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java
+++
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java
@@ -188,17 +188,6 @@ public abstract class BaseSparkInternalRecordContext
extends RecordContext<Inter
return (Comparable) value;
}
- @Override
- public Comparable convertPartitionValueToEngineType(Comparable value) {
- if (value instanceof String) {
- // Spark reads String field values as UTF8String.
- // To foster value comparison, if the value is of String type, e.g., from
- // the delete record, we convert it to UTF8String type.
- return UTF8String.fromString((String) value);
- }
- return value;
- }
-
@Override
public InternalRow getDeleteRow(String recordKey) {
UTF8String[] metaFields = new UTF8String[]{
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/engine/RecordContext.java
b/hudi-common/src/main/java/org/apache/hudi/common/engine/RecordContext.java
index 81cd60e8ec8d..633334c878f2 100644
--- a/hudi-common/src/main/java/org/apache/hudi/common/engine/RecordContext.java
+++ b/hudi-common/src/main/java/org/apache/hudi/common/engine/RecordContext.java
@@ -227,10 +227,6 @@ public abstract class RecordContext<T> implements
Serializable {
return value;
}
- public Comparable convertPartitionValueToEngineType(Comparable value) {
- return convertValueToEngineType(value);
- }
-
/**
* Converts the ordering value to the specific engine type, guaranteeing the
returned
* value supports comparison, e.g., by wrapping engine types that do not
implement
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
index b1a914a725ef..d9f4e63719e7 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
@@ -288,7 +288,7 @@ public final class HoodieFileGroupReader<T> implements
HoodieRecordReader<T> {
for (int i = 0; i < partitionFields.length; i++) {
String field = partitionFields[i];
if (dataSchema.getField(field).isPresent()) {
- filterFieldsAndValues.add(Pair.of(field,
readerContext.getRecordContext().convertPartitionValueToEngineType((Comparable)
partitionValues[i])));
+ filterFieldsAndValues.add(Pair.of(field,
readerContext.getRecordContext().convertValueToEngineType((Comparable)
partitionValues[i])));
}
}
return filterFieldsAndValues;
diff --git
a/hudi-spark-datasource/hudi-spark-common/src/test/scala/org/apache/spark/execution/datasources/parquet/TestSparkFileFormatInternalRowReaderContext.scala
b/hudi-spark-datasource/hudi-spark-common/src/test/scala/org/apache/spark/execution/datasources/parquet/TestSparkFileFormatInternalRowReaderContext.scala
index 6aaec1ca7427..5427df8965f5 100644
---
a/hudi-spark-datasource/hudi-spark-common/src/test/scala/org/apache/spark/execution/datasources/parquet/TestSparkFileFormatInternalRowReaderContext.scala
+++
b/hudi-spark-datasource/hudi-spark-common/src/test/scala/org/apache/spark/execution/datasources/parquet/TestSparkFileFormatInternalRowReaderContext.scala
@@ -19,23 +19,18 @@
package org.apache.spark.execution.datasources.parquet
-import org.apache.hudi.SparkFileFormatInternalRowReaderContext
+import org.apache.hudi.{SparkFileFormatInternalRecordContext,
SparkFileFormatInternalRowReaderContext}
import
org.apache.hudi.SparkFileFormatInternalRowReaderContext.{filterIsSafeForBootstrap,
filterIsSafeForPrimaryKey}
import org.apache.hudi.common.model.HoodieRecord
-import org.apache.hudi.common.table.HoodieTableConfig
import
org.apache.hudi.common.table.read.buffer.PositionBasedFileGroupRecordBuffer.ROW_INDEX_TEMPORARY_COLUMN_NAME
-import org.apache.hudi.testutils.SparkClientFunctionalTestHarness
-import org.apache.spark.sql.execution.datasources.SparkColumnarFileReader
import org.apache.spark.sql.sources.{And, IsNotNull, Or}
import org.apache.spark.sql.types.{LongType, StringType, StructField,
StructType}
import org.apache.spark.unsafe.types.UTF8String
import org.junit.jupiter.api.Assertions.{assertEquals, assertFalse, assertTrue}
import org.junit.jupiter.api.Test
-import org.mockito.Mockito
-import org.mockito.Mockito.when
-class TestSparkFileFormatInternalRowReaderContext extends
SparkClientFunctionalTestHarness {
+class TestSparkFileFormatInternalRowReaderContext {
@Test
def testBootstrapFilters(): Unit = {
@@ -104,19 +99,16 @@ class TestSparkFileFormatInternalRowReaderContext extends
SparkClientFunctionalT
@Test
def testConvertValueToEngineType(): Unit = {
- val reader = Mockito.mock(classOf[SparkColumnarFileReader])
val stringValue = "string_value"
- val tableConfig = Mockito.mock(classOf[HoodieTableConfig])
- when(tableConfig.populateMetaFields()).thenReturn(true)
- val sparkReaderContext = new
SparkFileFormatInternalRowReaderContext(reader, Seq.empty, Seq.empty,
storageConf(), tableConfig)
- assertEquals(1,
sparkReaderContext.getRecordContext().convertValueToEngineType(1))
- assertEquals(1L,
sparkReaderContext.getRecordContext().convertValueToEngineType(1L))
- assertEquals(1.1f,
sparkReaderContext.getRecordContext().convertValueToEngineType(1.1f))
- assertEquals(1.1d,
sparkReaderContext.getRecordContext().convertValueToEngineType(1.1d))
+ val recordContext = SparkFileFormatInternalRecordContext.apply()
+ assertEquals(1, recordContext.convertValueToEngineType(1))
+ assertEquals(1L, recordContext.convertValueToEngineType(1L))
+ assertEquals(1.1f, recordContext.convertValueToEngineType(1.1f))
+ assertEquals(1.1d, recordContext.convertValueToEngineType(1.1d))
assertEquals(UTF8String.fromString(stringValue),
-
sparkReaderContext.getRecordContext().convertPartitionValueToEngineType(stringValue))
+ recordContext.convertValueToEngineType(stringValue))
val utf8StringValue = UTF8String.fromString(stringValue)
assertEquals(utf8StringValue,
-
sparkReaderContext.getRecordContext().convertPartitionValueToEngineType(utf8StringValue))
+ recordContext.convertValueToEngineType(utf8StringValue))
}
}