This is an automated email from the ASF dual-hosted git repository.
sivabalan 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 6d3dbcc287d [HUDI-9255] Fix inferring correct merge behavior for few
scenarios (#13079)
6d3dbcc287d is described below
commit 6d3dbcc287d5082a3e5d5acb93ad705103c7a5ee
Author: Lokesh Jain <[email protected]>
AuthorDate: Fri Apr 4 11:24:47 2025 +0530
[HUDI-9255] Fix inferring correct merge behavior for few scenarios (#13079)
---------
Co-authored-by: sivabalan <[email protected]>
---
.../hudi/common/model/HoodieRecordMerger.java | 34 ++++++---
.../hudi/common/table/HoodieTableConfig.java | 2 +-
.../table/read/FileGroupReaderSchemaHandler.java | 3 +-
.../table/read/TestFileGroupRecordBuffer.java | 82 ++++++++++++++++------
.../hudi/common/table/TestHoodieTableConfig.java | 8 +--
.../hudi/utilities/streamer/HoodieStreamer.java | 5 +-
6 files changed, 93 insertions(+), 41 deletions(-)
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecordMerger.java
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecordMerger.java
index ac17959d5a8..ac302480065 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecordMerger.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecordMerger.java
@@ -24,6 +24,7 @@ import org.apache.hudi.common.config.RecordMergeMode;
import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType;
import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.HoodieTableVersion;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.StringUtils;
import org.apache.hudi.common.util.collection.Pair;
@@ -186,9 +187,8 @@ public interface HoodieRecordMerger extends Serializable {
*/
String getMergingStrategy();
- static String getRecordMergeStrategyId(RecordMergeMode mergeMode,
- String payloadClassName,
- String recordMergeStrategyId) {
+ static String getRecordMergeStrategyId(RecordMergeMode mergeMode, String
payloadClassName,
+ String recordMergeStrategyId,
HoodieTableVersion tableVersion) {
switch (mergeMode) {
case COMMIT_TIME_ORDERING:
return COMMIT_TIME_BASED_MERGE_STRATEGY_UUID;
@@ -196,13 +196,27 @@ public interface HoodieRecordMerger extends Serializable {
return EVENT_TIME_BASED_MERGE_STRATEGY_UUID;
case CUSTOM:
default:
- if (nonEmpty(recordMergeStrategyId)) {
- return recordMergeStrategyId;
- }
- if (nonEmpty(payloadClassName)) {
- return PAYLOAD_BASED_MERGE_STRATEGY_UUID;
- }
- return null;
+ return getCustomRecordMergeStrategyId(payloadClassName,
recordMergeStrategyId, tableVersion);
+ }
+ }
+
+ static String getCustomRecordMergeStrategyId(String payloadClassName, String
recordMergeStrategyId, HoodieTableVersion tableVersion) {
+ if (tableVersion.greaterThanOrEquals(HoodieTableVersion.EIGHT)) {
+ // For table version 8, we give preference to input
recordMergeStrategyId over payload based strategy
+ if (nonEmpty(recordMergeStrategyId)) {
+ return recordMergeStrategyId;
+ } else if (nonEmpty(payloadClassName)) {
+ return PAYLOAD_BASED_MERGE_STRATEGY_UUID;
+ }
+ return null;
+ } else {
+ // For table version 6, we give preference to payload based strategy
over input recordMergeStrategyId
+ if (nonEmpty(payloadClassName)) {
+ return PAYLOAD_BASED_MERGE_STRATEGY_UUID;
+ } else if (nonEmpty(recordMergeStrategyId)) {
+ return recordMergeStrategyId;
+ }
+ return null;
}
}
}
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java
b/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java
index 1ed8b68d80b..302645dc296 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java
@@ -827,7 +827,7 @@ public class HoodieTableConfig extends HoodieConfig {
inferredRecordMergeMode, payloadClassName);
// Inferring record merge strategy ID
inferredRecordMergeStrategyId =
HoodieRecordMerger.getRecordMergeStrategyId(
- inferredRecordMergeMode, inferredPayloadClassName,
recordMergeStrategyId);
+ inferredRecordMergeMode, inferredPayloadClassName,
recordMergeStrategyId, tableVersion);
// For custom merge mode, either payload class name or record merge
strategy ID must be configured
if (inferredRecordMergeMode == CUSTOM) {
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java
index a73398f2240..fa1048124fa 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java
@@ -26,7 +26,6 @@ import org.apache.hudi.common.engine.HoodieReaderContext;
import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.model.HoodieRecordMerger;
import org.apache.hudi.common.table.HoodieTableConfig;
-import org.apache.hudi.common.table.HoodieTableVersion;
import org.apache.hudi.common.util.LocalAvroSchemaCache;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.StringUtils;
@@ -216,7 +215,7 @@ public class FileGroupReaderSchemaHandler<T> {
cfg.getPayloadClass(),
cfg.getRecordMergeStrategyId(),
cfg.getPreCombineField(),
- HoodieTableVersion.current());
+ cfg.getTableVersion());
if (mergingConfigs.getLeft() == RecordMergeMode.CUSTOM) {
return recordMerger.get().getMandatoryFieldsForMerging(tableSchema, cfg,
props);
diff --git
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestFileGroupRecordBuffer.java
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestFileGroupRecordBuffer.java
index a4b912e33b5..db3846b459b 100644
---
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestFileGroupRecordBuffer.java
+++
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestFileGroupRecordBuffer.java
@@ -22,11 +22,15 @@ package org.apache.hudi.common.table.read;
import org.apache.hudi.common.config.RecordMergeMode;
import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.engine.HoodieReaderContext;
+import org.apache.hudi.common.model.DefaultHoodieRecordPayload;
import org.apache.hudi.common.model.DeleteRecord;
import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.model.HoodieRecordMerger;
+import org.apache.hudi.common.model.OverwriteNonDefaultsWithLatestAvroPayload;
+import org.apache.hudi.common.model.OverwriteWithLatestAvroPayload;
import org.apache.hudi.common.table.HoodieTableConfig;
import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.HoodieTableVersion;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.StringUtils;
import org.apache.hudi.common.util.collection.Pair;
@@ -110,33 +114,53 @@ class TestFileGroupRecordBuffer {
@ParameterizedTest
@CsvSource({
- "true, true, true, EVENT_TIME_ORDERING",
- "true, false, false, EVENT_TIME_ORDERING",
- "false, true, false, EVENT_TIME_ORDERING",
- "false, false, true, EVENT_TIME_ORDERING",
- "true, true, true, COMMIT_TIME_ORDERING",
- "true, false, false, COMMIT_TIME_ORDERING",
- "false, true, false, COMMIT_TIME_ORDERING",
- "false, false, true, COMMIT_TIME_ORDERING",
- "true, true, true, CUSTOM",
- "true, false, false, CUSTOM",
- "false, true, false, CUSTOM",
- "false, false, true, CUSTOM",
- "true, true, true,",
- "true, false, false,",
- "false, true, false,",
- "false, false, true,"
+ "true, true, true, EVENT_TIME_ORDERING, false, EIGHT,
eeb8d96f-b1e4-49fd-bbf8-28ac514178e5",
+ "true, false, false, EVENT_TIME_ORDERING, false, EIGHT,
eeb8d96f-b1e4-49fd-bbf8-28ac514178e5",
+ "false, true, false, EVENT_TIME_ORDERING, false, EIGHT,
eeb8d96f-b1e4-49fd-bbf8-28ac514178e5",
+ "false, false, true, EVENT_TIME_ORDERING, false, EIGHT,
eeb8d96f-b1e4-49fd-bbf8-28ac514178e5",
+ "true, true, true, COMMIT_TIME_ORDERING, false, EIGHT,
ce9acb64-bde0-424c-9b91-f6ebba25356d",
+ "true, false, false, COMMIT_TIME_ORDERING, false, EIGHT,
ce9acb64-bde0-424c-9b91-f6ebba25356d",
+ "false, true, false, COMMIT_TIME_ORDERING, false, EIGHT,
ce9acb64-bde0-424c-9b91-f6ebba25356d",
+ "false, false, true, COMMIT_TIME_ORDERING, false, EIGHT,
ce9acb64-bde0-424c-9b91-f6ebba25356d",
+ "true, true, true, CUSTOM, false, EIGHT,
00000000-0000-0000-0000-000000000000",
+ "true, false, false, CUSTOM, false, EIGHT,
00000000-0000-0000-0000-000000000000",
+ "false, true, false, CUSTOM, false, EIGHT,
00000000-0000-0000-0000-000000000000",
+ "false, false, true, CUSTOM, false, EIGHT,
00000000-0000-0000-0000-000000000000",
+ "true, true, true, , false, EIGHT, 00000000-0000-0000-0000-000000000000",
+ "true, false, false, , false, EIGHT,
00000000-0000-0000-0000-000000000000",
+ "false, true, false, , false, EIGHT,
00000000-0000-0000-0000-000000000000",
+ "false, false, true, , false, EIGHT,
00000000-0000-0000-0000-000000000000",
+ "true, true, true, EVENT_TIME_ORDERING, false, SIX,
eeb8d96f-b1e4-49fd-bbf8-28ac514178e5",
+ "true, false, false, EVENT_TIME_ORDERING, false, SIX,
eeb8d96f-b1e4-49fd-bbf8-28ac514178e5",
+ "false, true, false, EVENT_TIME_ORDERING, false, SIX,
eeb8d96f-b1e4-49fd-bbf8-28ac514178e5",
+ "false, false, true, EVENT_TIME_ORDERING, false, SIX,
eeb8d96f-b1e4-49fd-bbf8-28ac514178e5",
+ "true, true, true, COMMIT_TIME_ORDERING, false, SIX,
ce9acb64-bde0-424c-9b91-f6ebba25356d",
+ "true, false, false, COMMIT_TIME_ORDERING, false, SIX,
ce9acb64-bde0-424c-9b91-f6ebba25356d",
+ "false, true, false, COMMIT_TIME_ORDERING, false, SIX,
ce9acb64-bde0-424c-9b91-f6ebba25356d",
+ "false, false, true, COMMIT_TIME_ORDERING, false, SIX,
ce9acb64-bde0-424c-9b91-f6ebba25356d",
+ "true, true, true, CUSTOM, false, SIX,
00000000-0000-0000-0000-000000000000",
+ "true, false, false, CUSTOM, false, SIX,
00000000-0000-0000-0000-000000000000",
+ "false, true, false, CUSTOM, false, SIX,
00000000-0000-0000-0000-000000000000",
+ "false, false, true, CUSTOM, false, SIX,
00000000-0000-0000-0000-000000000000",
+ "true, true, true, , false, SIX, 00000000-0000-0000-0000-000000000000",
+ "true, false, false, , false, SIX, 00000000-0000-0000-0000-000000000000",
+ "false, true, false, , false, SIX, 00000000-0000-0000-0000-000000000000",
+ "false, false, true, , false, SIX, 00000000-0000-0000-0000-000000000000",
+ "true, true, true, COMMIT_TIME_ORDERING, true, SIX,
eeb8d96f-b1e4-49fd-bbf8-28ac514178e5", /// with table version 6, commit time
based merge mode can have event time based merge strategy id.
})
public void testSchemaForMandatoryFields(boolean setPrecombine,
boolean addHoodieIsDeleted,
boolean addCustomDeleteMarker,
- RecordMergeMode mergeMode) {
+ RecordMergeMode mergeMode,
+ boolean isProjectionCompatible,
+ HoodieTableVersion tableVersion,
+ String mergeStrategyId) {
HoodieReaderContext readerContext = mock(HoodieReaderContext.class);
when(readerContext.getHasBootstrapBaseFile()).thenReturn(false);
when(readerContext.getHasLogFiles()).thenReturn(true);
HoodieRecordMerger recordMerger = mock(HoodieRecordMerger.class);
when(readerContext.getRecordMerger()).thenReturn(Option.of(recordMerger));
- when(recordMerger.isProjectionCompatible()).thenReturn(false);
+
when(recordMerger.isProjectionCompatible()).thenReturn(isProjectionCompatible);
String preCombineField = "ts";
String customDeleteKey = "colC";
@@ -156,14 +180,26 @@ class TestFileGroupRecordBuffer {
when(tableConfig.getRecordMergeMode()).thenReturn(mergeMode);
when(tableConfig.populateMetaFields()).thenReturn(true);
when(tableConfig.getPreCombineField()).thenReturn(setPrecombine ?
preCombineField : StringUtils.EMPTY_STRING);
+ when(tableConfig.getTableVersion()).thenReturn(tableVersion);
+ if (tableConfig.getTableVersion() == HoodieTableVersion.SIX) {
+ if (mergeMode == RecordMergeMode.EVENT_TIME_ORDERING) {
+
when(tableConfig.getPayloadClass()).thenReturn(DefaultHoodieRecordPayload.class.getName());
+ } else if (mergeMode == RecordMergeMode.COMMIT_TIME_ORDERING) {
+
when(tableConfig.getPayloadClass()).thenReturn(OverwriteWithLatestAvroPayload.class.getName());
+ } else {
+
when(tableConfig.getPayloadClass()).thenReturn(OverwriteNonDefaultsWithLatestAvroPayload.class.getName());
+ }
+ }
+ if (mergeMode != null) {
+ when(tableConfig.getRecordMergeStrategyId()).thenReturn(mergeStrategyId);
+ }
TypedProperties props = new TypedProperties();
if (addCustomDeleteMarker) {
props.setProperty(DELETE_KEY, customDeleteKey);
props.setProperty(DELETE_MARKER, customDeleteValue);
}
- FileGroupReaderSchemaHandler fileGroupReaderSchemaHandler = new
FileGroupReaderSchemaHandler(readerContext,
- dataSchema, requestedSchema, Option.empty(), tableConfig, props);
+
List<String> expectedFields = new ArrayList();
expectedFields.add(HoodieRecord.RECORD_KEY_METADATA_FIELD);
expectedFields.add(HoodieRecord.PARTITION_PATH_METADATA_FIELD);
@@ -176,7 +212,11 @@ class TestFileGroupRecordBuffer {
if (addHoodieIsDeleted) {
expectedFields.add(HoodieRecord.HOODIE_IS_DELETED_FIELD);
}
- Schema expectedSchema = mergeMode == RecordMergeMode.CUSTOM ? dataSchema :
getSchema(expectedFields);
+ Schema expectedSchema = ((mergeMode == RecordMergeMode.CUSTOM) &&
!isProjectionCompatible) ? dataSchema : getSchema(expectedFields);
+ when(recordMerger.getMandatoryFieldsForMerging(dataSchema, tableConfig,
props)).thenReturn(expectedFields.toArray(new String[0]));
+
+ FileGroupReaderSchemaHandler fileGroupReaderSchemaHandler = new
FileGroupReaderSchemaHandler(readerContext,
+ dataSchema, requestedSchema, Option.empty(), tableConfig, props);
Schema actualSchema =
fileGroupReaderSchemaHandler.generateRequiredSchema();
assertEquals(expectedSchema, actualSchema);
assertEquals(addHoodieIsDeleted,
fileGroupReaderSchemaHandler.hasBuiltInDelete());
diff --git
a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieTableConfig.java
b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieTableConfig.java
index fbc778150cb..191c0e7ac84 100644
---
a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieTableConfig.java
+++
b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieTableConfig.java
@@ -377,7 +377,7 @@ public class TestHoodieTableConfig extends
HoodieCommonTestHarness {
arguments(CUSTOM, null, customStrategy, null,
"false", CUSTOM, defaultPayload, customStrategy),
arguments(CUSTOM, customPayload, customStrategy, null,
- "false", CUSTOM, customPayload, customStrategy),
+ "false", CUSTOM, customPayload, null),
//test legal configs that work but should not be used usually
arguments(CUSTOM, defaultPayload, customStrategy, null,
@@ -455,12 +455,10 @@ public class TestHoodieTableConfig extends
HoodieCommonTestHarness {
: Boolean.parseBoolean(shouldThrowString);
RecordMergeMode expectedMergeMode = outputMergeMode;
String expectedMergeStrategy = outputMergeStrategy;
- if (!shouldThrow && outputMergeMode == null) {
+ if (!shouldThrow && (outputMergeMode == null || outputMergeStrategy
== null)) {
expectedMergeMode =
tableVersion.greaterThanOrEquals(HoodieTableVersion.EIGHT)
? CUSTOM :
inferRecordMergeModeFromPayloadClass(outputPayloadClass);
- expectedMergeStrategy =
tableVersion.greaterThanOrEquals(HoodieTableVersion.EIGHT)
- ? PAYLOAD_BASED_MERGE_STRATEGY_UUID
- : getRecordMergeStrategyId(expectedMergeMode,
outputPayloadClass, null);
+ expectedMergeStrategy =
getRecordMergeStrategyId(expectedMergeMode, outputPayloadClass,
inputMergeStrategy, tableVersion);
}
if (shouldThrow) {
assertThrows(IllegalArgumentException.class,
diff --git
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamer.java
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamer.java
index 09cfe0a3bc1..c5fa5d7ffc8 100644
---
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamer.java
+++
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamer.java
@@ -44,6 +44,7 @@ import org.apache.hudi.common.table.HoodieTableVersion;
import org.apache.hudi.common.table.timeline.HoodieInstant;
import org.apache.hudi.common.util.ClusteringUtils;
import org.apache.hudi.common.util.CompactionUtils;
+import org.apache.hudi.common.util.ConfigUtils;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.StringUtils;
import org.apache.hudi.common.util.VisibleForTesting;
@@ -154,14 +155,14 @@ public class HoodieStreamer implements Serializable {
public HoodieStreamer(Config cfg, JavaSparkContext jssc, FileSystem fs,
Configuration conf,
Option<TypedProperties> propsOverride,
Option<SourceProfileSupplier> sourceProfileSupplier) throws IOException {
+ this.properties = combineProperties(cfg, propsOverride,
jssc.hadoopConfiguration());
Triple<RecordMergeMode, String, String> mergingConfigs =
HoodieTableConfig.inferCorrectMergingBehavior(
cfg.recordMergeMode, cfg.payloadClassName,
cfg.recordMergeStrategyId, cfg.sourceOrderingField,
- HoodieTableVersion.current());
+
HoodieTableVersion.fromVersionCode(ConfigUtils.getIntWithAltKeys(this.properties,
HoodieWriteConfig.WRITE_TABLE_VERSION)));
cfg.recordMergeMode = mergingConfigs.getLeft();
cfg.payloadClassName = mergingConfigs.getMiddle();
cfg.recordMergeStrategyId = mergingConfigs.getRight();
- this.properties = combineProperties(cfg, propsOverride,
jssc.hadoopConfiguration());
if (cfg.initialCheckpointProvider != null && cfg.checkpoint == null) {
InitialCheckPointProvider checkPointProvider =
UtilHelpers.createInitialCheckpointProvider(cfg.initialCheckpointProvider,
this.properties);