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(),