This is an automated email from the ASF dual-hosted git repository.
voonhous 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 cb0c7a493e07 refactor(core): decide the record representation from the
payload class alone (#19990)
cb0c7a493e07 is described below
commit cb0c7a493e07f13aeff2d50520ef0c08b9774561
Author: Y Ethan Guo <[email protected]>
AuthorDate: Fri Sep 18 00:06:29 2026 -0700
refactor(core): decide the record representation from the payload class
alone (#19990)
* refactor(core): decide the record representation from the payload class
alone
HoodieRecordUtils.createHoodieRecord chose between a payload-carrying
HoodieAvroRecord and a payload-free HoodieAvroIndexedRecord through
requiresPayload, which combined isChangingRecords with a subclass test on
the configured merge handle. That test is a poor proxy for whether a
payload-consuming merge path runs: it classifies by inheritance rather than
capability, it reads a config that most branches of HoodieMergeHandleFactory
override, and it is computed once on the driver while the handle is chosen
per file group with a runtime fallback.
Every merge and append handle on the changing-records path consumes incoming
records through BufferedRecords.fromHoodieRecord or an equivalent
abstraction
and never calls HoodieRecordPayload APIs on the incoming record; custom
merge
modes rebuild the payload from the buffered data in HoodieAvroRecordMerger.
Every caller outside the Spark write path and the streamer already passed a
literal false.
The representation is now decided by isPayloadClassDeprecated alone, and
requiresPayload and isFileGroupReaderBasedMergeHandle are removed.
* Add an UPDATE case with the merge handle pinned to HoodieWriteMergeHandle
The existing case runs with the default handle, which builds payload-free
records. Pinning HoodieWriteMergeHandle makes the write build
payload-carrying
records and merge on the ordering value each record carries, which is the
path
a prepped SQL UPDATE that leaves the ordering column unassigned exercises.
Both
cases now run.
* Scope the optimized-writes flag to the pinned-handle UPDATE case
Setting it through a SQL set statement leaked into the later tests in the
suite, which then ran their UPDATEs with the flag off instead of the
default.
---
.../org/apache/hudi/config/HoodieWriteConfig.java | 7 ----
.../org/apache/hudi/index/HoodieIndexUtils.java | 2 +-
.../apache/hudi/common/avro/AvroRecordContext.java | 2 +-
.../apache/hudi/common/util/HoodieRecordUtils.java | 16 ++++----
.../org/apache/hudi/HoodieCreateRecordUtils.scala | 6 +--
.../SparkFullBootstrapDataProviderBase.java | 2 +-
.../benchmark/CreateHandleBenchmark.scala | 5 ++-
.../sql/hudi/dml/others/TestUpdateTable.scala | 48 ++++++++++++++++++++++
.../utilities/streamer/HoodieStreamerUtils.java | 7 +---
9 files changed, 66 insertions(+), 29 deletions(-)
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java
index ffe2c0deb233..86a8cfefa0a8 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java
@@ -4090,11 +4090,4 @@ public class HoodieWriteConfig extends HoodieConfig {
}
}
- public boolean isFileGroupReaderBasedMergeHandle() {
- return isFileGroupReaderBasedMergeHandle(props);
- }
-
- public static boolean isFileGroupReaderBasedMergeHandle(TypedProperties
props) {
- return ReflectionUtils.isSubClass(ConfigUtils.getStringWithAltKeys(props,
HoodieWriteConfig.MERGE_HANDLE_CLASS_NAME, true),
FileGroupReaderBasedMergeHandle.class);
- }
}
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/index/HoodieIndexUtils.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/index/HoodieIndexUtils.java
index d9f6fe8c5fd5..9b1a6bd0d152 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/index/HoodieIndexUtils.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/index/HoodieIndexUtils.java
@@ -451,7 +451,7 @@ public class HoodieIndexUtils {
return Option.empty();
}
String partitionPath = inferPartitionPath(incoming, existing,
writeSchemaWithMetaFields, keyGenerator, existingRecordContext, mergeResult);
- if (config.isFileGroupReaderBasedMergeHandle() &&
HoodieRecordUtils.isPayloadClassDeprecated(ConfigUtils.getPayloadClass(properties)))
{
+ if
(HoodieRecordUtils.isPayloadClassDeprecated(ConfigUtils.getPayloadClass(properties)))
{
return
Option.of(existingRecordContext.constructHoodieRecord(mergeResult,
partitionPath));
} else {
HoodieRecord<R> result =
existingRecordContext.constructHoodieRecord(mergeResult, partitionPath);
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/avro/AvroRecordContext.java
b/hudi-common/src/main/java/org/apache/hudi/common/avro/AvroRecordContext.java
index 1dc7d4a5dbb3..8f4a5fa2ddf8 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/avro/AvroRecordContext.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/avro/AvroRecordContext.java
@@ -126,7 +126,7 @@ public class AvroRecordContext extends
RecordContext<IndexedRecord> {
}
return HoodieRecordUtils.createHoodieRecord((GenericRecord)
bufferedRecord.getRecord(), bufferedRecord.getOrderingValue(),
- hoodieKey, payloadClass, bufferedRecord.getHoodieOperation(),
Option.empty(), false, bufferedRecord.isDelete());
+ hoodieKey, payloadClass, bufferedRecord.getHoodieOperation(),
Option.empty(), bufferedRecord.isDelete());
}
@Override
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/util/HoodieRecordUtils.java
b/hudi-common/src/main/java/org/apache/hudi/common/util/HoodieRecordUtils.java
index 27337de679cb..4d2cb4b76849 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/util/HoodieRecordUtils.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/util/HoodieRecordUtils.java
@@ -159,14 +159,14 @@ public class HoodieRecordUtils {
}
}
- public static HoodieRecord createHoodieRecord(GenericRecord data, Comparable
orderingVal, HoodieKey hKey, String payloadClass, boolean requiresPayload,
Boolean isDelete) {
- return createHoodieRecord(data, orderingVal, hKey, payloadClass, null,
Option.empty(), requiresPayload, isDelete);
+ public static HoodieRecord createHoodieRecord(GenericRecord data, Comparable
orderingVal, HoodieKey hKey, String payloadClass, Boolean isDelete) {
+ return createHoodieRecord(data, orderingVal, hKey, payloadClass, null,
Option.empty(), isDelete);
}
public static HoodieRecord createHoodieRecord(GenericRecord data, Comparable
orderingVal, HoodieKey hKey,
- String payloadClass,
HoodieOperation hoodieOperation, Option<HoodieRecordLocation> recordLocation,
boolean requiresPayload, Boolean isDelete) {
+ String payloadClass,
HoodieOperation hoodieOperation, Option<HoodieRecordLocation> recordLocation,
Boolean isDelete) {
HoodieRecord record;
- if (!requiresPayload && isPayloadClassDeprecated(payloadClass)) {
+ if (isPayloadClassDeprecated(payloadClass)) {
record = new HoodieAvroIndexedRecord(hKey, data, orderingVal,
hoodieOperation, isDelete);
} else {
HoodieRecordPayload payload =
HoodieRecordUtils.loadPayload(payloadClass, data, orderingVal);
@@ -177,14 +177,14 @@ public class HoodieRecordUtils {
}
public static HoodieRecord createHoodieRecord(GenericRecord data, HoodieKey
hKey,
- String payloadClass, boolean
requiresPayload, Boolean isDelete) {
- return createHoodieRecord(data, hKey, payloadClass, Option.empty(),
requiresPayload, isDelete);
+ String payloadClass, Boolean
isDelete) {
+ return createHoodieRecord(data, hKey, payloadClass, Option.empty(),
isDelete);
}
public static HoodieRecord createHoodieRecord(GenericRecord data, HoodieKey
hKey,
- String payloadClass,
Option<HoodieRecordLocation> recordLocation, boolean requiresPayload, Boolean
isDelete) {
+ String payloadClass,
Option<HoodieRecordLocation> recordLocation, Boolean isDelete) {
HoodieRecord record;
- if (!requiresPayload && isPayloadClassDeprecated(payloadClass)) {
+ if (isPayloadClassDeprecated(payloadClass)) {
record = new HoodieAvroIndexedRecord(hKey, data, null, (HoodieOperation)
null, isDelete);
} else {
HoodieRecordPayload payload =
HoodieRecordUtils.loadPayload(payloadClass, data);
diff --git
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieCreateRecordUtils.scala
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieCreateRecordUtils.scala
index b8498c672559..d16db68b51d6 100644
---
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieCreateRecordUtils.scala
+++
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieCreateRecordUtils.scala
@@ -23,7 +23,6 @@ import org.apache.hudi.common.avro.{AvroRecordContext,
HoodieAvroUtils}
import org.apache.hudi.common.config.{RecordMergeMode, TypedProperties}
import org.apache.hudi.common.fs.FSUtils
import org.apache.hudi.common.model._
-import org.apache.hudi.common.model.WriteOperationType.isChangingRecords
import org.apache.hudi.common.schema.{HoodieSchema, HoodieSchemaCache}
import org.apache.hudi.common.table.HoodieTableConfig
import org.apache.hudi.common.table.read.DeleteContext
@@ -135,7 +134,6 @@ object HoodieCreateRecordUtils {
val consistentLogicalTimestampEnabled = parameters.getOrElse(
DataSourceWriteOptions.KEYGENERATOR_CONSISTENT_LOGICAL_TIMESTAMP_ENABLED.key(),
DataSourceWriteOptions.KEYGENERATOR_CONSISTENT_LOGICAL_TIMESTAMP_ENABLED.defaultValue()).toBoolean
- val requiresPayload = isChangingRecords(operation) &&
!config.isFileGroupReaderBasedMergeHandle
val mergeProps = ConfigUtils.getMergeProps(config.getProps,
args.tableConfig)
val deleteContext = new DeleteContext(mergeProps,
writerSchema).withReaderSchema(writerSchema);
@@ -159,10 +157,10 @@ object HoodieCreateRecordUtils {
val orderingVal = getOrderingValue(orderingFields, avroRec,
hoodieKey.getRecordKey,
consistentLogicalTimestampEnabled, requiresOrderingValue)
HoodieRecordUtils.createHoodieRecord(processedRecord,
orderingVal, hoodieKey,
- config.getPayloadClass, null, recordLocation, requiresPayload,
isDelete)
+ config.getPayloadClass, null, recordLocation, isDelete)
} else {
HoodieRecordUtils.createHoodieRecord(processedRecord, hoodieKey,
- config.getPayloadClass, recordLocation, requiresPayload,
isDelete)
+ config.getPayloadClass, recordLocation, isDelete)
}
hoodieRecord
}
diff --git
a/hudi-spark-datasource/hudi-spark/src/main/java/org/apache/hudi/bootstrap/SparkFullBootstrapDataProviderBase.java
b/hudi-spark-datasource/hudi-spark/src/main/java/org/apache/hudi/bootstrap/SparkFullBootstrapDataProviderBase.java
index 39d5d137de05..c2efb22edc68 100644
---
a/hudi-spark-datasource/hudi-spark/src/main/java/org/apache/hudi/bootstrap/SparkFullBootstrapDataProviderBase.java
+++
b/hudi-spark-datasource/hudi-spark/src/main/java/org/apache/hudi/bootstrap/SparkFullBootstrapDataProviderBase.java
@@ -81,7 +81,7 @@ public abstract class SparkFullBootstrapDataProviderBase
extends FullRecordBoots
gr, orderingFieldsStr, false, props.getBoolean(
KeyGeneratorOptions.KEYGENERATOR_CONSISTENT_LOGICAL_TIMESTAMP_ENABLED.key(),
Boolean.parseBoolean(KeyGeneratorOptions.KEYGENERATOR_CONSISTENT_LOGICAL_TIMESTAMP_ENABLED.defaultValue())));
- return HoodieRecordUtils.createHoodieRecord(gr, orderingVal,
keyGenerator.getKey(gr), ConfigUtils.getPayloadClass(props), false, null);
+ return HoodieRecordUtils.createHoodieRecord(gr, orderingVal,
keyGenerator.getKey(gr), ConfigUtils.getPayloadClass(props), null);
});
} else if (recordType == HoodieRecordType.SPARK) {
SparkKeyGeneratorInterface sparkKeyGenerator =
(SparkKeyGeneratorInterface) keyGenerator;
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/execution/benchmark/CreateHandleBenchmark.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/execution/benchmark/CreateHandleBenchmark.scala
index 4c4d8f93866c..47fbb25e017d 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/execution/benchmark/CreateHandleBenchmark.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/execution/benchmark/CreateHandleBenchmark.scala
@@ -27,7 +27,7 @@ import
org.apache.hudi.common.model.{DefaultHoodieRecordPayload, HoodieAvroIndex
import org.apache.hudi.common.schema.HoodieSchema
import org.apache.hudi.common.table.HoodieTableConfig
import org.apache.hudi.common.table.marker.MarkerType
-import org.apache.hudi.common.util.HoodieRecordUtils
+import org.apache.hudi.common.util.{HoodieRecordUtils, Option => HOption}
import org.apache.hudi.config.HoodieWriteConfig
import org.apache.hudi.io.HoodieCreateHandle
import org.apache.hudi.keygen.constant.KeyGeneratorOptions
@@ -128,7 +128,8 @@ object CreateHandleBenchmark extends HoodieBenchmarkBase {
it => {
it.map { genRec =>
val hoodieKey = new HoodieKey(genRec.get("key").toString, "")
- HoodieRecordUtils.createHoodieRecord(genRec, 0L, hoodieKey,
classOf[DefaultHoodieRecordPayload].getName, false, null)
+ HoodieRecordUtils.createHoodieRecord(genRec,
java.lang.Long.valueOf(0L), hoodieKey,
+ classOf[DefaultHoodieRecordPayload].getName, null,
HOption.empty(), null)
}
}).toJavaRDD().collect().stream().map[HoodieRecord[_]](hoodieRec => {
hoodieRec.asInstanceOf[HoodieAvroIndexedRecord].toIndexedRecord(schema,
dummpProps)
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/others/TestUpdateTable.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/others/TestUpdateTable.scala
index 9d67da0db5ce..bee62c001858 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/others/TestUpdateTable.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/others/TestUpdateTable.scala
@@ -25,6 +25,7 @@ import org.apache.hudi.HoodieSparkUtils.gteqSpark3_4
import org.apache.hudi.common.model.HoodieTableType
import org.apache.hudi.common.table.timeline.HoodieInstant
import org.apache.hudi.common.util.{Option => HOption}
+import org.apache.hudi.config.HoodieWriteConfig
import org.apache.hudi.testutils.HoodieClientTestUtils.createMetaClient
import org.apache.spark.sql.{AnalysisException, Row}
@@ -81,6 +82,53 @@ class TestUpdateTable extends HoodieSparkSqlTestBase {
})
}
+ test("Test Update Table With Row Merge Handle") {
+ // Same statements as above with the merge handle pinned to
HoodieWriteMergeHandle, so the write
+ // builds payload-carrying records and merges on the ordering value each
record carries.
+ withRecordType()(withTempDir { tmp =>
+ Seq(true, false).foreach { sparkSqlOptimizedWrites =>
+ Seq("cow", "mor").foreach { tableType =>
+ val tableName = generateTableName
+ spark.sql(
+ s"""
+ |create table $tableName (
+ | id int,
+ | name string,
+ | price double,
+ | ts long
+ |) using hudi
+ | location '${tmp.getCanonicalPath}/$tableName'
+ | tblproperties (
+ | type = '$tableType',
+ | primaryKey = 'id',
+ | preCombineField = 'ts',
+ | '${HoodieWriteConfig.MERGE_HANDLE_CLASS_NAME.key()}' =
'org.apache.hudi.io.HoodieWriteMergeHandle'
+ | )
+ """.stripMargin)
+
+ spark.sql(s"insert into $tableName select 1, 'a1', 10, 1000")
+ checkAnswer(s"select id, name, price, ts from $tableName")(
+ Seq(1, "a1", 10.0, 1000)
+ )
+
+ withSQLConf(SPARK_SQL_OPTIMIZED_WRITES.key() ->
sparkSqlOptimizedWrites.toString) {
+
+ // the ordering column is not assigned, so the update must take
effect
+ spark.sql(s"update $tableName set price = 20 where id = 1")
+ checkAnswer(s"select id, name, price, ts from $tableName")(
+ Seq(1, "a1", 20.0, 1000)
+ )
+
+ spark.sql(s"update $tableName set price = price * 2 where id = 1")
+ checkAnswer(s"select id, name, price, ts from $tableName")(
+ Seq(1, "a1", 40.0, 1000)
+ )
+ }
+ }
+ }
+ })
+ }
+
test("Test Update Table Without Primary Key") {
withRecordType()(withTempDir { tmp =>
Seq("cow", "mor").foreach { tableType =>
diff --git
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamerUtils.java
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamerUtils.java
index c0e6b118eae3..90eba1962405 100644
---
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamerUtils.java
+++
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamerUtils.java
@@ -41,7 +41,6 @@ import org.apache.hudi.common.util.OrderingValues;
import org.apache.hudi.common.util.StringUtils;
import org.apache.hudi.common.util.collection.ClosableIterator;
import org.apache.hudi.common.util.collection.CloseableMappingIterator;
-import org.apache.hudi.config.HoodieWriteConfig;
import org.apache.hudi.exception.HoodieException;
import org.apache.hudi.exception.HoodieIOException;
import org.apache.hudi.exception.HoodieKeyException;
@@ -71,7 +70,6 @@ import java.util.Iterator;
import java.util.Set;
import java.util.stream.Collectors;
-import static
org.apache.hudi.common.model.WriteOperationType.isChangingRecords;
import static
org.apache.hudi.common.table.HoodieTableConfig.DROP_PARTITION_COLUMNS;
import static
org.apache.hudi.config.HoodieErrorTableConfig.ERROR_ENABLE_VALIDATE_RECORD_CREATION;
@@ -99,7 +97,6 @@ public class HoodieStreamerUtils {
String payloadClassName = StringUtils.isNullOrEmpty(cfg.payloadClassName)
? HoodieRecordPayload.getAvroPayloadForMergeMode(cfg.recordMergeMode,
cfg.payloadClassName)
: cfg.payloadClassName;
- boolean requiresPayload = isChangingRecords(cfg.operation) &&
!HoodieWriteConfig.isFileGroupReaderBasedMergeHandle(props);
return avroRDDOptional.map(avroRDD -> {
HoodieSchema targetSchema = schemaProvider.getTargetHoodieSchema();
@@ -133,8 +130,8 @@ public class HoodieStreamerUtils {
? OrderingValues.create(orderingFieldsStr.split(","),
field -> (Comparable)
HoodieAvroUtils.getNestedFieldVal(gr, field, false,
useConsistentLogicalTimestamp))
: null;
- HoodieRecord record = shouldUseOrderingField ?
HoodieRecordUtils.createHoodieRecord(gr, orderingValue, hoodieKey,
payloadClassName, requiresPayload, isDelete)
- : HoodieRecordUtils.createHoodieRecord(gr, hoodieKey,
payloadClassName, requiresPayload, isDelete);
+ HoodieRecord record = shouldUseOrderingField ?
HoodieRecordUtils.createHoodieRecord(gr, orderingValue, hoodieKey,
payloadClassName, isDelete)
+ : HoodieRecordUtils.createHoodieRecord(gr, hoodieKey,
payloadClassName, isDelete);
return Either.left(record);
} catch (Exception e) {
return generateErrorRecordOrThrowException(genRec, e,
shouldErrorTable);