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

danny0405 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 6ea3e7dff95d [HUDI-8635] Fix stats for total records written and num 
inserts in FileGroupReader based compaction (#13556)
6ea3e7dff95d is described below

commit 6ea3e7dff95dd31a8523ba8bb4e21fb286b76f7b
Author: Tim Brown <[email protected]>
AuthorDate: Mon Jul 21 23:34:01 2025 -0400

    [HUDI-8635] Fix stats for total records written and num inserts in 
FileGroupReader based compaction (#13556)
---
 .../hudi/io/FileGroupReaderBasedAppendHandle.java  |  37 ++--
 .../hudi/io/FileGroupReaderBasedMergeHandle.java   |  32 +--
 .../org/apache/hudi/io/HoodieAppendHandle.java     |   4 +
 .../apache/hudi/io/HoodieMergeHandleFactory.java   |   6 +-
 .../hudi/table/action/compact/HoodieCompactor.java |  25 +--
 .../table/log/BaseHoodieLogRecordReader.java       |   3 +
 .../common/table/read/FileGroupRecordBuffer.java   |   6 +-
 .../common/table/read/HoodieFileGroupReader.java   |  11 +-
 .../table/read/KeyBasedFileGroupRecordBuffer.java  |   3 +-
 .../read/PositionBasedFileGroupRecordBuffer.java   |   3 +-
 .../read/SortedKeyBasedFileGroupRecordBuffer.java  |   3 +-
 .../table/read/UnmergedFileGroupRecordBuffer.java  |   4 +-
 .../hudi/common/table/read/UpdateProcessor.java    |  19 +-
 .../hudi/metadata/HoodieTableMetadataUtil.java     |   2 +-
 .../table/read/TestFileGroupRecordBuffer.java      |   3 -
 .../read/TestKeyBasedFileGroupRecordBuffer.java    |  19 +-
 .../TestSortedKeyBasedFileGroupRecordBuffer.java   |   4 +-
 .../common/testutils/HoodieTestDataGenerator.java  |   3 +-
 .../hudi/common/testutils/RawTripTestPayload.java  |   3 +-
 .../org/apache/hudi/cdc/CDCFileGroupIterator.scala |   2 +-
 .../TestPositionBasedFileGroupRecordBuffer.java    |   7 +-
 .../hudi/table/TestHoodieMergeOnReadTable.java     | 216 +++++++++++++++------
 22 files changed, 253 insertions(+), 162 deletions(-)

diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedAppendHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedAppendHandle.java
index 256bc27e6620..d8de67050ae1 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedAppendHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedAppendHandle.java
@@ -24,18 +24,22 @@ import org.apache.hudi.common.config.HoodieMemoryConfig;
 import org.apache.hudi.common.config.TypedProperties;
 import org.apache.hudi.common.engine.HoodieReaderContext;
 import org.apache.hudi.common.engine.TaskContextSupplier;
+import org.apache.hudi.common.fs.FSUtils;
 import org.apache.hudi.common.model.CompactionOperation;
-import org.apache.hudi.common.model.FileSlice;
+import org.apache.hudi.common.model.HoodieLogFile;
+import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.table.log.block.HoodieLogBlock;
 import org.apache.hudi.common.table.read.HoodieFileGroupReader;
 import org.apache.hudi.common.table.read.HoodieReadStats;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.StringUtils;
+import org.apache.hudi.common.util.collection.CloseableMappingIterator;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.exception.HoodieIOException;
 import org.apache.hudi.exception.HoodieUpsertException;
 import org.apache.hudi.internal.schema.InternalSchema;
 import org.apache.hudi.internal.schema.utils.SerDeHelper;
+import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.table.HoodieTable;
 import org.apache.hudi.table.action.compact.strategy.CompactionStrategy;
 
@@ -44,6 +48,7 @@ import javax.annotation.concurrent.NotThreadSafe;
 import java.io.IOException;
 import java.util.Collections;
 import java.util.List;
+import java.util.stream.Stream;
 
 import static 
org.apache.hudi.common.config.HoodieReaderConfig.MERGE_USE_RECORD_POSITIONS;
 
@@ -57,17 +62,14 @@ import static 
org.apache.hudi.common.config.HoodieReaderConfig.MERGE_USE_RECORD_
 @NotThreadSafe
 public class FileGroupReaderBasedAppendHandle<T, I, K, O> extends 
HoodieAppendHandle<T, I, K, O> {
   private final HoodieReaderContext<T> readerContext;
-  private final FileSlice fileSlice;
   private final CompactionOperation operation;
   private HoodieReadStats readStats;
 
   public FileGroupReaderBasedAppendHandle(HoodieWriteConfig config, String 
instantTime, HoodieTable<T, I, K, O> hoodieTable,
-                                          FileSlice fileSlice, 
CompactionOperation operation, TaskContextSupplier taskContextSupplier,
-                                          HoodieReaderContext<T> 
readerContext) {
+                                          CompactionOperation operation, 
TaskContextSupplier taskContextSupplier, HoodieReaderContext<T> readerContext) {
     super(config, instantTime, hoodieTable, operation.getPartitionPath(), 
operation.getFileId(), taskContextSupplier);
     this.operation = operation;
     this.readerContext = readerContext;
-    this.fileSlice = fileSlice;
   }
 
   @Override
@@ -77,23 +79,25 @@ public class FileGroupReaderBasedAppendHandle<T, I, K, O> 
extends HoodieAppendHa
     TypedProperties props = TypedProperties.copy(config.getProps());
     long maxMemoryPerCompaction = 
IOUtils.getMaxMemoryPerCompaction(taskContextSupplier, config);
     props.put(HoodieMemoryConfig.MAX_MEMORY_FOR_MERGE.key(), 
String.valueOf(maxMemoryPerCompaction));
-    // Initializes the record iterator
+    Stream<HoodieLogFile> logFiles = 
operation.getDeltaFileNames().stream().map(logFileName ->
+        new HoodieLogFile(new StoragePath(FSUtils.constructAbsolutePath(
+            config.getBasePath(), operation.getPartitionPath()), 
logFileName)));
+    // Initializes the record iterator, log compaction requires writing the 
deletes into the delete block of the resulting log file.
     try (HoodieFileGroupReader<T> fileGroupReader = 
HoodieFileGroupReader.<T>newBuilder().withReaderContext(readerContext).withHoodieTableMetaClient(hoodieTable.getMetaClient())
-        
.withLatestCommitTime(instantTime).withFileSlice(fileSlice).withDataSchema(writeSchemaWithMetaFields).withRequestedSchema(writeSchemaWithMetaFields).withEnableOptimizedLogBlockScan(true)
-        
.withInternalSchema(internalSchemaOption).withProps(props).withShouldUseRecordPosition(usePosition).withSortOutput(hoodieTable.requireSortedRecords()).build())
 {
-      recordItr = fileGroupReader.getClosableHoodieRecordIterator();
+        
.withLatestCommitTime(instantTime).withPartitionPath(partitionPath).withLogFiles(logFiles).withBaseFileOption(Option.empty()).withDataSchema(writeSchemaWithMetaFields)
+        
.withRequestedSchema(writeSchemaWithMetaFields).withEnableOptimizedLogBlockScan(true).withInternalSchema(internalSchemaOption).withProps(props).withEmitDelete(true)
+        
.withShouldUseRecordPosition(usePosition).withSortOutput(hoodieTable.requireSortedRecords()).build())
 {
+      recordItr = new 
CloseableMappingIterator<>(fileGroupReader.getLogRecordsOnly(), record -> {
+        HoodieRecord<T> hoodieRecord = 
readerContext.constructHoodieRecord(record);
+        hoodieRecord.setCurrentLocation(newRecordLocation);
+        return hoodieRecord;
+      });
       header.put(HoodieLogBlock.HeaderMetadataType.COMPACTED_BLOCK_TIMES,
           StringUtils.join(fileGroupReader.getValidBlockInstants(), ","));
       super.doAppend();
-      // The stats of inserts, updates, and deletes are updated once at the end
-      // These will be set in the write stat when closing the merge handle
       this.readStats = fileGroupReader.getStats();
-      this.insertRecordsWritten = readStats.getNumInserts();
-      this.updatedRecordsWritten = readStats.getNumUpdates();
-      this.recordsDeleted = readStats.getNumDeletes();
-      this.recordsWritten = readStats.getNumInserts() + 
readStats.getNumUpdates();
     } catch (IOException e) {
-      throw new HoodieIOException("Failed to initialize file group reader for 
" + fileSlice, e);
+      throw new HoodieIOException("Failed to initialize file group reader for 
" + fileId, e);
     }
   }
 
@@ -114,6 +118,7 @@ public class FileGroupReaderBasedAppendHandle<T, I, K, O> 
extends HoodieAppendHa
       if (writeStatus.getStat().getRuntimeStats() != null) {
         
writeStatus.getStat().getRuntimeStats().setTotalScanTime(readStats.getTotalLogReadTimeMs());
       }
+      writeStatus.getStat().setPrevCommit(operation.getBaseInstantTime());
       return Collections.singletonList(writeStatus);
     } catch (Exception e) {
       throw new HoodieUpsertException("Failed to close " + 
this.getClass().getSimpleName(), e);
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
index 9ce794347b79..a8cc7e6864e6 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
@@ -26,8 +26,7 @@ import org.apache.hudi.common.engine.HoodieReaderContext;
 import org.apache.hudi.common.engine.TaskContextSupplier;
 import org.apache.hudi.common.fs.FSUtils;
 import org.apache.hudi.common.model.CompactionOperation;
-import org.apache.hudi.common.model.FileSlice;
-import org.apache.hudi.common.model.HoodieBaseFile;
+import org.apache.hudi.common.model.HoodieLogFile;
 import org.apache.hudi.common.model.HoodiePartitionMetadata;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.HoodieWriteStat;
@@ -60,6 +59,7 @@ import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.function.UnaryOperator;
+import java.util.stream.Stream;
 
 import static 
org.apache.hudi.common.config.HoodieReaderConfig.MERGE_USE_RECORD_POSITIONS;
 import static org.apache.hudi.common.model.HoodieFileFormat.HFILE;
@@ -76,7 +76,6 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
   private static final Logger LOG = 
LoggerFactory.getLogger(FileGroupReaderBasedMergeHandle.class);
 
   private final HoodieReaderContext<T> readerContext;
-  private final FileSlice fileSlice;
   private final CompactionOperation operation;
   private final String maxInstantTime;
   private HoodieReadStats readStats;
@@ -84,14 +83,13 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
   private final Option<HoodieCDCLogger> cdcLogger;
 
   public FileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String 
instantTime, HoodieTable<T, I, K, O> hoodieTable,
-                                         FileSlice fileSlice, 
CompactionOperation operation, TaskContextSupplier taskContextSupplier,
+                                         CompactionOperation operation, 
TaskContextSupplier taskContextSupplier,
                                          HoodieReaderContext<T> readerContext, 
String maxInstantTime,
                                          HoodieRecord.HoodieRecordType 
enginRecordType) {
     super(config, instantTime, operation.getPartitionPath(), 
operation.getFileId(), hoodieTable, taskContextSupplier);
     this.maxInstantTime = maxInstantTime;
     this.keyToNewRecords = Collections.emptyMap();
     this.readerContext = readerContext;
-    this.fileSlice = fileSlice;
     this.operation = operation;
     // If the table is a metadata table or the base file is an HFile, we use 
AVRO record type, otherwise we use the engine record type.
     this.recordType = hoodieTable.isMetadataTable() || 
HFILE.getFileExtension().equals(hoodieTable.getBaseFileExtension()) ? 
HoodieRecord.HoodieRecordType.AVRO : enginRecordType;
@@ -108,21 +106,21 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
     } else {
       this.cdcLogger = Option.empty();
     }
-    init(operation, this.partitionPath, fileSlice.getBaseFile());
+    init(operation, this.partitionPath);
   }
 
-  private void init(CompactionOperation operation, String partitionPath, 
Option<HoodieBaseFile> baseFileToMerge) {
+  private void init(CompactionOperation operation, String partitionPath) {
     LOG.info("partitionPath:{}, fileId to be merged:{}", partitionPath, 
fileId);
-    this.baseFileToMerge = baseFileToMerge.orElse(null);
+    this.baseFileToMerge = operation.getBaseFile(config.getBasePath(), 
operation.getPartitionPath()).orElse(null);
     this.writtenRecordKeys = new HashSet<>();
     writeStatus.setStat(new HoodieWriteStat());
     writeStatus.getStat().setTotalLogSizeCompacted(
         
operation.getMetrics().get(CompactionStrategy.TOTAL_LOG_FILE_SIZE).longValue());
     try {
       Option<String> latestValidFilePath = Option.empty();
-      if (baseFileToMerge.isPresent()) {
-        latestValidFilePath = Option.of(baseFileToMerge.get().getFileName());
-        
writeStatus.getStat().setPrevCommit(baseFileToMerge.get().getCommitTime());
+      if (baseFileToMerge != null) {
+        latestValidFilePath = Option.of(baseFileToMerge.getFileName());
+        writeStatus.getStat().setPrevCommit(baseFileToMerge.getCommitTime());
         // At the moment, we only support SI for overwrite with latest 
payload. So, we don't need to embed entire file slice here.
         // HUDI-8518 will be taken up to fix it for any payload during which 
we might require entire file slice to be set here.
         // Already AppendHandle adds all logs file from current file slice to 
HoodieDeltaWriteStat.
@@ -176,10 +174,14 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
     TypedProperties props = TypedProperties.copy(config.getProps());
     long maxMemoryPerCompaction = 
IOUtils.getMaxMemoryPerCompaction(taskContextSupplier, config);
     props.put(HoodieMemoryConfig.MAX_MEMORY_FOR_MERGE.key(), 
String.valueOf(maxMemoryPerCompaction));
+    Stream<HoodieLogFile> logFiles = 
operation.getDeltaFileNames().stream().map(logFileName ->
+        new HoodieLogFile(new StoragePath(FSUtils.constructAbsolutePath(
+            config.getBasePath(), operation.getPartitionPath()), 
logFileName)));
     // Initializes file group reader
     try (HoodieFileGroupReader<T> fileGroupReader = 
HoodieFileGroupReader.<T>newBuilder().withReaderContext(readerContext).withHoodieTableMetaClient(hoodieTable.getMetaClient())
-        
.withLatestCommitTime(maxInstantTime).withFileSlice(fileSlice).withDataSchema(writeSchemaWithMetaFields).withRequestedSchema(writeSchemaWithMetaFields)
-        
.withInternalSchema(internalSchemaOption).withProps(props).withShouldUseRecordPosition(usePosition).withSortOutput(hoodieTable.requireSortedRecords())
+        
.withLatestCommitTime(maxInstantTime).withPartitionPath(partitionPath).withBaseFileOption(Option.ofNullable(baseFileToMerge)).withLogFiles(logFiles)
+        
.withDataSchema(writeSchemaWithMetaFields).withRequestedSchema(writeSchemaWithMetaFields).withInternalSchema(internalSchemaOption).withProps(props)
+        
.withShouldUseRecordPosition(usePosition).withSortOutput(hoodieTable.requireSortedRecords())
         .withFileGroupUpdateCallback(cdcLogger.map(logger -> new 
CDCCallback(logger, readerContext))).build()) {
       // Reads the records from the file slice
       try (ClosableIterator<HoodieRecord<T>> recordIterator = 
fileGroupReader.getClosableHoodieRecordIterator()) {
@@ -199,6 +201,7 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
             writeToFile(record.getKey(), record, writeSchemaWithMetaFields,
                 config.getPayloadConfig().getProps(), preserveMetadata);
             writeStatus.markSuccess(record, recordMetadata);
+            recordsWritten++;
           } catch (Exception e) {
             LOG.error("Error writing record {}", record, e);
             writeStatus.markFailure(record, e, recordMetadata);
@@ -211,10 +214,9 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
         this.insertRecordsWritten = readStats.getNumInserts();
         this.updatedRecordsWritten = readStats.getNumUpdates();
         this.recordsDeleted = readStats.getNumDeletes();
-        this.recordsWritten = readStats.getNumInserts() + 
readStats.getNumUpdates();
       }
     } catch (IOException e) {
-      throw new HoodieUpsertException("Failed to compact file slice: " + 
fileSlice, e);
+      throw new HoodieUpsertException("Failed to compact file group: " + 
fileId, e);
     }
   }
 
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
index c85262f9cc41..8dfa02d3bc8f 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
@@ -68,6 +68,7 @@ import org.apache.avro.Schema;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.io.Closeable;
 import java.io.IOException;
 import java.util.ArrayList;
 import java.util.Collections;
@@ -529,6 +530,9 @@ public class HoodieAppendHandle<T, I, K, O> extends 
HoodieWriteHandle<T, I, K, O
       markClosed();
       // flush any remaining records to disk
       appendDataAndDeleteBlocks(header, true);
+      if (recordItr instanceof Closeable) {
+        ((Closeable) recordItr).close();
+      }
       recordItr = null;
 
       if (writer != null) {
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleFactory.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleFactory.java
index 9b8ed2811b0a..a0fe1291a932 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleFactory.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleFactory.java
@@ -21,7 +21,6 @@ package org.apache.hudi.io;
 import org.apache.hudi.common.engine.HoodieReaderContext;
 import org.apache.hudi.common.engine.TaskContextSupplier;
 import org.apache.hudi.common.model.CompactionOperation;
-import org.apache.hudi.common.model.FileSlice;
 import org.apache.hudi.common.model.HoodieBaseFile;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.WriteOperationType;
@@ -114,7 +113,6 @@ public class HoodieMergeHandleFactory {
       HoodieWriteConfig config,
       String instantTime,
       HoodieTable<T, I, K, O> hoodieTable,
-      FileSlice fileSlice,
       CompactionOperation operation,
       TaskContextSupplier taskContextSupplier,
       HoodieReaderContext<T> readerContext,
@@ -128,13 +126,13 @@ public class HoodieMergeHandleFactory {
     LOG.info("Create HoodieMergeHandle implementation {} {}", 
mergeHandleClass, logContext);
 
     Class<?>[] constructorParamTypes = new Class<?>[] {
-        HoodieWriteConfig.class, String.class, HoodieTable.class, 
FileSlice.class, CompactionOperation.class,
+        HoodieWriteConfig.class, String.class, HoodieTable.class, 
CompactionOperation.class,
         TaskContextSupplier.class, HoodieReaderContext.class, String.class, 
HoodieRecord.HoodieRecordType.class
     };
 
     return instantiateMergeHandle(
         isFallbackEnabled, mergeHandleClass, 
COMPACT_MERGE_HANDLE_CLASS_NAME.defaultValue(), logContext, 
constructorParamTypes,
-        config, instantTime, hoodieTable, fileSlice, operation, 
taskContextSupplier, readerContext, maxInstantTime, recordType);
+        config, instantTime, hoodieTable, operation, taskContextSupplier, 
readerContext, maxInstantTime, recordType);
   }
 
   /**
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/HoodieCompactor.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/HoodieCompactor.java
index c9ddce099eeb..6e72397d541a 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/HoodieCompactor.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/HoodieCompactor.java
@@ -26,11 +26,7 @@ import org.apache.hudi.common.engine.HoodieEngineContext;
 import org.apache.hudi.common.engine.HoodieReaderContext;
 import org.apache.hudi.common.engine.ReaderContextFactory;
 import org.apache.hudi.common.engine.TaskContextSupplier;
-import org.apache.hudi.common.fs.FSUtils;
 import org.apache.hudi.common.model.CompactionOperation;
-import org.apache.hudi.common.model.FileSlice;
-import org.apache.hudi.common.model.HoodieBaseFile;
-import org.apache.hudi.common.model.HoodieLogFile;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.WriteOperationType;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
@@ -46,7 +42,6 @@ import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.io.FileGroupReaderBasedAppendHandle;
 import org.apache.hudi.io.HoodieMergeHandle;
 import org.apache.hudi.io.HoodieMergeHandleFactory;
-import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.table.HoodieTable;
 
 import org.apache.avro.Schema;
@@ -57,7 +52,6 @@ import org.slf4j.LoggerFactory;
 import java.io.IOException;
 import java.io.Serializable;
 import java.util.List;
-import java.util.stream.Collectors;
 
 import static java.util.stream.Collectors.toList;
 
@@ -159,7 +153,7 @@ public abstract class HoodieCompactor<T, I, K, O> 
implements Serializable {
                                    String maxInstantTime,
                                    TaskContextSupplier taskContextSupplier) 
throws IOException {
     HoodieMergeHandle<T, ?, ?, ?> mergeHandle = 
HoodieMergeHandleFactory.create(writeConfig,
-        instantTime, table, getFileSliceFromOperation(operation, 
writeConfig.getBasePath()), operation, taskContextSupplier, 
hoodieReaderContext, maxInstantTime, getEngineRecordType());
+        instantTime, table, operation, taskContextSupplier, 
hoodieReaderContext, maxInstantTime, getEngineRecordType());
     mergeHandle.doMerge();
     return mergeHandle.close();
   }
@@ -171,26 +165,11 @@ public abstract class HoodieCompactor<T, I, K, O> 
implements Serializable {
                                       HoodieTable table,
                                       TaskContextSupplier taskContextSupplier) 
throws IOException {
     HoodieReaderContext<IndexedRecord> readerContext = new 
HoodieAvroReaderContext(table.getStorageConf(), 
table.getMetaClient().getTableConfig(), instantRange, Option.empty());
-    FileGroupReaderBasedAppendHandle<IndexedRecord, ?, ?, ?> appendHandle = 
new FileGroupReaderBasedAppendHandle<>(writeConfig, instantTime, table, 
getFileSliceFromOperation(operation,
-        writeConfig.getBasePath()), operation,  taskContextSupplier, 
readerContext);
+    FileGroupReaderBasedAppendHandle<IndexedRecord, ?, ?, ?> appendHandle = 
new FileGroupReaderBasedAppendHandle<>(writeConfig, instantTime, table, 
operation,  taskContextSupplier, readerContext);
     appendHandle.doAppend();
     return appendHandle.close();
   }
 
-  private FileSlice getFileSliceFromOperation(CompactionOperation operation, 
String basePath) {
-    Option<HoodieBaseFile> baseFileOpt =
-        operation.getBaseFile(basePath, operation.getPartitionPath());
-    List<HoodieLogFile> logFiles = 
operation.getDeltaFileNames().stream().map(p ->
-            new HoodieLogFile(new StoragePath(FSUtils.constructAbsolutePath(
-                basePath, operation.getPartitionPath()), p)))
-        .collect(Collectors.toList());
-    return new FileSlice(
-        operation.getFileGroupId(),
-        operation.getBaseInstantTime(),
-        baseFileOpt.isPresent() ? baseFileOpt.get() : null,
-        logFiles);
-  }
-
   public String getMaxInstantTime(HoodieTableMetaClient metaClient) {
     String maxInstantTime = metaClient
         
.getActiveTimeline().getTimelineOfActions(CollectionUtils.createSet(HoodieTimeline.COMMIT_ACTION,
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/BaseHoodieLogRecordReader.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/BaseHoodieLogRecordReader.java
index 0362a619499f..6897375744c3 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/BaseHoodieLogRecordReader.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/BaseHoodieLogRecordReader.java
@@ -698,6 +698,9 @@ public abstract class BaseHoodieLogRecordReader<T> {
       }
       // Done
       progress = 1.0f;
+      if (recordBuffer != null) {
+        totalLogRecords.set(recordBuffer.getTotalLogRecords());
+      }
     } catch (IOException e) {
       LOG.error("Got IOException when reading log file", e);
       throw new HoodieIOException("IOException when reading log file ", e);
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupRecordBuffer.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupRecordBuffer.java
index a2a01f7ee0ee..67b68753c999 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupRecordBuffer.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupRecordBuffer.java
@@ -73,7 +73,6 @@ public abstract class FileGroupRecordBuffer<T> implements 
HoodieFileGroupRecordB
   protected final Option<String> payloadClass;
   protected final TypedProperties props;
   protected final ExternalSpillableMap<Serializable, BufferedRecord<T>> 
records;
-  protected final HoodieReadStats readStats;
   protected final boolean shouldCheckCustomDeleteMarker;
   protected final boolean shouldCheckBuiltInDeleteMarker;
   protected ClosableIterator<T> baseFileIterator;
@@ -91,7 +90,6 @@ public abstract class FileGroupRecordBuffer<T> implements 
HoodieFileGroupRecordB
                                   RecordMergeMode recordMergeMode,
                                   PartialUpdateMode partialUpdateMode,
                                   TypedProperties props,
-                                  HoodieReadStats readStats,
                                   Option<String> orderingFieldName,
                                   UpdateProcessor<T> updateProcessor) {
     this.readerContext = readerContext;
@@ -122,7 +120,6 @@ public abstract class FileGroupRecordBuffer<T> implements 
HoodieFileGroupRecordB
         
SPILLABLE_DISK_MAP_TYPE.defaultValue().name()).toUpperCase(Locale.ROOT));
     boolean isBitCaskDiskMapCompressionEnabled = 
props.getBoolean(DISK_MAP_BITCASK_COMPRESSION_ENABLED.key(),
         DISK_MAP_BITCASK_COMPRESSION_ENABLED.defaultValue());
-    this.readStats = readStats;
     try {
       // Store merged records for all versions for this log file, set the 
in-memory footprint to maxInMemoryMapSize
       this.records = new ExternalSpillableMap<>(maxMemorySizeInBytes, 
spillableMapBasePath, new DefaultSizeEstimator<>(),
@@ -324,12 +321,11 @@ public abstract class FileGroupRecordBuffer<T> implements 
HoodieFileGroupRecordB
 
     // Inserts
     nextRecord = readerContext.seal(baseRecord);
-    readStats.incrementNumInserts();
     return true;
   }
 
   protected void initializeLogRecordIterator() {
-    logRecordIterator = records.values().iterator();
+    logRecordIterator = records.iterator();
   }
 
   protected boolean hasNextLogRecord() {
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
index 59fd05fb46cf..73094129b14f 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
@@ -180,13 +180,13 @@ public final class HoodieFileGroupReader<T> implements 
Closeable {
           readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, readStats);
     } else if (sortOutput) {
       return new SortedKeyBasedFileGroupRecordBuffer<>(
-          readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, readStats, orderingFieldName, updateProcessor);
+          readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, orderingFieldName, updateProcessor);
     } else if (shouldUseRecordPosition && 
inputSplit.baseFileOption.isPresent()) {
       return new PositionBasedFileGroupRecordBuffer<>(
-          readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, inputSplit.baseFileOption.get().getCommitTime(), props, 
readStats, orderingFieldName, updateProcessor);
+          readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, inputSplit.baseFileOption.get().getCommitTime(), props, 
orderingFieldName, updateProcessor);
     } else {
       return new KeyBasedFileGroupRecordBuffer<>(
-          readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, readStats, orderingFieldName, updateProcessor);
+          readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, orderingFieldName, updateProcessor);
     }
   }
 
@@ -390,6 +390,11 @@ public final class HoodieFileGroupReader<T> implements 
Closeable {
         nextRecord -> readerContext.getRecordKey(nextRecord, 
readerContext.getSchemaHandler().getRequestedSchema()));
   }
 
+  public ClosableIterator<BufferedRecord<T>> getLogRecordsOnly() throws 
IOException {
+    initRecordIterators();
+    return recordBuffer.getLogRecordIterator();
+  }
+
   public static class HoodieFileGroupReaderIterator<T> implements 
ClosableIterator<T> {
     private HoodieFileGroupReader<T> reader;
 
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/KeyBasedFileGroupRecordBuffer.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/KeyBasedFileGroupRecordBuffer.java
index a5d8d1e5f7c7..f8614f431279 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/KeyBasedFileGroupRecordBuffer.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/KeyBasedFileGroupRecordBuffer.java
@@ -54,10 +54,9 @@ public class KeyBasedFileGroupRecordBuffer<T> extends 
FileGroupRecordBuffer<T> {
                                        RecordMergeMode recordMergeMode,
                                        PartialUpdateMode partialUpdateMode,
                                        TypedProperties props,
-                                       HoodieReadStats readStats,
                                        Option<String> orderingFieldName,
                                        UpdateProcessor<T> updateProcessor) {
-    super(readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, readStats, orderingFieldName, updateProcessor);
+    super(readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, orderingFieldName, updateProcessor);
   }
 
   @Override
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/PositionBasedFileGroupRecordBuffer.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/PositionBasedFileGroupRecordBuffer.java
index 2734a711afa4..2b96b028ef88 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/PositionBasedFileGroupRecordBuffer.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/PositionBasedFileGroupRecordBuffer.java
@@ -72,10 +72,9 @@ public class PositionBasedFileGroupRecordBuffer<T> extends 
KeyBasedFileGroupReco
                                             PartialUpdateMode 
partialUpdateMode,
                                             String baseFileInstantTime,
                                             TypedProperties props,
-                                            HoodieReadStats readStats,
                                             Option<String> orderingFieldName,
                                             UpdateProcessor<T> 
updateProcessor) {
-    super(readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, readStats, orderingFieldName, updateProcessor);
+    super(readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, orderingFieldName, updateProcessor);
     this.baseFileInstantTime = baseFileInstantTime;
   }
 
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/SortedKeyBasedFileGroupRecordBuffer.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/SortedKeyBasedFileGroupRecordBuffer.java
index adcd4419ffe7..ec7706a2ae0d 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/SortedKeyBasedFileGroupRecordBuffer.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/SortedKeyBasedFileGroupRecordBuffer.java
@@ -48,10 +48,9 @@ public class SortedKeyBasedFileGroupRecordBuffer<T> extends 
KeyBasedFileGroupRec
                                              RecordMergeMode recordMergeMode,
                                              PartialUpdateMode 
partialUpdateMode,
                                              TypedProperties props,
-                                             HoodieReadStats readStats,
                                              Option<String> orderingFieldName,
                                              UpdateProcessor<T> 
updateProcessor) {
-    super(readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, readStats, orderingFieldName, updateProcessor);
+    super(readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, orderingFieldName, updateProcessor);
   }
 
   @Override
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UnmergedFileGroupRecordBuffer.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UnmergedFileGroupRecordBuffer.java
index ee1f9c279e32..22d43ebfbcaf 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UnmergedFileGroupRecordBuffer.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UnmergedFileGroupRecordBuffer.java
@@ -43,6 +43,7 @@ import java.util.Deque;
 public class UnmergedFileGroupRecordBuffer<T> extends FileGroupRecordBuffer<T> 
{
 
   private final Deque<HoodieLogBlock> currentInstantLogBlocks;
+  private final HoodieReadStats readStats;
   private ClosableIterator<T> recordIterator;
 
   public UnmergedFileGroupRecordBuffer(
@@ -52,7 +53,8 @@ public class UnmergedFileGroupRecordBuffer<T> extends 
FileGroupRecordBuffer<T> {
       PartialUpdateMode partialUpdateMode,
       TypedProperties props,
       HoodieReadStats readStats) {
-    super(readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, readStats, Option.empty(), null);
+    super(readerContext, hoodieTableMetaClient, recordMergeMode, 
partialUpdateMode, props, Option.empty(), null);
+    this.readStats = readStats;
     this.currentInstantLogBlocks = new ArrayDeque<>();
   }
 
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java
index beea4b3a0d39..a0f3e0eecea0 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java
@@ -21,10 +21,6 @@ package org.apache.hudi.common.table.read;
 import org.apache.hudi.common.engine.HoodieReaderContext;
 import org.apache.hudi.common.util.Option;
 
-import org.apache.avro.Schema;
-
-import java.util.function.UnaryOperator;
-
 /**
  * Interface used within the {@link HoodieFileGroupReader<T>} for processing 
updates to records in Merge-on-Read tables.
  * Note that the updates are always relative to the base file's current state.
@@ -45,7 +41,7 @@ public interface UpdateProcessor<T> {
                                        boolean emitDeletes, 
Option<BaseFileUpdateCallback> updateCallback) {
     UpdateProcessor<T> handler = new StandardUpdateProcessor<>(readStats, 
readerContext, emitDeletes);
     if (updateCallback.isPresent()) {
-      return new CallbackProcessor<>(updateCallback.get(), handler, 
readerContext);
+      return new CallbackProcessor<>(updateCallback.get(), handler);
     }
     return handler;
   }
@@ -75,9 +71,9 @@ public interface UpdateProcessor<T> {
         }
         return null;
       } else {
-        if (previousRecord != null) {
+        if (previousRecord != null && previousRecord != mergedRecord) {
           readStats.incrementNumUpdates();
-        } else {
+        } else if (previousRecord == null) {
           readStats.incrementNumInserts();
         }
         return readerContext.seal(mergedRecord);
@@ -92,17 +88,10 @@ public interface UpdateProcessor<T> {
   class CallbackProcessor<T> implements UpdateProcessor<T> {
     private final BaseFileUpdateCallback<T> callback;
     private final UpdateProcessor<T> delegate;
-    private final HoodieReaderContext<T> readerContext;
-    private final Option<UnaryOperator<T>> outputConverter;
-    private final Schema requestedSchema;
 
-    public CallbackProcessor(BaseFileUpdateCallback callback, 
UpdateProcessor<T> delegate,
-                             HoodieReaderContext<T> readerContext) {
+    public CallbackProcessor(BaseFileUpdateCallback callback, 
UpdateProcessor<T> delegate) {
       this.callback = callback;
       this.delegate = delegate;
-      this.readerContext = readerContext;
-      this.outputConverter = 
readerContext.getSchemaHandler().getOutputConverter();
-      this.requestedSchema = 
readerContext.getSchemaHandler().getRequestedSchema();
     }
 
     @Override
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java
 
b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java
index fc5298dc76cc..e8429383179a 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java
@@ -1033,7 +1033,7 @@ public class HoodieTableMetadataUtil {
       readerContext.setSchemaHandler(new 
FileGroupReaderSchemaHandler<>(readerContext, writerSchemaOpt.get(), 
writerSchemaOpt.get(), Option.empty(), tableConfig, properties));
       HoodieReadStats readStats = new HoodieReadStats();
       KeyBasedFileGroupRecordBuffer<T> recordBuffer = new 
KeyBasedFileGroupRecordBuffer<>(readerContext, datasetMetaClient,
-          readerContext.getMergeMode(), PartialUpdateMode.NONE, properties, 
readStats, Option.ofNullable(tableConfig.getPreCombineField()),
+          readerContext.getMergeMode(), PartialUpdateMode.NONE, properties, 
Option.ofNullable(tableConfig.getPreCombineField()),
           UpdateProcessor.create(readStats, readerContext, true, 
Option.empty()));
 
       // CRITICAL: Ensure allowInflightInstants is set to true
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestFileGroupRecordBuffer.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestFileGroupRecordBuffer.java
index 70d73950fa23..f9cfd504bbdf 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestFileGroupRecordBuffer.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestFileGroupRecordBuffer.java
@@ -297,7 +297,6 @@ class TestFileGroupRecordBuffer {
             RecordMergeMode.COMMIT_TIME_ORDERING,
             PartialUpdateMode.NONE,
             props,
-            readStats,
             Option.empty(),
             updateProcessor
         );
@@ -312,7 +311,6 @@ class TestFileGroupRecordBuffer {
             RecordMergeMode.COMMIT_TIME_ORDERING,
             PartialUpdateMode.NONE,
             props,
-            readStats,
             Option.empty(),
             updateProcessor
     );
@@ -336,7 +334,6 @@ class TestFileGroupRecordBuffer {
             RecordMergeMode.COMMIT_TIME_ORDERING,
             PartialUpdateMode.NONE,
             props,
-            readStats,
             Option.empty(),
             updateProcessor
         );
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestKeyBasedFileGroupRecordBuffer.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestKeyBasedFileGroupRecordBuffer.java
index 7383d28ebea4..dd2a281b295d 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestKeyBasedFileGroupRecordBuffer.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestKeyBasedFileGroupRecordBuffer.java
@@ -100,6 +100,9 @@ class TestKeyBasedFileGroupRecordBuffer {
     List<IndexedRecord> actualRecords = 
getActualRecords(fileGroupRecordBuffer);
     // delete for record3 is ignored due to event time ordering
     assertEquals(Arrays.asList(testRecord1UpdateWithSameTime, 
testRecord2Update, testRecord3Update), actualRecords);
+    assertEquals(0, readStats.getNumInserts());
+    assertEquals(0, readStats.getNumDeletes());
+    assertEquals(3, readStats.getNumUpdates());
   }
 
   @Test
@@ -132,6 +135,9 @@ class TestKeyBasedFileGroupRecordBuffer {
 
     List<IndexedRecord> actualRecords = 
getActualRecords(fileGroupRecordBuffer);
     assertEquals(Arrays.asList(testRecord2Update, testRecord3Update), 
actualRecords);
+    assertEquals(0, readStats.getNumInserts());
+    assertEquals(1, readStats.getNumDeletes());
+    assertEquals(2, readStats.getNumUpdates());
   }
 
   @Test
@@ -159,6 +165,9 @@ class TestKeyBasedFileGroupRecordBuffer {
 
     List<IndexedRecord> actualRecords = 
getActualRecords(fileGroupRecordBuffer);
     assertEquals(Arrays.asList(testRecord1UpdateWithSameTime, 
testRecord2EarlierUpdate), actualRecords);
+    assertEquals(0, readStats.getNumInserts());
+    assertEquals(1, readStats.getNumDeletes());
+    assertEquals(2, readStats.getNumUpdates());
   }
 
   @Test
@@ -192,6 +201,10 @@ class TestKeyBasedFileGroupRecordBuffer {
 
     List<IndexedRecord> actualRecords = 
getActualRecords(fileGroupRecordBuffer);
     assertEquals(Collections.singletonList(testRecord1), actualRecords);
+
+    assertEquals(0, readStats.getNumInserts());
+    assertEquals(3, readStats.getNumDeletes());
+    assertEquals(0, readStats.getNumUpdates());
   }
 
   @Test
@@ -225,6 +238,10 @@ class TestKeyBasedFileGroupRecordBuffer {
 
     List<IndexedRecord> actualRecords = 
getActualRecords(fileGroupRecordBuffer);
     assertEquals(Collections.singletonList(testRecord1), actualRecords);
+
+    assertEquals(0, readStats.getNumInserts());
+    assertEquals(3, readStats.getNumDeletes());
+    assertEquals(0, readStats.getNumUpdates());
   }
 
   private static GenericRecord createTestRecord(String recordKey, int counter, 
long ts) {
@@ -255,7 +272,7 @@ class TestKeyBasedFileGroupRecordBuffer {
     TypedProperties props = new TypedProperties();
     UpdateProcessor<IndexedRecord> updateProcessor = 
UpdateProcessor.create(readStats, readerContext, false, Option.empty());
     return new KeyBasedFileGroupRecordBuffer<>(
-        readerContext, mockMetaClient, recordMergeMode, 
PartialUpdateMode.NONE, props, readStats, orderingFieldName, updateProcessor);
+        readerContext, mockMetaClient, recordMergeMode, 
PartialUpdateMode.NONE, props, orderingFieldName, updateProcessor);
   }
 
   private static List<IndexedRecord> 
getActualRecords(FileGroupRecordBuffer<IndexedRecord> fileGroupRecordBuffer) 
throws IOException {
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java
index ebd38ff5b26c..439bcf8f920e 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java
@@ -77,7 +77,7 @@ class TestSortedKeyBasedFileGroupRecordBuffer {
 
     List<TestRecord> actualRecords = getActualRecords(fileGroupRecordBuffer);
     assertEquals(Arrays.asList(testRecord1, testRecord2Update, testRecord4, 
testRecord5, testRecord6Update), actualRecords);
-    assertEquals(4, readStats.getNumInserts());
+    assertEquals(3, readStats.getNumInserts());
     assertEquals(1, readStats.getNumUpdates());
     assertEquals(1, readStats.getNumDeletes());
   }
@@ -129,7 +129,7 @@ class TestSortedKeyBasedFileGroupRecordBuffer {
     TypedProperties props = new TypedProperties();
     UpdateProcessor<TestRecord> updateProcessor = 
UpdateProcessor.create(readStats, mockReaderContext, false, Option.empty());
     return new SortedKeyBasedFileGroupRecordBuffer<>(
-        mockReaderContext, mockMetaClient, recordMergeMode, partialUpdateMode, 
props, readStats, Option.empty(), updateProcessor);
+        mockReaderContext, mockMetaClient, recordMergeMode, partialUpdateMode, 
props, Option.empty(), updateProcessor);
   }
 
   private static List<TestRecord> 
getActualRecords(SortedKeyBasedFileGroupRecordBuffer<TestRecord> 
fileGroupRecordBuffer) throws IOException {
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/testutils/HoodieTestDataGenerator.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/testutils/HoodieTestDataGenerator.java
index 97810a21d9f0..f7d4a4292c25 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/testutils/HoodieTestDataGenerator.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/testutils/HoodieTestDataGenerator.java
@@ -85,6 +85,7 @@ import java.util.stream.Collectors;
 import java.util.stream.IntStream;
 import java.util.stream.Stream;
 
+import static org.apache.hudi.common.model.HoodieRecord.DEFAULT_ORDERING_VALUE;
 import static 
org.apache.hudi.common.testutils.HoodieTestUtils.COMMIT_METADATA_SER_DE;
 import static 
org.apache.hudi.common.testutils.HoodieTestUtils.INSTANT_FILE_NAME_GENERATOR;
 import static org.apache.hudi.common.util.StringUtils.getUTF8Bytes;
@@ -913,7 +914,7 @@ Generate random record using TRIP_ENCODED_DECIMAL_SCHEMA
 
   public HoodieRecord generateDeleteRecord(HoodieKey key) throws IOException {
     RawTripTestPayload payload =
-        new RawTripTestPayload(Option.empty(), key.getRecordKey(), 
key.getPartitionPath(), null, true, 0L);
+        new RawTripTestPayload(Option.empty(), key.getRecordKey(), 
key.getPartitionPath(), null, true, DEFAULT_ORDERING_VALUE);
     return new HoodieAvroRecord(key, payload);
   }
 
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/testutils/RawTripTestPayload.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/testutils/RawTripTestPayload.java
index 47f3b550daf3..f23fb3eb9563 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/testutils/RawTripTestPayload.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/testutils/RawTripTestPayload.java
@@ -52,6 +52,7 @@ import java.util.zip.DeflaterOutputStream;
 import java.util.zip.InflaterInputStream;
 
 import static org.apache.hudi.avro.HoodieAvroUtils.createHoodieRecordFromAvro;
+import static org.apache.hudi.common.model.HoodieRecord.DEFAULT_ORDERING_VALUE;
 import static 
org.apache.hudi.common.testutils.HoodieTestDataGenerator.AVRO_SCHEMA;
 import static org.apache.hudi.common.util.StringUtils.getUTF8Bytes;
 
@@ -190,7 +191,7 @@ public class RawTripTestPayload implements 
HoodieRecordPayload<RawTripTestPayloa
 
   @Override
   public RawTripTestPayload preCombine(RawTripTestPayload oldValue) {
-    if (oldValue.orderingVal.compareTo(orderingVal) > 0) {
+    if (!orderingVal.equals(DEFAULT_ORDERING_VALUE) && 
oldValue.orderingVal.compareTo(orderingVal) > 0) {
       // pick the payload with greatest ordering value
       return oldValue;
     } else {
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala
index a921969f2e83..855a93c4dc8e 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala
@@ -501,7 +501,7 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit,
         Option.empty(), metaClient.getTableConfig, readerProperties))
     val stats = new HoodieReadStats
     val recordBuffer = new 
KeyBasedFileGroupRecordBuffer[InternalRow](readerContext, metaClient,
-      readerContext.getMergeMode, 
metaClient.getTableConfig.getPartialUpdateMode, readerProperties, stats,
+      readerContext.getMergeMode, 
metaClient.getTableConfig.getPartialUpdateMode, readerProperties,
       Option.ofNullable(metaClient.getTableConfig.getPreCombineField),
       UpdateProcessor.create(stats, readerContext, true, Option.empty()))
 
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestPositionBasedFileGroupRecordBuffer.java
 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestPositionBasedFileGroupRecordBuffer.java
index 2c3e5f7384c1..ebe383741056 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestPositionBasedFileGroupRecordBuffer.java
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestPositionBasedFileGroupRecordBuffer.java
@@ -35,12 +35,11 @@ import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.table.TableSchemaResolver;
 import org.apache.hudi.common.table.log.block.HoodieDeleteBlock;
 import org.apache.hudi.common.table.log.block.HoodieLogBlock;
-import org.apache.hudi.common.table.read.UpdateProcessor;
 import org.apache.hudi.common.table.read.CustomPayloadForTesting;
 import org.apache.hudi.common.table.read.FileGroupReaderSchemaHandler;
-import org.apache.hudi.common.table.read.HoodieReadStats;
-import org.apache.hudi.common.table.read.PositionBasedFileGroupRecordBuffer;
 import org.apache.hudi.common.table.read.ParquetRowIndexBasedSchemaHandler;
+import org.apache.hudi.common.table.read.PositionBasedFileGroupRecordBuffer;
+import org.apache.hudi.common.table.read.UpdateProcessor;
 import org.apache.hudi.common.testutils.HoodieTestDataGenerator;
 import org.apache.hudi.common.testutils.RawTripTestPayload;
 import org.apache.hudi.common.testutils.SchemaTestUtil;
@@ -152,7 +151,6 @@ public class TestPositionBasedFileGroupRecordBuffer extends 
SparkClientFunctiona
       writeConfigs.put(HoodieWriteConfig.WRITE_PAYLOAD_CLASS_NAME.key(), 
CustomPayloadForTesting.class.getName());
       writeConfigs.put(HoodieTableConfig.RECORD_MERGE_STRATEGY_ID.key(), 
HoodieRecordMerger.PAYLOAD_BASED_MERGE_STRATEGY_UUID);
     }
-    HoodieReadStats readStats = new HoodieReadStats();
     buffer = new PositionBasedFileGroupRecordBuffer<>(
         ctx,
         metaClient,
@@ -160,7 +158,6 @@ public class TestPositionBasedFileGroupRecordBuffer extends 
SparkClientFunctiona
         metaClient.getTableConfig().getPartialUpdateMode(),
         baseFileInstantTime,
         props,
-        readStats,
         Option.of("timestamp"),
         updateProcessor);
   }
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/table/TestHoodieMergeOnReadTable.java
 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/table/TestHoodieMergeOnReadTable.java
index 9262378c2192..4012293a2bca 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/table/TestHoodieMergeOnReadTable.java
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/table/TestHoodieMergeOnReadTable.java
@@ -18,8 +18,6 @@
 
 package org.apache.hudi.table;
 
-import org.apache.hudi.client.HoodieReadClient;
-import org.apache.hudi.client.SparkRDDReadClient;
 import org.apache.hudi.client.SparkRDDWriteClient;
 import org.apache.hudi.client.WriteClientTestUtils;
 import org.apache.hudi.client.WriteStatus;
@@ -36,6 +34,8 @@ import org.apache.hudi.common.model.HoodieWriteStat;
 import org.apache.hudi.common.model.TableServiceType;
 import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.log.HoodieLogFileReader;
+import org.apache.hudi.common.table.log.block.HoodieLogBlock;
 import org.apache.hudi.common.table.timeline.HoodieActiveTimeline;
 import org.apache.hudi.common.table.timeline.HoodieInstant;
 import org.apache.hudi.common.table.timeline.HoodieInstant.State;
@@ -48,10 +48,11 @@ import org.apache.hudi.config.HoodieClusteringConfig;
 import org.apache.hudi.config.HoodieCompactionConfig;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.data.HoodieJavaRDD;
-import org.apache.hudi.index.HoodieIndex;
+import org.apache.hudi.exception.HoodieIOException;
 import org.apache.hudi.index.HoodieIndex.IndexType;
 import org.apache.hudi.metadata.HoodieTableMetadataWriter;
 import org.apache.hudi.metadata.SparkHoodieBackedTableMetadataWriter;
+import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.storage.StoragePathInfo;
 import org.apache.hudi.table.action.HoodieWriteMetadata;
 import 
org.apache.hudi.table.action.deltacommit.BaseSparkDeltaCommitActionExecutor;
@@ -73,6 +74,7 @@ import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.CsvSource;
+import org.junit.jupiter.params.provider.EnumSource;
 import org.junit.jupiter.params.provider.ValueSource;
 
 import java.io.IOException;
@@ -89,9 +91,11 @@ import java.util.Set;
 import java.util.stream.Collectors;
 import java.util.stream.Stream;
 
+import static org.apache.hudi.common.model.HoodieWriteStat.NULL_COMMIT;
 import static 
org.apache.hudi.common.table.timeline.HoodieTimeline.DELTA_COMMIT_ACTION;
 import static 
org.apache.hudi.common.table.timeline.InstantComparison.GREATER_THAN;
 import static 
org.apache.hudi.common.table.timeline.InstantComparison.compareTimestamps;
+import static 
org.apache.hudi.common.testutils.HoodieTestDataGenerator.AVRO_SCHEMA;
 import static 
org.apache.hudi.common.testutils.HoodieTestUtils.INSTANT_GENERATOR;
 import static 
org.apache.hudi.common.testutils.RawTripTestPayload.recordsToStrings;
 import static org.apache.hudi.config.HoodieWriteConfig.WRITE_TABLE_VERSION;
@@ -352,37 +356,21 @@ public class TestHoodieMergeOnReadTable extends 
SparkClientFunctionalTestHarness
   }
 
   // TODO: Enable metadata virtual keys in this test once the feature 
HUDI-2593 is completed
-  @Test
-  public void testLogFileCountsAfterCompaction() throws Exception {
+  @ParameterizedTest
+  @EnumSource(value = IndexType.class, names = {"INMEMORY", "BLOOM"})
+  public void testLogFileCountsAfterCompaction(IndexType indexType) throws 
Exception {
     boolean populateMetaFields = true;
     // insert 100 records
-    HoodieWriteConfig.Builder cfgBuilder = getConfigBuilder(true, false, 
HoodieIndex.IndexType.BLOOM,
+    HoodieWriteConfig.Builder cfgBuilder = getConfigBuilder(true, false, 
indexType,
         1024 * 1024 * 1024L, HoodieClusteringConfig.newBuilder().build());
     addConfigsForPopulateMetaFields(cfgBuilder, populateMetaFields);
     HoodieWriteConfig config = cfgBuilder.build();
 
     setUp(config.getProps());
 
-    try (SparkRDDWriteClient writeClient = getHoodieWriteClient(config);) {
-      String newCommitTime = "100";
-      WriteClientTestUtils.startCommitWithTime(writeClient, newCommitTime);
-
-      List<HoodieRecord> records = dataGen.generateInserts(newCommitTime, 100);
-      JavaRDD<HoodieRecord> recordsRDD = jsc().parallelize(records, 1);
-      List<WriteStatus> statuses = writeClient.insert(recordsRDD, 
newCommitTime).collect();
-      writeClient.commit(newCommitTime, jsc().parallelize(statuses), 
Option.empty(), DELTA_COMMIT_ACTION, Collections.emptyMap());
-
-      // Update all the 100 records
-      newCommitTime = "101";
-      List<HoodieRecord> updatedRecords = 
dataGen.generateUpdates(newCommitTime, records);
-      JavaRDD<HoodieRecord> updatedRecordsRDD = 
jsc().parallelize(updatedRecords, 1);
-
-      SparkRDDReadClient readClient = new SparkRDDReadClient(context(), 
config);
-      JavaRDD<HoodieRecord> updatedTaggedRecordsRDD = 
readClient.tagLocation(updatedRecordsRDD);
-
-      WriteClientTestUtils.startCommitWithTime(writeClient, newCommitTime);
-      statuses = writeClient.upsertPreppedRecords(updatedTaggedRecordsRDD, 
newCommitTime).collect();
-      writeClient.commit(newCommitTime, jsc().parallelize(statuses), 
Option.empty(), DELTA_COMMIT_ACTION, Collections.emptyMap());
+    try (SparkRDDWriteClient writeClient = getHoodieWriteClient(config)) {
+      String firstCommitTime = "100";
+      LastCommit lastCommit = writeInsertUpdateAndDelete(firstCommitTime, 
writeClient);
 
       // Write them to corresponding avro logfiles
       metaClient = HoodieTableMetaClient.reload(metaClient);
@@ -392,20 +380,22 @@ public class TestHoodieMergeOnReadTable extends 
SparkClientFunctionalTestHarness
         HoodieSparkWriteableTestTable testTable = HoodieSparkWriteableTestTable
             .of(metaClient, 
HoodieTestDataGenerator.AVRO_SCHEMA_WITH_METADATA_FIELDS, metadataWriter);
 
-        Set<String> allPartitions = updatedRecords.stream()
+        Set<String> allPartitions = lastCommit.updatedRecords.stream()
             .map(record -> record.getPartitionPath())
             .collect(Collectors.groupingBy(partitionPath -> partitionPath))
             .keySet();
-        assertEquals(allPartitions.size(), 
testTable.listAllBaseFiles().size());
+        // In-Memory index will only create log files, so no base files 
expected if that is set
+        assertEquals(indexType != IndexType.INMEMORY ? allPartitions.size() : 
0, testTable.listAllBaseFiles().size());
 
         // Verify that all data file has one log file
         HoodieTable table = HoodieSparkTable.create(config, context(), 
metaClient);
         for (String partitionPath : dataGen.getPartitionPaths()) {
           List<FileSlice> groupedLogFiles =
               
table.getSliceView().getLatestFileSlices(partitionPath).collect(Collectors.toList());
+          int expected = indexType != IndexType.INMEMORY ? 2 : 3;
           for (FileSlice fileSlice : groupedLogFiles) {
-            assertEquals(1, fileSlice.getLogFiles().count(),
-                "There should be 1 log file written for the latest data file - 
" + fileSlice);
+            assertEquals(expected, fileSlice.getLogFiles().count(),
+                String.format("There should be %d log files written for the 
latest data file - %s", expected, fileSlice));
           }
         }
 
@@ -414,14 +404,29 @@ public class TestHoodieMergeOnReadTable extends 
SparkClientFunctionalTestHarness
         HoodieWriteMetadata<JavaRDD<WriteStatus>> result = 
writeClient.compact(compactionInstantTime);
         writeClient.commitCompaction(compactionInstantTime, result, 
Option.of(table));
         
assertTrue(metaClient.reloadActiveTimeline().filterCompletedInstants().containsInstant(compactionInstantTime));
+        HoodieCommitMetadata compactionMetadata = 
metaClient.getActiveTimeline().readCommitMetadata(metaClient.getActiveTimeline().reload().getCommitsAndCompactionTimeline().lastInstant().get());
+        // assert the compaction metadata counts
+        String previousCommit;
+        long expectedUpdates;
+        long expectedInserts;
+        if (indexType == IndexType.INMEMORY) {
+          previousCommit = NULL_COMMIT;
+          expectedInserts = 90;
+          expectedUpdates = 0;
+        } else {
+          previousCommit = firstCommitTime;
+          expectedInserts = 0;
+          expectedUpdates = 80;
+        }
+        validateCompactionMetadata(compactionMetadata, previousCommit, 90, 
expectedUpdates, expectedInserts, 10);
 
         // Verify that recently written compacted data file has no log file
         metaClient = HoodieTableMetaClient.reload(metaClient);
         table = HoodieSparkTable.create(config, context(), metaClient);
         HoodieActiveTimeline timeline = metaClient.getActiveTimeline();
 
-        
assertTrue(compareTimestamps(timeline.lastInstant().get().requestedTime(), 
GREATER_THAN, newCommitTime),
-            "Compaction commit should be > than last insert");
+        
assertTrue(compareTimestamps(timeline.lastInstant().get().requestedTime(), 
GREATER_THAN, lastCommit.finalDeleteTime),
+            "Compaction commit should be > than last delta-commit");
 
         for (String partitionPath : dataGen.getPartitionPaths()) {
           List<FileSlice> groupedLogFiles =
@@ -443,12 +448,20 @@ public class TestHoodieMergeOnReadTable extends 
SparkClientFunctionalTestHarness
         Dataset<Row> actual = HoodieClientTestUtils.read(
             jsc(), basePath(), sqlContext(), hoodieStorage(), 
fullPartitionPaths);
         List<Row> rows = actual.collectAsList();
-        assertEquals(updatedRecords.size(), rows.size());
+        assertEquals(90, rows.size());
+        int updatedCount = 0;
         for (Row row : rows) {
-          assertEquals(row.getAs(HoodieRecord.COMMIT_TIME_METADATA_FIELD), 
newCommitTime);
+          if 
(row.getAs(HoodieRecord.COMMIT_TIME_METADATA_FIELD).equals(lastCommit.finalUpdateTime))
 {
+            updatedCount++;
+          } else {
+            // check that the commit time is 100 for all records that are not 
updated
+            assertEquals(firstCommitTime, 
row.getAs(HoodieRecord.COMMIT_TIME_METADATA_FIELD));
+          }
           // check that file names metadata is updated
           
assertTrue(row.getString(HoodieRecord.FILENAME_META_FIELD_ORD).contains(compactionInstantTime));
         }
+        // check that 80 records are updated
+        assertEquals(80, updatedCount);
       }
     }
   }
@@ -476,29 +489,8 @@ public class TestHoodieMergeOnReadTable extends 
SparkClientFunctionalTestHarness
     setUp(config.getProps());
 
     try (SparkRDDWriteClient writeClient = getHoodieWriteClient(config)) {
-      String newCommitTime = "100";
-      WriteClientTestUtils.startCommitWithTime(writeClient, newCommitTime);
-
-      List<HoodieRecord> records = dataGen.generateInserts(newCommitTime, 100);
-      JavaRDD<HoodieRecord> recordsRDD = jsc().parallelize(records, 1);
-      List<WriteStatus> statuses = writeClient.insert(recordsRDD, 
newCommitTime).collect();
-      writeClient.commit(newCommitTime, jsc().parallelize(statuses), 
Option.empty(), DELTA_COMMIT_ACTION, Collections.emptyMap());
-      // Update all the 100 records
-      newCommitTime = "101";
-      List<HoodieRecord> updatedRecords = 
dataGen.generateUpdates(newCommitTime, records);
-      JavaRDD<HoodieRecord> updatedRecordsRDD = 
jsc().parallelize(updatedRecords, 1);
-
-      HoodieReadClient readClient = new HoodieReadClient(context(), config);
-      JavaRDD<HoodieRecord> updatedTaggedRecordsRDD = 
readClient.tagLocation(updatedRecordsRDD);
-
-      WriteClientTestUtils.startCommitWithTime(writeClient, newCommitTime);
-      statuses = writeClient.upsertPreppedRecords(updatedTaggedRecordsRDD, 
newCommitTime).collect();
-      writeClient.commit(newCommitTime, jsc().parallelize(statuses), 
Option.empty(), DELTA_COMMIT_ACTION, Collections.emptyMap());
-
-      newCommitTime = "102";
-      WriteClientTestUtils.startCommitWithTime(writeClient, newCommitTime);
-      statuses = writeClient.upsertPreppedRecords(updatedTaggedRecordsRDD, 
newCommitTime).collect();
-      writeClient.commit(newCommitTime, jsc().parallelize(statuses), 
Option.empty(), DELTA_COMMIT_ACTION, Collections.emptyMap());
+      String firstCommitTime = "100";
+      LastCommit lastCommit = writeInsertUpdateAndDelete(firstCommitTime, 
writeClient);
 
       // Write them to corresponding avro logfiles
       metaClient = HoodieTableMetaClient.reload(metaClient);
@@ -508,7 +500,7 @@ public class TestHoodieMergeOnReadTable extends 
SparkClientFunctionalTestHarness
         HoodieSparkWriteableTestTable testTable = HoodieSparkWriteableTestTable
             .of(metaClient, 
HoodieTestDataGenerator.AVRO_SCHEMA_WITH_METADATA_FIELDS, metadataWriter);
 
-        Set<String> allPartitions = updatedRecords.stream()
+        Set<String> allPartitions = lastCommit.updatedRecords.stream()
             .map(record -> record.getPartitionPath())
             .collect(Collectors.groupingBy(partitionPath -> partitionPath))
             .keySet();
@@ -528,13 +520,16 @@ public class TestHoodieMergeOnReadTable extends 
SparkClientFunctionalTestHarness
         // Do a log compaction
         String logCompactionInstantTime = 
writeClient.scheduleLogCompaction(Option.empty()).get().toString();
         HoodieWriteMetadata<JavaRDD<WriteStatus>> result = 
writeClient.logCompact(logCompactionInstantTime, true);
+        HoodieCommitMetadata compactionMetadata = 
metaClient.getActiveTimeline().readCommitMetadata(metaClient.getActiveTimeline().reload().getCommitsAndCompactionTimeline().lastInstant().get());
+        validateCompactionMetadata(compactionMetadata, firstCommitTime, 80, 
80, 0, 10);
+        validateLogCompactionMetadataHeaders(compactionMetadata, 
metaClient.getBasePath(), "102,101");
 
         // Verify that recently written compacted data file has no log file
         metaClient = HoodieTableMetaClient.reload(metaClient);
         table = HoodieSparkTable.create(config, context(), metaClient);
         HoodieActiveTimeline timeline = metaClient.getActiveTimeline();
 
-        
assertTrue(compareTimestamps(timeline.lastInstant().get().requestedTime(), 
GREATER_THAN, newCommitTime),
+        
assertTrue(compareTimestamps(timeline.lastInstant().get().requestedTime(), 
GREATER_THAN, lastCommit.finalDeleteTime),
             "Compaction commit should be > than last insert");
 
         for (String partitionPath : dataGen.getPartitionPaths()) {
@@ -551,6 +546,109 @@ public class TestHoodieMergeOnReadTable extends 
SparkClientFunctionalTestHarness
     }
   }
 
+  private LastCommit writeInsertUpdateAndDelete(String firstCommitTime, 
SparkRDDWriteClient writeClient) throws IOException {
+    WriteClientTestUtils.startCommitWithTime(writeClient, firstCommitTime);
+    List<HoodieRecord> records = dataGen.generateInserts(firstCommitTime, 100);
+    JavaRDD<HoodieRecord> recordsRDD = jsc().parallelize(records, 1);
+    List<WriteStatus> statuses = writeClient.insert(recordsRDD, 
firstCommitTime).collect();
+    validateWriteStatuses(statuses, 100, 0, 0);
+    writeClient.commit(firstCommitTime, jsc().parallelize(statuses), 
Option.empty(), DELTA_COMMIT_ACTION, Collections.emptyMap());
+
+    // Update 80 of 100 records
+    String newCommitTime = "101";
+    List<HoodieRecord> updatedRecords = dataGen.generateUpdates(newCommitTime, 
records.subList(0, 80));
+    JavaRDD<HoodieRecord> updatedRecordsRDD = 
jsc().parallelize(updatedRecords, 1);
+    WriteClientTestUtils.startCommitWithTime(writeClient, newCommitTime);
+    statuses = writeClient.upsert(updatedRecordsRDD, newCommitTime).collect();
+    validateWriteStatuses(statuses, 0, 80, 0);
+    writeClient.commit(newCommitTime, jsc().parallelize(statuses), 
Option.empty(), DELTA_COMMIT_ACTION, Collections.emptyMap());
+
+    // Delete 10 of the remaining records
+    List<HoodieRecord> deleteRecords = 
dataGen.generateDeletesFromExistingRecords(records.subList(90, 100));
+    JavaRDD<HoodieRecord> deleteRecordsRDD = jsc().parallelize(deleteRecords, 
1);
+    String deleteCommitTime = "102";
+    WriteClientTestUtils.startCommitWithTime(writeClient, deleteCommitTime);
+    statuses = writeClient.upsert(deleteRecordsRDD, 
deleteCommitTime).collect();
+    validateWriteStatuses(statuses, 0, 0, 10);
+    writeClient.commit(deleteCommitTime, jsc().parallelize(statuses), 
Option.empty(), DELTA_COMMIT_ACTION, Collections.emptyMap());
+    return new LastCommit(newCommitTime, deleteCommitTime, updatedRecords);
+  }
+
+  private static class LastCommit {
+    public final String finalUpdateTime;
+    public final String finalDeleteTime;
+    public final List<HoodieRecord> updatedRecords;
+
+    public LastCommit(String finalUpdateTime, String finalDeleteTime, 
List<HoodieRecord> updatedRecords) {
+      this.finalUpdateTime = finalUpdateTime;
+      this.finalDeleteTime = finalDeleteTime;
+      this.updatedRecords = updatedRecords;
+    }
+  }
+
+  private static void validateWriteStatuses(List<WriteStatus> statuses, long 
expectedInserts, long expectedUpdates, long expectedDeletes) {
+    long totalInserts = 0;
+    long totalUpdates = 0;
+    long totalDeletes = 0;
+    for (WriteStatus status : statuses) {
+      totalInserts += status.getStat().getNumInserts();
+      totalUpdates += status.getStat().getNumUpdateWrites();
+      totalDeletes += status.getStat().getNumDeletes();
+    }
+    assertEquals(expectedInserts, totalInserts);
+    assertEquals(expectedUpdates, totalUpdates);
+    assertEquals(expectedDeletes, totalDeletes);
+  }
+
+  private static void validateCompactionMetadata(HoodieCommitMetadata 
compactionMetadata, String previousCommit, long expectedTotalRecordsWritten, 
long expectedTotalUpdatedRecords,
+                                                 long 
expectedTotalInsertedRecords, long expectedTotalDeletedRecords) {
+    long totalRecordsWritten = 0;
+    long totalDeletedRecords = 0;
+    long totalUpdatedRecords = 0;
+    long totalInsertedRecords = 0;
+    for (HoodieWriteStat writeStat : compactionMetadata.getWriteStats()) {
+      totalRecordsWritten += writeStat.getNumWrites();
+      totalDeletedRecords += writeStat.getNumDeletes();
+      totalUpdatedRecords += writeStat.getNumUpdateWrites();
+      totalInsertedRecords += writeStat.getNumInserts();
+      assertEquals(previousCommit, writeStat.getPrevCommit());
+      assertNotNull(writeStat.getFileId());
+      assertNotNull(writeStat.getPath());
+      assertTrue(writeStat.getFileSizeInBytes() > 0);
+      assertTrue(writeStat.getTotalWriteBytes() > 0);
+      assertTrue(writeStat.getTotalLogBlocks() > 0);
+      assertTrue(writeStat.getTotalLogSizeCompacted() > 0);
+      assertTrue(writeStat.getTotalLogFilesCompacted() > 0);
+      assertTrue(writeStat.getTotalLogRecords() > 0);
+    }
+    assertEquals(expectedTotalRecordsWritten, totalRecordsWritten);
+    assertEquals(expectedTotalUpdatedRecords, totalUpdatedRecords);
+    assertEquals(expectedTotalInsertedRecords, totalInsertedRecords);
+    assertEquals(expectedTotalDeletedRecords, totalDeletedRecords);
+  }
+
+  private void validateLogCompactionMetadataHeaders(HoodieCommitMetadata 
compactionMetadata, StoragePath basePath, String expectedCompactedBlockTimes) {
+    compactionMetadata.getFileIdAndFullPaths(basePath).values().stream()
+        .map(StoragePath::new)
+        .filter(path -> FSUtils.isLogFile(path.getName()))
+        .forEach(logFilePath -> {
+          try {
+            HoodieLogFileReader reader = new 
HoodieLogFileReader(hoodieStorage(), new HoodieLogFile(logFilePath), 
AVRO_SCHEMA, 10000, false,
+                false, "_row_key", null);
+            Map<HoodieLogBlock.HeaderMetadataType, String> headers = 
Collections.emptyMap();
+            while (reader.hasNext()) {
+              // Get headers from the final block
+              headers = reader.next().getLogBlockHeader();
+            }
+            
headers.containsKey(HoodieLogBlock.HeaderMetadataType.INSTANT_TIME);
+            headers.containsKey(HoodieLogBlock.HeaderMetadataType.SCHEMA);
+            assertEquals(expectedCompactedBlockTimes, 
headers.get(HoodieLogBlock.HeaderMetadataType.COMPACTED_BLOCK_TIMES));
+          } catch (IOException ex) {
+            throw new HoodieIOException("Failed reading logs", ex);
+          }
+        });
+  }
+
   /**
    * Test to ensure metadata stats are correctly written to metadata file.
    */

Reply via email to