This is an automated email from the ASF dual-hosted git repository. voonhous pushed a commit to branch release-1.2.1 in repository https://gitbox.apache.org/repos/asf/hudi.git
commit 72fee44b3e28daf8bff6d0cf23cf08ea6b80653a Author: Shuo Cheng <[email protected]> AuthorDate: Tue Jul 7 19:24:17 2026 +0800 fix(reader): Remove redundant partition value conversion in RecordContext (#19201) (cherry picked from commit b84f61823610647fbf56dba53d260741a3e38d23) --- .../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 6bd5420af014..2944eeb8affa 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 @@ -182,17 +182,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 d63c9bd63408..2996f9a3cb19 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 @@ -217,10 +217,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. */ 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 0b536b511db1..b4b1bc6b74cb 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 @@ -289,7 +289,7 @@ public final class HoodieFileGroupReader<T> implements Closeable { 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)) } }
