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 e017d85d76b [HUDI-9527] Update HoodieTableMetadataUtil to use 
HoodieMergedLogRecordReader and FileGroupReader (#13470)
e017d85d76b is described below

commit e017d85d76b5a2332e96ce0b7e4b2a552f98dadc
Author: Tim Brown <[email protected]>
AuthorDate: Fri Jun 20 21:14:08 2025 -0500

    [HUDI-9527] Update HoodieTableMetadataUtil to use 
HoodieMergedLogRecordReader and FileGroupReader (#13470)
---
 .../hudi/common/engine/HoodieReaderContext.java    |  34 ++-
 .../common/table/read/HoodieFileGroupReader.java   |  22 +-
 .../hudi/metadata/HoodieTableMetadataUtil.java     | 292 ++++++++++++---------
 .../hudi/metadata/TestHoodieTableMetadataUtil.java |  15 +-
 .../TestMetadataUtilRLIandSIRecordGeneration.java  |   5 +-
 5 files changed, 212 insertions(+), 156 deletions(-)

diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/engine/HoodieReaderContext.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/engine/HoodieReaderContext.java
index e226f26c339..dd137af8d19 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/engine/HoodieReaderContext.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/engine/HoodieReaderContext.java
@@ -20,12 +20,14 @@
 package org.apache.hudi.common.engine;
 
 import org.apache.hudi.common.config.RecordMergeMode;
+import org.apache.hudi.common.config.TypedProperties;
 import org.apache.hudi.common.model.HoodieFileFormat;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.HoodieRecordMerger;
 import org.apache.hudi.common.serialization.CustomSerializer;
 import org.apache.hudi.common.serialization.DefaultSerializer;
 import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.HoodieTableVersion;
 import org.apache.hudi.common.table.log.InstantRange;
 import org.apache.hudi.common.table.read.BufferedRecord;
 import org.apache.hudi.common.table.read.FileGroupReaderSchemaHandler;
@@ -36,6 +38,7 @@ import org.apache.hudi.common.util.SizeEstimator;
 import org.apache.hudi.common.util.collection.ClosableIterator;
 import org.apache.hudi.common.util.collection.CloseableFilterIterator;
 import org.apache.hudi.common.util.collection.Pair;
+import org.apache.hudi.common.util.collection.Triple;
 import org.apache.hudi.expression.Predicate;
 import org.apache.hudi.keygen.KeyGenerator;
 import org.apache.hudi.metadata.HoodieTableMetadata;
@@ -57,6 +60,8 @@ import java.util.Map;
 import java.util.function.BiFunction;
 import java.util.function.UnaryOperator;
 
+import static 
org.apache.hudi.common.config.HoodieReaderConfig.RECORD_MERGE_IMPL_CLASSES_DEPRECATED_WRITE_CONFIG_KEY;
+import static 
org.apache.hudi.common.config.HoodieReaderConfig.RECORD_MERGE_IMPL_CLASSES_WRITE_CONFIG_KEY;
 import static org.apache.hudi.common.model.HoodieRecord.DEFAULT_ORDERING_VALUE;
 import static 
org.apache.hudi.common.model.HoodieRecord.RECORD_KEY_METADATA_FIELD;
 
@@ -77,6 +82,7 @@ public abstract class HoodieReaderContext<T> {
   protected final HoodieFileFormat baseFileFormat;
   // For general predicate pushdown.
   protected final Option<Predicate> keyFilterOpt;
+  protected final HoodieTableConfig tableConfig;
   private FileGroupReaderSchemaHandler<T> schemaHandler = null;
   private String tablePath = null;
   private String latestCommitTime = null;
@@ -87,6 +93,7 @@ public abstract class HoodieReaderContext<T> {
   private Boolean shouldMergeUseRecordPosition = null;
   protected String partitionPath;
   protected Option<InstantRange> instantRangeOpt = Option.empty();
+  private RecordMergeMode mergeMode;
 
   // for encoding and decoding schemas to the spillable map
   private final LocalAvroSchemaCache localAvroSchemaCache = 
LocalAvroSchemaCache.getInstance();
@@ -95,6 +102,7 @@ public abstract class HoodieReaderContext<T> {
                                 HoodieTableConfig tableConfig,
                                 Option<InstantRange> instantRangeOpt,
                                 Option<Predicate> keyFilterOpt) {
+    this.tableConfig = tableConfig;
     this.storageConfiguration = storageConfiguration;
     this.recordKeyExtractor = tableConfig.populateMetaFields() ? 
metadataKeyExtractor() : virtualKeyExtractor(tableConfig.getRecordKeyFields()
         .orElseThrow(() -> new IllegalArgumentException("No record keys 
specified and meta fields are not populated")));
@@ -263,7 +271,31 @@ public abstract class HoodieReaderContext<T> {
    *
    * @return {@link HoodieRecordMerger} to use.
    */
-  public abstract Option<HoodieRecordMerger> getRecordMerger(RecordMergeMode 
mergeMode, String mergeStrategyId, String mergeImplClasses);
+  protected abstract Option<HoodieRecordMerger> 
getRecordMerger(RecordMergeMode mergeMode, String mergeStrategyId, String 
mergeImplClasses);
+
+  /**
+   * Initializes the record merger based on the table configuration and 
properties.
+   * @param properties the properties for the reader.
+   */
+  public void initRecordMerger(TypedProperties properties) {
+    RecordMergeMode recordMergeMode = tableConfig.getRecordMergeMode();
+    String mergeStrategyId = tableConfig.getRecordMergeStrategyId();
+    if 
(!tableConfig.getTableVersion().greaterThanOrEquals(HoodieTableVersion.EIGHT)) {
+      Triple<RecordMergeMode, String, String> triple = 
HoodieTableConfig.inferCorrectMergingBehavior(
+          recordMergeMode, tableConfig.getPayloadClass(),
+          mergeStrategyId, null, tableConfig.getTableVersion());
+      recordMergeMode = triple.getLeft();
+      mergeStrategyId = triple.getRight();
+    }
+    this.mergeMode = recordMergeMode;
+    this.recordMerger = getRecordMerger(recordMergeMode, mergeStrategyId,
+        properties.getString(RECORD_MERGE_IMPL_CLASSES_WRITE_CONFIG_KEY,
+            
properties.getString(RECORD_MERGE_IMPL_CLASSES_DEPRECATED_WRITE_CONFIG_KEY, 
"")));
+  }
+
+  public RecordMergeMode getMergeMode() {
+    return mergeMode;
+  }
 
   /**
    * Gets the field value.
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 0f395434850..ee9b6fa2436 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
@@ -31,7 +31,6 @@ import org.apache.hudi.common.model.HoodieLogFile;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
-import org.apache.hudi.common.table.HoodieTableVersion;
 import org.apache.hudi.common.table.PartitionPathParser;
 import org.apache.hudi.common.table.log.HoodieMergedLogRecordReader;
 import org.apache.hudi.common.util.ConfigUtils;
@@ -43,7 +42,6 @@ import 
org.apache.hudi.common.util.collection.ClosableIterator;
 import org.apache.hudi.common.util.collection.CloseableMappingIterator;
 import org.apache.hudi.common.util.collection.EmptyIterator;
 import org.apache.hudi.common.util.collection.Pair;
-import org.apache.hudi.common.util.collection.Triple;
 import org.apache.hudi.exception.HoodieIOException;
 import org.apache.hudi.internal.schema.InternalSchema;
 import org.apache.hudi.storage.HoodieStorage;
@@ -60,8 +58,6 @@ import java.util.List;
 import java.util.function.UnaryOperator;
 import java.util.stream.Collectors;
 
-import static 
org.apache.hudi.common.config.HoodieReaderConfig.RECORD_MERGE_IMPL_CLASSES_DEPRECATED_WRITE_CONFIG_KEY;
-import static 
org.apache.hudi.common.config.HoodieReaderConfig.RECORD_MERGE_IMPL_CLASSES_WRITE_CONFIG_KEY;
 import static org.apache.hudi.common.fs.FSUtils.getRelativePartitionPath;
 import static org.apache.hudi.common.util.ConfigUtils.getIntWithAltKeys;
 
@@ -131,19 +127,7 @@ public final class HoodieFileGroupReader<T> implements 
Closeable {
     HoodieTableConfig tableConfig = hoodieTableMetaClient.getTableConfig();
     this.partitionPath = fileSlice.getPartitionPath();
     this.partitionPathFields = tableConfig.getPartitionFields();
-    RecordMergeMode recordMergeMode = tableConfig.getRecordMergeMode();
-    String mergeStrategyId = tableConfig.getRecordMergeStrategyId();
-    if 
(!tableConfig.getTableVersion().greaterThanOrEquals(HoodieTableVersion.EIGHT)) {
-      Triple<RecordMergeMode, String, String> triple = 
HoodieTableConfig.inferCorrectMergingBehavior(
-          recordMergeMode, tableConfig.getPayloadClass(),
-          mergeStrategyId, null, tableConfig.getTableVersion());
-      recordMergeMode = triple.getLeft();
-      mergeStrategyId = triple.getRight();
-    }
-    readerContext.setRecordMerger(readerContext.getRecordMerger(
-        recordMergeMode, mergeStrategyId,
-        props.getString(RECORD_MERGE_IMPL_CLASSES_WRITE_CONFIG_KEY,
-            
props.getString(RECORD_MERGE_IMPL_CLASSES_DEPRECATED_WRITE_CONFIG_KEY, ""))));
+    readerContext.initRecordMerger(props);
     readerContext.setTablePath(tablePath);
     readerContext.setLatestCommitTime(latestCommitTime);
     boolean isSkipMerge = ConfigUtils.getStringWithAltKeys(props, 
HoodieReaderConfig.MERGE_TYPE, 
true).equalsIgnoreCase(HoodieReaderConfig.REALTIME_SKIP_MERGE);
@@ -158,7 +142,7 @@ public final class HoodieFileGroupReader<T> implements 
Closeable {
         ? new PositionBasedSchemaHandler<>(readerContext, dataSchema, 
requestedSchema, internalSchemaOpt, tableConfig, props)
         : new FileGroupReaderSchemaHandler<>(readerContext, dataSchema, 
requestedSchema, internalSchemaOpt, tableConfig, props));
     this.outputConverter = 
readerContext.getSchemaHandler().getOutputConverter();
-    this.orderingFieldName = recordMergeMode == 
RecordMergeMode.COMMIT_TIME_ORDERING
+    this.orderingFieldName = readerContext.getMergeMode() == 
RecordMergeMode.COMMIT_TIME_ORDERING
         ? Option.empty()
         : Option.ofNullable(ConfigUtils.getOrderingField(props))
         .or(() -> {
@@ -170,7 +154,7 @@ public final class HoodieFileGroupReader<T> implements 
Closeable {
         });
     this.readStats = new HoodieReadStats();
     this.recordBuffer = getRecordBuffer(readerContext, hoodieTableMetaClient,
-        recordMergeMode, props, hoodieBaseFileOption, this.logFiles.isEmpty(),
+        readerContext.getMergeMode(), props, hoodieBaseFileOption, 
this.logFiles.isEmpty(),
         isSkipMerge, shouldUseRecordPosition, readStats, emitDelete, 
sortOutput);
     this.allowInflightInstants = allowInflightInstants;
   }
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 f507a115979..966f084e26c 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
@@ -51,6 +51,7 @@ import org.apache.hudi.common.engine.EngineType;
 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.engine.ReaderContextFactory;
 import org.apache.hudi.common.fs.FSUtils;
 import org.apache.hudi.common.function.SerializableBiFunction;
 import org.apache.hudi.common.function.SerializablePairFunction;
@@ -69,15 +70,19 @@ import org.apache.hudi.common.model.HoodiePartitionMetadata;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType;
 import org.apache.hudi.common.model.HoodieRecordGlobalLocation;
-import org.apache.hudi.common.model.HoodieRecordMerger;
 import org.apache.hudi.common.model.HoodieReplaceCommitMetadata;
 import org.apache.hudi.common.model.HoodieWriteStat;
 import org.apache.hudi.common.model.WriteOperationType;
 import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.table.TableSchemaResolver;
-import org.apache.hudi.common.table.log.HoodieMergedLogRecordScanner;
+import org.apache.hudi.common.table.log.HoodieMergedLogRecordReader;
+import org.apache.hudi.common.table.read.BufferedRecord;
+import org.apache.hudi.common.table.read.FileGroupReaderSchemaHandler;
+import org.apache.hudi.common.table.read.FileGroupRecordBuffer;
 import org.apache.hudi.common.table.read.HoodieFileGroupReader;
+import org.apache.hudi.common.table.read.HoodieReadStats;
+import org.apache.hudi.common.table.read.KeyBasedFileGroupRecordBuffer;
 import org.apache.hudi.common.table.timeline.HoodieActiveTimeline;
 import org.apache.hudi.common.table.timeline.HoodieInstant;
 import org.apache.hudi.common.table.timeline.HoodieInstantTimeGenerator;
@@ -87,8 +92,6 @@ import org.apache.hudi.common.table.timeline.TimelineFactory;
 import org.apache.hudi.common.table.view.HoodieTableFileSystemView;
 import org.apache.hudi.common.util.CollectionUtils;
 import org.apache.hudi.common.util.FileFormatUtils;
-import org.apache.hudi.common.util.FileIOUtils;
-import org.apache.hudi.common.util.HoodieRecordUtils;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.StringUtils;
 import org.apache.hudi.common.util.ValidationUtils;
@@ -164,7 +167,6 @@ import static 
org.apache.hudi.common.config.HoodieCommonConfig.DISK_MAP_BITCASK_
 import static 
org.apache.hudi.common.config.HoodieCommonConfig.MAX_MEMORY_FOR_COMPACTION;
 import static 
org.apache.hudi.common.config.HoodieCommonConfig.SPILLABLE_DISK_MAP_TYPE;
 import static 
org.apache.hudi.common.config.HoodieMemoryConfig.MAX_MEMORY_FOR_MERGE;
-import static 
org.apache.hudi.common.config.HoodieReaderConfig.ENABLE_OPTIMIZED_LOG_BLOCKS_SCAN;
 import static 
org.apache.hudi.common.config.HoodieReaderConfig.REALTIME_SKIP_MERGE;
 import static org.apache.hudi.common.fs.FSUtils.getFileNameFromPath;
 import static 
org.apache.hudi.common.model.HoodieRecord.COMMIT_TIME_METADATA_FIELD;
@@ -821,13 +823,13 @@ public class HoodieTableMetadataUtil {
   }
 
   @VisibleForTesting
-  public static HoodieData<HoodieRecord> 
convertMetadataToRecordIndexRecords(HoodieEngineContext engineContext,
-                                                                             
HoodieCommitMetadata commitMetadata,
-                                                                             
HoodieMetadataConfig metadataConfig,
-                                                                             
HoodieTableMetaClient dataTableMetaClient,
-                                                                             
int writesFileIdEncoding,
-                                                                             
String instantTime,
-                                                                             
EngineType engineType) {
+  public static <T> HoodieData<HoodieRecord> 
convertMetadataToRecordIndexRecords(HoodieEngineContext engineContext,
+                                                                               
  HoodieCommitMetadata commitMetadata,
+                                                                               
  HoodieMetadataConfig metadataConfig,
+                                                                               
  HoodieTableMetaClient dataTableMetaClient,
+                                                                               
  int writesFileIdEncoding,
+                                                                               
  String instantTime,
+                                                                               
  EngineType engineType) {
     List<HoodieWriteStat> allWriteStats = 
commitMetadata.getPartitionToWriteStats().values().stream()
         .flatMap(Collection::stream).collect(Collectors.toList());
     // Return early if there are no write stats, or if the operation is a 
compaction.
@@ -850,6 +852,7 @@ public class HoodieTableMetadataUtil {
       StorageConfiguration storageConfiguration = 
dataTableMetaClient.getStorageConf();
       Option<Schema> writerSchemaOpt = 
tryResolveSchemaForTable(dataTableMetaClient);
       Option<Schema> finalWriterSchemaOpt = writerSchemaOpt;
+      ReaderContextFactory<T> readerContextFactory = 
engineContext.getReaderContextFactory(dataTableMetaClient);
       HoodieData<HoodieRecord> recordIndexRecords = 
engineContext.parallelize(new ArrayList<>(writeStatsByFileId.entrySet()), 
parallelism)
           .flatMap(writeStatsByFileIdEntry -> {
             String fileId = writeStatsByFileIdEntry.getKey();
@@ -889,8 +892,8 @@ public class HoodieTableMetadataUtil {
                   })
                   .collect(Collectors.toList());
               // Extract revived and deleted keys
-              Pair<Set<String>, Set<String>> revivedAndDeletedKeys =
-                  getRevivedAndDeletedKeysFromMergedLogs(dataTableMetaClient, 
instantTime, engineType, allLogFilePaths, finalWriterSchemaOpt, 
currentLogFilePaths);
+              Pair<Set<String>, Set<String>> revivedAndDeletedKeys = 
getRevivedAndDeletedKeysFromMergedLogs(dataTableMetaClient, instantTime, 
allLogFilePaths, finalWriterSchemaOpt,
+                  currentLogFilePaths, partitionPath, 
readerContextFactory.getContext());
               Set<String> revivedKeys = revivedAndDeletedKeys.getLeft();
               Set<String> deletedKeys = revivedAndDeletedKeys.getRight();
               // Process revived keys to create updates
@@ -940,109 +943,141 @@ public class HoodieTableMetadataUtil {
    *
    * @param dataTableMetaClient  data table meta client
    * @param instantTime          timestamp of the commit
-   * @param engineType           engine type (SPARK, FLINK, JAVA)
+   * @param partitionPath        partition path of the log files
+   * @param readerContext        the reader context for the engine
    * @param logFilePaths         list of log file paths including current and 
previous file slices
    * @param finalWriterSchemaOpt records schema
    * @param currentLogFilePaths  list of log file paths for the current instant
    * @return pair of revived and deleted keys
    */
   @VisibleForTesting
-  public static Pair<Set<String>, Set<String>> 
getRevivedAndDeletedKeysFromMergedLogs(HoodieTableMetaClient 
dataTableMetaClient,
-                                                                               
       String instantTime,
-                                                                               
       EngineType engineType,
-                                                                               
       List<String> logFilePaths,
-                                                                               
       Option<Schema> finalWriterSchemaOpt,
-                                                                               
       List<String> currentLogFilePaths) {
+  public static <T> Pair<Set<String>, Set<String>> 
getRevivedAndDeletedKeysFromMergedLogs(HoodieTableMetaClient 
dataTableMetaClient,
+                                                                               
           String instantTime,
+                                                                               
           List<String> logFilePaths,
+                                                                               
           Option<Schema> finalWriterSchemaOpt,
+                                                                               
           List<String> currentLogFilePaths,
+                                                                               
           String partitionPath,
+                                                                               
           HoodieReaderContext<T> readerContext) {
     // Separate out the current log files
     List<String> logFilePathsWithoutCurrentLogFiles = logFilePaths.stream()
         .filter(logFilePath -> !currentLogFilePaths.contains(logFilePath))
         .collect(toList());
     if (logFilePathsWithoutCurrentLogFiles.isEmpty()) {
       // Only current log file is present, so we can directly get the deleted 
record keys from it and return the RLI records.
-      Map<String, HoodieRecord> currentLogRecords =
-          getLogRecords(currentLogFilePaths, dataTableMetaClient, 
finalWriterSchemaOpt, instantTime, engineType);
-      Set<String> deletedKeys = currentLogRecords.entrySet().stream()
-          .filter(entry -> isDeleteRecord(dataTableMetaClient, 
finalWriterSchemaOpt, entry.getValue()))
-          .map(Map.Entry::getKey)
-          .collect(Collectors.toSet());
-      return Pair.of(Collections.emptySet(), deletedKeys);
+      try (ClosableIterator<BufferedRecord<T>> currentLogRecords =
+               getLogRecords(currentLogFilePaths, dataTableMetaClient, 
finalWriterSchemaOpt, instantTime, partitionPath, readerContext)) {
+        Set<String> deletedKeys = new HashSet<>();
+        currentLogRecords.forEachRemaining(record -> {
+          if (record.isDelete()) {
+            deletedKeys.add(record.getRecordKey());
+          }
+        });
+        return Pair.of(Collections.emptySet(), deletedKeys);
+      }
     }
-    return getRevivedAndDeletedKeys(dataTableMetaClient, instantTime, 
engineType, logFilePaths, finalWriterSchemaOpt, 
logFilePathsWithoutCurrentLogFiles);
+    return getRevivedAndDeletedKeys(dataTableMetaClient, instantTime, 
partitionPath, readerContext, logFilePaths, finalWriterSchemaOpt, 
logFilePathsWithoutCurrentLogFiles);
   }
 
-  private static Pair<Set<String>, Set<String>> 
getRevivedAndDeletedKeys(HoodieTableMetaClient dataTableMetaClient, String 
instantTime, EngineType engineType, List<String> logFilePaths,
-                                                                         
Option<Schema> finalWriterSchemaOpt, List<String> 
logFilePathsWithoutCurrentLogFiles) {
+  private static <T> Pair<Set<String>, Set<String>> 
getRevivedAndDeletedKeys(HoodieTableMetaClient dataTableMetaClient, String 
instantTime, String partitionPath, HoodieReaderContext<T> readerContext,
+                                                                             
List<String> logFilePaths, Option<Schema> finalWriterSchemaOpt, List<String> 
logFilePathsWithoutCurrentLogFiles) {
+    // Partition valid (non-deleted) and deleted keys from all log files, 
including current, in a single pass
+    Set<String> validKeysForAllLogs = new HashSet<>();
+    Set<String> deletedKeysForAllLogs = new HashSet<>();
     // Fetch log records for all log files
-    Map<String, HoodieRecord> allLogRecords =
-        getLogRecords(logFilePaths, dataTableMetaClient, finalWriterSchemaOpt, 
instantTime, engineType);
-
-    // Fetch log records for previous log files (excluding the current log 
files)
-    Map<String, HoodieRecord> previousLogRecords =
-        getLogRecords(logFilePathsWithoutCurrentLogFiles, dataTableMetaClient, 
finalWriterSchemaOpt, instantTime, engineType);
+    try (ClosableIterator<BufferedRecord<T>> allLogRecords =
+             getLogRecords(logFilePaths, dataTableMetaClient, 
finalWriterSchemaOpt, instantTime, partitionPath, readerContext)) {
+      allLogRecords.forEachRemaining(record -> {
+        if (record.isDelete()) {
+          deletedKeysForAllLogs.add(record.getRecordKey());
+        } else {
+          validKeysForAllLogs.add(record.getRecordKey());
+        }
+      });
+    }
 
     // Partition valid (non-deleted) and deleted keys from previous log files 
in a single pass
-    Map<Boolean, Set<String>> partitionedKeysForPreviousLogs = 
previousLogRecords.entrySet().stream()
-        .collect(Collectors.partitioningBy(
-            entry -> !isDeleteRecord(dataTableMetaClient, 
finalWriterSchemaOpt, entry.getValue()),
-            Collectors.mapping(Map.Entry::getKey, Collectors.toSet())
-        ));
-    Set<String> validKeysForPreviousLogs = 
partitionedKeysForPreviousLogs.get(true);
-    Set<String> deletedKeysForPreviousLogs = 
partitionedKeysForPreviousLogs.get(false);
-
-    // Partition valid (non-deleted) and deleted keys from all log files, 
including current, in a single pass
-    Map<Boolean, Set<String>> partitionedKeysForAllLogs = 
allLogRecords.entrySet().stream()
-        .collect(Collectors.partitioningBy(
-            entry -> !isDeleteRecord(dataTableMetaClient, 
finalWriterSchemaOpt, entry.getValue()),
-            Collectors.mapping(Map.Entry::getKey, Collectors.toSet())
-        ));
-    Set<String> validKeysForAllLogs = partitionedKeysForAllLogs.get(true);
-    Set<String> deletedKeysForAllLogs = partitionedKeysForAllLogs.get(false);
+    Set<String> validKeysForPreviousLogs = new HashSet<>();
+    Set<String> deletedKeysForPreviousLogs = new HashSet<>();
+    // Fetch log records for previous log files (excluding the current log 
files)
+    try (ClosableIterator<BufferedRecord<T>> previousLogRecords =
+             getLogRecords(logFilePathsWithoutCurrentLogFiles, 
dataTableMetaClient, finalWriterSchemaOpt, instantTime, partitionPath, 
readerContext)) {
+      previousLogRecords.forEachRemaining(record -> {
+        if (record.isDelete()) {
+          deletedKeysForPreviousLogs.add(record.getRecordKey());
+        } else {
+          validKeysForPreviousLogs.add(record.getRecordKey());
+        }
+      });
+    }
 
     return computeRevivedAndDeletedKeys(validKeysForPreviousLogs, 
deletedKeysForPreviousLogs, validKeysForAllLogs, deletedKeysForAllLogs);
   }
 
-  private static boolean isDeleteRecord(HoodieTableMetaClient 
dataTableMetaClient, Option<Schema> finalWriterSchemaOpt, HoodieRecord record) {
-    try {
-      return record.isDelete(finalWriterSchemaOpt.get(), 
dataTableMetaClient.getTableConfig().getProps());
-    } catch (IOException e) {
-      throw new HoodieException("Failed to check if record is delete", e);
-    }
-  }
-
-  private static Map<String, HoodieRecord> getLogRecords(List<String> 
logFilePaths,
-                                                         HoodieTableMetaClient 
datasetMetaClient,
-                                                         Option<Schema> 
writerSchemaOpt,
-                                                         String 
latestCommitTimestamp,
-                                                         EngineType 
engineType) {
-    if (writerSchemaOpt.isPresent()) {
+  private static <T> ClosableIterator<BufferedRecord<T>> 
getLogRecords(List<String> logFilePaths,
+                                                                       
HoodieTableMetaClient datasetMetaClient,
+                                                                       
Option<Schema> writerSchemaOpt,
+                                                                       String 
latestCommitTimestamp,
+                                                                       String 
partitionPath,
+                                                                       
HoodieReaderContext<T> readerContext) {
+    if (writerSchemaOpt.isPresent() && !logFilePaths.isEmpty()) {
+      List<HoodieLogFile> logFiles = 
logFilePaths.stream().map(HoodieLogFile::new).collect(Collectors.toList());
+      FileSlice fileSlice = new FileSlice(partitionPath, 
logFiles.get(0).getFileId(), logFiles.get(0).getDeltaCommitTime());
+      logFiles.forEach(fileSlice::addLogFile);
       final StorageConfiguration<?> storageConf = 
datasetMetaClient.getStorageConf();
-      HoodieRecordMerger recordMerger = HoodieRecordUtils.createRecordMerger(
-          datasetMetaClient.getBasePath().toString(),
-          engineType,
-          Collections.emptyList(),
-          datasetMetaClient.getTableConfig().getRecordMergeStrategyId());
-
-      // CRITICAL: Ensure allowInflightInstants is set to true while replacing 
the scanner with *LogRecordReader or HoodieFileGroupReader
-      HoodieMergedLogRecordScanner mergedLogRecordScanner = 
HoodieMergedLogRecordScanner.newBuilder()
+      TypedProperties properties = 
getFileGroupReaderPropertiesFromStorageConf(storageConf);
+      readerContext.setLatestCommitTime(latestCommitTimestamp);
+      readerContext.setHasBootstrapBaseFile(false);
+      readerContext.setHasLogFiles(true);
+      HoodieTableConfig tableConfig = datasetMetaClient.getTableConfig();
+      readerContext.initRecordMerger(properties);
+      readerContext.setSchemaHandler(new 
FileGroupReaderSchemaHandler<>(readerContext, writerSchemaOpt.get(), 
writerSchemaOpt.get(), Option.empty(), tableConfig, properties));
+      KeyBasedFileGroupRecordBuffer<T> recordBuffer = new 
KeyBasedFileGroupRecordBuffer<>(readerContext, datasetMetaClient,
+          readerContext.getMergeMode(), properties, new HoodieReadStats(), 
Option.ofNullable(tableConfig.getPreCombineField()), true);
+
+      // CRITICAL: Ensure allowInflightInstants is set to true
+      HoodieMergedLogRecordReader<T> mergedLogRecordReader = 
HoodieMergedLogRecordReader.<T>newBuilder()
           .withStorage(datasetMetaClient.getStorage())
-          .withBasePath(datasetMetaClient.getBasePath())
-          .withLogFilePaths(logFilePaths)
-          .withReaderSchema(writerSchemaOpt.get())
-          .withLatestInstantTime(latestCommitTimestamp)
+          .withHoodieReaderContext(readerContext)
+          
.withLogFiles(logFilePaths.stream().map(HoodieLogFile::new).collect(toList()))
           .withReverseReader(false)
-          
.withMaxMemorySizeInBytes(storageConf.getLong(MAX_MEMORY_FOR_COMPACTION.key(), 
DEFAULT_MAX_MEMORY_FOR_SPILLABLE_MAP_IN_BYTES))
           
.withBufferSize(HoodieMetadataConfig.MAX_READER_BUFFER_SIZE_PROP.defaultValue())
-          
.withSpillableMapBasePath(FileIOUtils.getDefaultSpillableMapBasePath())
-          .withOptimizedLogBlocksScan(storageConf.getBoolean("hoodie" + 
HoodieMetadataConfig.OPTIMIZED_LOG_BLOCKS_SCAN, false))
-          .withDiskMapType(storageConf.getEnum(SPILLABLE_DISK_MAP_TYPE.key(), 
SPILLABLE_DISK_MAP_TYPE.defaultValue()))
-          
.withBitCaskDiskMapCompressionEnabled(storageConf.getBoolean(DISK_MAP_BITCASK_COMPRESSION_ENABLED.key(),
 DISK_MAP_BITCASK_COMPRESSION_ENABLED.defaultValue()))
-          .withRecordMerger(recordMerger)
-          .withTableMetaClient(datasetMetaClient)
+          .withPartition(partitionPath)
+          .withAllowInflightInstants(true)
+          .withMetaClient(datasetMetaClient)
           .withAllowInflightInstants(true)
+          .withRecordBuffer(recordBuffer)
           .build();
-      return mergedLogRecordScanner.getRecords();
+      return new CloseableLogRecordsIterator<>(mergedLogRecordReader, 
recordBuffer);
+    }
+    return ClosableIterator.wrap(Collections.emptyIterator());
+  }
+
+  private static class CloseableLogRecordsIterator<T> implements 
ClosableIterator<BufferedRecord<T>> {
+    private final HoodieMergedLogRecordReader<T> mergedLogRecordReader;
+    private final FileGroupRecordBuffer<T> recordBuffer;
+    private final Iterator<BufferedRecord<T>> iterator;
+
+    public CloseableLogRecordsIterator(HoodieMergedLogRecordReader<T> 
mergedLogRecordReader, FileGroupRecordBuffer<T> recordBuffer) {
+      this.mergedLogRecordReader = mergedLogRecordReader;
+      this.recordBuffer = recordBuffer;
+      this.iterator = mergedLogRecordReader.iterator();
+    }
+
+    @Override
+    public void close() {
+      mergedLogRecordReader.close();
+      recordBuffer.close();
+    }
+
+    @Override
+    public boolean hasNext() {
+      return iterator.hasNext();
+    }
+
+    @Override
+    public BufferedRecord<T> next() {
+      return iterator.next();
     }
-    return Collections.emptyMap();
   }
 
   @VisibleForTesting
@@ -2042,7 +2077,7 @@ public class HoodieTableMetadataUtil {
         // Restore is made up of several rollbacks
         HoodieRestoreMetadata restoreMetadata = 
timeline.readRestoreMetadata(instant);
         restoreMetadata.getHoodieRestoreMetadata().values()
-                .forEach(rms -> rms.forEach(rm -> 
rollbackedCommits.addAll(rm.getCommitsRollback())));
+            .forEach(rms -> rms.forEach(rm -> 
rollbackedCommits.addAll(rm.getCommitsRollback())));
       }
       return rollbackedCommits;
     } catch (IOException e) {
@@ -2327,7 +2362,7 @@ public class HoodieTableMetadataUtil {
 
   /**
    * Reads the record keys from the base files and returns a {@link 
HoodieData} of {@link HoodieRecord} to be updated in the metadata table.
-   * Use {@link #readRecordKeysFromFileSlices(HoodieEngineContext, List, 
boolean, int, String, HoodieTableMetaClient, EngineType)} instead.
+   * Use {@link #readRecordKeysFromFileSlices} instead.
    */
   @Deprecated
   public static HoodieData<HoodieRecord> 
readRecordKeysFromBaseFiles(HoodieEngineContext engineContext,
@@ -2363,11 +2398,11 @@ public class HoodieTableMetadataUtil {
    * Reads the record keys from the given file slices and returns a {@link 
HoodieData} of {@link HoodieRecord} to be updated in the metadata table.
    * If file slice does not have any base file, then iterates over the log 
files to get the record keys.
    */
-  public static HoodieData<HoodieRecord> 
readRecordKeysFromFileSlices(HoodieEngineContext engineContext,
-                                                                      
List<Pair<String, FileSlice>> partitionFileSlicePairs,
-                                                                      boolean 
forDelete,
-                                                                      int 
recordIndexMaxParallelism,
-                                                                      String 
activeModule, HoodieTableMetaClient metaClient, EngineType engineType) {
+  public static <T> HoodieData<HoodieRecord> 
readRecordKeysFromFileSlices(HoodieEngineContext engineContext,
+                                                                          
List<Pair<String, FileSlice>> partitionFileSlicePairs,
+                                                                          int 
recordIndexMaxParallelism,
+                                                                          
String activeModule,
+                                                                          
HoodieTableMetaClient metaClient) {
     if (partitionFileSlicePairs.isEmpty()) {
       return engineContext.emptyHoodieData();
     }
@@ -2376,36 +2411,30 @@ public class HoodieTableMetadataUtil {
     final int parallelism = Math.min(partitionFileSlicePairs.size(), 
recordIndexMaxParallelism);
     final StoragePath basePath = metaClient.getBasePath();
     final StorageConfiguration<?> storageConf = metaClient.getStorageConf();
+    final Schema tableSchema;
+    try {
+      tableSchema = new TableSchemaResolver(metaClient).getTableAvroSchema();
+    } catch (Exception e) {
+      throw new HoodieException("Unable to resolve table schema for table", e);
+    }
+    ReaderContextFactory<T> readerContextFactory = 
engineContext.getReaderContextFactory(metaClient);
+    String latestCommitTime = 
metaClient.getActiveTimeline().filterCompletedInstants().lastInstant().map(HoodieInstant::requestedTime).orElse("");
     return engineContext.parallelize(partitionFileSlicePairs, 
parallelism).flatMap(partitionAndBaseFile -> {
       final String partition = partitionAndBaseFile.getKey();
       final FileSlice fileSlice = partitionAndBaseFile.getValue();
       if (!fileSlice.getBaseFile().isPresent()) {
-        List<String> logFilePaths = 
fileSlice.getLogFiles().sorted(HoodieLogFile.getLogFileComparator())
-            .map(l -> l.getPath().toString()).collect(Collectors.toList());
-        HoodieMergedLogRecordScanner mergedLogRecordScanner = 
HoodieMergedLogRecordScanner.newBuilder()
-            .withStorage(metaClient.getStorage())
-            .withBasePath(basePath)
-            .withLogFilePaths(logFilePaths)
-            .withReaderSchema(HoodieAvroUtils.getRecordKeySchema())
-            
.withLatestInstantTime(metaClient.getActiveTimeline().filterCompletedInstants().lastInstant().map(HoodieInstant::requestedTime).orElse(""))
-            .withReverseReader(false)
-            .withMaxMemorySizeInBytes(storageConf.getLong(
-                MAX_MEMORY_FOR_COMPACTION.key(), 
DEFAULT_MAX_MEMORY_FOR_SPILLABLE_MAP_IN_BYTES))
-            
.withSpillableMapBasePath(FileIOUtils.getDefaultSpillableMapBasePath())
-            .withPartition(fileSlice.getPartitionPath())
-            
.withOptimizedLogBlocksScan(storageConf.getBoolean(ENABLE_OPTIMIZED_LOG_BLOCKS_SCAN.key(),
 false))
-            
.withDiskMapType(storageConf.getEnum(SPILLABLE_DISK_MAP_TYPE.key(), 
SPILLABLE_DISK_MAP_TYPE.defaultValue()))
-            .withBitCaskDiskMapCompressionEnabled(storageConf.getBoolean(
-                DISK_MAP_BITCASK_COMPRESSION_ENABLED.key(), 
DISK_MAP_BITCASK_COMPRESSION_ENABLED.defaultValue()))
-            .withRecordMerger(HoodieRecordUtils.createRecordMerger(
-                metaClient.getBasePath().toString(),
-                engineType,
-                Collections.emptyList(), // TODO: support different merger 
classes, which is currently only known to write config
-                metaClient.getTableConfig().getRecordMergeStrategyId()))
-            .withTableMetaClient(metaClient)
+        HoodieFileGroupReader fileGroupReader = 
HoodieFileGroupReader.<T>newBuilder()
+            .withReaderContext(readerContextFactory.getContext())
+            .withHoodieTableMetaClient(metaClient)
+            .withFileSlice(fileSlice)
+            .withDataSchema(tableSchema)
+            .withRequestedSchema(HoodieAvroUtils.getRecordKeySchema())
+            .withLatestCommitTime(latestCommitTime)
+            
.withProps(getFileGroupReaderPropertiesFromStorageConf(storageConf))
             .build();
-        ClosableIterator<String> recordKeyIterator = 
ClosableIterator.wrap(mergedLogRecordScanner.getRecords().keySet().iterator());
-        return getHoodieRecordIterator(recordKeyIterator, forDelete, 
partition, fileSlice.getFileId(), fileSlice.getBaseInstantTime());
+
+        ClosableIterator<String> recordKeyIterator = 
fileGroupReader.getClosableKeyIterator();
+        return getHoodieRecordIterator(recordKeyIterator, false, partition, 
fileSlice.getFileId(), fileSlice.getBaseInstantTime());
       }
       final HoodieBaseFile baseFile = fileSlice.getBaseFile().get();
       final String filename = baseFile.getFileName();
@@ -2417,7 +2446,7 @@ public class HoodieTableMetadataUtil {
       HoodieFileReader reader = 
HoodieIOFactory.getIOFactory(metaClient.getStorage())
           .getReaderFactory(HoodieRecord.HoodieRecordType.AVRO)
           .getFileReader(hoodieConfig, dataFilePath);
-      return getHoodieRecordIterator(reader.getRecordKeyIterator(), forDelete, 
partition, fileId, instantTime);
+      return getHoodieRecordIterator(reader.getRecordKeyIterator(), false, 
partition, fileId, instantTime);
     });
   }
 
@@ -2474,16 +2503,16 @@ public class HoodieTableMetadataUtil {
       @Override
       public HoodieRecord next() {
         return forDelete
-                ? 
HoodieMetadataPayload.createRecordIndexDelete(recordKeyIterator.next())
-                : 
HoodieMetadataPayload.createRecordIndexUpdate(recordKeyIterator.next(), 
partition, fileId, instantTime, 0);
+            ? 
HoodieMetadataPayload.createRecordIndexDelete(recordKeyIterator.next())
+            : 
HoodieMetadataPayload.createRecordIndexUpdate(recordKeyIterator.next(), 
partition, fileId, instantTime, 0);
       }
     };
   }
 
   private static Stream<HoodieRecord> collectAndProcessColumnMetadata(
-          List<List<HoodieColumnRangeMetadata<Comparable>>> fileColumnMetadata,
-          String partitionPath, boolean isTightBound,
-          Map<String, Schema> colsToIndexSchemaMap
+      List<List<HoodieColumnRangeMetadata<Comparable>>> fileColumnMetadata,
+      String partitionPath, boolean isTightBound,
+      Map<String, Schema> colsToIndexSchemaMap
   ) {
     return collectAndProcessColumnMetadata(partitionPath, isTightBound, 
Option.empty(), fileColumnMetadata.stream().flatMap(List::stream), 
colsToIndexSchemaMap);
   }
@@ -3063,4 +3092,15 @@ public class HoodieTableMetadataUtil {
       return filenameToSizeMap;
     }
   }
+
+  private static TypedProperties 
getFileGroupReaderPropertiesFromStorageConf(StorageConfiguration<?> 
storageConf) {
+    TypedProperties properties = new TypedProperties();
+    properties.setProperty(MAX_MEMORY_FOR_MERGE.key(),
+        Long.toString(storageConf.getLong(MAX_MEMORY_FOR_COMPACTION.key(), 
DEFAULT_MAX_MEMORY_FOR_SPILLABLE_MAP_IN_BYTES)));
+    properties.setProperty(SPILLABLE_DISK_MAP_TYPE.key(),
+        storageConf.getEnum(SPILLABLE_DISK_MAP_TYPE.key(), 
SPILLABLE_DISK_MAP_TYPE.defaultValue()).toString());
+    properties.setProperty(DISK_MAP_BITCASK_COMPRESSION_ENABLED.key(),
+        
Boolean.toString(storageConf.getBoolean(DISK_MAP_BITCASK_COMPRESSION_ENABLED.key(),
 DISK_MAP_BITCASK_COMPRESSION_ENABLED.defaultValue())));
+    return properties;
+  }
 }
diff --git 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestHoodieTableMetadataUtil.java
 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestHoodieTableMetadataUtil.java
index 3e02dae081f..6ec7676274f 100644
--- 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestHoodieTableMetadataUtil.java
+++ 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestHoodieTableMetadataUtil.java
@@ -22,7 +22,6 @@ package org.apache.hudi.metadata;
 import org.apache.hudi.avro.HoodieAvroUtils;
 import org.apache.hudi.common.config.HoodieMetadataConfig;
 import org.apache.hudi.common.data.HoodieData;
-import org.apache.hudi.common.engine.EngineType;
 import org.apache.hudi.common.engine.HoodieLocalEngineContext;
 import org.apache.hudi.common.model.FileSlice;
 import org.apache.hudi.common.model.HoodieBaseFile;
@@ -67,6 +66,7 @@ import java.util.stream.Collectors;
 import static org.apache.hudi.avro.AvroSchemaUtils.createNullableSchema;
 import static 
org.apache.hudi.avro.TestHoodieAvroUtils.SCHEMA_WITH_AVRO_TYPES_STR;
 import static 
org.apache.hudi.avro.TestHoodieAvroUtils.SCHEMA_WITH_NESTED_FIELD_STR;
+import static 
org.apache.hudi.common.testutils.HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA;
 import static 
org.apache.hudi.metadata.HoodieTableMetadataUtil.computeRevivedAndDeletedKeys;
 import static 
org.apache.hudi.metadata.HoodieTableMetadataUtil.getFileIDForFileGroup;
 import static 
org.apache.hudi.metadata.HoodieTableMetadataUtil.validateDataTypeForSecondaryOrExpressionIndex;
@@ -106,11 +106,9 @@ public class TestHoodieTableMetadataUtil extends 
HoodieCommonTestHarness {
     HoodieData<HoodieRecord> result = 
HoodieTableMetadataUtil.readRecordKeysFromFileSlices(
         engineContext,
         partitionFileSlicePairs,
-        false,
         1,
         "activeModule",
-        metaClient,
-        EngineType.SPARK
+        metaClient
     );
     assertTrue(result.isEmpty());
   }
@@ -178,7 +176,10 @@ public class TestHoodieTableMetadataUtil extends 
HoodieCommonTestHarness {
   public void testReadRecordKeysFromBaseFilesWithValidRecords() throws 
Exception {
     HoodieLocalEngineContext engineContext = new 
HoodieLocalEngineContext(metaClient.getStorageConf());
     String instant = "20230918120000000";
-    hoodieTestTable = hoodieTestTable.addCommit(instant);
+    HoodieCommitMetadata commitMetadata = new HoodieCommitMetadata();
+    commitMetadata.setOperationType(WriteOperationType.INSERT);
+    commitMetadata.addMetadata(HoodieCommitMetadata.SCHEMA_KEY, 
TRIP_EXAMPLE_SCHEMA);
+    hoodieTestTable = hoodieTestTable.addCommit(instant, 
Option.of(commitMetadata));
     Set<String> recordKeys = new HashSet<>();
     final List<Pair<String, FileSlice>> partitionFileSlicePairs = new 
ArrayList<>();
     // Generate 10 inserts for each partition and populate 
partitionBaseFilePairs and recordKeys.
@@ -206,11 +207,9 @@ public class TestHoodieTableMetadataUtil extends 
HoodieCommonTestHarness {
     HoodieData<HoodieRecord> result = 
HoodieTableMetadataUtil.readRecordKeysFromFileSlices(
         engineContext,
         partitionFileSlicePairs,
-        false,
         1,
         "activeModule",
-        metaClient,
-        EngineType.SPARK
+        metaClient
     );
     // Validate the result.
     List<HoodieRecord> records = result.collectAsList();
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/client/functional/TestMetadataUtilRLIandSIRecordGeneration.java
 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/client/functional/TestMetadataUtilRLIandSIRecordGeneration.java
index e3740353192..9e4b8f169f1 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/client/functional/TestMetadataUtilRLIandSIRecordGeneration.java
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/client/functional/TestMetadataUtilRLIandSIRecordGeneration.java
@@ -531,8 +531,9 @@ public class TestMetadataUtilRLIandSIRecordGeneration 
extends HoodieClientTestBa
             HoodieWriteStat writeStat = writeStatus.getStat();
             StoragePath fullFilePath = new StoragePath(basePath, 
writeStat.getPath());
             // used for RLI
-            
finalActualDeletes.addAll(getRevivedAndDeletedKeysFromMergedLogs(metaClient, 
latestCommitTimestamp, EngineType.SPARK, 
Collections.singletonList(fullFilePath.toString()), writerSchemaOpt,
-                
Collections.singletonList(fullFilePath.toString())).getValue());
+            HoodieReaderContext<?> readerContext = 
context.getReaderContextFactory(metaClient).getContext();
+            
finalActualDeletes.addAll(getRevivedAndDeletedKeysFromMergedLogs(metaClient, 
latestCommitTimestamp, Collections.singletonList(fullFilePath.toString()), 
writerSchemaOpt,
+                Collections.singletonList(fullFilePath.toString()), 
writeStat.getPartitionPath(), readerContext).getValue());
 
             // used in SI flow
             
actualUpdatesAndDeletes.addAll(getRecordKeys(writeStat.getPartitionPath(), 
writeStat.getPrevCommit(), writeStat.getFileId(),


Reply via email to