nsivabalan commented on code in PR #13519:
URL: https://github.com/apache/hudi/pull/13519#discussion_r2279728437


##########
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/NineToEightDowngradeHandler.java:
##########
@@ -20,24 +20,96 @@
 package org.apache.hudi.table.upgrade;
 
 import org.apache.hudi.common.config.ConfigProperty;
+import org.apache.hudi.common.config.RecordMergeMode;
 import org.apache.hudi.common.engine.HoodieEngineContext;
-import org.apache.hudi.common.util.collection.Pair;
+import org.apache.hudi.common.model.AWSDmsAvroPayload;
+import org.apache.hudi.common.model.DefaultHoodieRecordPayload;
+import org.apache.hudi.common.model.HoodieTableType;
+import org.apache.hudi.common.model.debezium.PostgresDebeziumAvroPayload;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.model.OverwriteWithLatestAvroPayload;
+import org.apache.hudi.common.util.StringUtils;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.table.HoodieTable;
 
-import java.util.Collections;
-import java.util.List;
+import java.util.HashMap;
+import java.util.HashSet;
 import java.util.Map;
+import java.util.Set;
 
+import static 
org.apache.hudi.common.model.DefaultHoodieRecordPayload.DELETE_KEY;
+import static 
org.apache.hudi.common.model.DefaultHoodieRecordPayload.DELETE_MARKER;
+import static 
org.apache.hudi.common.model.HoodieRecordMerger.PAYLOAD_BASED_MERGE_STRATEGY_UUID;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.LEGACY_PAYLOAD_CLASS_NAME;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PARTIAL_UPDATE_CUSTOM_MARKER;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PARTIAL_UPDATE_MODE;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PAYLOAD_CLASS_NAME;
+import static org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_MODE;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_PROPERTY_PREFIX;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_STRATEGY_ID;
+import static 
org.apache.hudi.table.upgrade.UpgradeDowngradeUtils.PAYLOAD_CLASSES_TO_HANDLE;
+
+/**
+ * Version 8 is the placeholder version from 1.0.0 to 1.0.2.
+ * Version 9 is the placeholder version >= 1.1.0.
+ * The major change introduced in version 9 is two table configurations for 
payload deprecation.
+ * The main downgrade logic:
+ *   for all tables:
+ *     remove hoodie.table.partial.update.mode from table_configs
+ *   for table with payload class defined in RFC-97,
+ *     remove hoodie.legacy.payload.class from table_configs
+ *     set hoodie.compaction.payload.class=payload
+ *     set hoodie.record.merge.mode=CUSTOM
+ *     set hoodie.record.merge.strategy.id accordingly
+ */
 public class NineToEightDowngradeHandler implements DowngradeHandler {
   @Override
-  public Pair<Map<ConfigProperty, String>, List<ConfigProperty>> 
downgrade(HoodieWriteConfig config,
-                                                                           
HoodieEngineContext context,
-                                                                           
String instantTime,
-                                                                           
SupportsUpgradeDowngrade upgradeDowngradeHelper) {
+  public UpgradeDowngrade.TableConfigChangeSet downgrade(HoodieWriteConfig 
config,
+                                                         HoodieEngineContext 
context,
+                                                         String instantTime,
+                                                         
SupportsUpgradeDowngrade upgradeDowngradeHelper) {
     final HoodieTable table = upgradeDowngradeHelper.getTable(config, context);
+    HoodieTableMetaClient metaClient = table.getMetaClient();
+    // Handle secondary index.
     UpgradeDowngradeUtils.dropNonV1SecondaryIndexPartitions(
         config, context, table, upgradeDowngradeHelper, "downgrading from 
table version nine to eight");
-    return Pair.of(Collections.emptyMap(), Collections.emptyList());
+    // Update table properties.
+    Set<ConfigProperty> propertiesToRemove = new HashSet<>();
+    Map<ConfigProperty, String> propertiesToAdd = new HashMap<>();
+    // TODO: handle COW table after write path is done.
+    if (metaClient.getTableConfig().getTableType() == 
HoodieTableType.MERGE_ON_READ) {
+      updateMergeRelatedConfigs(propertiesToAdd, propertiesToRemove, 
metaClient);
+    }
+    return new UpgradeDowngrade.TableConfigChangeSet(propertiesToAdd, 
propertiesToRemove);
+  }
+
+  private void updateMergeRelatedConfigs(Map<ConfigProperty, String> 
propertiesToAdd,
+                                         Set<ConfigProperty> 
propertiesToRemove,
+                                         HoodieTableMetaClient metaClient) {
+    // Update table properties.
+    propertiesToRemove.add(PARTIAL_UPDATE_MODE);
+    // For specified payload classes, add strategy id and custom merge mode.
+    HoodieTableConfig tableConfig = metaClient.getTableConfig();
+    String payloadClass = tableConfig.getLegacyPayloadClass();
+    if (!StringUtils.isNullOrEmpty(payloadClass) && 
(PAYLOAD_CLASSES_TO_HANDLE.contains(payloadClass))) {
+      propertiesToRemove.add(LEGACY_PAYLOAD_CLASS_NAME);
+      propertiesToAdd.put(PAYLOAD_CLASS_NAME, payloadClass);
+      if (!payloadClass.equals(OverwriteWithLatestAvroPayload.class.getName())

Review Comment:
   @linliu-code : don't we need to fix merge strategy Id for 
OverwriteWithLatestAvroPayload and DefaultHoodieRecordPayload ? 



##########
hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/upgrade/TestNineToEightDowngradeHandler.java:
##########
@@ -19,75 +19,202 @@
 
 package org.apache.hudi.table.upgrade;
 
-import org.apache.hudi.client.BaseHoodieWriteClient;
+import org.apache.hudi.common.config.RecordMergeMode;
 import org.apache.hudi.common.engine.HoodieEngineContext;
-import org.apache.hudi.common.model.HoodieIndexDefinition;
-import org.apache.hudi.common.model.HoodieIndexMetadata;
+import org.apache.hudi.common.model.AWSDmsAvroPayload;
+import org.apache.hudi.common.model.DefaultHoodieRecordPayload;
+import org.apache.hudi.common.model.HoodieTableType;
+import org.apache.hudi.common.model.OverwriteNonDefaultsWithLatestAvroPayload;
+import org.apache.hudi.common.model.OverwriteWithLatestAvroPayload;
+import org.apache.hudi.common.model.PartialUpdateAvroPayload;
+import org.apache.hudi.common.model.debezium.MySqlDebeziumAvroPayload;
+import org.apache.hudi.common.model.debezium.PostgresDebeziumAvroPayload;
 import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.client.BaseHoodieWriteClient;
+import org.apache.hudi.common.model.HoodieIndexDefinition;
+import org.apache.hudi.common.model.HoodieIndexMetadata;
 import org.apache.hudi.common.table.HoodieTableVersion;
 import org.apache.hudi.common.util.Option;
-import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.metadata.HoodieIndexVersion;
 import org.apache.hudi.metadata.MetadataPartitionType;
 import org.apache.hudi.table.HoodieTable;
 
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.extension.ExtendWith;
-import org.mockito.ArgumentCaptor;
-import org.mockito.Mock;
-import org.mockito.MockedStatic;
-import org.mockito.junit.jupiter.MockitoExtension;
-import org.mockito.junit.jupiter.MockitoSettings;
-import org.mockito.quality.Strictness;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
 
 import java.util.Arrays;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.stream.Stream;
 
+import static 
org.apache.hudi.common.model.HoodieRecordMerger.PAYLOAD_BASED_MERGE_STRATEGY_UUID;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.LEGACY_PAYLOAD_CLASS_NAME;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PARTIAL_UPDATE_MODE;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PAYLOAD_CLASS_NAME;
+import static org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_MODE;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_STRATEGY_ID;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyBoolean;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.mockStatic;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
-@ExtendWith(MockitoExtension.class)
-@MockitoSettings(strictness = Strictness.LENIENT)
-class TestNineToEightDowngradeHandler {
+import org.mockito.ArgumentCaptor;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
 
-  @Mock
-  private HoodieWriteConfig config;
-  @Mock
-  private HoodieTableConfig tblConfig;
-  @Mock
-  private HoodieTableMetaClient metaClient;
-  @Mock
-  private HoodieEngineContext context;
-  @Mock
-  private SupportsUpgradeDowngrade upgradeDowngradeHelper;
-  @Mock
-  private HoodieTable table;
-
-  private NineToEightDowngradeHandler downgradeHandler;
+class TestNineToEightDowngradeHandler {
+  private final NineToEightDowngradeHandler handler = new 
NineToEightDowngradeHandler();
+  private final HoodieWriteConfig config = mock(HoodieWriteConfig.class);
+  private final HoodieEngineContext context = mock(HoodieEngineContext.class);
+  private final HoodieTable table = mock(HoodieTable.class);
+  private HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+  private HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+  private SupportsUpgradeDowngrade upgradeDowngradeHelper = 
mock(SupportsUpgradeDowngrade.class);
 
   @BeforeEach
-  void setUp() {
-    downgradeHandler = new NineToEightDowngradeHandler();
-    when(upgradeDowngradeHelper.getTable(config, context)).thenReturn(table);
+  public void setUp() {
+    when(upgradeDowngradeHelper.getTable(any(), any())).thenReturn(table);
+    when(table.getMetaClient()).thenReturn(metaClient);
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+  }
+
+  static Stream<Arguments> payloadClassTestCases() {
+    return Stream.of(
+        // AWSDmsAvroPayload - requires RECORD_MERGE_MODE and 
RECORD_MERGE_STRATEGY_ID
+        Arguments.of(
+            AWSDmsAvroPayload.class.getName(),
+            4, // propertiesToRemove size

Review Comment:
   can we assert the properties to remove as well in addition to just the size 
check



##########
hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/upgrade/TestNineToEightDowngradeHandler.java:
##########
@@ -19,75 +19,202 @@
 
 package org.apache.hudi.table.upgrade;
 
-import org.apache.hudi.client.BaseHoodieWriteClient;
+import org.apache.hudi.common.config.RecordMergeMode;
 import org.apache.hudi.common.engine.HoodieEngineContext;
-import org.apache.hudi.common.model.HoodieIndexDefinition;
-import org.apache.hudi.common.model.HoodieIndexMetadata;
+import org.apache.hudi.common.model.AWSDmsAvroPayload;
+import org.apache.hudi.common.model.DefaultHoodieRecordPayload;
+import org.apache.hudi.common.model.HoodieTableType;
+import org.apache.hudi.common.model.OverwriteNonDefaultsWithLatestAvroPayload;
+import org.apache.hudi.common.model.OverwriteWithLatestAvroPayload;
+import org.apache.hudi.common.model.PartialUpdateAvroPayload;
+import org.apache.hudi.common.model.debezium.MySqlDebeziumAvroPayload;
+import org.apache.hudi.common.model.debezium.PostgresDebeziumAvroPayload;
 import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.client.BaseHoodieWriteClient;
+import org.apache.hudi.common.model.HoodieIndexDefinition;
+import org.apache.hudi.common.model.HoodieIndexMetadata;
 import org.apache.hudi.common.table.HoodieTableVersion;
 import org.apache.hudi.common.util.Option;
-import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.metadata.HoodieIndexVersion;
 import org.apache.hudi.metadata.MetadataPartitionType;
 import org.apache.hudi.table.HoodieTable;
 
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.extension.ExtendWith;
-import org.mockito.ArgumentCaptor;
-import org.mockito.Mock;
-import org.mockito.MockedStatic;
-import org.mockito.junit.jupiter.MockitoExtension;
-import org.mockito.junit.jupiter.MockitoSettings;
-import org.mockito.quality.Strictness;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
 
 import java.util.Arrays;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.stream.Stream;
 
+import static 
org.apache.hudi.common.model.HoodieRecordMerger.PAYLOAD_BASED_MERGE_STRATEGY_UUID;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.LEGACY_PAYLOAD_CLASS_NAME;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PARTIAL_UPDATE_MODE;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PAYLOAD_CLASS_NAME;
+import static org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_MODE;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_STRATEGY_ID;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyBoolean;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.mockStatic;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
-@ExtendWith(MockitoExtension.class)
-@MockitoSettings(strictness = Strictness.LENIENT)
-class TestNineToEightDowngradeHandler {
+import org.mockito.ArgumentCaptor;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
 
-  @Mock
-  private HoodieWriteConfig config;
-  @Mock
-  private HoodieTableConfig tblConfig;
-  @Mock
-  private HoodieTableMetaClient metaClient;
-  @Mock
-  private HoodieEngineContext context;
-  @Mock
-  private SupportsUpgradeDowngrade upgradeDowngradeHelper;
-  @Mock
-  private HoodieTable table;
-
-  private NineToEightDowngradeHandler downgradeHandler;
+class TestNineToEightDowngradeHandler {
+  private final NineToEightDowngradeHandler handler = new 
NineToEightDowngradeHandler();
+  private final HoodieWriteConfig config = mock(HoodieWriteConfig.class);
+  private final HoodieEngineContext context = mock(HoodieEngineContext.class);
+  private final HoodieTable table = mock(HoodieTable.class);
+  private HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+  private HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+  private SupportsUpgradeDowngrade upgradeDowngradeHelper = 
mock(SupportsUpgradeDowngrade.class);
 
   @BeforeEach
-  void setUp() {
-    downgradeHandler = new NineToEightDowngradeHandler();
-    when(upgradeDowngradeHelper.getTable(config, context)).thenReturn(table);
+  public void setUp() {
+    when(upgradeDowngradeHelper.getTable(any(), any())).thenReturn(table);
+    when(table.getMetaClient()).thenReturn(metaClient);
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+  }
+
+  static Stream<Arguments> payloadClassTestCases() {
+    return Stream.of(
+        // AWSDmsAvroPayload - requires RECORD_MERGE_MODE and 
RECORD_MERGE_STRATEGY_ID
+        Arguments.of(
+            AWSDmsAvroPayload.class.getName(),
+            4, // propertiesToRemove size
+            3, // propertiesToAdd size
+            true, // hasRecordMergeMode
+            true, // hasRecordMergeStrategyId
+            "AWSDmsAvroPayload"
+        ),
+        // OverwriteNonDefaultsWithLatestAvroPayload - requires 
RECORD_MERGE_MODE and RECORD_MERGE_STRATEGY_ID
+        Arguments.of(
+            OverwriteNonDefaultsWithLatestAvroPayload.class.getName(),
+            2,
+            3,
+            true,
+            true,
+            "OverwriteNonDefaultsWithLatestAvroPayload"
+        ),
+        // PartialUpdateAvroPayload - requires RECORD_MERGE_MODE and 
RECORD_MERGE_STRATEGY_ID
+        Arguments.of(
+            PartialUpdateAvroPayload.class.getName(),
+            2,
+            3,
+            true,
+            true,
+            "PartialUpdateAvroPayload"
+        ),
+        // MySqlDebeziumAvroPayload - requires RECORD_MERGE_MODE and 
RECORD_MERGE_STRATEGY_ID
+        Arguments.of(
+            MySqlDebeziumAvroPayload.class.getName(),
+            2,
+            3,
+            true,
+            true,
+            "MySqlDebeziumAvroPayload"
+        ),
+        // PostgresDebeziumAvroPayload - requires RECORD_MERGE_MODE and 
RECORD_MERGE_STRATEGY_ID
+        Arguments.of(
+            PostgresDebeziumAvroPayload.class.getName(),
+            3,
+            3,
+            true,
+            true,
+            "PostgresDebeziumAvroPayload"
+        ),
+        // OverwriteWithLatestAvroPayload - only requires PAYLOAD_CLASS_NAME
+        Arguments.of(
+            OverwriteWithLatestAvroPayload.class.getName(),
+            2,
+            1,
+            false,
+            false,
+            "OverwriteWithLatestAvroPayload"
+        ),
+        // DefaultHoodieRecordPayload - only requires PAYLOAD_CLASS_NAME
+        Arguments.of(

Review Comment:
   EventTimeAvroPayload? 



##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPayloadDeprecationFlow.scala:
##########
@@ -110,46 +127,104 @@ class TestPayloadDeprecationFlow extends 
SparkClientFunctionalTestHarness {
           option(HoodieCompactionConfig.INLINE_COMPACT.key(), "false").
           
option(HoodieCompactionConfig.INLINE_COMPACT_NUM_DELTA_COMMITS.key(), "1").
           option(HoodieTableConfig.PAYLOAD_CLASS_NAME.key(),
-            classOf[MySqlDebeziumAvroPayload].getName). // Position is 
important.
+            classOf[MySqlDebeziumAvroPayload].getName).
           mode(SaveMode.Append).
           save(basePath)
       }
     }
+
     // 5. Validate.
+    // Validate table configs.
+    tableConfig = metaClient.getTableConfig
+    expectedConfigs.foreach { case (key, expectedValue) =>
+      if (expectedValue != null) {
+        assertEquals(expectedValue, tableConfig.getString(key), s"Config $key 
should be $expectedValue")
+      } else {
+        assertFalse(tableConfig.contains(key), s"Config $key should not be 
present")
+      }
+    }
+    // Validate snapshot query.
     val df = spark.read.format("hudi").load(basePath)
-    val finalDf = df.select("ts", "key", "rider", "driver", "fare", 
"Op").sort("key")
+    val finalDf = df.select("ts", "_event_lsn", "rider", "driver", "fare", 
"Op", "_event_seq").sort("_event_lsn")
+    val expectedData = getExpectedResultForSnapshotQuery(payloadClazz)
+    val expectedDf = 
spark.createDataFrame(spark.sparkContext.parallelize(expectedData)).toDF(columns:
 _*).sort("_event_lsn")
+    expectedDf.show(false)
+    finalDf.show(false)
+    assertTrue(expectedDf.except(finalDf).isEmpty && 
finalDf.except(expectedDf).isEmpty)
+    // Validate time travel query.
+    val timeTravelDf = spark.read.format("hudi")
+      .option("as.of.instant", firstUpdateInstantTime).load(basePath)
+      .select("ts", "_event_lsn", "rider", "driver", "fare", "Op", 
"_event_seq").sort("_event_lsn")
+    timeTravelDf.show(false)
+    val expectedTimeTravelData = 
getExpectedResultForTimeTravelQuery(payloadClazz)
+    val expectedTimeTravelDf = spark.createDataFrame(
+      spark.sparkContext.parallelize(expectedTimeTravelData)).toDF(columns: 
_*).sort("_event_lsn")
+    expectedTimeTravelDf.show(false)
+    timeTravelDf.show(false)
+    assertTrue(
+      expectedTimeTravelDf.except(timeTravelDf).isEmpty
+      && timeTravelDf.except(expectedTimeTravelDf).isEmpty)
+  }
+
+  def getWriteConfig(hudiOpts: Map[String, String]): HoodieWriteConfig = {
+    val props = TypedProperties.fromMap(hudiOpts.asJava)
+    HoodieWriteConfig.newBuilder()
+      .withProps(props)
+      .withPath(basePath())
+      .build()
+  }
 
-    val expectedData = if 
(!payloadClazz.equals(classOf[AWSDmsAvroPayload].getName)) {
-      if 
(HoodieTableConfig.EVENT_TIME_ORDERING_PAYLOADS.contains(payloadClazz)) {
+  def getExpectedResultForSnapshotQuery(payloadClazz: String): Seq[(Int, Long, 
String, String, Double, String, String)] = {
+    if (!payloadClazz.equals(classOf[AWSDmsAvroPayload].getName)) {
+      if (payloadClazz.equals(classOf[PartialUpdateAvroPayload].getName)
+        || payloadClazz.equals(classOf[EventTimeAvroPayload].getName)
+        || payloadClazz.equals(classOf[DefaultHoodieRecordPayload].getName)
+        || payloadClazz.equals(classOf[PostgresDebeziumAvroPayload].getName)
+        || payloadClazz.equals(classOf[MySqlDebeziumAvroPayload].getName)) {
         Seq(
-          (11, "1", "rider-X", "driver-X", 19.10, "D"),
-          (11, "2", "rider-Y", "driver-Y", 27.70, "u"),
-          (12, "3", "rider-CC", "driver-CC", 33.90, "i"),
-          (10, "4", "rider-D", "driver-D", 34.15, "i"),
-          (12, "5", "rider-EE", "driver-EE", 17.85, "i"))
+          (11, 1, "rider-X", "driver-X", 19.10, "D", "11.1"),
+          (11, 2, "rider-Y", "driver-Y", 27.70, "u", "11.1"),
+          (12, 3, "rider-CC", "driver-CC", 33.90, "i", "12.1"),
+          (10, 4, "rider-D", "driver-D", 34.15, "i", "10.1"),
+          (12, 5, "rider-EE", "driver-EE", 17.85, "i", "12.1"))
       } else {
         Seq(
-          (11, "1", "rider-X", "driver-X", 19.10, "D"),
-          (11, "2", "rider-Y", "driver-Y", 27.70, "u"),
-          (12, "3", "rider-CC", "driver-CC", 33.90, "i"),
-          (9, "4", "rider-DD", "driver-DD", 34.15, "i"),
-          (12, "5", "rider-EE", "driver-EE", 17.85, "i"))
+          (11, 1, "rider-X", "driver-X", 19.10, "D", "11.1"),
+          (11, 2, "rider-Y", "driver-Y", 27.70, "u", "11.1"),
+          (12, 3, "rider-CC", "driver-CC", 33.90, "i", "12.1"),
+          (9, 4, "rider-DD", "driver-DD", 34.15, "i", "9.1"),
+          (12, 5, "rider-EE", "driver-EE", 17.85, "i", "12.1"))
       }
     } else {
       Seq(
-        (11, "2", "rider-Y", "driver-Y", 27.70, "u"),
-        (12, "3", "rider-CC", "driver-CC", 33.90, "i"),
-        (9, "4", "rider-DD", "driver-DD", 34.15, "i"),
-        (12, "5", "rider-EE", "driver-EE", 17.85, "i"))
+        (11, 2, "rider-Y", "driver-Y", 27.70, "u", "11.1"),
+        (12, 3, "rider-CC", "driver-CC", 33.90, "i", "12.1"),
+        (9, 4, "rider-DD", "driver-DD", 34.15, "i", "9.1"),
+        (12, 5, "rider-EE", "driver-EE", 17.85, "i", "12.1"))
+    }
+  }
+
+  def getExpectedResultForTimeTravelQuery(payloadClazz: String):
+  Seq[(Int, Long, String, String, Double, String, String)] = {
+    if (!payloadClazz.equals(classOf[AWSDmsAvroPayload].getName)) {
+      Seq(
+        (11, 1, "rider-X", "driver-X", 19.10, "D", "11.1"),
+        (11, 2, "rider-Y", "driver-Y", 27.70, "u", "11.1"),
+        (10, 3, "rider-C", "driver-C", 33.90, "i", "10.1"),
+        (10, 4, "rider-D", "driver-D", 34.15, "i", "10.1"),
+        (10, 5, "rider-E", "driver-E", 17.85, "i", "10.1"))
+    } else {
+      Seq(
+        (11, 2, "rider-Y", "driver-Y", 27.70, "u", "11.1"),
+        (10, 3, "rider-C", "driver-C", 33.90, "i", "10.1"),
+        (10, 4, "rider-D", "driver-D", 34.15, "i", "10.1"),
+        (10, 5, "rider-E", "driver-E", 17.85, "i", "10.1"))
     }
-    val expectedDf = spark.createDataFrame(
-      spark.sparkContext.parallelize(expectedData)).toDF(columns: 
_*).sort("key")
-    assertTrue(
-      expectedDf.except(finalDf).isEmpty && finalDf.except(expectedDf).isEmpty)
   }
 }
 
 // TODO: Add COPY_ON_WRITE table type tests when write path is updated 
accordingly.
+// TODO: Add Test for MySqlDebeziumAvroPayload.

Review Comment:
   I already see we are upgrading this payload as well. So, whats pending to 
enable this for tests? 



##########
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/EightToNineUpgradeHandler.java:
##########
@@ -19,41 +19,188 @@
 package org.apache.hudi.table.upgrade;
 
 import org.apache.hudi.common.config.ConfigProperty;
+import org.apache.hudi.common.config.RecordMergeMode;
 import org.apache.hudi.common.engine.HoodieEngineContext;
+import org.apache.hudi.common.model.AWSDmsAvroPayload;
+import org.apache.hudi.common.model.EventTimeAvroPayload;
 import org.apache.hudi.common.model.HoodieIndexMetadata;
+import org.apache.hudi.common.model.HoodieTableType;
+import org.apache.hudi.common.model.OverwriteNonDefaultsWithLatestAvroPayload;
+import org.apache.hudi.common.model.PartialUpdateAvroPayload;
+import org.apache.hudi.common.model.debezium.MySqlDebeziumAvroPayload;
+import org.apache.hudi.common.model.debezium.PostgresDebeziumAvroPayload;
+import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.PartialUpdateMode;
 import org.apache.hudi.common.util.Option;
+import org.apache.hudi.common.util.StringUtils;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.metadata.HoodieIndexVersion;
 import org.apache.hudi.table.HoodieTable;
 
+import java.util.Arrays;
 import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
 import java.util.Map;
+import java.util.Set;
 
+import static 
org.apache.hudi.common.model.DefaultHoodieRecordPayload.DELETE_KEY;
+import static 
org.apache.hudi.common.model.DefaultHoodieRecordPayload.DELETE_MARKER;
+import static 
org.apache.hudi.common.model.HoodieRecordMerger.COMMIT_TIME_BASED_MERGE_STRATEGY_UUID;
+import static 
org.apache.hudi.common.model.HoodieRecordMerger.CUSTOM_MERGE_STRATEGY_UUID;
+import static 
org.apache.hudi.common.model.HoodieRecordMerger.EVENT_TIME_BASED_MERGE_STRATEGY_UUID;
+import static 
org.apache.hudi.common.model.HoodieRecordMerger.PAYLOAD_BASED_MERGE_STRATEGY_UUID;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.DEBEZIUM_UNAVAILABLE_VALUE;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.LEGACY_PAYLOAD_CLASS_NAME;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PARTIAL_UPDATE_CUSTOM_MARKER;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PARTIAL_UPDATE_MODE;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PAYLOAD_CLASS_NAME;
+import static org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_MODE;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_PROPERTY_PREFIX;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_STRATEGY_ID;
+import static 
org.apache.hudi.table.upgrade.UpgradeDowngradeUtils.PAYLOAD_CLASSES_TO_HANDLE;
+
+/**
+ * Version 8 presents Hudi version from 1.0.0 to 1.0.2.
+ * Version 9 presents Hudi version >= 1.1.0.
+ * Major upgrade logic:
+ *  Deprecate a given set of payload classes to prefer merge mode. That is,
+ *   for table with payload class defined in RFC-97,
+ *     remove hoodie.compaction.payload.class from table_configs
+ *     add hoodie.legacy.payload.class=payload to table_configs
+ *     set hoodie.table.partial.update.mode based on RFC-97
+ *     set hoodie.table.merge.properties based on RFC-97
+ *     set hoodie.record.merge.mode based on RFC-97
+ *     set hoodie.record.merge.strategy.id based on RFC-97
+ *   for table with event_time/commit_time merge mode,
+ *     set hoodie.table.partial.update.mode to default value
+ *     set hoodie.table.merge.properties to default value
+ *   for table with custom merger or payload,
+ *     set hoodie.table.partial.update.mode to default value
+ *     set hoodie.table.merge.properties to default value
+ */
 public class EightToNineUpgradeHandler implements UpgradeHandler {
+  private static final Set<String> PAYLOADS_UPGRADE_TO_EVENT_TIME_MERGE_MODE = 
new HashSet<>(Arrays.asList(
+      EventTimeAvroPayload.class.getName(),
+      MySqlDebeziumAvroPayload.class.getName(),
+      PartialUpdateAvroPayload.class.getName(),
+      PostgresDebeziumAvroPayload.class.getName()));
+  private static final Set<String> PAYLOADS_UPGRADE_TO_COMMIT_TIME_MERGE_MODE 
= new HashSet<>(Arrays.asList(
+      AWSDmsAvroPayload.class.getName(),
+      OverwriteNonDefaultsWithLatestAvroPayload.class.getName()));
+  public static final Set<String> BUILTIN_MERGE_STRATEGIES = 
Collections.unmodifiableSet(
+      new HashSet<>(Arrays.asList(
+          COMMIT_TIME_BASED_MERGE_STRATEGY_UUID,
+          CUSTOM_MERGE_STRATEGY_UUID,
+          EVENT_TIME_BASED_MERGE_STRATEGY_UUID,
+          PAYLOAD_BASED_MERGE_STRATEGY_UUID)));
 
   @Override
-  public Map<ConfigProperty, String> upgrade(HoodieWriteConfig config,
-                                             HoodieEngineContext context,
-                                             String instantTime,
-                                             SupportsUpgradeDowngrade 
upgradeDowngradeHelper) {
+  public UpgradeDowngrade.TableConfigChangeSet upgrade(HoodieWriteConfig 
config,
+                                                       HoodieEngineContext 
context,
+                                                       String instantTime,
+                                                       
SupportsUpgradeDowngrade upgradeDowngradeHelper) {
     final HoodieTable table = upgradeDowngradeHelper.getTable(config, context);
-
+    Map<ConfigProperty, String> tablePropsToAdd = new HashMap<>();
     HoodieTableMetaClient metaClient = table.getMetaClient();
-
+    HoodieTableConfig tableConfig = metaClient.getTableConfig();
     // Populate missing index versions indexes
     Option<HoodieIndexMetadata> indexMetadataOpt = 
metaClient.getIndexMetadata();
     if (indexMetadataOpt.isPresent()) {
       populateIndexVersionIfMissing(indexMetadataOpt);
-
       // Write the updated index metadata back to storage
       HoodieTableMetaClient.writeIndexMetadataToStorage(
           metaClient.getStorage(),
           metaClient.getIndexDefinitionPath(),
           indexMetadataOpt.get(),
           metaClient.getTableConfig().getTableVersion());
     }
-    return Collections.emptyMap();
+    Set<ConfigProperty> tablePropsToRemove = new HashSet<>();
+    // TODO: make it work for COW after write path is ready.

Review Comment:
   wrt lines 113 to 116. 
   We are not idempotent here. 
   for eg, we can crash after upgrading the index defn, but before doing 
anything else. i.e actually upgrading the table. 
   
   CC @Davis-Zhang-Onehouse @linliu-code @rahil-c @yihua  : We might need to 
fix this properly to handle crashes and retries. 
   can we file a follow up blocking ticket. we can decide who can pick it up. 



##########
hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/upgrade/TestEightToNineUpgradeHandler.java:
##########
@@ -50,70 +59,187 @@
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
-
+import java.util.Set;
+import java.util.stream.Stream;
+
+import static 
org.apache.hudi.common.config.RecordMergeMode.COMMIT_TIME_ORDERING;
+import static 
org.apache.hudi.common.config.RecordMergeMode.EVENT_TIME_ORDERING;
+import static 
org.apache.hudi.common.model.DefaultHoodieRecordPayload.DELETE_KEY;
+import static 
org.apache.hudi.common.model.DefaultHoodieRecordPayload.DELETE_MARKER;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.DEBEZIUM_UNAVAILABLE_VALUE;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.LEGACY_PAYLOAD_CLASS_NAME;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PARTIAL_UPDATE_CUSTOM_MARKER;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PARTIAL_UPDATE_MODE;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.PAYLOAD_CLASS_NAME;
+import static org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_MODE;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_PROPERTY_PREFIX;
+import static org.apache.hudi.common.table.PartialUpdateMode.IGNORE_MARKERS;
 import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyBoolean;
 import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.mockStatic;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
-@ExtendWith(MockitoExtension.class)
-@MockitoSettings(strictness = Strictness.LENIENT)
 class TestEightToNineUpgradeHandler {
-
-  @TempDir
-  private Path tempDir;
-
-  @Mock
-  private HoodieWriteConfig config;
-  @Mock
-  private HoodieEngineContext context;
-  @Mock
-  private SupportsUpgradeDowngrade upgradeDowngradeHelper;
-  @Mock
-  private HoodieTable table;
-  @Mock
-  private HoodieTableMetaClient metaClient;
-  @Mock
-  private HoodieTableConfig tableConfig;
-  @Mock
-  private HoodieStorage storage;
-
-  private EightToNineUpgradeHandler upgradeHandler;
+  private final EightToNineUpgradeHandler handler = new 
EightToNineUpgradeHandler();
+  private final HoodieStorage storage = mock(HoodieStorage.class);
+  private final HoodieEngineContext context = mock(HoodieEngineContext.class);
+  private final HoodieTable table = mock(HoodieTable.class);
+  private final HoodieTableMetaClient metaClient = 
mock(HoodieTableMetaClient.class);
+  private final HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+  private final SupportsUpgradeDowngrade upgradeDowngradeHelper =
+      mock(SupportsUpgradeDowngrade.class);
+  private final HoodieWriteConfig config = mock(HoodieWriteConfig.class);
+  private static final Map<ConfigProperty, String> DEFAULT_CONFIG_UPDATED = 
Collections.emptyMap();
+  private static final Set<ConfigProperty> DEFAULT_CONFIG_REMOVED = 
Collections.emptySet();
+  private static final UpgradeDowngrade.TableConfigChangeSet 
DEFAULT_UPGRADE_RESULT =
+      new UpgradeDowngrade.TableConfigChangeSet(DEFAULT_CONFIG_UPDATED, 
DEFAULT_CONFIG_REMOVED);
   private static final String INSTANT_TIME = "20231201120000";
   private StoragePath indexDefPath;
+  @TempDir
+  private Path tempDir;
 
   @BeforeEach
-  void setUp() throws IOException {
-    upgradeHandler = new EightToNineUpgradeHandler();
-    
+  public void setUp() throws IOException {
+    when(upgradeDowngradeHelper.getTable(any(), any())).thenReturn(table);
+    when(table.getMetaClient()).thenReturn(metaClient);
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+    when(config.autoUpgrade()).thenReturn(true);
+
     // Setup common mocks
     when(upgradeDowngradeHelper.getTable(config, context)).thenReturn(table);
     when(table.getMetaClient()).thenReturn(metaClient);
     when(metaClient.getTableConfig()).thenReturn(tableConfig);
     when(metaClient.getStorage()).thenReturn(storage);
     when(tableConfig.getTableVersion()).thenReturn(HoodieTableVersion.EIGHT);
-    
+
     // Use a temp file for index definition path
     indexDefPath = new StoragePath(tempDir.resolve("index.json").toString());
     
when(metaClient.getIndexDefinitionPath()).thenReturn(indexDefPath.toString());
-    
+
     // Mock storage methods for file creation
     when(storage.exists(any(StoragePath.class))).thenReturn(false);
     when(storage.createNewFile(any(StoragePath.class))).thenReturn(true);
-    
+
     // Mock create method to capture written content
     ByteArrayOutputStream capturedContent = new ByteArrayOutputStream();
     when(storage.create(any(StoragePath.class), 
anyBoolean())).thenReturn(capturedContent);
-    
+
     // Mock autoUpgrade to return true
     when(config.autoUpgrade()).thenReturn(true);
   }
 
+  static Stream<Arguments> payloadClassTestCases() {
+    return Stream.of(
+        Arguments.of(
+            DefaultHoodieRecordPayload.class.getName(),
+            "",
+            null,
+            PartialUpdateMode.NONE.name(),
+            "DefaultHoodieRecordPayload"
+        ),
+        Arguments.of(
+            OverwriteWithLatestAvroPayload.class.getName(),
+            "",
+            null,
+            PartialUpdateMode.NONE.name(),
+            "OverwriteWithLatestAvroPayload"
+        ),
+        Arguments.of(
+            AWSDmsAvroPayload.class.getName(),
+            RECORD_MERGE_PROPERTY_PREFIX + DELETE_KEY + "=Op,"
+                + RECORD_MERGE_PROPERTY_PREFIX + DELETE_MARKER + "=D", // 
mergeProperties
+            COMMIT_TIME_ORDERING.name(),
+            PartialUpdateMode.NONE.name(),
+            "AWSDmsAvroPayload"
+        ),
+        Arguments.of(
+            PostgresDebeziumAvroPayload.class.getName(),
+            RECORD_MERGE_PROPERTY_PREFIX + PARTIAL_UPDATE_CUSTOM_MARKER
+                + "=" + DEBEZIUM_UNAVAILABLE_VALUE,
+            EVENT_TIME_ORDERING.name(),
+            IGNORE_MARKERS.name(),
+            "PostgresDebeziumAvroPayload"
+        ),
+        Arguments.of(
+            PartialUpdateAvroPayload.class.getName(),
+            "",
+            EVENT_TIME_ORDERING.name(),
+            PartialUpdateMode.IGNORE_DEFAULTS.name(),
+            "PartialUpdateAvroPayload"
+        ),
+        Arguments.of(
+            MySqlDebeziumAvroPayload.class.getName(),
+            "",
+            EVENT_TIME_ORDERING.name(),
+            PartialUpdateMode.NONE.name(),
+            "MySqlDebeziumAvroPayload"
+        ),
+        Arguments.of(
+            OverwriteNonDefaultsWithLatestAvroPayload.class.getName(),

Review Comment:
   EventTimeAvroPayload



##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestEightToNineUpgrade.scala:
##########
@@ -0,0 +1,174 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.hudi.functional
+
+import 
org.apache.hudi.DataSourceWriteOptions.{INSERT_OVERWRITE_OPERATION_OPT_VAL, 
PARTITIONPATH_FIELD, PAYLOAD_CLASS_NAME, RECORD_MERGE_IMPL_CLASSES, TABLE_TYPE, 
UPSERT_OPERATION_OPT_VAL}
+import org.apache.hudi.common.config.{HoodieStorageConfig, RecordMergeMode}
+import org.apache.hudi.common.model.{AWSDmsAvroPayload, HoodieRecordMerger, 
HoodieTableType, OverwriteNonDefaultsWithLatestAvroPayload, 
PartialUpdateAvroPayload}
+import org.apache.hudi.common.model.DefaultHoodieRecordPayload.{DELETE_KEY, 
DELETE_MARKER}
+import org.apache.hudi.common.model.debezium.PostgresDebeziumAvroPayload
+import org.apache.hudi.common.table.{HoodieTableConfig, HoodieTableMetaClient, 
HoodieTableVersion, PartialUpdateMode}
+import 
org.apache.hudi.common.table.HoodieTableConfig.{DEBEZIUM_UNAVAILABLE_VALUE, 
PARTIAL_UPDATE_CUSTOM_MARKER, RECORD_MERGE_PROPERTY_PREFIX}
+import org.apache.hudi.common.testutils.HoodieTestDataGenerator
+import org.apache.hudi.config.HoodieWriteConfig
+import org.apache.hudi.table.upgrade.{SparkUpgradeDowngradeHelper, 
UpgradeDowngrade}
+
+import org.apache.spark.sql.SaveMode
+import org.junit.jupiter.api.Assertions.assertEquals
+import org.junit.jupiter.params.ParameterizedTest
+import org.junit.jupiter.params.provider.{Arguments, MethodSource}
+
+class TestEightToNineUpgrade extends RecordLevelIndexTestBase {
+  @ParameterizedTest
+  @MethodSource(Array("payloadConfigs"))
+  def testUpgradeDowngradeBetweenEightAndNine(tableType: HoodieTableType,
+                                              payloadClass: String): Unit = {
+    val partitionFields = "partition:simple"
+    val mergerClasses = "org.apache.hudi.DefaultSparkRecordMerger," +
+      "org.apache.hudi.OverwriteWithLatestSparkRecordMerger," +
+      "org.apache.hudi.common.model.HoodieAvroRecordMerger"
+    var hudiOpts= commonOpts ++ Map(
+      TABLE_TYPE.key -> tableType.name(),
+      PARTITIONPATH_FIELD.key -> partitionFields,
+      PAYLOAD_CLASS_NAME.key -> payloadClass,
+      RECORD_MERGE_IMPL_CLASSES.key -> mergerClasses,
+      HoodieWriteConfig.WRITE_TABLE_VERSION.key -> "8",
+      HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet"
+    )
+
+    // Create a table in table version 8.
+    doWriteAndValidateDataAndRecordIndex(hudiOpts,
+      operation = INSERT_OVERWRITE_OPERATION_OPT_VAL,
+      saveMode = SaveMode.Overwrite,
+      schemaStr = 
HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA_WITH_PAYLOAD_SPECIFIC_COLS)
+    metaClient = getLatestMetaClient(true)
+    // Assert table version is 8.
+    checkResultForVersion8(payloadClass)
+    checkResultForVersion8(payloadClass)
+    // Add an extra commit.
+    doWriteAndValidateDataAndRecordIndex(hudiOpts,
+      operation = INSERT_OVERWRITE_OPERATION_OPT_VAL,
+      saveMode = SaveMode.Append,
+      schemaStr = 
HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA_WITH_PAYLOAD_SPECIFIC_COLS)
+    // Do validations.
+    checkResultForVersion8(payloadClass)
+
+    // Upgrade to version 9.
+    // Remove the write table version config, such that an upgrade could be 
triggered.
+    hudiOpts = hudiOpts ++ Map(HoodieWriteConfig.WRITE_TABLE_VERSION.key -> 
"9")
+    doWriteAndValidateDataAndRecordIndex(hudiOpts,
+      operation = INSERT_OVERWRITE_OPERATION_OPT_VAL,
+      saveMode = SaveMode.Append,
+      schemaStr = 
HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA_WITH_PAYLOAD_SPECIFIC_COLS)
+    // Table should be automatically upgraded to version 9.
+    // Do validations for table version 9.
+    checkResultForVersion9(partitionFields, payloadClass)
+    // Add an extra commit.
+    doWriteAndValidateDataAndRecordIndex(hudiOpts,
+      operation = INSERT_OVERWRITE_OPERATION_OPT_VAL,
+      saveMode = SaveMode.Append,
+      schemaStr = 
HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA_WITH_PAYLOAD_SPECIFIC_COLS)
+    // Do validations for table version 9.
+    checkResultForVersion9(partitionFields, payloadClass)
+
+    // Downgrade to table version 8 explicitly.
+    // Note that downgrade is NOT automatic.
+    // It has to be triggered explicitly.
+    hudiOpts = hudiOpts ++ Map(HoodieWriteConfig.WRITE_TABLE_VERSION.key -> 
"8")
+    new UpgradeDowngrade(metaClient, getWriteConfig(hudiOpts), context, 
SparkUpgradeDowngradeHelper.getInstance)
+      .run(HoodieTableVersion.EIGHT, null)
+    doWriteAndValidateDataAndRecordIndex(hudiOpts,
+      operation = INSERT_OVERWRITE_OPERATION_OPT_VAL,
+      saveMode = SaveMode.Append,
+      schemaStr = 
HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA_WITH_PAYLOAD_SPECIFIC_COLS)
+    checkResultForVersion8(payloadClass)
+    // Add an extra commit.
+    doWriteAndValidateDataAndRecordIndex(hudiOpts,
+      operation = INSERT_OVERWRITE_OPERATION_OPT_VAL,
+      saveMode = SaveMode.Append,
+      schemaStr = 
HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA_WITH_PAYLOAD_SPECIFIC_COLS)
+    // Do validations.
+    checkResultForVersion8(payloadClass)
+  }
+
+  def checkResultForVersion8(payloadClass: String): Unit = {
+    metaClient = HoodieTableMetaClient.reload(metaClient)
+    assertEquals(HoodieTableVersion.EIGHT, 
metaClient.getTableConfig.getTableVersion)
+    // The payload class should be maintained.
+    assertEquals(payloadClass, metaClient.getTableConfig.getPayloadClass)
+    // The partial update mode should be NONE.
+    assertEquals(PartialUpdateMode.NONE, 
metaClient.getTableConfig.getPartialUpdateMode)
+    // The merge mode should be CUSTOM.
+    assertEquals(
+      HoodieRecordMerger.PAYLOAD_BASED_MERGE_STRATEGY_UUID,
+      metaClient.getTableConfig.getRecordMergeStrategyId)
+  }
+
+  def checkResultForVersion9(partitionFields: String, payloadClass: String): 
Unit = {
+    metaClient = HoodieTableMetaClient.reload(metaClient)
+    assertEquals(HoodieTableVersion.NINE, 
metaClient.getTableConfig.getTableVersion)
+    assertEquals(
+      partitionFields,
+      
HoodieTableConfig.getPartitionFieldPropForKeyGenerator(metaClient.getTableConfig).get())
+    assertEquals(payloadClass, metaClient.getTableConfig.getLegacyPayloadClass)
+    // Based on the payload and table type, the merge mode is updated 
accordingly.
+    if (payloadClass.equals(classOf[PartialUpdateAvroPayload].getName)) {
+      assertEquals(
+        HoodieRecordMerger.EVENT_TIME_BASED_MERGE_STRATEGY_UUID,
+        metaClient.getTableConfig.getRecordMergeStrategyId)
+      assertEquals(RecordMergeMode.EVENT_TIME_ORDERING, 
metaClient.getTableConfig.getRecordMergeMode)
+      assertEquals(PartialUpdateMode.IGNORE_DEFAULTS, 
metaClient.getTableConfig.getPartialUpdateMode)
+    } else if 
(payloadClass.equals(classOf[OverwriteNonDefaultsWithLatestAvroPayload].getName))
 {
+      assertEquals(
+        HoodieRecordMerger.COMMIT_TIME_BASED_MERGE_STRATEGY_UUID,
+        metaClient.getTableConfig.getRecordMergeStrategyId)
+      assertEquals(RecordMergeMode.COMMIT_TIME_ORDERING, 
metaClient.getTableConfig.getRecordMergeMode)
+      assertEquals(PartialUpdateMode.IGNORE_DEFAULTS, 
metaClient.getTableConfig.getPartialUpdateMode)
+    } else if 
(payloadClass.equals(classOf[PostgresDebeziumAvroPayload].getName)) {
+      assertEquals(
+        HoodieRecordMerger.EVENT_TIME_BASED_MERGE_STRATEGY_UUID,
+        metaClient.getTableConfig.getRecordMergeStrategyId)
+      assertEquals(RecordMergeMode.EVENT_TIME_ORDERING, 
metaClient.getTableConfig.getRecordMergeMode)
+      assertEquals(PartialUpdateMode.IGNORE_MARKERS, 
metaClient.getTableConfig.getPartialUpdateMode)
+      val customMarker = 
metaClient.getTableConfig.getString(s"${RECORD_MERGE_PROPERTY_PREFIX}${PARTIAL_UPDATE_CUSTOM_MARKER}")
+      assertEquals(DEBEZIUM_UNAVAILABLE_VALUE, customMarker)
+    } else if (payloadClass.equals(classOf[AWSDmsAvroPayload].getName)) {
+      assertEquals(
+        HoodieRecordMerger.COMMIT_TIME_BASED_MERGE_STRATEGY_UUID,
+        metaClient.getTableConfig.getRecordMergeStrategyId)
+      assertEquals(RecordMergeMode.COMMIT_TIME_ORDERING, 
metaClient.getTableConfig.getRecordMergeMode)
+      val deleteField = 
metaClient.getTableConfig.getString(s"${RECORD_MERGE_PROPERTY_PREFIX}${DELETE_KEY}")
+      assertEquals(AWSDmsAvroPayload.OP_FIELD, deleteField)
+      val deleteMarker = 
metaClient.getTableConfig.getString(s"${RECORD_MERGE_PROPERTY_PREFIX}${DELETE_MARKER}")
+      assertEquals(AWSDmsAvroPayload.DELETE_OPERATION_VALUE, deleteMarker)
+    }
+  }
+}
+
+object TestEightToNineUpgrade {
+  def payloadConfigs(): java.util.stream.Stream[Arguments] = {
+    java.util.stream.Stream.of(
+      Arguments.of("MERGE_ON_READ", classOf[PartialUpdateAvroPayload].getName),

Review Comment:
   Lets expand this for COW table type as well. 
   and also all payloads as much as possible. 
   



##########
hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieWriterUtils.scala:
##########
@@ -255,6 +255,8 @@ object HoodieWriterUtils {
       val diffConfigs = StringBuilder.newBuilder
       params.foreach { case (key, value) =>
         if (!shouldIgnoreConfig(key, value, params, tableConfig)) {
+          // TODO: To disable payload class overwrite during writes,

Review Comment:
   hey @linliu-code : we have a jira for this right. can you link that in java 
docs here. 



##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPayloadDeprecationFlow.scala:
##########
@@ -185,24 +260,17 @@ object TestPayloadDeprecationFlow {
           HoodieTableConfig.LEGACY_PAYLOAD_CLASS_NAME.key() -> 
classOf[PostgresDebeziumAvroPayload].getName,
           HoodieTableConfig.RECORD_MERGE_STRATEGY_ID.key() -> 
HoodieRecordMerger.EVENT_TIME_BASED_MERGE_STRATEGY_UUID),
           HoodieTableConfig.PARTIAL_UPDATE_MODE.key() -> "IGNORE_MARKERS",
-        HoodieTableConfig.RECORD_MERGE_PROPERTY_PREFIX
-            + HoodieTableConfig.PARTIAL_UPDATE_CUSTOM_MARKER -> 
"__debezium_unavailable_value"),
+          HoodieTableConfig.RECORD_MERGE_PROPERTY_PREFIX + 
HoodieTableConfig.PARTIAL_UPDATE_CUSTOM_MARKER

Review Comment:
   lets enable these for COW as well. 



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to