This is an automated email from the ASF dual-hosted git repository.

yihua 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 276133b814e [HUDI-8003] Add OverwriteWithLatestHiveRecordMerger and 
restructure existing merger classes (#11649)
276133b814e is described below

commit 276133b814ef97aa427260fd2d0feb70723ae1b2
Author: Jon Vexler <[email protected]>
AuthorDate: Tue Sep 10 22:05:21 2024 -0400

    [HUDI-8003] Add OverwriteWithLatestHiveRecordMerger and restructure 
existing merger classes (#11649)
    
    Co-authored-by: Jonathan Vexler <=>
---
 hudi-client/hudi-java-client/pom.xml               |   8 +
 .../hadoop/TestHoodieFileGroupReaderOnHive.java    | 318 +++++++++++++++++++++
 .../hudi/testutils/ArrayWritableTestUtil.java      | 194 +++++++++++++
 .../testutils/HoodieJavaClientTestHarness.java     |   2 +-
 ...rdMerger.java => DefaultSparkRecordMerger.java} |  10 +-
 .../org/apache/hudi/HoodieSparkRecordMerger.java   | 107 +------
 ...a => OverwriteWithLatestSparkRecordMerger.java} |   7 +-
 .../hudi/BaseSparkInternalRowReaderContext.java    |  15 +-
 .../table/read/TestHoodieFileGroupReaderBase.java  |  16 +-
 ...ordMerger.java => DefaultHiveRecordMerger.java} |  10 +-
 .../hudi/hadoop/HiveHoodieReaderContext.java       |  66 ++---
 .../HoodieFileGroupReaderBasedRecordReader.java    |  42 ++-
 .../apache/hudi/hadoop/HoodieHiveRecordMerger.java |  51 +---
 .../hudi/hadoop/HoodieParquetInputFormat.java      |   8 +-
 .../OverwriteWithLatestHiveRecordMerger.java       |   8 +-
 .../hadoop/utils/HoodieArrayWritableAvroUtils.java |  12 -
 .../utils/HoodieRealtimeRecordReaderUtils.java     |   2 +-
 .../hudi/hadoop/utils/ObjectInspectorCache.java    |  16 +-
 ...odieSparkValidateDuplicateKeyRecordMerger.scala |  11 +-
 .../hudi/TestHoodieMergeHandleWithSparkMerger.java |  10 +-
 ...stHoodiePositionBasedFileGroupRecordBuffer.java |   7 +-
 .../read/TestHoodieFileGroupReaderOnSpark.scala    |  11 +-
 .../TestSpark35RecordPositionMetadataColumn.scala  |   2 +-
 .../apache/hudi/functional/CommonOptionUtils.scala |   4 +-
 .../apache/hudi/functional/TestCOWDataSource.scala |   2 +-
 .../functional/TestDataSourceForBootstrap.scala    |   4 +-
 .../TestHoodieMultipleBaseFileFormat.scala         |   4 +-
 .../apache/hudi/functional/TestMORDataSource.scala |   8 +-
 .../TestSparkDataSourceDAGExecution.scala          |   4 +-
 .../ReadAndWriteWithoutAvroBenchmark.scala         |  12 +-
 .../sql/hudi/common/HoodieSparkSqlTestBase.scala   |   6 +-
 .../apache/spark/sql/hudi/ddl/TestSpark3DDL.scala  |   4 +-
 .../deltastreamer/TestHoodieDeltaStreamer.java     |   4 +-
 33 files changed, 688 insertions(+), 297 deletions(-)

diff --git a/hudi-client/hudi-java-client/pom.xml 
b/hudi-client/hudi-java-client/pom.xml
index c6b36bd6bed..95bb71e54e4 100644
--- a/hudi-client/hudi-java-client/pom.xml
+++ b/hudi-client/hudi-java-client/pom.xml
@@ -69,6 +69,14 @@
             <type>test-jar</type>
             <scope>test</scope>
         </dependency>
+        <dependency>
+            <groupId>org.apache.hudi</groupId>
+            <artifactId>hudi-io</artifactId>
+            <version>${project.version}</version>
+            <classifier>tests</classifier>
+            <type>test-jar</type>
+            <scope>test</scope>
+        </dependency>
         <dependency>
             <groupId>org.apache.hudi</groupId>
             <artifactId>hudi-hadoop-common</artifactId>
diff --git 
a/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/hadoop/TestHoodieFileGroupReaderOnHive.java
 
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/hadoop/TestHoodieFileGroupReaderOnHive.java
new file mode 100644
index 00000000000..b29668d0d53
--- /dev/null
+++ 
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/hadoop/TestHoodieFileGroupReaderOnHive.java
@@ -0,0 +1,318 @@
+/*
+ * 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.hadoop;
+
+import org.apache.hudi.avro.HoodieAvroUtils;
+import org.apache.hudi.client.HoodieJavaWriteClient;
+import org.apache.hudi.client.common.HoodieJavaEngineContext;
+import org.apache.hudi.common.config.HoodieMemoryConfig;
+import org.apache.hudi.common.config.HoodieReaderConfig;
+import org.apache.hudi.common.config.RecordMergeMode;
+import org.apache.hudi.common.engine.EngineType;
+import org.apache.hudi.common.engine.HoodieReaderContext;
+import org.apache.hudi.common.model.DefaultHoodieRecordPayload;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.model.HoodieRecordMerger;
+import org.apache.hudi.common.model.OverwriteWithLatestAvroPayload;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.read.CustomPayloadForTesting;
+import org.apache.hudi.common.table.read.TestHoodieFileGroupReaderBase;
+import org.apache.hudi.common.testutils.HoodieTestDataGenerator;
+import org.apache.hudi.common.testutils.HoodieTestUtils;
+import org.apache.hudi.common.testutils.minicluster.HdfsTestService;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.hadoop.hive.HoodieCombineHiveInputFormat;
+import org.apache.hudi.hadoop.realtime.HoodieParquetRealtimeInputFormat;
+import org.apache.hudi.testutils.ArrayWritableTestUtil;
+import org.apache.hudi.hadoop.utils.ObjectInspectorCache;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StorageConfiguration;
+import org.apache.hudi.storage.hadoop.HoodieHadoopStorage;
+import org.apache.hudi.testutils.HoodieJavaClientTestHarness;
+
+import org.apache.avro.Schema;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.hive.metastore.api.hive_metastoreConstants;
+import org.apache.hadoop.hive.ql.exec.Utilities;
+import org.apache.hadoop.hive.ql.exec.mr.ExecMapper;
+import org.apache.hadoop.hive.ql.io.parquet.MapredParquetInputFormat;
+import org.apache.hadoop.hive.ql.plan.MapredWork;
+import org.apache.hadoop.hive.ql.plan.PartitionDesc;
+import org.apache.hadoop.hive.ql.plan.TableDesc;
+import org.apache.hadoop.hive.serde2.ColumnProjectionUtils;
+import org.apache.hadoop.io.ArrayWritable;
+import org.apache.hadoop.io.NullWritable;
+import org.apache.hadoop.mapred.FileInputFormat;
+import org.apache.hadoop.mapred.InputSplit;
+import org.apache.hadoop.mapred.JobConf;
+import org.apache.hadoop.mapred.RecordReader;
+import org.apache.hadoop.mapred.Reporter;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Disabled;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+import static org.apache.hadoop.hive.ql.exec.Utilities.HAS_MAP_WORK;
+import static org.apache.hadoop.hive.ql.exec.Utilities.MAPRED_MAPPER_CLASS;
+import static 
org.apache.hudi.hadoop.HoodieFileGroupReaderBasedRecordReader.getRecordKeyField;
+import static 
org.apache.hudi.hadoop.HoodieFileGroupReaderBasedRecordReader.getStoredPartitionFieldNames;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+public class TestHoodieFileGroupReaderOnHive extends 
TestHoodieFileGroupReaderBase<ArrayWritable> {
+
+  @Override
+  @Disabled("[HUDI-8072]")
+  public void testReadLogFilesOnlyInMergeOnReadTable(RecordMergeMode 
recordMergeMode, String logDataBlockFormat) throws Exception {
+  }
+
+  private static final String PARTITION_COLUMN = "datestr";
+  private static JobConf baseJobConf;
+  private static HdfsTestService hdfsTestService;
+  private static HoodieStorage storage;
+  private static FileSystem fs;
+  private static StorageConfiguration<Configuration> storageConf;
+
+  //currently always true. If we ever have a test with a nonpartitioned table, 
the usages of this should be tied together
+  private static final boolean USE_FAKE_PARTITION = true;
+
+  @BeforeAll
+  public static void setUpClass() throws IOException, InterruptedException {
+    // Append is not supported in LocalFileSystem. HDFS needs to be setup.
+    hdfsTestService = new HdfsTestService();
+    fs = hdfsTestService.start(true).getFileSystem();
+    storageConf = HoodieTestUtils.getDefaultStorageConf();
+    baseJobConf = new JobConf(storageConf.unwrap());
+    baseJobConf.set(HoodieMemoryConfig.MAX_DFS_STREAM_BUFFER_SIZE.key(), 
String.valueOf(1024 * 1024));
+    fs.setConf(baseJobConf);
+    storage = new HoodieHadoopStorage(fs);
+  }
+
+  @AfterAll
+  public static void tearDownClass() throws IOException {
+    hdfsTestService.stop();
+    if (fs != null) {
+      fs.close();
+      storage.close();
+    }
+  }
+
+  @Override
+  public StorageConfiguration<?> getStorageConf() {
+    return storageConf;
+  }
+
+  @Override
+  public String getBasePath() {
+    return tempDir.toAbsolutePath() + "/myTable";
+  }
+
+  @Override
+  public HoodieReaderContext<ArrayWritable> getHoodieReaderContext(String 
tablePath, Schema avroSchema, StorageConfiguration<?> storageConf) {
+    HoodieFileGroupReaderBasedRecordReader.HiveReaderCreator readerCreator = 
(inputSplit, jobConf) -> new 
MapredParquetInputFormat().getRecordReader(inputSplit, jobConf, null);
+    HoodieTableMetaClient metaClient = 
HoodieTableMetaClient.builder().setConf(storageConf).setBasePath(tablePath).build();
+    JobConf jobConf = new JobConf(storageConf.unwrapAs(Configuration.class));
+    setupJobconf(jobConf);
+    return new HiveHoodieReaderContext(readerCreator, 
getRecordKeyField(metaClient),
+        getStoredPartitionFieldNames(new 
JobConf(storageConf.unwrapAs(Configuration.class)), avroSchema),
+        new ObjectInspectorCache(avroSchema, jobConf));
+  }
+
+  @Override
+  public String getRecordPayloadForMergeMode(RecordMergeMode mergeMode) {
+    switch (mergeMode) {
+      case EVENT_TIME_ORDERING:
+        return DefaultHoodieRecordPayload.class.getName();
+      case OVERWRITE_WITH_LATEST:
+        return OverwriteWithLatestAvroPayload.class.getName();
+      case CUSTOM:
+      default:
+        return CustomPayloadForTesting.class.getName();
+    }
+  }
+
+  @Override
+  public void commitToTable(List<HoodieRecord> recordList, String operation, 
Map<String, String> writeConfigs) {
+    HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+        .withEngineType(EngineType.JAVA)
+        .withEmbeddedTimelineServerEnabled(false)
+        .withProps(writeConfigs)
+        .withPath(getBasePath())
+        .withSchema(HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA)
+        .build();
+
+    HoodieJavaClientTestHarness.TestJavaTaskContextSupplier 
taskContextSupplier = new 
HoodieJavaClientTestHarness.TestJavaTaskContextSupplier();
+    HoodieJavaEngineContext context = new 
HoodieJavaEngineContext(getStorageConf(), taskContextSupplier);
+    //init table if not exists
+    Path basePath = new Path(getBasePath());
+    try {
+      try (FileSystem lfs = basePath.getFileSystem(baseJobConf)) {
+        boolean basepathExists = lfs.exists(basePath);
+        boolean operationIsInsert = operation.equalsIgnoreCase("insert");
+        if (!basepathExists || operationIsInsert) {
+          if (basepathExists) {
+            lfs.delete(new Path(getBasePath()), true);
+          }
+          String recordMergerStrategy = "";
+          if 
(RecordMergeMode.valueOf(writeConfigs.get("hoodie.record.merge.mode")).equals(RecordMergeMode.OVERWRITE_WITH_LATEST))
 {
+            recordMergerStrategy = 
HoodieRecordMerger.OVERWRITE_MERGER_STRATEGY_UUID;
+          } else if 
(RecordMergeMode.valueOf(writeConfigs.get("hoodie.record.merge.mode")).equals(RecordMergeMode.EVENT_TIME_ORDERING))
 {
+            recordMergerStrategy = 
HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
+          } else if 
(RecordMergeMode.valueOf(writeConfigs.get("hoodie.record.merge.mode")).equals(RecordMergeMode.CUSTOM))
 {
+            //match the behavior of spark for now, but this should be a config
+            recordMergerStrategy = 
HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
+          }
+          Map<String, Object> initConfigs = new HashMap<>(writeConfigs);
+          HoodieTableMetaClient.withPropertyBuilder()
+              
.setTableType(writeConfigs.getOrDefault("hoodie.datasource.write.table.type", 
"MERGE_ON_READ"))
+              .setTableName(writeConfigs.get("hoodie.table.name"))
+              
.setPartitionFields(writeConfigs.getOrDefault("hoodie.datasource.write.partitionpath.field",
 ""))
+              
.setRecordMergeMode(RecordMergeMode.valueOf(writeConfigs.get("hoodie.record.merge.mode")))
+              .setRecordMergerStrategy(recordMergerStrategy)
+              .set(initConfigs).initTable(storageConf, getBasePath());
+        }
+      }
+    } catch (IOException e) {
+      throw new RuntimeException(e);
+    }
+
+    HoodieJavaWriteClient writeClient = new HoodieJavaWriteClient(context, 
writeConfig);
+    String instantTime = writeClient.createNewInstantTime();
+    writeClient.startCommitWithTime(instantTime);
+    if (operation.toLowerCase().equals("insert")) {
+      writeClient.insert(recordList, instantTime);
+    } else {
+      writeClient.upsert(recordList, instantTime);
+    }
+  }
+
+  @Override
+  public void validateRecordsInFileGroup(String tablePath, List<ArrayWritable> 
actualRecordList, Schema schema, String fileGroupId) {
+    
assertEquals(HoodieAvroUtils.addMetadataFields(HoodieTestDataGenerator.AVRO_SCHEMA),
 schema);
+    try {
+      //prepare fg reader records to be compared to the baseline reader
+      HoodieReaderContext<ArrayWritable> readerContext = 
getHoodieReaderContext(tablePath, schema, storageConf);
+      Map<String, ArrayWritable> recordMap = new HashMap<>();
+      for (ArrayWritable record : actualRecordList) {
+        recordMap.put(readerContext.getRecordKey(record, schema), record);
+      }
+
+      RecordReader<NullWritable, ArrayWritable> reader = 
createRecordReader(tablePath);
+      // use reader to read log file.
+      NullWritable key = reader.createKey();
+      ArrayWritable value = reader.createValue();
+      while (reader.next(key, value)) {
+        if (readerContext.getValue(value, schema, 
HoodieRecord.FILENAME_METADATA_FIELD).toString().contains(fileGroupId)) {
+          //only evaluate records from the specified filegroup. Maybe there is 
a way to get
+          //hive to do this?
+          ArrayWritable compVal = 
recordMap.remove(readerContext.getRecordKey(value, schema));
+          assertNotNull(compVal);
+          ArrayWritableTestUtil.assertArrayWritableEqual(schema, value, 
compVal, USE_FAKE_PARTITION);
+        }
+        key = reader.createKey();
+        value = reader.createValue();
+      }
+      reader.close();
+      assertEquals(0, recordMap.size());
+    } catch (IOException e) {
+      throw new RuntimeException(e);
+    }
+  }
+
+  private RecordReader<NullWritable, ArrayWritable> createRecordReader(String 
tablePath) throws IOException {
+    JobConf jobConf = new JobConf(baseJobConf);
+    jobConf.set(HoodieReaderConfig.FILE_GROUP_READER_ENABLED.key(), "false");
+
+    TableDesc tblDesc = Utilities.defaultTd;
+    // Set the input format
+    tblDesc.setInputFileFormatClass(HoodieParquetRealtimeInputFormat.class);
+    LinkedHashMap<Path, PartitionDesc> pt = new LinkedHashMap<>();
+    LinkedHashMap<Path, ArrayList<String>> talias = new LinkedHashMap<>();
+
+    PartitionDesc partDesc = new PartitionDesc(tblDesc, null);
+
+    pt.put(new Path(tablePath), partDesc);
+
+    ArrayList<String> arrayList = new ArrayList<>();
+    arrayList.add(tablePath);
+    talias.put(new Path(tablePath), arrayList);
+
+    MapredWork mrwork = new MapredWork();
+    mrwork.getMapWork().setPathToPartitionInfo(pt);
+    mrwork.getMapWork().setPathToAliases(talias);
+
+    Path mapWorkPath = new Path(tablePath);
+    Utilities.setMapRedWork(jobConf, mrwork, mapWorkPath);
+
+    // Add three partition path to InputPaths
+    Path[] partitionDirArray = new 
Path[HoodieTestDataGenerator.DEFAULT_PARTITION_PATHS.length];
+    Arrays.stream(HoodieTestDataGenerator.DEFAULT_PARTITION_PATHS).map(s -> 
new Path(tablePath, s)).collect(Collectors.toList()).toArray(partitionDirArray);
+    FileInputFormat.setInputPaths(jobConf, partitionDirArray);
+    jobConf.set(HAS_MAP_WORK, "true");
+    // The following config tells Hive to choose ExecMapper to read the 
MAP_WORK
+    jobConf.set(MAPRED_MAPPER_CLASS, ExecMapper.class.getName());
+    // setting the split size to be 3 to create one split for 3 file groups
+    
jobConf.set(org.apache.hadoop.mapreduce.lib.input.FileInputFormat.SPLIT_MAXSIZE,
 "128000000");
+    setupJobconf(jobConf);
+
+    HoodieCombineHiveInputFormat combineHiveInputFormat = new 
HoodieCombineHiveInputFormat();
+    InputSplit[] splits = combineHiveInputFormat.getSplits(jobConf, 1);
+
+    assertEquals(1, splits.length);
+    return  combineHiveInputFormat.getRecordReader(splits[0], jobConf, 
Reporter.NULL);
+  }
+
+  private void setupJobconf(JobConf jobConf) {
+    Schema schema = 
HoodieAvroUtils.addMetadataFields(HoodieTestDataGenerator.AVRO_SCHEMA);
+    List<Schema.Field> fields = schema.getFields();
+    setHiveColumnNameProps(fields, jobConf, USE_FAKE_PARTITION);
+    jobConf.set("columns.types","string,string,string,string,string," + 
HoodieTestDataGenerator.TRIP_HIVE_COLUMN_TYPES + ",string");
+  }
+
+  private void setHiveColumnNameProps(List<Schema.Field> fields, JobConf 
jobConf, boolean isPartitioned) {
+    String names = 
fields.stream().map(Schema.Field::name).collect(Collectors.joining(","));
+    String positions = fields.stream().map(f -> 
String.valueOf(f.pos())).collect(Collectors.joining(","));
+    jobConf.set(ColumnProjectionUtils.READ_COLUMN_NAMES_CONF_STR, names);
+    jobConf.set(ColumnProjectionUtils.READ_COLUMN_IDS_CONF_STR, positions);
+
+    String hiveOrderedColumnNames = fields.stream().filter(field -> 
!field.name().equalsIgnoreCase(PARTITION_COLUMN))
+        .map(Schema.Field::name).collect(Collectors.joining(","));
+    if (isPartitioned) {
+      hiveOrderedColumnNames += "," + PARTITION_COLUMN;
+      jobConf.set(hive_metastoreConstants.META_TABLE_PARTITION_COLUMNS, 
PARTITION_COLUMN);
+    }
+    jobConf.set(hive_metastoreConstants.META_TABLE_COLUMNS, 
hiveOrderedColumnNames);
+  }
+
+  @Override
+  public Comparable getComparableUTF8String(String value) {
+    return value;
+  }
+}
diff --git 
a/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/ArrayWritableTestUtil.java
 
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/ArrayWritableTestUtil.java
new file mode 100644
index 00000000000..bbca2671127
--- /dev/null
+++ 
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/ArrayWritableTestUtil.java
@@ -0,0 +1,194 @@
+/*
+ * 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.testutils;
+
+import org.apache.avro.LogicalTypes;
+import org.apache.avro.Schema;
+import org.apache.hadoop.hive.serde2.io.DateWritable;
+import org.apache.hadoop.hive.serde2.io.HiveDecimalWritable;
+import org.apache.hadoop.io.ArrayWritable;
+import org.apache.hadoop.io.BooleanWritable;
+import org.apache.hadoop.io.BytesWritable;
+import org.apache.hadoop.io.DoubleWritable;
+import org.apache.hadoop.io.FloatWritable;
+import org.apache.hadoop.io.IntWritable;
+import org.apache.hadoop.io.LongWritable;
+import org.apache.hadoop.io.NullWritable;
+import org.apache.hadoop.io.Text;
+import org.apache.hadoop.io.Writable;
+
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+public class ArrayWritableTestUtil {
+  public static void assertArrayWritableEqual(Schema schema, ArrayWritable 
expected, ArrayWritable actual, boolean isPartitioned) {
+    assertArrayWritableEqualInternal(schema, expected, actual, isPartitioned);
+  }
+
+  private static void assertArrayWritableEqualInternal(Schema schema, Writable 
expected, Writable actual, boolean ignoreOneExtraCol) {
+    switch (schema.getType()) {
+      case RECORD: {
+        assertInstanceOf(ArrayWritable.class, expected);
+        assertInstanceOf(ArrayWritable.class, actual);
+        //adjust for fake partition
+        int expectedLen = ((ArrayWritable) expected).get().length - 
(ignoreOneExtraCol ? 1 : 0);
+        int actualLen = ((ArrayWritable) actual).get().length;
+        assertEquals(expectedLen, actualLen);
+        assertEquals(schema.getFields().size(), expectedLen);
+        for (Schema.Field field : schema.getFields()) {
+          assertArrayWritableEqualInternal(field.schema(), ((ArrayWritable) 
expected).get()[field.pos()], ((ArrayWritable) actual).get()[field.pos()], 
false);
+        }
+        break;
+      }
+      case ARRAY: {
+        assertInstanceOf(ArrayWritable.class, expected);
+        assertInstanceOf(ArrayWritable.class, actual);
+        int expectedLen = ((ArrayWritable) expected).get().length;
+        int actualLen = ((ArrayWritable) actual).get().length;
+        assertEquals(expectedLen, actualLen);
+        for (int i = 0; i < expectedLen; i++) {
+          assertArrayWritableEqualInternal(schema.getElementType(), 
((ArrayWritable) expected).get()[i], ((ArrayWritable) expected).get()[i], 
false);
+        }
+        break;
+      }
+      case MAP: {
+        assertInstanceOf(ArrayWritable.class, expected);
+        assertInstanceOf(ArrayWritable.class, actual);
+        int expectedLen = ((ArrayWritable) expected).get().length;
+        int actualLen = ((ArrayWritable) actual).get().length;
+        assertEquals(expectedLen, actualLen);
+        Map<Writable, Writable> expectedMap = new HashMap<>(expectedLen);
+        Map<Writable, Writable> actualMap = new HashMap<>(actualLen);
+        for (int i = 0; i < expectedLen; i++) {
+          Writable expectedKV = ((ArrayWritable) expected).get()[i];
+          assertInstanceOf(ArrayWritable.class, expectedKV);
+          assertEquals(2, ((ArrayWritable) expectedKV).get().length);
+          expectedMap.put(((ArrayWritable) expectedKV).get()[0], 
((ArrayWritable) expectedKV).get()[1]);
+          Writable actualKV = ((ArrayWritable) actual).get()[i];
+          assertInstanceOf(ArrayWritable.class, actualKV);
+          assertEquals(2, ((ArrayWritable) actualKV).get().length);
+          actualMap.put(((ArrayWritable) actualKV).get()[0], ((ArrayWritable) 
actualKV).get()[1]);
+        }
+
+        for (Writable key : expectedMap.keySet()) {
+          Writable expectedValue = expectedMap.get(key);
+          assertNotNull(expectedValue);
+          Writable actualValue = actualMap.remove(key);
+          assertNotNull(actualValue);
+          assertArrayWritableEqualInternal(schema.getValueType(), 
expectedValue, actualValue, false);
+        }
+        assertEquals(0, actualMap.size());
+        break;
+      }
+      case UNION:
+        if (schema.getTypes().size() == 2
+            && schema.getTypes().get(0).getType() == Schema.Type.NULL) {
+          assertArrayWritableEqualInternal(schema.getTypes().get(1), expected, 
actual, false);
+        } else if (schema.getTypes().size() == 2
+            && schema.getTypes().get(1).getType() == Schema.Type.NULL) {
+          assertArrayWritableEqualInternal(schema.getTypes().get(0), expected, 
actual, false);
+        } else if (schema.getTypes().size() == 1) {
+          assertArrayWritableEqualInternal(schema.getTypes().get(0), expected, 
actual, false);
+        } else {
+          throw new IllegalStateException("Union has more than 2 types or one 
type is not null: " + schema);
+        }
+        break;
+
+      default:
+        assertWritablePrimaryType(schema, expected, actual);
+    }
+  }
+
+  private static void assertWritablePrimaryType(Schema schema, Writable 
expected, Writable actual) {
+    switch (schema.getType()) {
+      case NULL:
+        assertInstanceOf(NullWritable.class, expected);
+        assertInstanceOf(NullWritable.class, actual);
+        assertEquals(expected, actual);
+        break;
+
+      case BOOLEAN:
+        assertInstanceOf(BooleanWritable.class, expected);
+        assertInstanceOf(BooleanWritable.class, actual);
+        assertEquals(expected, actual);
+        break;
+
+      case INT:
+        if (schema.getLogicalType() instanceof LogicalTypes.Date) {
+          assertInstanceOf(DateWritable.class, expected);
+          assertInstanceOf(DateWritable.class, actual);
+        } else {
+          assertInstanceOf(IntWritable.class, expected);
+          assertInstanceOf(IntWritable.class, actual);
+        }
+        assertEquals(expected, actual);
+        break;
+
+      case LONG:
+        assertInstanceOf(LongWritable.class, expected);
+        assertInstanceOf(LongWritable.class, actual);
+        assertEquals(expected, actual);
+        break;
+
+      case FLOAT:
+        assertInstanceOf(FloatWritable.class, expected);
+        assertInstanceOf(FloatWritable.class, actual);
+        assertEquals(expected, actual);
+        break;
+
+      case DOUBLE:
+        assertInstanceOf(DoubleWritable.class, expected);
+        assertInstanceOf(DoubleWritable.class, actual);
+        assertEquals(expected, actual);
+        break;
+
+      case BYTES:
+      case ENUM:
+        assertInstanceOf(BytesWritable.class, expected);
+        assertInstanceOf(BytesWritable.class, actual);
+        assertEquals(expected, actual);
+        break;
+
+      case STRING:
+        assertInstanceOf(Text.class, expected);
+        assertInstanceOf(Text.class, actual);
+        assertEquals(expected, actual);
+        break;
+
+      case FIXED:
+        if (schema.getLogicalType() instanceof LogicalTypes.Decimal) {
+          assertInstanceOf(HiveDecimalWritable.class, expected);
+          assertInstanceOf(HiveDecimalWritable.class, actual);
+        } else {
+          assertEquals(expected.getClass(), actual.getClass());
+        }
+        assertEquals(expected, actual);
+        break;
+
+      default:
+        assertEquals(expected.getClass(), actual.getClass());
+        assertEquals(expected, actual);
+    }
+  }
+}
diff --git 
a/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/HoodieJavaClientTestHarness.java
 
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/HoodieJavaClientTestHarness.java
index 9195ba9f165..eeebd7960a3 100644
--- 
a/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/HoodieJavaClientTestHarness.java
+++ 
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/HoodieJavaClientTestHarness.java
@@ -154,7 +154,7 @@ public abstract class HoodieJavaClientTestHarness extends 
HoodieWriterClientTest
     cleanupExecutorService();
   }
 
-  public class TestJavaTaskContextSupplier extends TaskContextSupplier {
+  public static class TestJavaTaskContextSupplier extends TaskContextSupplier {
     int partitionId = 0;
     int stageId = 0;
     long attemptId = 0;
diff --git 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/DefaultSparkRecordMerger.java
similarity index 96%
copy from 
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
copy to 
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/DefaultSparkRecordMerger.java
index a488e7174e5..1b0197a3330 100644
--- 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
+++ 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/DefaultSparkRecordMerger.java
@@ -33,7 +33,10 @@ import org.apache.avro.Schema;
 
 import java.io.IOException;
 
-public class HoodieSparkRecordMerger implements HoodieRecordMerger {
+/**
+ * Record merger for spark that implements the default merger strategy
+ */
+public class DefaultSparkRecordMerger extends HoodieSparkRecordMerger {
 
   @Override
   public String getMergingStrategy() {
@@ -115,9 +118,4 @@ public class HoodieSparkRecordMerger implements 
HoodieRecordMerger {
           (HoodieSparkRecord) older, oldSchema, (HoodieSparkRecord) newer, 
newSchema, readerSchema, props));
     }
   }
-
-  @Override
-  public HoodieRecordType getRecordType() {
-    return HoodieRecordType.SPARK;
-  }
 }
\ No newline at end of file
diff --git 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
index a488e7174e5..18fa76004af 100644
--- 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
+++ 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
@@ -19,105 +19,24 @@
 
 package org.apache.hudi;
 
-import org.apache.hudi.common.config.TypedProperties;
 import org.apache.hudi.common.model.HoodieRecord;
-import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType;
 import org.apache.hudi.common.model.HoodieRecordMerger;
-import org.apache.hudi.common.model.HoodieSparkRecord;
-import org.apache.hudi.common.util.Option;
-import org.apache.hudi.common.util.ValidationUtils;
-import org.apache.hudi.common.util.collection.Pair;
-import org.apache.hudi.merge.SparkRecordMergingUtils;
-
-import org.apache.avro.Schema;
-
-import java.io.IOException;
-
-public class HoodieSparkRecordMerger implements HoodieRecordMerger {
+import org.apache.hudi.exception.HoodieException;
 
+public abstract class HoodieSparkRecordMerger implements HoodieRecordMerger {
   @Override
-  public String getMergingStrategy() {
-    return HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
+  public HoodieRecord.HoodieRecordType getRecordType() {
+    return HoodieRecord.HoodieRecordType.SPARK;
   }
 
-  @Override
-  public Option<Pair<HoodieRecord, Schema>> merge(HoodieRecord older, Schema 
oldSchema, HoodieRecord newer, Schema newSchema, TypedProperties props) throws 
IOException {
-    ValidationUtils.checkArgument(older.getRecordType() == 
HoodieRecordType.SPARK);
-    ValidationUtils.checkArgument(newer.getRecordType() == 
HoodieRecordType.SPARK);
-
-    if (newer instanceof HoodieSparkRecord) {
-      HoodieSparkRecord newSparkRecord = (HoodieSparkRecord) newer;
-      if (newSparkRecord.isDeleted()) {
-        // Delete record
-        return Option.empty();
-      }
-    } else {
-      if (newer.getData() == null) {
-        // Delete record
-        return Option.empty();
-      }
+  static HoodieRecordMerger getRecordMerger(String mergerStrategy) {
+    switch (mergerStrategy) {
+      case DEFAULT_MERGER_STRATEGY_UUID:
+        return new DefaultSparkRecordMerger();
+      case OVERWRITE_MERGER_STRATEGY_UUID:
+        return new OverwriteWithLatestSparkRecordMerger();
+      default:
+        throw new HoodieException("This merger strategy UUID is not supported: 
" + mergerStrategy);
     }
-
-    if (older instanceof HoodieSparkRecord) {
-      HoodieSparkRecord oldSparkRecord = (HoodieSparkRecord) older;
-      if (oldSparkRecord.isDeleted()) {
-        // use natural order for delete record
-        return Option.of(Pair.of(newer, newSchema));
-      }
-    } else {
-      if (older.getData() == null) {
-        // use natural order for delete record
-        return Option.of(Pair.of(newer, newSchema));
-      }
-    }
-    if (older.getOrderingValue(oldSchema, 
props).compareTo(newer.getOrderingValue(newSchema, props)) > 0) {
-      return Option.of(Pair.of(older, oldSchema));
-    } else {
-      return Option.of(Pair.of(newer, newSchema));
-    }
-  }
-
-  @Override
-  public Option<Pair<HoodieRecord, Schema>> partialMerge(HoodieRecord older, 
Schema oldSchema, HoodieRecord newer, Schema newSchema, Schema readerSchema, 
TypedProperties props) throws IOException {
-    ValidationUtils.checkArgument(older.getRecordType() == 
HoodieRecordType.SPARK);
-    ValidationUtils.checkArgument(newer.getRecordType() == 
HoodieRecordType.SPARK);
-
-    if (newer instanceof HoodieSparkRecord) {
-      HoodieSparkRecord newSparkRecord = (HoodieSparkRecord) newer;
-      if (newSparkRecord.isDeleted()) {
-        // Delete record
-        return Option.empty();
-      }
-    } else {
-      if (newer.getData() == null) {
-        // Delete record
-        return Option.empty();
-      }
-    }
-
-    if (older instanceof HoodieSparkRecord) {
-      HoodieSparkRecord oldSparkRecord = (HoodieSparkRecord) older;
-      if (oldSparkRecord.isDeleted()) {
-        // use natural order for delete record
-        return Option.of(Pair.of(newer, newSchema));
-      }
-    } else {
-      if (older.getData() == null) {
-        // use natural order for delete record
-        return Option.of(Pair.of(newer, newSchema));
-      }
-    }
-    if (older.getOrderingValue(oldSchema, 
props).compareTo(newer.getOrderingValue(newSchema, props)) > 0) {
-      return Option.of(SparkRecordMergingUtils.mergePartialRecords(
-          (HoodieSparkRecord) newer, newSchema, (HoodieSparkRecord) older, 
oldSchema, readerSchema, props));
-    } else {
-      return Option.of(SparkRecordMergingUtils.mergePartialRecords(
-          (HoodieSparkRecord) older, oldSchema, (HoodieSparkRecord) newer, 
newSchema, readerSchema, props));
-    }
-  }
-
-  @Override
-  public HoodieRecordType getRecordType() {
-    return HoodieRecordType.SPARK;
   }
-}
\ No newline at end of file
+}
diff --git 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkMerger.java
 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkRecordMerger.java
similarity index 89%
copy from 
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkMerger.java
copy to 
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkRecordMerger.java
index 611f045f645..4cd052e8356 100644
--- 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkMerger.java
+++ 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkRecordMerger.java
@@ -31,7 +31,7 @@ import java.io.IOException;
 /**
  * Spark merger that always chooses the newer record
  */
-public class OverwriteWithLatestSparkMerger extends HoodieSparkRecordMerger {
+public class OverwriteWithLatestSparkRecordMerger extends 
HoodieSparkRecordMerger {
 
   @Override
   public String getMergingStrategy() {
@@ -43,4 +43,9 @@ public class OverwriteWithLatestSparkMerger extends 
HoodieSparkRecordMerger {
     return Option.of(Pair.of(newer, newSchema));
   }
 
+  @Override
+  public HoodieRecord.HoodieRecordType getRecordType() {
+    return null;
+  }
+
 }
diff --git 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
index 72b6276b458..817e83fce44 100644
--- 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
+++ 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
@@ -28,7 +28,9 @@ import org.apache.hudi.common.model.HoodieRecordMerger;
 import org.apache.hudi.common.model.HoodieSparkRecord;
 import org.apache.hudi.common.util.ConfigUtils;
 import org.apache.hudi.common.util.Option;
-import org.apache.hudi.exception.HoodieException;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.HoodieStorageUtils;
+import org.apache.hudi.storage.StorageConfiguration;
 
 import org.apache.avro.Schema;
 import org.apache.spark.sql.HoodieInternalRowUtils;
@@ -44,8 +46,6 @@ import java.util.function.UnaryOperator;
 import scala.Function1;
 
 import static 
org.apache.hudi.common.model.HoodieRecord.RECORD_KEY_METADATA_FIELD;
-import static 
org.apache.hudi.common.model.HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
-import static 
org.apache.hudi.common.model.HoodieRecordMerger.OVERWRITE_MERGER_STRATEGY_UUID;
 import static org.apache.spark.sql.HoodieInternalRowUtils.getCachedSchema;
 
 /**
@@ -56,14 +56,7 @@ public abstract class BaseSparkInternalRowReaderContext 
extends HoodieReaderCont
 
   @Override
   public HoodieRecordMerger getRecordMerger(String mergerStrategy) {
-    switch (mergerStrategy) {
-      case DEFAULT_MERGER_STRATEGY_UUID:
-        return new HoodieSparkRecordMerger();
-      case OVERWRITE_MERGER_STRATEGY_UUID:
-        return new OverwriteWithLatestSparkMerger();
-      default:
-        throw new HoodieException("The merger strategy UUID is not supported: 
" + mergerStrategy);
-    }
+    return HoodieSparkRecordMerger.getRecordMerger(mergerStrategy);
   }
 
   @Override
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderBase.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderBase.java
index a53583e63da..0a6b8f3937d 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderBase.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderBase.java
@@ -29,6 +29,7 @@ import org.apache.hudi.common.engine.HoodieEngineContext;
 import org.apache.hudi.common.engine.HoodieLocalEngineContext;
 import org.apache.hudi.common.engine.HoodieReaderContext;
 import org.apache.hudi.common.model.FileSlice;
+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.HoodieTableMetaClient;
@@ -68,7 +69,6 @@ import static 
org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGER_STRAT
 import static org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_MODE;
 import static 
org.apache.hudi.common.table.read.HoodieBaseFileGroupRecordBuffer.compareTo;
 import static 
org.apache.hudi.common.testutils.HoodieTestUtils.getLogFileListFromFileSlice;
-import static 
org.apache.hudi.common.testutils.RawTripTestPayload.recordsToStrings;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.params.provider.Arguments.arguments;
@@ -90,7 +90,7 @@ public abstract class TestHoodieFileGroupReaderBase<T> {
 
   public abstract String getRecordPayloadForMergeMode(RecordMergeMode 
mergeMode);
 
-  public abstract void commitToTable(List<String> recordList, String operation,
+  public abstract void commitToTable(List<HoodieRecord> recordList, String 
operation,
                                      Map<String, String> writeConfigs);
 
   public abstract void validateRecordsInFileGroup(String tablePath,
@@ -105,7 +105,7 @@ public abstract class TestHoodieFileGroupReaderBase<T> {
     Map<String, String> writeConfigs = new 
HashMap<>(getCommonConfigs(RecordMergeMode.EVENT_TIME_ORDERING));
     // Prepare a table for initializing reader context
     try (HoodieTestDataGenerator dataGen = new 
HoodieTestDataGenerator(0xDEEF)) {
-      commitToTable(recordsToStrings(dataGen.generateInserts("001", 1)), 
BULK_INSERT.value(), writeConfigs);
+      commitToTable(dataGen.generateInserts("001", 1), BULK_INSERT.value(), 
writeConfigs);
     }
     StorageConfiguration<?> storageConf = getStorageConf();
     String tablePath = getBasePath();
@@ -184,17 +184,17 @@ public abstract class TestHoodieFileGroupReaderBase<T> {
 
     try (HoodieTestDataGenerator dataGen = new 
HoodieTestDataGenerator(0xDEEF)) {
       // One commit; reading one file group containing a base file only
-      commitToTable(recordsToStrings(dataGen.generateInserts("001", 100)), 
INSERT.value(), writeConfigs);
+      commitToTable(dataGen.generateInserts("001", 100), INSERT.value(), 
writeConfigs);
       validateOutputFromFileGroupReader(
           getStorageConf(), getBasePath(), dataGen.getPartitionPaths(), true, 
0, recordMergeMode);
 
       // Two commits; reading one file group containing a base file and a log 
file
-      commitToTable(recordsToStrings(dataGen.generateUpdates("002", 100)), 
UPSERT.value(), writeConfigs);
+      commitToTable(dataGen.generateUpdates("002", 100), UPSERT.value(), 
writeConfigs);
       validateOutputFromFileGroupReader(
           getStorageConf(), getBasePath(), dataGen.getPartitionPaths(), true, 
1, recordMergeMode);
 
       // Three commits; reading one file group containing a base file and two 
log files
-      commitToTable(recordsToStrings(dataGen.generateUpdates("003", 100)), 
UPSERT.value(), writeConfigs);
+      commitToTable(dataGen.generateUpdates("003", 100), UPSERT.value(), 
writeConfigs);
       validateOutputFromFileGroupReader(
           getStorageConf(), getBasePath(), dataGen.getPartitionPaths(), true, 
2, recordMergeMode);
     }
@@ -210,12 +210,12 @@ public abstract class TestHoodieFileGroupReaderBase<T> {
 
     try (HoodieTestDataGenerator dataGen = new 
HoodieTestDataGenerator(0xDEEF)) {
       // One commit; reading one file group containing a base file only
-      commitToTable(recordsToStrings(dataGen.generateInserts("001", 100)), 
INSERT.value(), writeConfigs);
+      commitToTable(dataGen.generateInserts("001", 100), INSERT.value(), 
writeConfigs);
       validateOutputFromFileGroupReader(
           getStorageConf(), getBasePath(), dataGen.getPartitionPaths(), false, 
1, recordMergeMode);
 
       // Two commits; reading one file group containing a base file and a log 
file
-      commitToTable(recordsToStrings(dataGen.generateUpdates("002", 100)), 
UPSERT.value(), writeConfigs);
+      commitToTable(dataGen.generateUpdates("002", 100), UPSERT.value(), 
writeConfigs);
       validateOutputFromFileGroupReader(
           getStorageConf(), getBasePath(), dataGen.getPartitionPaths(), false, 
2, recordMergeMode);
     }
diff --git 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/DefaultHiveRecordMerger.java
similarity index 92%
copy from 
hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
copy to 
hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/DefaultHiveRecordMerger.java
index 17a4738569e..7d353f7ee10 100644
--- 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
+++ 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/DefaultHiveRecordMerger.java
@@ -30,7 +30,10 @@ import org.apache.avro.Schema;
 
 import java.io.IOException;
 
-public class HoodieHiveRecordMerger implements HoodieRecordMerger {
+/**
+ * Record merger for hive that implements the default merger strategy
+ */
+public class DefaultHiveRecordMerger extends HoodieHiveRecordMerger {
   @Override
   public Option<Pair<HoodieRecord, Schema>> merge(HoodieRecord older, Schema 
oldSchema, HoodieRecord newer, Schema newSchema, TypedProperties props) throws 
IOException {
     ValidationUtils.checkArgument(older.getRecordType() == 
HoodieRecord.HoodieRecordType.HIVE);
@@ -59,11 +62,6 @@ public class HoodieHiveRecordMerger implements 
HoodieRecordMerger {
     }
   }
 
-  @Override
-  public HoodieRecord.HoodieRecordType getRecordType() {
-    return HoodieRecord.HoodieRecordType.HIVE;
-  }
-
   @Override
   public String getMergingStrategy() {
     return HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
diff --git 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveHoodieReaderContext.java
 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveHoodieReaderContext.java
index d2b6ccd8ee8..e31b1756ec9 100644
--- 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveHoodieReaderContext.java
+++ 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveHoodieReaderContext.java
@@ -25,13 +25,10 @@ import org.apache.hudi.common.model.HoodieEmptyRecord;
 import org.apache.hudi.common.model.HoodieKey;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.HoodieRecordMerger;
-import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.StringUtils;
-import org.apache.hudi.common.util.ValidationUtils;
 import org.apache.hudi.common.util.collection.ClosableIterator;
 import org.apache.hudi.common.util.collection.CloseableMappingIterator;
-import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.hadoop.utils.HoodieArrayWritableAvroUtils;
 import org.apache.hudi.hadoop.utils.HoodieRealtimeRecordReaderUtils;
 import org.apache.hudi.hadoop.utils.ObjectInspectorCache;
@@ -41,7 +38,10 @@ import org.apache.hudi.storage.StoragePathInfo;
 
 import org.apache.avro.Schema;
 import org.apache.avro.generic.IndexedRecord;
+import org.apache.hadoop.conf.Configuration;
 import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.hive.ql.io.sarg.ConvertAstToSearchArg;
+import org.apache.hadoop.hive.ql.plan.TableScanDesc;
 import org.apache.hadoop.hive.serde.serdeConstants;
 import org.apache.hadoop.hive.serde2.ColumnProjectionUtils;
 import org.apache.hadoop.hive.serde2.typeinfo.TypeInfo;
@@ -52,7 +52,6 @@ import org.apache.hadoop.mapred.FileSplit;
 import org.apache.hadoop.mapred.InputSplit;
 import org.apache.hadoop.mapred.JobConf;
 import org.apache.hadoop.mapred.RecordReader;
-import org.apache.hadoop.mapred.Reporter;
 
 import java.io.IOException;
 import java.util.Arrays;
@@ -65,18 +64,11 @@ import java.util.function.UnaryOperator;
 import java.util.stream.Collectors;
 import java.util.stream.Stream;
 
-import static 
org.apache.hudi.common.model.HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
-import static 
org.apache.hudi.hadoop.utils.HoodieInputFormatUtils.getPartitionFieldNames;
-
 /**
  * {@link HoodieReaderContext} for Hive-specific {@link 
HoodieFileGroupReaderBasedRecordReader}.
  */
 public class HiveHoodieReaderContext extends 
HoodieReaderContext<ArrayWritable> {
   protected final HoodieFileGroupReaderBasedRecordReader.HiveReaderCreator 
readerCreator;
-  protected final InputSplit split;
-  protected final JobConf jobConf;
-  protected final Reporter reporter;
-  protected final Schema writerSchema;
   protected final Map<String, TypeInfo> columnTypeMap;
   private final ObjectInspectorCache objectInspectorCache;
   private RecordReader<NullWritable, ArrayWritable> firstRecordReader = null;
@@ -87,39 +79,17 @@ public class HiveHoodieReaderContext extends 
HoodieReaderContext<ArrayWritable>
   private final String recordKeyField;
 
   protected 
HiveHoodieReaderContext(HoodieFileGroupReaderBasedRecordReader.HiveReaderCreator
 readerCreator,
-                                    InputSplit split,
-                                    JobConf jobConf,
-                                    Reporter reporter,
-                                    Schema writerSchema,
-                                    HoodieTableMetaClient metaClient) {
+                                    String recordKeyField,
+                                    List<String> partitionCols,
+                                    ObjectInspectorCache objectInspectorCache) 
{
     this.readerCreator = readerCreator;
-    this.split = split;
-    this.jobConf = jobConf;
-    this.reporter = reporter;
-    this.writerSchema = writerSchema;
-    this.partitionCols = getPartitionFieldNames(jobConf).stream().filter(n -> 
writerSchema.getField(n) != null).collect(Collectors.toList());
+    this.partitionCols = partitionCols;
     this.partitionColSet = new HashSet<>(this.partitionCols);
-    String tableName = metaClient.getTableConfig().getTableName();
-    recordKeyField = getRecordKeyField(metaClient);
-    this.objectInspectorCache = 
HoodieArrayWritableAvroUtils.getCacheForTable(tableName, writerSchema, jobConf);
+    this.recordKeyField = recordKeyField;
+    this.objectInspectorCache = objectInspectorCache;
     this.columnTypeMap = objectInspectorCache.getColumnTypeMap();
   }
 
-  /**
-   * If populate meta fields is false, then getRecordKeyFields()
-   * should return exactly 1 recordkey field.
-   */
-  private static String getRecordKeyField(HoodieTableMetaClient metaClient) {
-    if (metaClient.getTableConfig().populateMetaFields()) {
-      return HoodieRecord.RECORD_KEY_METADATA_FIELD;
-    }
-
-    Option<String[]> recordKeyFieldsOpt = 
metaClient.getTableConfig().getRecordKeyFields();
-    ValidationUtils.checkArgument(recordKeyFieldsOpt.isPresent(), "No record 
key field set in table config, but populateMetaFields is disabled");
-    ValidationUtils.checkArgument(recordKeyFieldsOpt.get().length == 1, "More 
than 1 record key set in table config, but populateMetaFields is disabled");
-    return recordKeyFieldsOpt.get()[0];
-  }
-
   private void setSchemas(JobConf jobConf, Schema dataSchema, Schema 
requiredSchema) {
     List<String> dataColumnNameList = dataSchema.getFields().stream().map(f -> 
f.name().toLowerCase(Locale.ROOT)).collect(Collectors.toList());
     List<TypeInfo> dataColumnTypeList = 
dataColumnNameList.stream().map(fieldName -> {
@@ -153,14 +123,23 @@ public class HiveHoodieReaderContext extends 
HoodieReaderContext<ArrayWritable>
 
   private ClosableIterator<ArrayWritable> getFileRecordIterator(StoragePath 
filePath, String[] hosts, long start, long length, Schema dataSchema,
                                                                 Schema 
requiredSchema, HoodieStorage storage) throws IOException {
-    JobConf jobConfCopy = new JobConf(jobConf);
+    JobConf jobConfCopy = new 
JobConf(storage.getConf().unwrapAs(Configuration.class));
+    if (getNeedsBootstrapMerge()) {
+      // Hive PPD works at row-group level and only enabled when 
hive.optimize.index.filter=true;
+      // The above config is disabled by default. But when enabled, would 
cause misalignment between
+      // skeleton and bootstrap file. We will disable them specifically when 
query needs bootstrap and skeleton
+      // file to be stitched.
+      // This disables row-group filtering
+      jobConfCopy.unset(TableScanDesc.FILTER_EXPR_CONF_STR);
+      jobConfCopy.unset(ConvertAstToSearchArg.SARG_PUSHDOWN);
+    }
     //move the partition cols to the end, because in some cases it has issues 
if we don't do that
     Schema modifiedDataSchema = 
HoodieAvroUtils.generateProjectionSchema(dataSchema, 
Stream.concat(dataSchema.getFields().stream()
             .map(f -> f.name().toLowerCase(Locale.ROOT)).filter(n -> 
!partitionColSet.contains(n)),
         partitionCols.stream().filter(c -> dataSchema.getField(c) != 
null)).collect(Collectors.toList()));
     setSchemas(jobConfCopy, modifiedDataSchema, requiredSchema);
     InputSplit inputSplit = new FileSplit(new Path(filePath.toString()), 
start, length, hosts);
-    RecordReader<NullWritable, ArrayWritable> recordReader = 
readerCreator.getRecordReader(inputSplit, jobConfCopy, reporter);
+    RecordReader<NullWritable, ArrayWritable> recordReader = 
readerCreator.getRecordReader(inputSplit, jobConfCopy);
     if (firstRecordReader == null) {
       firstRecordReader = recordReader;
     }
@@ -179,10 +158,7 @@ public class HiveHoodieReaderContext extends 
HoodieReaderContext<ArrayWritable>
 
   @Override
   public HoodieRecordMerger getRecordMerger(String mergerStrategy) {
-    if (mergerStrategy.equals(DEFAULT_MERGER_STRATEGY_UUID)) {
-      return new HoodieHiveRecordMerger();
-    }
-    throw new HoodieException(String.format("The merger strategy UUID is not 
supported, Default: %s, Passed: %s", mergerStrategy, 
DEFAULT_MERGER_STRATEGY_UUID));
+    return HoodieHiveRecordMerger.getRecordMerger(mergerStrategy);
   }
 
   @Override
diff --git 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieFileGroupReaderBasedRecordReader.java
 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieFileGroupReaderBasedRecordReader.java
index f72b79b56e1..93b6ca8ad13 100644
--- 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieFileGroupReaderBasedRecordReader.java
+++ 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieFileGroupReaderBasedRecordReader.java
@@ -25,15 +25,19 @@ import org.apache.hudi.common.model.BaseFile;
 import org.apache.hudi.common.model.FileSlice;
 import org.apache.hudi.common.model.HoodieBaseFile;
 import org.apache.hudi.common.model.HoodieFileGroupId;
+import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.table.TableSchemaResolver;
 import org.apache.hudi.common.table.read.HoodieFileGroupReader;
 import org.apache.hudi.common.table.timeline.HoodieInstant;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.StringUtils;
+import org.apache.hudi.common.util.ValidationUtils;
+import org.apache.hudi.common.util.VisibleForTesting;
 import org.apache.hudi.hadoop.realtime.RealtimeSplit;
 import org.apache.hudi.hadoop.utils.HoodieRealtimeInputFormatUtils;
 import org.apache.hudi.hadoop.utils.HoodieRealtimeRecordReaderUtils;
+import org.apache.hudi.hadoop.utils.ObjectInspectorCache;
 
 import org.apache.avro.Schema;
 import org.apache.hadoop.fs.FileSystem;
@@ -47,7 +51,6 @@ import org.apache.hadoop.mapred.FileSplit;
 import org.apache.hadoop.mapred.InputSplit;
 import org.apache.hadoop.mapred.JobConf;
 import org.apache.hadoop.mapred.RecordReader;
-import org.apache.hadoop.mapred.Reporter;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -55,6 +58,7 @@ import java.io.IOException;
 import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashSet;
+import java.util.List;
 import java.util.Locale;
 import java.util.Set;
 import java.util.function.UnaryOperator;
@@ -84,8 +88,7 @@ public class HoodieFileGroupReaderBasedRecordReader 
implements RecordReader<Null
   public interface HiveReaderCreator {
     org.apache.hadoop.mapred.RecordReader<NullWritable, ArrayWritable> 
getRecordReader(
         final org.apache.hadoop.mapred.InputSplit split,
-        final org.apache.hadoop.mapred.JobConf job,
-        final org.apache.hadoop.mapred.Reporter reporter
+        final org.apache.hadoop.mapred.JobConf job
     ) throws IOException;
   }
 
@@ -99,8 +102,7 @@ public class HoodieFileGroupReaderBasedRecordReader 
implements RecordReader<Null
 
   public HoodieFileGroupReaderBasedRecordReader(HiveReaderCreator 
readerCreator,
                                                 final InputSplit split,
-                                                final JobConf jobConf,
-                                                final Reporter reporter) 
throws IOException {
+                                                final JobConf jobConf) throws 
IOException {
     this.jobConfCopy = new JobConf(jobConf);
     HoodieRealtimeInputFormatUtils.cleanProjectionColumnIds(jobConfCopy);
     Set<String> partitionColumns = new 
HashSet<>(getPartitionFieldNames(jobConfCopy));
@@ -115,7 +117,9 @@ public class HoodieFileGroupReaderBasedRecordReader 
implements RecordReader<Null
     String latestCommitTime = getLatestCommitTime(split, metaClient);
     Schema tableSchema = getLatestTableSchema(metaClient, jobConfCopy, 
latestCommitTime);
     Schema requestedSchema = createRequestedSchema(tableSchema, jobConfCopy);
-    this.readerContext = new HiveHoodieReaderContext(readerCreator, split, 
jobConfCopy, reporter, tableSchema, metaClient);
+    this.readerContext = new HiveHoodieReaderContext(readerCreator, 
getRecordKeyField(metaClient),
+        getStoredPartitionFieldNames(jobConfCopy, tableSchema),
+        new ObjectInspectorCache(tableSchema, jobConfCopy));
     this.arrayWritable = new ArrayWritable(Writable.class, new 
Writable[requestedSchema.getFields().size()]);
     TypedProperties props = metaClient.getTableConfig().getProps();
     jobConf.forEach(e -> {
@@ -171,6 +175,30 @@ public class HoodieFileGroupReaderBasedRecordReader 
implements RecordReader<Null
     return readerContext.getProgress();
   }
 
+  /**
+   * If populate meta fields is false, then getRecordKeyFields()
+   * should return exactly 1 recordkey field.
+   */
+  @VisibleForTesting
+  static String getRecordKeyField(HoodieTableMetaClient metaClient) {
+    if (metaClient.getTableConfig().populateMetaFields()) {
+      return HoodieRecord.RECORD_KEY_METADATA_FIELD;
+    }
+
+    Option<String[]> recordKeyFieldsOpt = 
metaClient.getTableConfig().getRecordKeyFields();
+    ValidationUtils.checkArgument(recordKeyFieldsOpt.isPresent(), "No record 
key field set in table config, but populateMetaFields is disabled");
+    ValidationUtils.checkArgument(recordKeyFieldsOpt.get().length == 1, "More 
than 1 record key set in table config, but populateMetaFields is disabled");
+    return recordKeyFieldsOpt.get()[0];
+  }
+
+  /**
+   * List of partition fields that are actually written to the file
+   */
+  @VisibleForTesting
+  static List<String> getStoredPartitionFieldNames(JobConf jobConf, Schema 
writerSchema) {
+    return getPartitionFieldNames(jobConf).stream().filter(n -> 
writerSchema.getField(n) != null).collect(Collectors.toList());
+  }
+
   public RealtimeSplit getSplit() {
     return (RealtimeSplit) inputSplit;
   }
@@ -203,7 +231,7 @@ public class HoodieFileGroupReaderBasedRecordReader 
implements RecordReader<Null
   }
 
   /**
-   * Convert FileSplit to FileSlice, but save the locations in 'hosts' because 
that data is otherwise lost.
+   * Convert FileSplit to FileSlice
    */
   private static FileSlice getFileSliceFromSplit(FileSplit split, FileSystem 
fs, String tableBasePath) throws IOException {
     BaseFile bootstrapBaseFile = createBootstrapBaseFile(split, fs);
diff --git 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
index 17a4738569e..e35bc68fada 100644
--- 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
+++ 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
@@ -19,53 +19,24 @@
 
 package org.apache.hudi.hadoop;
 
-import org.apache.hudi.common.config.TypedProperties;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.HoodieRecordMerger;
-import org.apache.hudi.common.util.Option;
-import org.apache.hudi.common.util.ValidationUtils;
-import org.apache.hudi.common.util.collection.Pair;
-
-import org.apache.avro.Schema;
-
-import java.io.IOException;
-
-public class HoodieHiveRecordMerger implements HoodieRecordMerger {
-  @Override
-  public Option<Pair<HoodieRecord, Schema>> merge(HoodieRecord older, Schema 
oldSchema, HoodieRecord newer, Schema newSchema, TypedProperties props) throws 
IOException {
-    ValidationUtils.checkArgument(older.getRecordType() == 
HoodieRecord.HoodieRecordType.HIVE);
-    ValidationUtils.checkArgument(newer.getRecordType() == 
HoodieRecord.HoodieRecordType.HIVE);
-    if (newer instanceof HoodieHiveRecord) {
-      HoodieHiveRecord newHiveRecord = (HoodieHiveRecord) newer;
-      if (newHiveRecord.isDeleted()) {
-        return Option.empty();
-      }
-    } else if (newer.getData() == null) {
-      return Option.empty();
-    }
-
-    if (older instanceof HoodieHiveRecord) {
-      HoodieHiveRecord oldHiveRecord = (HoodieHiveRecord) older;
-      if (oldHiveRecord.isDeleted()) {
-        return Option.of(Pair.of(newer, newSchema));
-      }
-    } else if (older.getData() == null) {
-      return Option.empty();
-    }
-    if (older.getOrderingValue(oldSchema, 
props).compareTo(newer.getOrderingValue(newSchema, props)) > 0) {
-      return Option.of(Pair.of(older, oldSchema));
-    } else {
-      return Option.of(Pair.of(newer, newSchema));
-    }
-  }
+import org.apache.hudi.exception.HoodieException;
 
+abstract class HoodieHiveRecordMerger implements HoodieRecordMerger {
   @Override
   public HoodieRecord.HoodieRecordType getRecordType() {
     return HoodieRecord.HoodieRecordType.HIVE;
   }
 
-  @Override
-  public String getMergingStrategy() {
-    return HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
+  static HoodieRecordMerger getRecordMerger(String mergerStrategy) {
+    switch (mergerStrategy) {
+      case DEFAULT_MERGER_STRATEGY_UUID:
+        return new DefaultHiveRecordMerger();
+      case OVERWRITE_MERGER_STRATEGY_UUID:
+        return new OverwriteWithLatestHiveRecordMerger();
+      default:
+        throw new HoodieException("This merger strategy UUID is not supported: 
" + mergerStrategy);
+    }
   }
 }
diff --git 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieParquetInputFormat.java
 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieParquetInputFormat.java
index 18b9e221978..ed044ef7ddb 100644
--- 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieParquetInputFormat.java
+++ 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieParquetInputFormat.java
@@ -121,15 +121,15 @@ public class HoodieParquetInputFormat extends 
HoodieParquetInputFormatBase {
           return super.getRecordReader(split, job, reporter);
         }
         if (supportAvroRead && 
HoodieColumnProjectionUtils.supportTimestamp(job)) {
-          return new HoodieFileGroupReaderBasedRecordReader((s, j, r) -> {
+          return new HoodieFileGroupReaderBasedRecordReader((s, j) -> {
             try {
-              return new ParquetRecordReaderWrapper(new 
HoodieTimestampAwareParquetInputFormat(), s, j, r);
+              return new ParquetRecordReaderWrapper(new 
HoodieTimestampAwareParquetInputFormat(), s, j, reporter);
             } catch (InterruptedException e) {
               throw new RuntimeException(e);
             }
-          }, split, job, reporter);
+          }, split, job);
         } else {
-          return new 
HoodieFileGroupReaderBasedRecordReader(super::getRecordReader, split, job, 
reporter);
+          return new HoodieFileGroupReaderBasedRecordReader((s, j) -> 
super.getRecordReader(s, j, reporter), split, job);
         }
       } catch (final IOException e) {
         throw new RuntimeException("Cannot create a RecordReaderWrapper", e);
diff --git 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkMerger.java
 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/OverwriteWithLatestHiveRecordMerger.java
similarity index 89%
rename from 
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkMerger.java
rename to 
hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/OverwriteWithLatestHiveRecordMerger.java
index 611f045f645..4a4acb18a36 100644
--- 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkMerger.java
+++ 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/OverwriteWithLatestHiveRecordMerger.java
@@ -17,7 +17,7 @@
  * under the License.
  */
 
-package org.apache.hudi;
+package org.apache.hudi.hadoop;
 
 import org.apache.hudi.common.config.TypedProperties;
 import org.apache.hudi.common.model.HoodieRecord;
@@ -29,10 +29,9 @@ import org.apache.avro.Schema;
 import java.io.IOException;
 
 /**
- * Spark merger that always chooses the newer record
+ * Hive merger that always chooses the newer record
  */
-public class OverwriteWithLatestSparkMerger extends HoodieSparkRecordMerger {
-
+public class OverwriteWithLatestHiveRecordMerger extends 
HoodieHiveRecordMerger {
   @Override
   public String getMergingStrategy() {
     return OVERWRITE_MERGER_STRATEGY_UUID;
@@ -42,5 +41,4 @@ public class OverwriteWithLatestSparkMerger extends 
HoodieSparkRecordMerger {
   public Option<Pair<HoodieRecord, Schema>> merge(HoodieRecord older, Schema 
oldSchema, HoodieRecord newer, Schema newSchema, TypedProperties props) throws 
IOException {
     return Option.of(Pair.of(newer, newSchema));
   }
-
 }
diff --git 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieArrayWritableAvroUtils.java
 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieArrayWritableAvroUtils.java
index a2da796c6f7..3657c6bec3e 100644
--- 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieArrayWritableAvroUtils.java
+++ 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieArrayWritableAvroUtils.java
@@ -27,24 +27,12 @@ import org.apache.avro.Schema;
 import org.apache.hadoop.hive.ql.io.parquet.serde.ArrayWritableObjectInspector;
 import org.apache.hadoop.io.ArrayWritable;
 import org.apache.hadoop.io.Writable;
-import org.apache.hadoop.mapred.JobConf;
 
 import java.util.List;
 import java.util.function.UnaryOperator;
 
 public class HoodieArrayWritableAvroUtils {
 
-  private static final Cache<String, ObjectInspectorCache>
-      OBJECT_INSPECTOR_TABLE_CACHE = 
Caffeine.newBuilder().maximumSize(1000).build();
-
-  public static ObjectInspectorCache getCacheForTable(String table, Schema 
tableSchema, JobConf jobConf) {
-    ObjectInspectorCache cache = 
OBJECT_INSPECTOR_TABLE_CACHE.getIfPresent(table);
-    if (cache == null) {
-      cache = new ObjectInspectorCache(tableSchema, jobConf);
-    }
-    return cache;
-  }
-
   private static final Cache<Pair<Schema, Schema>, int[]>
       PROJECTION_CACHE = Caffeine.newBuilder().maximumSize(1000).build();
 
diff --git 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieRealtimeRecordReaderUtils.java
 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieRealtimeRecordReaderUtils.java
index 07863ac5b97..32dd8920742 100644
--- 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieRealtimeRecordReaderUtils.java
+++ 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieRealtimeRecordReaderUtils.java
@@ -216,7 +216,7 @@ public class HoodieRealtimeRecordReaderUtils {
         }
         return new ArrayWritable(Writable.class, recordValues);
       case ENUM:
-        return new Text(value.toString());
+        return new BytesWritable(value.toString().getBytes());
       case ARRAY:
         GenericArray arrayValue = (GenericArray) value;
         Writable[] arrayValues = new Writable[arrayValue.size()];
diff --git 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/ObjectInspectorCache.java
 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/ObjectInspectorCache.java
index ddcc28851df..c5bb8b4fa6f 100644
--- 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/ObjectInspectorCache.java
+++ 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/ObjectInspectorCache.java
@@ -19,8 +19,6 @@
 
 package org.apache.hudi.hadoop.utils;
 
-import com.github.benmanes.caffeine.cache.Cache;
-import com.github.benmanes.caffeine.cache.Caffeine;
 import org.apache.avro.Schema;
 import org.apache.hadoop.hive.ql.io.parquet.serde.ArrayWritableObjectInspector;
 import org.apache.hadoop.hive.serde.serdeConstants;
@@ -46,8 +44,7 @@ import java.util.stream.IntStream;
  */
 public class ObjectInspectorCache {
   private final Map<String, TypeInfo> columnTypeMap = new HashMap<>();
-  private final Cache<Schema, ArrayWritableObjectInspector>
-      objectInspectorCache = Caffeine.newBuilder().maximumSize(1000).build();
+  private final Map<Schema, ArrayWritableObjectInspector> objectInspectorCache 
= new HashMap<>();
 
   public Map<String, TypeInfo> getColumnTypeMap() {
     return columnTypeMap;
@@ -82,7 +79,7 @@ public class ObjectInspectorCache {
   }
 
   public ArrayWritableObjectInspector getObjectInspector(Schema schema) {
-    return objectInspectorCache.get(schema, s -> {
+    return objectInspectorCache.computeIfAbsent(schema, s -> {
       List<String> columnNameList = 
s.getFields().stream().map(Schema.Field::name).collect(Collectors.toList());
       List<TypeInfo> columnTypeList = 
columnNameList.stream().map(columnTypeMap::get).collect(Collectors.toList());
       StructTypeInfo rowTypeInfo = (StructTypeInfo) 
TypeInfoFactory.getStructTypeInfo(columnNameList, columnTypeList);
@@ -92,12 +89,7 @@ public class ObjectInspectorCache {
   }
 
   public Object getValue(ArrayWritable record, Schema schema, String 
fieldName) {
-    try {
-      ArrayWritableObjectInspector objectInspector = 
getObjectInspector(schema);
-      return objectInspector.getStructFieldData(record, 
objectInspector.getStructFieldRef(fieldName));
-    } catch (Exception e) {
-      throw e;
-    }
-
+    ArrayWritableObjectInspector objectInspector = getObjectInspector(schema);
+    return objectInspector.getStructFieldData(record, 
objectInspector.getStructFieldRef(fieldName));
   }
 }
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/command/HoodieSparkValidateDuplicateKeyRecordMerger.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/command/HoodieSparkValidateDuplicateKeyRecordMerger.scala
index 8a64007f053..c157f58f516 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/command/HoodieSparkValidateDuplicateKeyRecordMerger.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/command/HoodieSparkValidateDuplicateKeyRecordMerger.scala
@@ -18,10 +18,10 @@
 package org.apache.spark.sql.hudi.command
 
 import org.apache.avro.Schema
-import org.apache.hudi.HoodieSparkRecordMerger
+import org.apache.hudi.{DefaultSparkRecordMerger, HoodieSparkRecordMerger}
 import org.apache.hudi.common.config.TypedProperties
 import org.apache.hudi.common.model.{HoodieRecord, HoodieRecordMerger, 
OperationModeAwareness}
-import org.apache.hudi.common.util.{collection, HoodieRecordUtils, Option => 
HOption}
+import org.apache.hudi.common.util.{HoodieRecordUtils, collection, Option => 
HOption}
 import org.apache.hudi.exception.HoodieDuplicateKeyException
 
 /**
@@ -37,6 +37,11 @@ class HoodieSparkValidateDuplicateKeyRecordMerger extends 
HoodieSparkRecordMerge
   }
 
   override def asPreCombiningMode(): HoodieRecordMerger = {
-    
HoodieRecordUtils.loadRecordMerger(classOf[HoodieSparkRecordMerger].getName)
+    
HoodieRecordUtils.loadRecordMerger(classOf[DefaultSparkRecordMerger].getName)
   }
+
+  /**
+   * The kind of merging strategy this recordMerger belongs to. An UUID 
represents merging strategy.
+   */
+  override def getMergingStrategy: String = 
HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID
 }
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodieMergeHandleWithSparkMerger.java
 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodieMergeHandleWithSparkMerger.java
index 620bf4beb13..b9309632743 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodieMergeHandleWithSparkMerger.java
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodieMergeHandleWithSparkMerger.java
@@ -160,7 +160,7 @@ public class TestHoodieMergeHandleWithSparkMerger extends 
SparkClientFunctionalT
     Properties extraProperties = new Properties();
     extraProperties.setProperty(
         RECORD_MERGER_IMPLS.key(),
-        "org.apache.hudi.HoodieSparkRecordMerger");
+        "org.apache.hudi.DefaultSparkRecordMerger");
     extraProperties.setProperty(
         LOGFILE_DATA_BLOCK_FORMAT.key(),
         "parquet");
@@ -229,7 +229,7 @@ public class TestHoodieMergeHandleWithSparkMerger extends 
SparkClientFunctionalT
     Map<String, String> properties = new HashMap<>();
     properties.put(
         RECORD_MERGER_IMPLS.key(),
-        "org.apache.hudi.HoodieSparkRecordMerger");
+        "org.apache.hudi.DefaultSparkRecordMerger");
     properties.put(
         LOGFILE_DATA_BLOCK_FORMAT.key(),
         "parquet");
@@ -361,21 +361,21 @@ public class TestHoodieMergeHandleWithSparkMerger extends 
SparkClientFunctionalT
     }
   }
 
-  public static class DefaultMerger extends HoodieSparkRecordMerger {
+  public static class DefaultMerger extends DefaultSparkRecordMerger {
     @Override
     public boolean shouldFlush(HoodieRecord record, Schema schema, 
TypedProperties props) {
       return true;
     }
   }
 
-  public static class NoFlushMerger extends HoodieSparkRecordMerger {
+  public static class NoFlushMerger extends DefaultSparkRecordMerger {
     @Override
     public boolean shouldFlush(HoodieRecord record, Schema schema, 
TypedProperties props) {
       return false;
     }
   }
 
-  public static class CustomMerger extends HoodieSparkRecordMerger {
+  public static class CustomMerger extends DefaultSparkRecordMerger {
     @Override
     public boolean shouldFlush(HoodieRecord record, Schema schema, 
TypedProperties props) throws IOException {
       return !((HoodieSparkRecord) 
record).getData().getString(0).equals("001");
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodiePositionBasedFileGroupRecordBuffer.java
 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodiePositionBasedFileGroupRecordBuffer.java
index edcc1f6dad2..c80a2f28ded 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodiePositionBasedFileGroupRecordBuffer.java
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodiePositionBasedFileGroupRecordBuffer.java
@@ -62,7 +62,6 @@ import java.util.stream.Collectors;
 import static 
org.apache.hudi.common.engine.HoodieReaderContext.INTERNAL_META_RECORD_KEY;
 import static org.apache.hudi.common.model.WriteOperationType.INSERT;
 import static 
org.apache.hudi.common.testutils.HoodieTestUtils.createMetaClient;
-import static 
org.apache.hudi.common.testutils.RawTripTestPayload.recordsToStrings;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNull;
@@ -90,7 +89,7 @@ public class TestHoodiePositionBasedFileGroupRecordBuffer 
extends TestHoodieFile
     writeConfigs.put("hoodie.compact.inline", "false");
     writeConfigs.put(HoodieWriteConfig.WRITE_RECORD_POSITIONS.key(), "true");
     writeConfigs.put(HoodieWriteConfig.WRITE_PAYLOAD_CLASS_NAME.key(), 
getRecordPayloadForMergeMode(mergeMode));
-    commitToTable(recordsToStrings(dataGen.generateInserts("001", 100)), 
INSERT.value(), writeConfigs);
+    commitToTable(dataGen.generateInserts("001", 100), INSERT.value(), 
writeConfigs);
 
     String[] partitionPaths = dataGen.getPartitionPaths();
     String[] partitionValues = new String[1];
@@ -115,11 +114,11 @@ public class TestHoodiePositionBasedFileGroupRecordBuffer 
extends TestHoodieFile
         ctx.setRecordMerger(new CustomMerger());
         break;
       case EVENT_TIME_ORDERING:
-        ctx.setRecordMerger(new HoodieSparkRecordMerger());
+        ctx.setRecordMerger(new DefaultSparkRecordMerger());
         break;
       case OVERWRITE_WITH_LATEST:
       default:
-        ctx.setRecordMerger(new OverwriteWithLatestSparkMerger());
+        ctx.setRecordMerger(new OverwriteWithLatestSparkRecordMerger());
         break;
     }
     ctx.setSchemaHandler(new HoodiePositionBasedSchemaHandler<>(ctx, 
avroSchema, avroSchema,
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderOnSpark.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderOnSpark.scala
index 23013db5d33..85cf2fc1cb9 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderOnSpark.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderOnSpark.scala
@@ -24,9 +24,9 @@ import org.apache.hudi.common.config.RecordMergeMode
 import org.apache.hudi.common.engine.HoodieReaderContext
 import org.apache.hudi.common.model.{DefaultHoodieRecordPayload, HoodieRecord, 
OverwriteWithLatestAvroPayload, WriteOperationType}
 import org.apache.hudi.common.table.HoodieTableMetaClient
-import org.apache.hudi.common.testutils.HoodieTestUtils
+import org.apache.hudi.common.testutils.{HoodieTestUtils, RawTripTestPayload}
 import org.apache.hudi.storage.StorageConfiguration
-import org.apache.hudi.{HoodieSparkRecordMerger, SparkAdapterSupport, 
SparkFileFormatInternalRowReaderContext}
+import org.apache.hudi.{DefaultSparkRecordMerger, SparkAdapterSupport, 
SparkFileFormatInternalRowReaderContext}
 
 import org.apache.avro.Schema
 import org.apache.hadoop.conf.Configuration
@@ -89,12 +89,13 @@ class TestHoodieFileGroupReaderOnSpark extends 
TestHoodieFileGroupReaderBase[Int
   override def getHoodieReaderContext(tablePath: String, avroSchema: Schema, 
storageConf: StorageConfiguration[_]): HoodieReaderContext[InternalRow] = {
     val reader = sparkAdapter.createParquetFileReader(vectorized = false, 
spark.sessionState.conf, Map.empty, 
storageConf.unwrapAs(classOf[Configuration]))
     val metaClient = 
HoodieTableMetaClient.builder().setConf(storageConf).setBasePath(tablePath).build
-    val recordKeyField = new 
HoodieSparkRecordMerger().getMandatoryFieldsForMerging(metaClient.getTableConfig)(0)
+    val recordKeyField = new 
DefaultSparkRecordMerger().getMandatoryFieldsForMerging(metaClient.getTableConfig)(0)
     new SparkFileFormatInternalRowReaderContext(reader, recordKeyField, 
Seq.empty, Seq.empty)
   }
 
-  override def commitToTable(recordList: util.List[String], operation: String, 
options: util.Map[String, String]): Unit = {
-    val inputDF: Dataset[Row] = 
spark.read.json(spark.sparkContext.parallelize(recordList.asScala.toList, 2))
+  override def commitToTable(recordList: util.List[HoodieRecord[_]], 
operation: String, options: util.Map[String, String]): Unit = {
+    val recs = RawTripTestPayload.recordsToStrings(recordList)
+    val inputDF: Dataset[Row] = 
spark.read.json(spark.sparkContext.parallelize(recs.asScala.toList, 2))
 
     inputDF.write.format("hudi")
       .options(options)
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestSpark35RecordPositionMetadataColumn.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestSpark35RecordPositionMetadataColumn.scala
index 59dc89c5a21..87bdd4d3685 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestSpark35RecordPositionMetadataColumn.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestSpark35RecordPositionMetadataColumn.scala
@@ -39,7 +39,7 @@ import org.junit.jupiter.api.{BeforeEach, Test}
 
 class TestSpark35RecordPositionMetadataColumn extends 
SparkClientFunctionalTestHarness {
   private val PARQUET_FORMAT = "parquet"
-  private val SPARK_MERGER = "org.apache.hudi.HoodieSparkRecordMerger"
+  private val SPARK_MERGER = "org.apache.hudi.DefaultSparkRecordMerger"
 
   @BeforeEach
   def setUp(): Unit = {
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/CommonOptionUtils.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/CommonOptionUtils.scala
index dffe06d4611..45248bbc9ae 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/CommonOptionUtils.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/CommonOptionUtils.scala
@@ -23,7 +23,7 @@ import org.apache.hudi.common.config.HoodieMetadataConfig
 import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType
 import org.apache.hudi.common.table.HoodieTableConfig
 import org.apache.hudi.config.HoodieWriteConfig
-import org.apache.hudi.{DataSourceReadOptions, DataSourceWriteOptions, 
HoodieSparkRecordMerger}
+import org.apache.hudi.{DataSourceReadOptions, DataSourceWriteOptions, 
DefaultSparkRecordMerger}
 
 object CommonOptionUtils {
 
@@ -39,7 +39,7 @@ object CommonOptionUtils {
     HoodieWriteConfig.TBL_NAME.key -> "hoodie_test",
     HoodieMetadataConfig.COMPACT_NUM_DELTA_COMMITS.key -> "1"
   )
-  val sparkOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[HoodieSparkRecordMerger].getName)
+  val sparkOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[DefaultSparkRecordMerger].getName)
 
   def getWriterReaderOpts(recordType: HoodieRecordType,
                           opt: Map[String, String] = commonOpts,
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestCOWDataSource.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestCOWDataSource.scala
index 593034865a3..af98f85db6a 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestCOWDataSource.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestCOWDataSource.scala
@@ -1762,7 +1762,7 @@ class TestCOWDataSource extends HoodieSparkClientTestBase 
with ScalaAssertionSup
       HoodieWriteConfig.KEYGENERATOR_CLASS_NAME.key() -> 
"org.apache.hudi.keygen.ComplexKeyGenerator",
       KeyGeneratorOptions.HIVE_STYLE_PARTITIONING_ENABLE.key() -> "true",
       HiveSyncConfigHolder.HIVE_SYNC_ENABLED.key() -> "false",
-      HoodieWriteConfig.RECORD_MERGER_IMPLS.key() -> 
"org.apache.hudi.HoodieSparkRecordMerger"
+      HoodieWriteConfig.RECORD_MERGER_IMPLS.key() -> 
"org.apache.hudi.DefaultSparkRecordMerger"
     )
     
df1.write.format("hudi").options(hudiOptions).mode(SaveMode.Append).save(basePath)
 
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestDataSourceForBootstrap.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestDataSourceForBootstrap.scala
index fe3373f1a9f..5209cc3a80c 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestDataSourceForBootstrap.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestDataSourceForBootstrap.scala
@@ -28,7 +28,7 @@ import 
org.apache.hudi.functional.TestDataSourceForBootstrap.{dropMetaCols, sort
 import org.apache.hudi.hadoop.fs.HadoopFSUtils
 import org.apache.hudi.keygen.{NonpartitionedKeyGenerator, SimpleKeyGenerator}
 import org.apache.hudi.testutils.HoodieClientTestUtils
-import org.apache.hudi.{DataSourceReadOptions, DataSourceWriteOptions, 
HoodieDataSourceHelpers, HoodieSparkRecordMerger}
+import org.apache.hudi.{DataSourceReadOptions, DataSourceWriteOptions, 
HoodieDataSourceHelpers, DefaultSparkRecordMerger}
 
 import org.apache.hadoop.fs.{FileSystem, Path}
 import org.apache.spark.api.java.JavaSparkContext
@@ -62,7 +62,7 @@ class TestDataSourceForBootstrap {
   )
 
   val sparkRecordTypeOpts = Map(
-    HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[HoodieSparkRecordMerger].getName,
+    HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[DefaultSparkRecordMerger].getName,
     HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet"
   )
 
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestHoodieMultipleBaseFileFormat.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestHoodieMultipleBaseFileFormat.scala
index 5b4e7950708..f67ac29f8a7 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestHoodieMultipleBaseFileFormat.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestHoodieMultipleBaseFileFormat.scala
@@ -26,7 +26,7 @@ import 
org.apache.hudi.common.testutils.HoodieTestDataGenerator.{DEFAULT_FIRST_P
 import org.apache.hudi.common.testutils.RawTripTestPayload.recordsToStrings
 import org.apache.hudi.config.HoodieWriteConfig
 import org.apache.hudi.testutils.HoodieSparkClientTestBase
-import org.apache.hudi.{DataSourceWriteOptions, HoodieSparkRecordMerger, 
SparkDatasetMixin}
+import org.apache.hudi.{DataSourceWriteOptions, DefaultSparkRecordMerger, 
SparkDatasetMixin}
 
 import org.apache.spark.sql.{Dataset, Row, SaveMode, SparkSession}
 import org.junit.jupiter.api.Assertions.assertEquals
@@ -52,7 +52,7 @@ class TestHoodieMultipleBaseFileFormat extends 
HoodieSparkClientTestBase with Sp
     HoodieWriteConfig.TBL_NAME.key -> "hoodie_test"
   )
   val sparkOpts = Map(
-    HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[HoodieSparkRecordMerger].getName,
+    HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[DefaultSparkRecordMerger].getName,
     HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet"
   )
 
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestMORDataSource.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestMORDataSource.scala
index f9625f4aa8f..5ca998a5e3a 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestMORDataSource.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestMORDataSource.scala
@@ -35,7 +35,7 @@ import org.apache.hudi.storage.StoragePath
 import org.apache.hudi.table.action.compact.CompactionTriggerStrategy
 import org.apache.hudi.testutils.{DataSourceTestUtils, 
HoodieSparkClientTestBase}
 import org.apache.hudi.util.JFunction
-import org.apache.hudi.{DataSourceReadOptions, DataSourceUtils, 
DataSourceWriteOptions, HoodieDataSourceHelpers, HoodieSparkRecordMerger, 
SparkDatasetMixin}
+import org.apache.hudi.{DataSourceReadOptions, DataSourceUtils, 
DataSourceWriteOptions, HoodieDataSourceHelpers, DefaultSparkRecordMerger, 
SparkDatasetMixin}
 
 import org.apache.hadoop.fs.Path
 import org.apache.spark.sql._
@@ -68,7 +68,7 @@ class TestMORDataSource extends HoodieSparkClientTestBase 
with SparkDatasetMixin
     HoodieWriteConfig.TBL_NAME.key -> "hoodie_test"
   )
   val sparkOpts = Map(
-    HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[HoodieSparkRecordMerger].getName,
+    HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[DefaultSparkRecordMerger].getName,
     HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet"
   )
 
@@ -1019,7 +1019,7 @@ class TestMORDataSource extends HoodieSparkClientTestBase 
with SparkDatasetMixin
       "hoodie.datasource.write.row.writer.enable" -> "false"
     )
     if (recordType.equals(HoodieRecordType.SPARK)) {
-      writeOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[HoodieSparkRecordMerger].getName,
+      writeOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[DefaultSparkRecordMerger].getName,
         HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet") ++ 
writeOpts
     }
     val records1 = recordsToStrings(dataGen.generateInserts("001", 
10)).asScala.toSeq
@@ -1061,7 +1061,7 @@ class TestMORDataSource extends HoodieSparkClientTestBase 
with SparkDatasetMixin
       "hoodie.datasource.write.row.writer.enable" -> "false"
     )
     if (recordType.equals(HoodieRecordType.SPARK)) {
-      writeOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[HoodieSparkRecordMerger].getName,
+      writeOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[DefaultSparkRecordMerger].getName,
         HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet") ++ 
writeOpts
     }
     val records1 = recordsToStrings(dataGen.generateInserts("001", 
10)).asScala.toSeq
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestSparkDataSourceDAGExecution.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestSparkDataSourceDAGExecution.scala
index fff0046d56f..8c1591e6f78 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestSparkDataSourceDAGExecution.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestSparkDataSourceDAGExecution.scala
@@ -25,7 +25,7 @@ import org.apache.hudi.common.util.Option
 import org.apache.hudi.config.HoodieWriteConfig
 import org.apache.hudi.testutils.HoodieSparkClientTestBase
 import org.apache.hudi.util.JFunction
-import org.apache.hudi.{DataSourceWriteOptions, HoodieSparkRecordMerger, 
ScalaAssertionSupport}
+import org.apache.hudi.{DataSourceWriteOptions, DefaultSparkRecordMerger, 
ScalaAssertionSupport}
 
 import org.apache.hadoop.fs.FileSystem
 import org.apache.spark.scheduler.{SparkListener, SparkListenerStageCompleted}
@@ -58,7 +58,7 @@ class TestSparkDataSourceDAGExecution extends 
HoodieSparkClientTestBase with Sca
     HoodieWriteConfig.TBL_NAME.key -> "hoodie_test",
     HoodieMetadataConfig.ENABLE.key -> "false"
   )
-  val sparkOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[HoodieSparkRecordMerger].getName)
+  val sparkOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[DefaultSparkRecordMerger].getName)
 
   val verificationCol: String = "driver"
   val updatedVerificationVal: String = "driver_update"
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/execution/benchmark/ReadAndWriteWithoutAvroBenchmark.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/execution/benchmark/ReadAndWriteWithoutAvroBenchmark.scala
index d7f7cf6a9f3..266740a4877 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/execution/benchmark/ReadAndWriteWithoutAvroBenchmark.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/execution/benchmark/ReadAndWriteWithoutAvroBenchmark.scala
@@ -21,7 +21,7 @@ package org.apache.spark.sql.execution.benchmark
 import org.apache.hudi.common.config.HoodieStorageConfig
 import org.apache.hudi.common.model.HoodieAvroRecordMerger
 import org.apache.hudi.config.{HoodieCompactionConfig, HoodieWriteConfig}
-import org.apache.hudi.{HoodieSparkRecordMerger, HoodieSparkUtils}
+import org.apache.hudi.{DefaultSparkRecordMerger, HoodieSparkUtils}
 
 import org.apache.hadoop.fs.Path
 import org.apache.spark.SparkConf
@@ -116,12 +116,12 @@ object ReadAndWriteWithoutAvroBenchmark extends 
HoodieBenchmarkBase {
    *  pref insert overwrite:                               Best Time(ms)   Avg 
Time(ms)   Stdev(ms)    Rate(M/s)   Per Row(ns)   Relative
    *  
-----------------------------------------------------------------------------------------------------------------------------------
    *  org.apache.hudi.common.model.HoodieAvroRecordMerger          16714       
   17107         353          0.1       16714.5       1.0X
-   *  org.apache.hudi.HoodieSparkRecordMerger                      12654       
   13924        1100          0.1       12653.8       1.3X
+   *  org.apache.hudi.DefaultSparkRecordMerger                      12654      
    13924        1100          0.1       12653.8       1.3X
    */
   private def overwriteBenchmark(): Unit = {
     val df = createComplexDataFrame(1000000)
     val benchmark = new HoodieBenchmark("pref insert overwrite", 1000000, 3)
-    Seq(classOf[HoodieAvroRecordMerger].getName, 
classOf[HoodieSparkRecordMerger].getName).zip(Seq(avroTable, 
sparkTable)).foreach {
+    Seq(classOf[HoodieAvroRecordMerger].getName, 
classOf[DefaultSparkRecordMerger].getName).zip(Seq(avroTable, 
sparkTable)).foreach {
       case (merger, tableName) => benchmark.addCase(merger) { _ =>
         withTempDir { f =>
           prepareHoodieTable(tableName, new Path(f.getCanonicalPath, 
tableName).toUri.toString, "mor", merger, df)
@@ -137,18 +137,18 @@ object ReadAndWriteWithoutAvroBenchmark extends 
HoodieBenchmarkBase {
    * pref upsert:                                         Best Time(ms)   Avg 
Time(ms)   Stdev(ms)    Rate(M/s)   Per Row(ns)   Relative
    * 
-----------------------------------------------------------------------------------------------------------------------------------
    * org.apache.hudi.common.model.HoodieAvroRecordMerger           6108        
   6383         257          0.0      610785.6       1.0X
-   * org.apache.hudi.HoodieSparkRecordMerger                       4833        
   5468         614          0.0      483300.0       1.3X
+   * org.apache.hudi.DefaultSparkRecordMerger                       4833       
    5468         614          0.0      483300.0       1.3X
    *
    * Java HotSpot(TM) 64-Bit Server VM 1.8.0_211-b12 on Mac OS X 10.16
    * Intel(R) Core(TM) i7-9750H CPU @ 2.60GHz
    * pref read:                                           Best Time(ms)   Avg 
Time(ms)   Stdev(ms)    Rate(M/s)   Per Row(ns)   Relative
    * 
-----------------------------------------------------------------------------------------------------------------------------------
    * org.apache.hudi.common.model.HoodieAvroRecordMerger            813        
    818           8          0.0       81302.1       1.0X
-   * org.apache.hudi.HoodieSparkRecordMerger                        604        
    616          18          0.0       60430.1       1.3X
+   * org.apache.hudi.DefaultSparkRecordMerger                        604       
     616          18          0.0       60430.1       1.3X
    */
   private def upsertThenReadBenchmark(): Unit = {
     val avroMergerImpl = classOf[HoodieAvroRecordMerger].getName
-    val sparkMergerImpl = classOf[HoodieSparkRecordMerger].getName
+    val sparkMergerImpl = classOf[DefaultSparkRecordMerger].getName
     val df = createComplexDataFrame(10000)
     withTempDir { avroPath =>
       withTempDir { sparkPath =>
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/HoodieSparkSqlTestBase.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/HoodieSparkSqlTestBase.scala
index 432d2ab767e..cb558b66b94 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/HoodieSparkSqlTestBase.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/HoodieSparkSqlTestBase.scala
@@ -17,7 +17,7 @@
 
 package org.apache.spark.sql.hudi.common
 
-import org.apache.hudi.HoodieSparkRecordMerger
+import org.apache.hudi.DefaultSparkRecordMerger
 import org.apache.hudi.common.config.HoodieStorageConfig
 import org.apache.hudi.common.model.HoodieAvroRecordMerger
 import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType
@@ -224,7 +224,7 @@ class HoodieSparkSqlTestBase extends FunSuite with 
BeforeAndAfterAll {
     // TODO HUDI-5264 Test parquet log with avro record in spark sql test
     recordTypes.foreach { recordType =>
       val (merger, format) = recordType match {
-        case HoodieRecordType.SPARK => 
(classOf[HoodieSparkRecordMerger].getName, "parquet")
+        case HoodieRecordType.SPARK => 
(classOf[DefaultSparkRecordMerger].getName, "parquet")
         case _ => (classOf[HoodieAvroRecordMerger].getName, "avro")
       }
       val config = Map(
@@ -240,7 +240,7 @@ class HoodieSparkSqlTestBase extends FunSuite with 
BeforeAndAfterAll {
 
   protected def getRecordType(): HoodieRecordType = {
     val merger = 
spark.sessionState.conf.getConfString(HoodieWriteConfig.RECORD_MERGER_IMPLS.key,
 HoodieWriteConfig.RECORD_MERGER_IMPLS.defaultValue())
-    if (merger.equals(classOf[HoodieSparkRecordMerger].getName)) {
+    if (merger.equals(classOf[DefaultSparkRecordMerger].getName)) {
       HoodieRecordType.SPARK
     } else {
       HoodieRecordType.AVRO
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/ddl/TestSpark3DDL.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/ddl/TestSpark3DDL.scala
index fb39e8eaef7..c4b67799530 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/ddl/TestSpark3DDL.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/ddl/TestSpark3DDL.scala
@@ -26,7 +26,7 @@ import org.apache.hudi.config.HoodieWriteConfig
 import org.apache.hudi.index.inmemory.HoodieInMemoryHashIndex
 import org.apache.hudi.testutils.DataSourceTestUtils
 import org.apache.hudi.testutils.HoodieClientTestUtils.createMetaClient
-import org.apache.hudi.{DataSourceWriteOptions, HoodieSparkRecordMerger, 
HoodieSparkUtils, QuickstartUtils}
+import org.apache.hudi.{DataSourceWriteOptions, DefaultSparkRecordMerger, 
HoodieSparkUtils, QuickstartUtils}
 
 import org.apache.hadoop.fs.Path
 import org.apache.spark.sql.catalyst.TableIdentifier
@@ -698,7 +698,7 @@ class TestSpark3DDL extends HoodieSparkSqlTestBase {
   }
 
   val sparkOpts = Map(
-    HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[HoodieSparkRecordMerger].getName,
+    HoodieWriteConfig.RECORD_MERGER_IMPLS.key -> 
classOf[DefaultSparkRecordMerger].getName,
     HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet"
   )
 
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java
index e7f3acd6fc5..cd93bdbd9ce 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java
@@ -21,7 +21,7 @@ package org.apache.hudi.utilities.deltastreamer;
 
 import org.apache.hudi.DataSourceReadOptions;
 import org.apache.hudi.DataSourceWriteOptions;
-import org.apache.hudi.HoodieSparkRecordMerger;
+import org.apache.hudi.DefaultSparkRecordMerger;
 import org.apache.hudi.HoodieSparkUtils$;
 import org.apache.hudi.client.SparkRDDWriteClient;
 import org.apache.hudi.client.transaction.lock.InProcessLockProvider;
@@ -183,7 +183,7 @@ public class TestHoodieDeltaStreamer extends 
HoodieDeltaStreamerTestBase {
   private void addRecordMerger(HoodieRecordType type, List<String> 
hoodieConfig) {
     if (type == HoodieRecordType.SPARK) {
       Map<String, String> opts = new HashMap<>();
-      opts.put(HoodieWriteConfig.RECORD_MERGER_IMPLS.key(), 
HoodieSparkRecordMerger.class.getName());
+      opts.put(HoodieWriteConfig.RECORD_MERGER_IMPLS.key(), 
DefaultSparkRecordMerger.class.getName());
       opts.put(HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key(), "parquet");
       for (Map.Entry<String, String> entry : opts.entrySet()) {
         hoodieConfig.add(String.format("%s=%s", entry.getKey(), 
entry.getValue()));

Reply via email to