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.
*/