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);

Reply via email to