1
0

[HUDI-2285][HUDI-2476] Metadata table synchronous design. Rebased and Squashed from pull/3426 (#3590)

* [HUDI-2285] Adding Synchronous updates to metadata before completion of commits in data timelime.

- This patch adds synchronous updates to metadata table. In other words, every write is first committed to metadata table followed by data table. While reading metadata table, we ignore any delta commits that are present only in metadata table and not in data table timeline.
- Compaction of metadata table is fenced by the condition that we trigger compaction only when there are no inflight requests in datatable. This ensures that all base files in metadata table is always in sync with data table(w/o any holes) and only there could be some extra invalid commits among delta log files in metadata table.
- Due to this, archival of data table also fences itself up until compacted instant in metadata table.
All writes to metadata table happens within the datatable lock. So, metadata table works in one writer mode only. This might be tough to loosen since all writers write to same FILES partition and so, will result in a conflict anyways.
- As part of this, have added acquiring locks in data table for those operations which were not before while committing (rollback, clean, compaction, cluster). To note, we were not doing any conflict resolution. All we are doing here is to commit by taking a lock. So that all writes to metadata table is always a single writer. 
- Also added building block to add buckets for partitions, which will be leveraged by other indexes like record level index, etc. For now, FILES partition has only one bucket. In general, any number of buckets per partition is allowed and each partition has a fixed fileId prefix with incremental suffix for each bucket within each partition.
Have fixed [HUDI-2476]. This fix is about retrying a failed compaction if it succeeded in metadata for first time, but failed w/ data table.
- Enabling metadata table by default.
- Adding more tests for metadata table

Co-authored-by: Prashant Wason <pwason@uber.com>
This commit is contained in:
Sivabalan Narayanan
2021-10-06 00:17:52 -04:00
committed by GitHub
parent 46808dcb1f
commit 5f32162a2f
101 changed files with 3329 additions and 2069 deletions

View File

@@ -23,6 +23,7 @@ import org.apache.hudi.avro.model.HoodieCleanMetadata;
import org.apache.hudi.avro.model.HoodieCleanerPlan;
import org.apache.hudi.avro.model.HoodieCompactionPlan;
import org.apache.hudi.avro.model.HoodieRequestedReplaceMetadata;
import org.apache.hudi.avro.model.HoodieRestoreMetadata;
import org.apache.hudi.avro.model.HoodieRollbackMetadata;
import org.apache.hudi.common.fs.FSUtils;
import org.apache.hudi.common.model.HoodieCommitMetadata;
@@ -60,6 +61,7 @@ import static org.apache.hudi.common.table.timeline.TimelineMetadataUtils.serial
import static org.apache.hudi.common.table.timeline.TimelineMetadataUtils.serializeCleanerPlan;
import static org.apache.hudi.common.table.timeline.TimelineMetadataUtils.serializeCompactionPlan;
import static org.apache.hudi.common.table.timeline.TimelineMetadataUtils.serializeRequestedReplaceMetadata;
import static org.apache.hudi.common.table.timeline.TimelineMetadataUtils.serializeRestoreMetadata;
import static org.apache.hudi.common.table.timeline.TimelineMetadataUtils.serializeRollbackMetadata;
public class FileCreateUtils {
@@ -130,6 +132,14 @@ public class FileCreateUtils {
}
}
private static void deleteMetaFile(String basePath, String instantTime, String suffix) throws IOException {
Path parentPath = Paths.get(basePath, HoodieTableMetaClient.METAFOLDER_NAME);
Path metaFilePath = parentPath.resolve(instantTime + suffix);
if (Files.exists(metaFilePath)) {
Files.delete(metaFilePath);
}
}
public static void createCommit(String basePath, String instantTime) throws IOException {
createMetaFile(basePath, instantTime, HoodieTimeline.COMMIT_EXTENSION);
}
@@ -150,6 +160,10 @@ public class FileCreateUtils {
createMetaFile(basePath, instantTime, HoodieTimeline.INFLIGHT_COMMIT_EXTENSION);
}
public static void createDeltaCommit(String basePath, String instantTime, HoodieCommitMetadata metadata) throws IOException {
createMetaFile(basePath, instantTime, HoodieTimeline.DELTA_COMMIT_EXTENSION, metadata.toJsonString().getBytes(StandardCharsets.UTF_8));
}
public static void createDeltaCommit(String basePath, String instantTime) throws IOException {
createMetaFile(basePath, instantTime, HoodieTimeline.DELTA_COMMIT_EXTENSION);
}
@@ -166,6 +180,10 @@ public class FileCreateUtils {
createMetaFile(basePath, instantTime, HoodieTimeline.INFLIGHT_DELTA_COMMIT_EXTENSION);
}
public static void createInflightReplaceCommit(String basePath, String instantTime) throws IOException {
createMetaFile(basePath, instantTime, HoodieTimeline.INFLIGHT_REPLACE_COMMIT_EXTENSION);
}
public static void createReplaceCommit(String basePath, String instantTime, HoodieReplaceCommitMetadata metadata) throws IOException {
createMetaFile(basePath, instantTime, HoodieTimeline.REPLACE_COMMIT_EXTENSION, metadata.toJsonString().getBytes(StandardCharsets.UTF_8));
}
@@ -210,6 +228,10 @@ public class FileCreateUtils {
createMetaFile(basePath, instantTime, HoodieTimeline.ROLLBACK_EXTENSION, serializeRollbackMetadata(hoodieRollbackMetadata).get());
}
public static void createRestoreFile(String basePath, String instantTime, HoodieRestoreMetadata hoodieRestoreMetadata) throws IOException {
createMetaFile(basePath, instantTime, HoodieTimeline.RESTORE_ACTION, serializeRestoreMetadata(hoodieRestoreMetadata).get());
}
private static void createAuxiliaryMetaFile(String basePath, String instantTime, String suffix) throws IOException {
Path parentPath = Paths.get(basePath, HoodieTableMetaClient.AUXILIARYFOLDER_NAME);
Files.createDirectories(parentPath);
@@ -224,7 +246,7 @@ public class FileCreateUtils {
}
public static void createInflightCompaction(String basePath, String instantTime) throws IOException {
createAuxiliaryMetaFile(basePath, instantTime, HoodieTimeline.REQUESTED_COMPACTION_EXTENSION);
createAuxiliaryMetaFile(basePath, instantTime, HoodieTimeline.INFLIGHT_COMPACTION_EXTENSION);
}
public static void createPartitionMetaFile(String basePath, String partitionPath) throws IOException {
@@ -309,6 +331,10 @@ public class FileCreateUtils {
removeMetaFile(basePath, instantTime, HoodieTimeline.DELTA_COMMIT_EXTENSION);
}
public static void deleteReplaceCommit(String basePath, String instantTime) throws IOException {
removeMetaFile(basePath, instantTime, HoodieTimeline.REPLACE_COMMIT_EXTENSION);
}
public static long getTotalMarkerFileCount(String basePath, String partitionPath, String instantTime, IOType ioType) throws IOException {
Path parentPath = Paths.get(basePath, HoodieTableMetaClient.TEMPFOLDER_NAME, instantTime, partitionPath);
if (Files.notExists(parentPath)) {

View File

@@ -19,15 +19,13 @@
package org.apache.hudi.common.testutils;
import org.apache.hadoop.fs.FileStatus;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hudi.avro.model.HoodieActionInstant;
import org.apache.hudi.avro.model.HoodieCleanMetadata;
import org.apache.hudi.avro.model.HoodieCleanerPlan;
import org.apache.hudi.avro.model.HoodieCompactionPlan;
import org.apache.hudi.avro.model.HoodieInstantInfo;
import org.apache.hudi.avro.model.HoodieRequestedReplaceMetadata;
import org.apache.hudi.avro.model.HoodieRestoreMetadata;
import org.apache.hudi.avro.model.HoodieRollbackMetadata;
import org.apache.hudi.avro.model.HoodieRollbackPartitionMetadata;
import org.apache.hudi.avro.model.HoodieSavepointMetadata;
@@ -55,6 +53,10 @@ import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.ValidationUtils;
import org.apache.hudi.common.util.collection.Pair;
import org.apache.hudi.exception.HoodieIOException;
import org.apache.hadoop.fs.FileStatus;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.log4j.LogManager;
import org.apache.log4j.Logger;
@@ -102,6 +104,7 @@ import static org.apache.hudi.common.testutils.FileCreateUtils.createRequestedCo
import static org.apache.hudi.common.testutils.FileCreateUtils.createRequestedCompaction;
import static org.apache.hudi.common.testutils.FileCreateUtils.createRequestedDeltaCommit;
import static org.apache.hudi.common.testutils.FileCreateUtils.createRequestedReplaceCommit;
import static org.apache.hudi.common.testutils.FileCreateUtils.createRestoreFile;
import static org.apache.hudi.common.testutils.FileCreateUtils.createRollbackFile;
import static org.apache.hudi.common.testutils.FileCreateUtils.logFileName;
import static org.apache.hudi.common.util.CleanerUtils.convertCleanMetadata;
@@ -114,7 +117,7 @@ public class HoodieTestTable {
private static final Logger LOG = LogManager.getLogger(HoodieTestTable.class);
private static final Random RANDOM = new Random();
private static HoodieTestTableState testTableState;
protected static HoodieTestTableState testTableState;
private final List<String> inflightCommits = new ArrayList<>();
protected final String basePath;
@@ -152,16 +155,19 @@ public class HoodieTestTable {
}
public static List<String> makeIncrementalCommitTimes(int num, int firstOffsetSeconds) {
return makeIncrementalCommitTimes(num, firstOffsetSeconds, 0);
}
public static List<String> makeIncrementalCommitTimes(int num, int firstOffsetSeconds, int deltaSecs) {
final Instant now = Instant.now();
return IntStream.range(0, num)
.mapToObj(i -> makeNewCommitTime(now.plus(firstOffsetSeconds + i, SECONDS)))
.mapToObj(i -> makeNewCommitTime(now.plus(deltaSecs == 0 ? (firstOffsetSeconds + i) : (i == 0 ? (firstOffsetSeconds) : (i * deltaSecs) + i), SECONDS)))
.collect(Collectors.toList());
}
public HoodieTestTable addRequestedCommit(String instantTime) throws Exception {
createRequestedCommit(basePath, instantTime);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
@@ -170,7 +176,14 @@ public class HoodieTestTable {
createInflightCommit(basePath, instantTime);
inflightCommits.add(instantTime);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
public HoodieTestTable addInflightDeltaCommit(String instantTime) throws Exception {
createRequestedDeltaCommit(basePath, instantTime);
createInflightDeltaCommit(basePath, instantTime);
inflightCommits.add(instantTime);
currentInstantTime = instantTime;
return this;
}
@@ -179,7 +192,6 @@ public class HoodieTestTable {
createInflightCommit(basePath, instantTime);
createCommit(basePath, instantTime);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
@@ -210,15 +222,17 @@ public class HoodieTestTable {
createInflightCommit(basePath, instantTime);
createCommit(basePath, instantTime, metadata);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
public HoodieTestTable moveInflightCommitToComplete(String instantTime, HoodieCommitMetadata metadata) throws IOException {
createCommit(basePath, instantTime, metadata);
if (metaClient.getTableType() == HoodieTableType.COPY_ON_WRITE) {
createCommit(basePath, instantTime, metadata);
} else {
createDeltaCommit(basePath, instantTime, metadata);
}
inflightCommits.remove(instantTime);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
@@ -227,7 +241,14 @@ public class HoodieTestTable {
createInflightDeltaCommit(basePath, instantTime);
createDeltaCommit(basePath, instantTime);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
public HoodieTestTable addDeltaCommit(String instantTime, HoodieCommitMetadata metadata) throws Exception {
createRequestedDeltaCommit(basePath, instantTime);
createInflightDeltaCommit(basePath, instantTime);
createDeltaCommit(basePath, instantTime, metadata);
currentInstantTime = instantTime;
return this;
}
@@ -240,14 +261,12 @@ public class HoodieTestTable {
createInflightReplaceCommit(basePath, instantTime, inflightReplaceMetadata);
createReplaceCommit(basePath, instantTime, completeReplaceMetadata);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
public HoodieTestTable addRequestedReplace(String instantTime, Option<HoodieRequestedReplaceMetadata> requestedReplaceMetadata) throws Exception {
createRequestedReplaceCommit(basePath, instantTime, requestedReplaceMetadata);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
@@ -255,7 +274,6 @@ public class HoodieTestTable {
createRequestedCleanFile(basePath, instantTime, cleanerPlan);
createInflightCleanFile(basePath, instantTime, cleanerPlan);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
@@ -264,7 +282,6 @@ public class HoodieTestTable {
createInflightCleanFile(basePath, instantTime, cleanerPlan);
createCleanFile(basePath, instantTime, metadata);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
@@ -296,7 +313,6 @@ public class HoodieTestTable {
public HoodieTestTable addInflightRollback(String instantTime) throws IOException {
createInflightRollbackFile(basePath, instantTime);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
@@ -304,7 +320,12 @@ public class HoodieTestTable {
createInflightRollbackFile(basePath, instantTime);
createRollbackFile(basePath, instantTime, rollbackMetadata);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
public HoodieTestTable addRestore(String instantTime, HoodieRestoreMetadata restoreMetadata) throws IOException {
createRestoreFile(basePath, instantTime, restoreMetadata);
currentInstantTime = instantTime;
return this;
}
@@ -319,7 +340,11 @@ public class HoodieTestTable {
rollbackPartitionMetadata.setSuccessDeleteFiles(entry.getValue());
rollbackPartitionMetadata.setFailedDeleteFiles(new ArrayList<>());
rollbackPartitionMetadata.setWrittenLogFiles(getWrittenLogFiles(instantTimeToDelete, entry));
rollbackPartitionMetadata.setRollbackLogFiles(createImmutableMap(logFileName(instantTimeToDelete, UUID.randomUUID().toString(), 0), (long) (100 + RANDOM.nextInt(500))));
long rollbackLogFileSize = 50 + RANDOM.nextInt(500);
String fileId = UUID.randomUUID().toString();
String logFileName = logFileName(instantTimeToDelete, fileId, 0);
FileCreateUtils.createLogFile(basePath, entry.getKey(), instantTimeToDelete, fileId, 0, (int) rollbackLogFileSize);
rollbackPartitionMetadata.setRollbackLogFiles(createImmutableMap(logFileName, rollbackLogFileSize));
partitionMetadataMap.put(entry.getKey(), rollbackPartitionMetadata);
}
rollbackMetadata.setPartitionMetadata(partitionMetadataMap);
@@ -335,7 +360,7 @@ public class HoodieTestTable {
for (String fileName : entry.getValue()) {
if (FSUtils.isLogFile(new Path(fileName))) {
if (testTableState.getPartitionToLogFileInfoMap(instant) != null
&& testTableState.getPartitionToLogFileInfoMap(instant).containsKey(entry.getKey())) {
&& testTableState.getPartitionToLogFileInfoMap(instant).containsKey(entry.getKey())) {
List<Pair<String, Integer[]>> fileInfos = testTableState.getPartitionToLogFileInfoMap(instant).get(entry.getKey());
for (Pair<String, Integer[]> fileInfo : fileInfos) {
if (fileName.equals(logFileName(instant, fileInfo.getLeft(), fileInfo.getRight()[0]))) {
@@ -366,7 +391,6 @@ public class HoodieTestTable {
public HoodieTestTable addRequestedCompaction(String instantTime) throws IOException {
createRequestedCompaction(basePath, instantTime);
currentInstantTime = instantTime;
metaClient = HoodieTableMetaClient.reload(metaClient);
return this;
}
@@ -384,11 +408,31 @@ public class HoodieTestTable {
return addRequestedCompaction(instantTime, plan);
}
public HoodieTestTable addInflightCompaction(String instantTime, HoodieCommitMetadata commitMetadata) throws Exception {
List<FileSlice> fileSlices = new ArrayList<>();
for (Map.Entry<String, List<HoodieWriteStat>> entry : commitMetadata.getPartitionToWriteStats().entrySet()) {
for (HoodieWriteStat stat: entry.getValue()) {
fileSlices.add(new FileSlice(entry.getKey(), instantTime, stat.getPath()));
}
}
this.addRequestedCompaction(instantTime, fileSlices.toArray(new FileSlice[0]));
createInflightCompaction(basePath, instantTime);
inflightCommits.add(instantTime);
currentInstantTime = instantTime;
return this;
}
public HoodieTestTable addCompaction(String instantTime, HoodieCommitMetadata commitMetadata) throws Exception {
createRequestedCompaction(basePath, instantTime);
createInflightCompaction(basePath, instantTime);
return HoodieTestTable.of(metaClient)
.addCommit(instantTime, commitMetadata);
return addCommit(instantTime, commitMetadata);
}
public HoodieTestTable moveInflightCompactionToComplete(String instantTime, HoodieCommitMetadata metadata) throws IOException {
createCommit(basePath, instantTime, metadata);
inflightCommits.remove(instantTime);
currentInstantTime = instantTime;
return this;
}
public HoodieTestTable forCommit(String instantTime) {
@@ -648,6 +692,7 @@ public class HoodieTestTable {
}
public HoodieTestTable doRollback(String commitTimeToRollback, String commitTime) throws Exception {
metaClient = HoodieTableMetaClient.reload(metaClient);
Option<HoodieCommitMetadata> commitMetadata = getMetadataForInstant(commitTimeToRollback);
if (!commitMetadata.isPresent()) {
throw new IllegalArgumentException("Instant to rollback not present in timeline: " + commitTimeToRollback);
@@ -660,7 +705,32 @@ public class HoodieTestTable {
return addRollback(commitTime, rollbackMetadata);
}
public HoodieTestTable doCluster(String commitTime, Map<String, List<String>> partitionToReplaceFileIds) throws Exception {
public HoodieTestTable doRestore(String commitToRestoreTo, String restoreTime) throws Exception {
metaClient = HoodieTableMetaClient.reload(metaClient);
List<HoodieInstant> commitsToRollback = metaClient.getActiveTimeline().getCommitsTimeline()
.filterCompletedInstants().findInstantsAfter(commitToRestoreTo).getReverseOrderedInstants().collect(Collectors.toList());
Map<String, List<HoodieRollbackMetadata>> rollbackMetadataMap = new HashMap<>();
for (HoodieInstant commitInstantToRollback: commitsToRollback) {
Option<HoodieCommitMetadata> commitMetadata = getCommitMeta(commitInstantToRollback);
if (!commitMetadata.isPresent()) {
throw new IllegalArgumentException("Instant to rollback not present in timeline: " + commitInstantToRollback.getTimestamp());
}
Map<String, List<String>> partitionFiles = getPartitionFiles(commitMetadata.get());
rollbackMetadataMap.put(commitInstantToRollback.getTimestamp(),
Collections.singletonList(getRollbackMetadata(commitInstantToRollback.getTimestamp(), partitionFiles)));
for (Map.Entry<String, List<String>> entry : partitionFiles.entrySet()) {
deleteFilesInPartition(entry.getKey(), entry.getValue());
}
}
HoodieRestoreMetadata restoreMetadata = TimelineMetadataUtils.convertRestoreMetadata(restoreTime,1000L,
commitsToRollback, rollbackMetadataMap);
return addRestore(restoreTime, restoreMetadata);
}
public HoodieReplaceCommitMetadata doCluster(String commitTime, Map<String, List<String>> partitionToReplaceFileIds, List<String> partitions, int filesPerPartition) throws Exception {
HoodieTestTableState testTableState = getTestTableStateWithPartitionFileInfo(CLUSTER, metaClient.getTableType(), commitTime, partitions, filesPerPartition);
this.currentInstantTime = commitTime;
Map<String, List<Pair<String, Integer>>> partitionToReplaceFileIdsWithLength = new HashMap<>();
for (Map.Entry<String, List<String>> entry : partitionToReplaceFileIds.entrySet()) {
String partition = entry.getKey();
@@ -670,10 +740,15 @@ public class HoodieTestTable {
partitionToReplaceFileIdsWithLength.get(partition).add(Pair.of(fileId, length));
}
}
List<HoodieWriteStat> writeStats = generateHoodieWriteStatForPartition(partitionToReplaceFileIdsWithLength, commitTime, false);
List<HoodieWriteStat> writeStats = generateHoodieWriteStatForPartition(testTableState.getPartitionToBaseFileInfoMap(commitTime), commitTime, false);
for (String partition : testTableState.getPartitionToBaseFileInfoMap(commitTime).keySet()) {
this.withBaseFilesInPartition(partition, testTableState.getPartitionToBaseFileInfoMap(commitTime).get(partition));
}
HoodieReplaceCommitMetadata replaceMetadata =
(HoodieReplaceCommitMetadata) buildMetadata(writeStats, partitionToReplaceFileIds, Option.empty(), CLUSTER, EMPTY_STRING, REPLACE_COMMIT_ACTION);
return addReplaceCommit(commitTime, Option.empty(), Option.empty(), replaceMetadata);
(HoodieReplaceCommitMetadata) buildMetadata(writeStats, partitionToReplaceFileIds, Option.empty(), CLUSTER, EMPTY_STRING,
REPLACE_COMMIT_ACTION);
addReplaceCommit(commitTime, Option.empty(), Option.empty(), replaceMetadata);
return replaceMetadata;
}
public HoodieCleanMetadata doClean(String commitTime, Map<String, Integer> partitionFileCountsToDelete) throws IOException {
@@ -718,7 +793,11 @@ public class HoodieTestTable {
return savepointMetadata;
}
public HoodieTestTable doCompaction(String commitTime, List<String> partitions) throws Exception {
public HoodieCommitMetadata doCompaction(String commitTime, List<String> partitions) throws Exception {
return doCompaction(commitTime, partitions, false);
}
public HoodieCommitMetadata doCompaction(String commitTime, List<String> partitions, boolean inflight) throws Exception {
this.currentInstantTime = commitTime;
if (partitions.isEmpty()) {
partitions = Collections.singletonList(EMPTY_STRING);
@@ -728,7 +807,12 @@ public class HoodieTestTable {
for (String partition : partitions) {
this.withBaseFilesInPartition(partition, testTableState.getPartitionToBaseFileInfoMap(commitTime).get(partition));
}
return addCompaction(commitTime, commitMetadata);
if (inflight) {
this.addInflightCompaction(commitTime, commitMetadata);
} else {
this.addCompaction(commitTime, commitMetadata);
}
return commitMetadata;
}
public HoodieCommitMetadata doWriteOperation(String commitTime, WriteOperationType operationType,
@@ -765,9 +849,17 @@ public class HoodieTestTable {
this.withPartitionMetaFiles(str);
}
if (createInflightCommit) {
this.addInflightCommit(commitTime);
if (metaClient.getTableType() == HoodieTableType.COPY_ON_WRITE) {
this.addInflightCommit(commitTime);
} else {
this.addInflightDeltaCommit(commitTime);
}
} else {
this.addCommit(commitTime, commitMetadata);
if (metaClient.getTableType() == HoodieTableType.COPY_ON_WRITE) {
this.addCommit(commitTime, commitMetadata);
} else {
this.addDeltaCommit(commitTime, commitMetadata);
}
}
for (String partition : partitions) {
this.withBaseFilesInPartition(partition, testTableState.getPartitionToBaseFileInfoMap(commitTime).get(partition));
@@ -779,23 +871,12 @@ public class HoodieTestTable {
}
private Option<HoodieCommitMetadata> getMetadataForInstant(String instantTime) {
metaClient = HoodieTableMetaClient.reload(metaClient);
Option<HoodieInstant> hoodieInstant = metaClient.getActiveTimeline().getCommitsTimeline()
.filterCompletedInstants().filter(i -> i.getTimestamp().equals(instantTime)).firstInstant();
try {
if (hoodieInstant.isPresent()) {
switch (hoodieInstant.get().getAction()) {
case HoodieTimeline.REPLACE_COMMIT_ACTION:
HoodieReplaceCommitMetadata replaceCommitMetadata = HoodieReplaceCommitMetadata
.fromBytes(metaClient.getActiveTimeline().getInstantDetails(hoodieInstant.get()).get(), HoodieReplaceCommitMetadata.class);
return Option.of(replaceCommitMetadata);
case HoodieTimeline.DELTA_COMMIT_ACTION:
case HoodieTimeline.COMMIT_ACTION:
HoodieCommitMetadata commitMetadata = HoodieCommitMetadata
.fromBytes(metaClient.getActiveTimeline().getInstantDetails(hoodieInstant.get()).get(), HoodieCommitMetadata.class);
return Option.of(commitMetadata);
default:
throw new IllegalArgumentException("Unknown instant action" + hoodieInstant.get().getAction());
}
return getCommitMeta(hoodieInstant.get());
} else {
return Option.empty();
}
@@ -804,6 +885,22 @@ public class HoodieTestTable {
}
}
private Option<HoodieCommitMetadata> getCommitMeta(HoodieInstant hoodieInstant) throws IOException {
switch (hoodieInstant.getAction()) {
case HoodieTimeline.REPLACE_COMMIT_ACTION:
HoodieReplaceCommitMetadata replaceCommitMetadata = HoodieReplaceCommitMetadata
.fromBytes(metaClient.getActiveTimeline().getInstantDetails(hoodieInstant).get(), HoodieReplaceCommitMetadata.class);
return Option.of(replaceCommitMetadata);
case HoodieTimeline.DELTA_COMMIT_ACTION:
case HoodieTimeline.COMMIT_ACTION:
HoodieCommitMetadata commitMetadata = HoodieCommitMetadata
.fromBytes(metaClient.getActiveTimeline().getInstantDetails(hoodieInstant).get(), HoodieCommitMetadata.class);
return Option.of(commitMetadata);
default:
throw new IllegalArgumentException("Unknown instant action" + hoodieInstant.getAction());
}
}
private static Map<String, List<String>> getPartitionFiles(HoodieCommitMetadata commitMetadata) {
Map<String, List<String>> partitionFilesToDelete = new HashMap<>();
Map<String, List<HoodieWriteStat>> partitionToWriteStats = commitMetadata.getPartitionToWriteStats();
@@ -815,7 +912,7 @@ public class HoodieTestTable {
}
private static HoodieTestTableState getTestTableStateWithPartitionFileInfo(WriteOperationType operationType, HoodieTableType tableType, String commitTime,
List<String> partitions, int filesPerPartition) {
List<String> partitions, int filesPerPartition) {
for (String partition : partitions) {
Stream<Integer> fileLengths = IntStream.range(0, filesPerPartition).map(i -> 100 + RANDOM.nextInt(500)).boxed();
if (MERGE_ON_READ.equals(tableType) && UPSERT.equals(operationType)) {
@@ -861,7 +958,7 @@ public class HoodieTestTable {
for (Pair<String, Integer[]> fileIdInfo : entry.getValue()) {
HoodieWriteStat writeStat = new HoodieWriteStat();
String fileName = bootstrap ? fileIdInfo.getKey() :
FileCreateUtils.logFileName(commitTime, fileIdInfo.getKey(), fileIdInfo.getValue()[0]);
FileCreateUtils.logFileName(commitTime, fileIdInfo.getKey(), fileIdInfo.getValue()[0]);
writeStat.setFileId(fileName);
writeStat.setPartitionPath(partition);
writeStat.setPath(partition + "/" + fileName);