This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 169b92eb54 [core] Refactor sort compact to produce COMPACT commits
(#8631)
169b92eb54 is described below
commit 169b92eb541628106b9e0d8b52270e3387dcfbd9
Author: Bingeng Huang <[email protected]>
AuthorDate: Thu Sep 24 17:35:21 2026 +0800
[core] Refactor sort compact to produce COMPACT commits (#8631)
---
.../append/SortCompactCommitMessageRewriter.java | 718 ++++++++++
.../paimon/append/SortCompactPlanMetadata.java | 165 +++
.../java/org/apache/paimon/io/DataFileMeta.java | 2 +
.../org/apache/paimon/io/PojoDataFileMeta.java | 26 +
.../apache/paimon/io/ProjectedDataFileMeta.java | 5 +
.../SortCompactCommitMessageRewriterTest.java | 1370 ++++++++++++++++++++
6 files changed, 2286 insertions(+)
diff --git
a/paimon-core/src/main/java/org/apache/paimon/append/SortCompactCommitMessageRewriter.java
b/paimon-core/src/main/java/org/apache/paimon/append/SortCompactCommitMessageRewriter.java
new file mode 100644
index 0000000000..77c5d537e4
--- /dev/null
+++
b/paimon-core/src/main/java/org/apache/paimon/append/SortCompactCommitMessageRewriter.java
@@ -0,0 +1,718 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.append;
+
+import org.apache.paimon.Snapshot;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.deletionvectors.append.AppendDeleteFileMaintainer;
+import org.apache.paimon.deletionvectors.append.BaseAppendDeleteFileMaintainer;
+import org.apache.paimon.index.DeletionVectorMeta;
+import org.apache.paimon.index.IndexFileHandler;
+import org.apache.paimon.index.IndexFileMeta;
+import org.apache.paimon.io.CompactIncrement;
+import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.io.DataIncrement;
+import org.apache.paimon.manifest.FileKind;
+import org.apache.paimon.manifest.IndexManifestEntry;
+import org.apache.paimon.manifest.ManifestEntry;
+import org.apache.paimon.operation.FileStoreScan;
+import org.apache.paimon.table.BucketMode;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.sink.BatchWriteBuilder;
+import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.sink.CommitMessageImpl;
+import org.apache.paimon.table.sink.TableCommit;
+import org.apache.paimon.table.source.DataSplit;
+
+import javax.annotation.Nullable;
+
+import java.io.FileNotFoundException;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+
+import static
org.apache.paimon.deletionvectors.DeletionVectorsIndexFile.DELETION_VECTORS_INDEX;
+
+/**
+ * Rewrites the {@link CommitMessage}s produced by a sort compact write into
compact commit
+ * messages, so that the commit is a {@link Snapshot.CommitKind#COMPACT}
commit instead of an {@link
+ * Snapshot.CommitKind#OVERWRITE} commit.
+ *
+ * <p>The sort compact write only creates new data files (the sorted output).
The old files which
+ * are replaced by the sort compact are captured by the planned input {@link
DataSplit}s. This
+ * helper rewrites them into compact changes: all planned old files become
{@code compactBefore} for
+ * their original (partition, bucket), and all newly written files become
{@code compactAfter} for
+ * their output (partition, bucket). The result is a {@link CommitMessageImpl}
with an empty {@link
+ * DataIncrement} and a populated {@link CompactIncrement}.
+ *
+ * <p>For deletion-vector enabled append tables, only the deletion-vector
index entries captured
+ * from the base snapshot are cleaned up, mirroring {@link
+ * org.apache.paimon.append.AppendCompactTask}. Concurrent deletion-vector
writes after the base
+ * snapshot are <b>not</b> merged into cleanup; if deletion vectors on input
files changed, rewrite
+ * or commit conflict detection fails with an explicit error so the job can be
retried.
+ */
+public class SortCompactCommitMessageRewriter {
+
+ private static final String ABORT_COMMIT_USER = "sort-compact-abort";
+ private static final int DV_DRIFT_SAMPLE_LIMIT = 5;
+
+ private final FileStoreTable table;
+ private final long baseSnapshotId;
+
+ /** Old files grouped by partition then bucket, captured from the planned
input splits. */
+ private final Map<BinaryRow, Map<Integer, List<DataFileMeta>>>
compactBeforeFiles;
+
+ /**
+ * Total bucket counts grouped like {@link #compactBeforeFiles}, when
carried by input splits.
+ */
+ private final Map<BinaryRow, Map<Integer, Integer>>
compactBeforeTotalBuckets;
+
+ /**
+ * Deletion-vector index entries captured from the base snapshot at
planning time, grouped by
+ * partition.
+ */
+ private final Map<BinaryRow, List<IndexManifestEntry>>
baseDeletionVectorEntries;
+
+ /**
+ * Whether {@link #baseDeletionVectorEntries} is a known base-snapshot
state (including a known
+ * empty map). False only when the base snapshot was already missing at
construction and no
+ * {@link SortCompactPlanMetadata} was provided.
+ */
+ private final boolean baseDeletionVectorStateKnown;
+
+ public SortCompactCommitMessageRewriter(
+ FileStoreTable table, long baseSnapshotId, List<DataSplit>
compactInputSplits) {
+ this(table, baseSnapshotId, compactInputSplits, null);
+ }
+
+ public SortCompactCommitMessageRewriter(
+ FileStoreTable table,
+ long baseSnapshotId,
+ List<DataSplit> compactInputSplits,
+ @Nullable SortCompactPlanMetadata planMetadata) {
+ this.table = table;
+ this.baseSnapshotId = baseSnapshotId;
+ this.compactBeforeFiles = new HashMap<>();
+ this.compactBeforeTotalBuckets = new HashMap<>();
+ this.baseDeletionVectorEntries = new HashMap<>();
+ Set<BinaryRow> partitions = new HashSet<>();
+ for (DataSplit split : compactInputSplits) {
+ partitions.add(split.partition());
+ compactBeforeFiles
+ .computeIfAbsent(split.partition(), k -> new HashMap<>())
+ .computeIfAbsent(split.bucket(), k -> new ArrayList<>())
+ .addAll(split.dataFiles());
+ if (split.totalBuckets() != null) {
+ Integer previous =
+ compactBeforeTotalBuckets
+ .computeIfAbsent(split.partition(), k -> new
HashMap<>())
+ .putIfAbsent(split.bucket(),
split.totalBuckets());
+ if (previous != null &&
!previous.equals(split.totalBuckets())) {
+ throw new IllegalArgumentException(
+ String.format(
+ "Conflicting total bucket counts for
partition %s bucket %s: "
+ + "%s and %s.",
+ split.partition(),
+ split.bucket(),
+ previous,
+ split.totalBuckets()));
+ }
+ }
+ }
+ if (planMetadata != null) {
+ planMetadata.copyInto(baseDeletionVectorEntries);
+ this.baseDeletionVectorStateKnown =
planMetadata.baseSnapshotCaptured();
+ } else {
+ this.baseDeletionVectorStateKnown =
+ SortCompactPlanMetadata.captureInto(
+ table, baseSnapshotId, partitions,
baseDeletionVectorEntries);
+ }
+ }
+
+ /**
+ * Rewrite the given written append commit messages into compact commit
messages.
+ *
+ * <p>Both the planned input splits and the written messages are grouped
by (partition, bucket).
+ * Planned input groups are always emitted, even if the sort compact write
produces no files for
+ * that group. Written output groups which were not present in the planned
input are emitted as
+ * add-only compact messages.
+ *
+ * @param writtenMessages commit messages produced by the sort compact
write stage (only new
+ * files in {@link DataIncrement})
+ * @return rewritten commit messages carrying {@link CompactIncrement}s
+ */
+ public List<CommitMessage> rewrite(List<CommitMessage> writtenMessages) {
+ validateWriteOnlyMessages(writtenMessages);
+ validateNoDeletionVectorDrift(writtenMessages);
+
+ // group written messages by (partition, bucket)
+ Map<BinaryRow, Map<Integer, List<CommitMessageImpl>>> grouped = new
HashMap<>();
+ for (CommitMessage written : writtenMessages) {
+ CommitMessageImpl impl = (CommitMessageImpl) written;
+ grouped.computeIfAbsent(impl.partition(), k -> new HashMap<>())
+ .computeIfAbsent(impl.bucket(), k -> new ArrayList<>())
+ .add(impl);
+ }
+
+ List<CommitMessage> result = new ArrayList<>();
+ try {
+ for (Map.Entry<BinaryRow, Map<Integer, List<DataFileMeta>>>
partitionEntry :
+ compactBeforeFiles.entrySet()) {
+ BinaryRow partition = partitionEntry.getKey();
+ Map<Integer, List<CommitMessageImpl>> writtenInPartition =
grouped.get(partition);
+ for (Integer bucket : partitionEntry.getValue().keySet()) {
+ List<CommitMessageImpl> group = Collections.emptyList();
+ if (writtenInPartition != null) {
+ List<CommitMessageImpl> writtenGroup =
writtenInPartition.remove(bucket);
+ if (writtenGroup != null) {
+ group = writtenGroup;
+ }
+ }
+ result.add(rewriteGroup(partition, bucket, group));
+ }
+ if (writtenInPartition != null &&
writtenInPartition.isEmpty()) {
+ grouped.remove(partition);
+ }
+ }
+
+ for (Map.Entry<BinaryRow, Map<Integer, List<CommitMessageImpl>>>
partitionEntry :
+ grouped.entrySet()) {
+ BinaryRow partition = partitionEntry.getKey();
+ for (Map.Entry<Integer, List<CommitMessageImpl>> bucketEntry :
+ partitionEntry.getValue().entrySet()) {
+ result.add(
+ rewriteGroup(partition, bucketEntry.getKey(),
bucketEntry.getValue()));
+ }
+ }
+ } catch (RuntimeException e) {
+ // A partial rewrite may have already persisted new DV index files
in result.
+ // The sorted output files in writtenMessages must also be aborted
so callers
+ // (for example Flink's rewriteAll outside commit try/catch) do
not leave orphans.
+ abortQuietly(result, e);
+ abortQuietly(writtenMessages, e);
+ throw e;
+ }
+ return result;
+ }
+
+ /**
+ * Validate that the written messages only contain write-only append
output. Sort compact must
+ * not run inline compaction in the write stage; otherwise compact output
would be dropped and
+ * orphan files would be left on disk.
+ */
+ private void validateWriteOnlyMessages(List<CommitMessage>
writtenMessages) {
+ for (CommitMessage written : writtenMessages) {
+ CommitMessageImpl impl = (CommitMessageImpl) written;
+ CompactIncrement compactIncrement = impl.compactIncrement();
+ if (!compactIncrement.compactBefore().isEmpty()
+ || !compactIncrement.compactAfter().isEmpty()) {
+ abortAndFail(
+ writtenMessages,
+ String.format(
+ "Sort compact write produced inline compaction
changes for "
+ + "partition %s bucket %s
(compactBefore = %s, "
+ + "compactAfter = %s). The write stage
must run in "
+ + "write-only mode without waiting for
compaction.",
+ impl.partition(),
+ impl.bucket(),
+ compactIncrement.compactBefore(),
+ compactIncrement.compactAfter()));
+ }
+ }
+ }
+
+ /**
+ * Fail fast when deletion vectors on compact-before files changed after
the base snapshot.
+ *
+ * <p>Sort compact reads rows from the base snapshot. Committing after a
concurrent DV write
+ * would drop the newer deletion vectors and restore deleted rows. This
check only validates;
+ * cleanup still uses {@link #baseDeletionVectorEntries} only.
+ */
+ private void validateNoDeletionVectorDrift(List<CommitMessage>
writtenMessages) {
+ if (!table.coreOptions().deletionVectorsEnabled()
+ || table.bucketMode() != BucketMode.BUCKET_UNAWARE
+ || !hasInput()) {
+ return;
+ }
+
+ // Without a known base DV state we cannot tell whether DVs changed;
skip the proactive
+ // check and rely on commit conflict detection.
+ if (!baseDeletionVectorStateKnown) {
+ return;
+ }
+
+ Snapshot latestSnapshot = tryLatestSnapshot();
+ // No readable latest snapshot (e.g. the only snapshot was expired):
nothing to compare.
+ if (latestSnapshot == null) {
+ return;
+ }
+
+ Map<BinaryRow, List<IndexManifestEntry>> latestDeletionVectorEntries =
+ scanDeletionVectorEntries(latestSnapshot);
+ Long latestSnapshotId = latestSnapshot.id();
+
+ BinaryRow firstChangedPartition = null;
+ List<String> changedSamples = new ArrayList<>();
+ for (Map.Entry<BinaryRow, Map<Integer, List<DataFileMeta>>>
partitionEntry :
+ compactBeforeFiles.entrySet()) {
+ BinaryRow partition = partitionEntry.getKey();
+ Map<String, String> baseDvByDataFile =
+ dataFileToDvIndexFileName(
+ baseDeletionVectorEntries.getOrDefault(
+ partition, Collections.emptyList()));
+ Map<String, String> latestDvByDataFile =
+ dataFileToDvIndexFileName(
+ latestDeletionVectorEntries.getOrDefault(
+ partition, Collections.emptyList()));
+ for (List<DataFileMeta> files :
partitionEntry.getValue().values()) {
+ for (DataFileMeta file : files) {
+ String baseDv = baseDvByDataFile.get(file.fileName());
+ String latestDv = latestDvByDataFile.get(file.fileName());
+ if (Objects.equals(baseDv, latestDv)) {
+ continue;
+ }
+ if (firstChangedPartition == null) {
+ firstChangedPartition = partition;
+ }
+ if (changedSamples.size() < DV_DRIFT_SAMPLE_LIMIT) {
+ changedSamples.add(
+ String.format(
+ "%s (baseDv=%s -> latestDv=%s)",
+ file.fileName(), baseDv, latestDv));
+ }
+ }
+ }
+ }
+
+ if (firstChangedPartition != null) {
+ abortAndFail(
+ writtenMessages,
+ deletionVectorDriftMessage(
+ firstChangedPartition, changedSamples,
latestSnapshotId));
+ }
+ }
+
+ private String deletionVectorDriftMessage(
+ BinaryRow partition, List<String> changedSamples, @Nullable Long
latestSnapshotId) {
+ return "Sort compact cannot commit because deletion vectors on input
files changed after the base snapshot. "
+ + "Sort compact reads data from the base snapshot, so
committing would drop newer deletion vectors and restore deleted rows. "
+ + "Changed files (partition="
+ + partition
+ + ", sample): "
+ + String.join(", ", changedSamples)
+ + ". baseSnapshotId="
+ + baseSnapshotId
+ + ", latestSnapshotId="
+ + latestSnapshotId
+ + ". Please retry the sort compact job after concurrent
deletes/updates have finished.";
+ }
+
+ /**
+ * Abort newly written files when sort compact rewrite or commit fails.
+ *
+ * <p>This only covers files tracked by the original append write
messages. The new
+ * deletion-vector index files produced by {@code dvMaintainer.persist()}
during {@link
+ * #rewrite} live in the rewritten compact messages and must be cleaned up
via {@link
+ * #abortCompactMessages}.
+ */
+ public void abortWrittenMessages(List<CommitMessage> writtenMessages) {
+ abortMessages(writtenMessages);
+ }
+
+ /**
+ * Abort the rewritten compact messages, including the new deletion-vector
index files produced
+ * by {@code dvMaintainer.persist()} during {@link #rewrite}.
+ *
+ * <p>These index files are not referenced by the original written
messages, so aborting only
+ * the written messages (as failure cleanup used to do) orphans them. For
delete-only compact
+ * the written messages are empty, so the rewritten compact messages are
the only place the new
+ * DV index files are tracked. {@code commit.abort} only deletes new files
({@code
+ * compactAfter}, {@code newIndexFiles}); planned {@code compactBefore}
files and {@code
+ * deletedIndexFiles} (still referenced by the latest snapshot) are left
untouched.
+ */
+ public void abortCompactMessages(List<CommitMessage> compactMessages) {
+ abortMessages(compactMessages);
+ }
+
+ private void abortMessages(List<CommitMessage> messages) {
+ if (messages.isEmpty()) {
+ return;
+ }
+ try (TableCommit commit = table.newCommit(ABORT_COMMIT_USER)) {
+ commit.abort(messages);
+ } catch (Exception e) {
+ throw new IllegalStateException(
+ "Failed to clean up sort compact write output before
aborting commit.", e);
+ }
+ }
+
+ private void abortQuietly(List<CommitMessage> messages, RuntimeException
cause) {
+ try {
+ abortMessages(messages);
+ } catch (Exception abortException) {
+ cause.addSuppressed(abortException);
+ }
+ }
+
+ private void abortAndFail(List<CommitMessage> writtenMessages, String
message) {
+ abortWrittenMessages(writtenMessages);
+ throw new IllegalStateException(message);
+ }
+
+ private CommitMessage rewriteGroup(
+ BinaryRow partition, int bucket, List<CommitMessageImpl> group) {
+ List<DataFileMeta> compactBefore = compactBefore(partition, bucket);
+
+ // merge all newly written sorted files of this (partition, bucket) as
compact output
+ List<DataFileMeta> compactAfter = new ArrayList<>();
+ List<IndexFileMeta> newIndexFiles = new ArrayList<>();
+ List<IndexFileMeta> deletedIndexFiles = new ArrayList<>();
+ Integer totalBuckets = null;
+ for (CommitMessageImpl impl : group) {
+
compactAfter.addAll(toCompactAfter(impl.newFilesIncrement().newFiles()));
+ newIndexFiles.addAll(impl.compactIncrement().newIndexFiles());
+ newIndexFiles.addAll(impl.newFilesIncrement().newIndexFiles());
+
deletedIndexFiles.addAll(impl.compactIncrement().deletedIndexFiles());
+
deletedIndexFiles.addAll(impl.newFilesIncrement().deletedIndexFiles());
+ if (totalBuckets == null) {
+ totalBuckets = impl.totalBuckets();
+ }
+ }
+ if (totalBuckets == null) {
+ totalBuckets = compactBeforeTotalBuckets(partition, bucket);
+ }
+
+ // for deletion-vector append tables, clean up only the base-snapshot
DV index entries of
+ // removed old files (do not merge concurrent latest-snapshot DVs)
+ if (table.coreOptions().deletionVectorsEnabled()
+ && table.bucketMode() == BucketMode.BUCKET_UNAWARE
+ && !compactBefore.isEmpty()) {
+ AppendDeleteFileMaintainer dvMaintainer =
+ BaseAppendDeleteFileMaintainer.forUnawareAppend(
+ table.store().newIndexFileHandler(),
+ partition,
+ baseDeletionVectorEntries.getOrDefault(
+ partition, Collections.emptyList()));
+ for (DataFileMeta oldFile : compactBefore) {
+ dvMaintainer.notifyRemovedDeletionVector(oldFile.fileName());
+ }
+ for (IndexManifestEntry entry : dvMaintainer.persist()) {
+ if (entry.kind() == FileKind.ADD) {
+ newIndexFiles.add(entry.indexFile());
+ } else {
+ deletedIndexFiles.add(entry.indexFile());
+ }
+ }
+ }
+
+ CompactIncrement compactIncrement =
+ new CompactIncrement(
+ compactBefore,
+ compactAfter,
+ Collections.emptyList(),
+ newIndexFiles,
+ deletedIndexFiles);
+ return new CommitMessageImpl(
+ partition, bucket, totalBuckets,
DataIncrement.emptyIncrement(), compactIncrement);
+ }
+
+ private Map<BinaryRow, List<IndexManifestEntry>> scanDeletionVectorEntries(
+ @Nullable Snapshot snapshot) {
+ Map<BinaryRow, List<IndexManifestEntry>> entries = new HashMap<>();
+ if (snapshot == null) {
+ return entries;
+ }
+ IndexFileHandler indexFileHandler =
table.store().newIndexFileHandler();
+ for (IndexManifestEntry entry : indexFileHandler.scan(snapshot,
DELETION_VECTORS_INDEX)) {
+ entries.computeIfAbsent(entry.partition(), k -> new
ArrayList<>()).add(entry);
+ }
+ return entries;
+ }
+
+ private static Map<String, String>
dataFileToDvIndexFileName(List<IndexManifestEntry> entries) {
+ Map<String, String> result = new HashMap<>();
+ for (IndexManifestEntry entry : entries) {
+ LinkedHashMap<String, DeletionVectorMeta> dvRanges =
entry.indexFile().dvRanges();
+ if (dvRanges == null) {
+ continue;
+ }
+ String indexFileName = entry.indexFile().fileName();
+ for (String dataFileName : dvRanges.keySet()) {
+ result.put(dataFileName, indexFileName);
+ }
+ }
+ return result;
+ }
+
+ @Nullable
+ private Snapshot tryLatestSnapshot() {
+ Long latestId = table.snapshotManager().latestSnapshotId();
+ if (latestId == null) {
+ return null;
+ }
+ try {
+ return table.snapshotManager().tryGetSnapshot(latestId);
+ } catch (FileNotFoundException e) {
+ return null;
+ }
+ }
+
+ private List<DataFileMeta> toCompactAfter(List<DataFileMeta> newFiles) {
+ if (newFiles.isEmpty()) {
+ return newFiles;
+ }
+ List<DataFileMeta> result = new ArrayList<>(newFiles.size());
+ for (DataFileMeta newFile : newFiles) {
+ if (newFile.fileSource().orElse(null)
+ == org.apache.paimon.manifest.FileSource.COMPACT) {
+ result.add(newFile);
+ continue;
+ }
+
result.add(newFile.assignFileSource(org.apache.paimon.manifest.FileSource.COMPACT));
+ }
+ return result;
+ }
+
+ private List<DataFileMeta> compactBefore(BinaryRow partition, int bucket) {
+ return compactBeforeFiles
+ .getOrDefault(partition, Collections.emptyMap())
+ .getOrDefault(bucket, Collections.emptyList());
+ }
+
+ @Nullable
+ private Integer compactBeforeTotalBuckets(BinaryRow partition, int bucket)
{
+ return compactBeforeTotalBuckets
+ .getOrDefault(partition, Collections.emptyMap())
+ .get(bucket);
+ }
+
+ /** Latest snapshot id, or 0 when the table has no snapshot yet. */
+ public long latestSnapshotIdOrZero() {
+ Long latestId = table.snapshotManager().latestSnapshotId();
+ return latestId == null ? 0L : latestId;
+ }
+
+ /**
+ * Whether the given compact commit messages were already committed after
{@code
+ * snapshotIdBeforeCommit}.
+ *
+ * <p>Used to avoid aborting sort compact write output when {@link
+ * org.apache.paimon.table.sink.TableCommitImpl} fails after the snapshot
is already visible.
+ * Matching is based on the unique new files written by this rewrite: the
compact output data
+ * files, or the new deletion-vector index files for delete-only commits
(whose input files may
+ * also be deleted by a concurrent compaction, so removed input files
alone prove nothing).
+ */
+ public boolean isBatchCompactCommitSucceeded(
+ long snapshotIdBeforeCommit, List<CommitMessage> compactMessages) {
+ Long latestId = table.snapshotManager().latestSnapshotId();
+ if (latestId == null || latestId <= snapshotIdBeforeCommit) {
+ return false;
+ }
+
+ CompactCommitFingerprint fingerprint =
CompactCommitFingerprint.from(compactMessages);
+ Snapshot latestSnapshot = table.snapshotManager().snapshot(latestId);
+ if (matchesCompactCommit(latestSnapshot, fingerprint)) {
+ return true;
+ }
+
+ // Check older snapshots in reverse order. The compact commit is
usually near latest, and
+ // scoped scans below avoid full-table manifest reads when probing
history.
+ for (long id = latestId - 1; id > snapshotIdBeforeCommit; id--) {
+ Snapshot snapshot;
+ try {
+ snapshot = table.snapshotManager().tryGetSnapshot(id);
+ } catch (FileNotFoundException e) {
+ // Expired snapshots create gaps. Keep scanning older
snapshots instead of treating
+ // the commit as failed.
+ continue;
+ }
+ if (matchesCompactCommit(snapshot, fingerprint)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private boolean matchesCompactCommit(Snapshot snapshot,
CompactCommitFingerprint fingerprint) {
+ // Output data file names and new deletion-vector index file names are
unique, and
+ // snapshot commit is atomic. Finding any of them in a surviving
snapshot proves THIS
+ // batch compact commit succeeded, even when concurrent compaction has
replaced other
+ // outputs or the COMPACT snapshot expired.
+ if (!fingerprint.compactAfterFileNames.isEmpty()
+ && snapshotContainsAnyFile(
+ snapshot, fingerprint,
fingerprint.compactAfterFileNames)) {
+ return true;
+ }
+ if (!fingerprint.newIndexFileNames.isEmpty()
+ && snapshotContainsAnyIndexFile(snapshot,
fingerprint.newIndexFileNames)) {
+ return true;
+ }
+ if (!fingerprint.compactAfterFileNames.isEmpty()
+ || !fingerprint.newIndexFileNames.isEmpty()) {
+ return false;
+ }
+ // The commit produces no new files at all (delete-only compact
without deletion-vector
+ // rewrite), so there is nothing to abort and this result is
inconsequential. Fall back
+ // to the weak heuristic: input files already removed by a COMPACT
snapshot.
+ if (snapshot.commitIdentifier() != BatchWriteBuilder.COMMIT_IDENTIFIER
+ || snapshot.commitKind() != Snapshot.CommitKind.COMPACT) {
+ return false;
+ }
+ if (!fingerprint.compactBeforeFileNames.isEmpty()) {
+ return !snapshotContainsAnyFile(
+ snapshot, fingerprint, fingerprint.compactBeforeFileNames);
+ }
+ return true;
+ }
+
+ private boolean snapshotContainsAnyIndexFile(Snapshot snapshot,
Set<String> indexFileNames) {
+ // A snapshot's index manifest only contains live entries, so finding
a new index file
+ // name proves the commit that added it is visible in this snapshot.
+ for (IndexManifestEntry entry :
+ table.store().newIndexFileHandler().scan(snapshot,
DELETION_VECTORS_INDEX)) {
+ if (indexFileNames.contains(entry.indexFile().fileName())) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private boolean snapshotContainsAnyFile(
+ Snapshot snapshot, CompactCommitFingerprint fingerprint,
Set<String> fileNames) {
+ if (fileNames.isEmpty()) {
+ return false;
+ }
+ for (ManifestEntry entry :
+ createScopedScan(snapshot, fingerprint,
fileNames).plan().files()) {
+ if (fileNames.contains(entry.file().fileName())) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private FileStoreScan createScopedScan(
+ Snapshot snapshot, CompactCommitFingerprint fingerprint,
Set<String> fileNames) {
+ FileStoreScan scan =
table.store().newScan().withSnapshot(snapshot).dropStats();
+ if (!fingerprint.partitions.isEmpty()) {
+ scan = scan.withPartitionFilter(fingerprint.partitions);
+ }
+ if (!fingerprint.buckets.isEmpty()) {
+ scan = scan.withBucketFilter(fingerprint.buckets::contains);
+ }
+ if (!fileNames.isEmpty()) {
+ scan = scan.withDataFileNameFilter(fileNames::contains);
+ }
+ return scan;
+ }
+
+ private static final class CompactCommitFingerprint {
+ private final Set<String> compactAfterFileNames;
+ private final Set<String> compactBeforeFileNames;
+ private final Set<String> newIndexFileNames;
+ private final List<BinaryRow> partitions;
+ private final Set<Integer> buckets;
+
+ private CompactCommitFingerprint(
+ Set<String> compactAfterFileNames,
+ Set<String> compactBeforeFileNames,
+ Set<String> newIndexFileNames,
+ List<BinaryRow> partitions,
+ Set<Integer> buckets) {
+ this.compactAfterFileNames = compactAfterFileNames;
+ this.compactBeforeFileNames = compactBeforeFileNames;
+ this.newIndexFileNames = newIndexFileNames;
+ this.partitions = partitions;
+ this.buckets = buckets;
+ }
+
+ private static CompactCommitFingerprint from(List<CommitMessage>
compactMessages) {
+ Set<String> compactAfterFileNames = new HashSet<>();
+ Set<String> compactBeforeFileNames = new HashSet<>();
+ Set<String> newIndexFileNames = new HashSet<>();
+ Set<BinaryRow> partitionSet = new HashSet<>();
+ Set<Integer> buckets = new HashSet<>();
+ for (CommitMessage message : compactMessages) {
+ CommitMessageImpl impl = (CommitMessageImpl) message;
+ partitionSet.add(impl.partition());
+ buckets.add(impl.bucket());
+ for (DataFileMeta file :
impl.compactIncrement().compactAfter()) {
+ compactAfterFileNames.add(file.fileName());
+ }
+ for (DataFileMeta file :
impl.compactIncrement().compactBefore()) {
+ compactBeforeFileNames.add(file.fileName());
+ }
+ for (IndexFileMeta indexFile :
impl.compactIncrement().newIndexFiles()) {
+ newIndexFileNames.add(indexFile.fileName());
+ }
+ }
+ return new CompactCommitFingerprint(
+ compactAfterFileNames,
+ compactBeforeFileNames,
+ newIndexFileNames,
+ new ArrayList<>(partitionSet),
+ buckets);
+ }
+ }
+
+ /** Whether all planned compact-before files are already absent from the
latest snapshot. */
+ public boolean isPlannedInputAlreadyCommitted() {
+ if (!hasInput()) {
+ return true;
+ }
+ Long latestId = table.snapshotManager().latestSnapshotId();
+ if (latestId == null) {
+ return false;
+ }
+ Snapshot snapshot = table.snapshotManager().snapshot(latestId);
+ Set<String> beforeFileNames = new HashSet<>();
+ Set<BinaryRow> partitionSet = new HashSet<>();
+ Set<Integer> buckets = new HashSet<>();
+ for (Map.Entry<BinaryRow, Map<Integer, List<DataFileMeta>>>
partitionEntry :
+ compactBeforeFiles.entrySet()) {
+ partitionSet.add(partitionEntry.getKey());
+ for (Map.Entry<Integer, List<DataFileMeta>> bucketEntry :
+ partitionEntry.getValue().entrySet()) {
+ buckets.add(bucketEntry.getKey());
+ for (DataFileMeta file : bucketEntry.getValue()) {
+ beforeFileNames.add(file.fileName());
+ }
+ }
+ }
+ CompactCommitFingerprint fingerprint =
+ new CompactCommitFingerprint(
+ Collections.emptySet(),
+ beforeFileNames,
+ Collections.emptySet(),
+ new ArrayList<>(partitionSet),
+ buckets);
+ return !snapshotContainsAnyFile(snapshot, fingerprint,
beforeFileNames);
+ }
+
+ /** Whether this rewriter has captured any old files to compact. */
+ public boolean hasInput() {
+ return !compactBeforeFiles.isEmpty();
+ }
+}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/append/SortCompactPlanMetadata.java
b/paimon-core/src/main/java/org/apache/paimon/append/SortCompactPlanMetadata.java
new file mode 100644
index 0000000000..ebd7802fc3
--- /dev/null
+++
b/paimon-core/src/main/java/org/apache/paimon/append/SortCompactPlanMetadata.java
@@ -0,0 +1,165 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.append;
+
+import org.apache.paimon.Snapshot;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.index.IndexFileHandler;
+import org.apache.paimon.manifest.IndexManifestEntry;
+import org.apache.paimon.manifest.IndexManifestEntrySerializer;
+import org.apache.paimon.table.BucketMode;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.source.DataSplit;
+
+import javax.annotation.Nullable;
+
+import java.io.IOException;
+import java.io.Serializable;
+import java.io.UncheckedIOException;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+import static
org.apache.paimon.deletionvectors.DeletionVectorsIndexFile.DELETION_VECTORS_INDEX;
+
+/**
+ * Base-snapshot index metadata captured at sort compact planning time.
+ *
+ * <p>Flink carries this object in the job graph so commit recovery can still
clean up deletion
+ * vectors even if the base snapshot has expired before the committer runs
again. Index metadata is
+ * stored as Paimon byte arrays instead of non-{@link Serializable} POJOs.
+ */
+public final class SortCompactPlanMetadata implements Serializable {
+
+ private static final long serialVersionUID = 1L;
+
+ @Nullable private final byte[] serializedDeletionVectorEntries;
+
+ /**
+ * Whether the base snapshot was readable when this metadata was captured.
Distinguishes a
+ * successful capture of an empty deletion-vector state from a failed
capture.
+ */
+ private final boolean baseSnapshotCaptured;
+
+ private SortCompactPlanMetadata(
+ @Nullable byte[] serializedDeletionVectorEntries, boolean
baseSnapshotCaptured) {
+ this.serializedDeletionVectorEntries = serializedDeletionVectorEntries;
+ this.baseSnapshotCaptured = baseSnapshotCaptured;
+ }
+
+ /** Capture index metadata from the base snapshot for the planned compact
input. */
+ public static SortCompactPlanMetadata capture(
+ FileStoreTable table, long baseSnapshotId, List<DataSplit>
compactInputSplits) {
+ Set<BinaryRow> partitions = new HashSet<>();
+ for (DataSplit split : compactInputSplits) {
+ partitions.add(split.partition());
+ }
+
+ Map<BinaryRow, List<IndexManifestEntry>> baseDeletionVectorEntries =
new HashMap<>();
+ boolean captured =
+ captureInto(table, baseSnapshotId, partitions,
baseDeletionVectorEntries);
+ return fromCapturedMap(baseDeletionVectorEntries, captured);
+ }
+
+ static boolean captureInto(
+ FileStoreTable table,
+ long baseSnapshotId,
+ Set<BinaryRow> partitions,
+ Map<BinaryRow, List<IndexManifestEntry>>
baseDeletionVectorEntries) {
+ Snapshot snapshot = resolveBaseSnapshot(table, baseSnapshotId);
+ if (snapshot == null) {
+ return false;
+ }
+
+ IndexFileHandler indexFileHandler =
table.store().newIndexFileHandler();
+ if (table.coreOptions().deletionVectorsEnabled()
+ && table.bucketMode() == BucketMode.BUCKET_UNAWARE) {
+ for (IndexManifestEntry entry :
+ indexFileHandler.scan(snapshot, DELETION_VECTORS_INDEX)) {
+ if (partitions.contains(entry.partition())) {
+ baseDeletionVectorEntries
+ .computeIfAbsent(entry.partition(), k -> new
ArrayList<>())
+ .add(entry);
+ }
+ }
+ }
+ return true;
+ }
+
+ void copyInto(Map<BinaryRow, List<IndexManifestEntry>>
baseDeletionVectorEntries) {
+ if (serializedDeletionVectorEntries != null) {
+ IndexManifestEntrySerializer entrySerializer = new
IndexManifestEntrySerializer();
+ try {
+ for (IndexManifestEntry entry :
+
entrySerializer.deserializeList(serializedDeletionVectorEntries)) {
+ baseDeletionVectorEntries
+ .computeIfAbsent(entry.partition(), k -> new
ArrayList<>())
+ .add(entry);
+ }
+ } catch (IOException e) {
+ throw new UncheckedIOException(
+ "Failed to deserialize captured deletion vector
metadata.", e);
+ }
+ }
+ }
+
+ /**
+ * Whether the base snapshot was readable at capture time. An empty
deletion-vector payload with
+ * {@code true} means known-empty; with {@code false} means capture failed.
+ */
+ boolean baseSnapshotCaptured() {
+ return baseSnapshotCaptured;
+ }
+
+ private static SortCompactPlanMetadata fromCapturedMap(
+ Map<BinaryRow, List<IndexManifestEntry>> baseDeletionVectorEntries,
+ boolean baseSnapshotCaptured) {
+ byte[] serializedDeletionVectorEntries = null;
+ if (!baseDeletionVectorEntries.isEmpty()) {
+ List<IndexManifestEntry> entries = new ArrayList<>();
+ for (List<IndexManifestEntry> partitionEntries :
baseDeletionVectorEntries.values()) {
+ entries.addAll(partitionEntries);
+ }
+ IndexManifestEntrySerializer entrySerializer = new
IndexManifestEntrySerializer();
+ try {
+ serializedDeletionVectorEntries =
entrySerializer.serializeList(entries);
+ } catch (IOException e) {
+ throw new UncheckedIOException(
+ "Failed to serialize captured deletion vector
metadata.", e);
+ }
+ }
+
+ return new SortCompactPlanMetadata(serializedDeletionVectorEntries,
baseSnapshotCaptured);
+ }
+
+ @Nullable
+ private static Snapshot resolveBaseSnapshot(FileStoreTable table, long
snapshotId) {
+ if (snapshotId <= 0) {
+ return table.snapshotManager().latestSnapshot();
+ }
+ try {
+ return table.snapshotManager().tryGetSnapshot(snapshotId);
+ } catch (java.io.FileNotFoundException e) {
+ return null;
+ }
+ }
+}
diff --git a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java
b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java
index 955750613a..48aaf07809 100644
--- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java
+++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java
@@ -387,6 +387,8 @@ public interface DataFileMeta {
DataFileMeta assignSequenceNumber(long minSequenceNumber, long
maxSequenceNumber);
+ DataFileMeta assignFileSource(FileSource fileSource);
+
DataFileMeta withWriteColsSequences(long[] writeColsSequences);
DataFileMeta assignFirstRowId(long firstRowId);
diff --git
a/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java
b/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java
index af65b0800e..4632b8332d 100644
--- a/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java
+++ b/paimon-core/src/main/java/org/apache/paimon/io/PojoDataFileMeta.java
@@ -381,6 +381,32 @@ public class PojoDataFileMeta implements DataFileMeta {
writeColsSequences);
}
+ @Override
+ public PojoDataFileMeta assignFileSource(FileSource fileSource) {
+ return new PojoDataFileMeta(
+ fileName,
+ fileSize,
+ rowCount,
+ minKey,
+ maxKey,
+ keyStats,
+ valueStats,
+ minSequenceNumber,
+ maxSequenceNumber,
+ schemaId,
+ level,
+ extraFiles,
+ creationTime,
+ deleteRowCount,
+ embeddedIndex,
+ fileSource,
+ valueStatsCols,
+ externalPath,
+ firstRowId,
+ writeCols,
+ writeColsSequences);
+ }
+
@Override
public PojoDataFileMeta assignFirstRowId(long firstRowId) {
return new PojoDataFileMeta(
diff --git
a/paimon-core/src/main/java/org/apache/paimon/io/ProjectedDataFileMeta.java
b/paimon-core/src/main/java/org/apache/paimon/io/ProjectedDataFileMeta.java
index 1eaaf9af39..00451ed1be 100644
--- a/paimon-core/src/main/java/org/apache/paimon/io/ProjectedDataFileMeta.java
+++ b/paimon-core/src/main/java/org/apache/paimon/io/ProjectedDataFileMeta.java
@@ -307,6 +307,11 @@ public final class ProjectedDataFileMeta implements
DataFileMeta {
throw unsupportedOperation("assignSequenceNumber(long, long)");
}
+ @Override
+ public DataFileMeta assignFileSource(FileSource fileSource) {
+ throw unsupportedOperation("assignFileSource(FileSource)");
+ }
+
@Override
public DataFileMeta withWriteColsSequences(long[] writeColsSequences) {
throw unsupportedOperation("withWriteColsSequences(long[])");
diff --git
a/paimon-core/src/test/java/org/apache/paimon/append/SortCompactCommitMessageRewriterTest.java
b/paimon-core/src/test/java/org/apache/paimon/append/SortCompactCommitMessageRewriterTest.java
new file mode 100644
index 0000000000..974b8e3b26
--- /dev/null
+++
b/paimon-core/src/test/java/org/apache/paimon/append/SortCompactCommitMessageRewriterTest.java
@@ -0,0 +1,1370 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.append;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.Snapshot;
+import org.apache.paimon.TestAppendFileStore;
+import org.apache.paimon.TestKeyValueGenerator;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.deletionvectors.BitmapDeletionVector;
+import org.apache.paimon.deletionvectors.DeletionVector;
+import org.apache.paimon.deletionvectors.append.BaseAppendDeleteFileMaintainer;
+import org.apache.paimon.fs.FileIO;
+import org.apache.paimon.fs.FileIOFinder;
+import org.apache.paimon.fs.Path;
+import org.apache.paimon.fs.local.LocalFileIO;
+import org.apache.paimon.index.IndexFileMeta;
+import org.apache.paimon.index.IndexPathFactory;
+import org.apache.paimon.io.CompactIncrement;
+import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.io.DataIncrement;
+import org.apache.paimon.manifest.FileKind;
+import org.apache.paimon.manifest.FileSource;
+import org.apache.paimon.manifest.IndexManifestEntry;
+import org.apache.paimon.schema.FileSystemSchemaManager;
+import org.apache.paimon.schema.Schema;
+import org.apache.paimon.schema.SchemaManager;
+import org.apache.paimon.schema.SchemaUtils;
+import org.apache.paimon.schema.TableSchema;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.FileStoreTableFactory;
+import org.apache.paimon.table.sink.BatchTableCommit;
+import org.apache.paimon.table.sink.BatchWriteBuilder;
+import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.sink.CommitMessageImpl;
+import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.TraceableFileIO;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
+import java.io.ObjectInputStream;
+import java.io.ObjectOutputStream;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static org.apache.paimon.io.DataFileTestUtils.newFile;
+import static org.apache.paimon.table.BucketMode.UNAWARE_BUCKET;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+/** Test for {@link SortCompactCommitMessageRewriter}. */
+public class SortCompactCommitMessageRewriterTest {
+
+ @TempDir java.nio.file.Path tempDir;
+
+ private static DataFileMeta asCompactAfter(DataFileMeta file) {
+ return file.assignFileSource(FileSource.COMPACT);
+ }
+
+ @Test
+ public void testRewriteToCompactMessages() throws Exception {
+ FileStoreTable table = createAppendTable(Collections.emptyMap());
+
+ DataFileMeta old0 = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta old1 = newFile("data-1.orc", 0, 101, 200, 200);
+ DataFileMeta sorted = newFile("sorted-0.orc", 0, 0, 200, 200);
+
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Arrays.asList(old0, old1))
+ .build();
+
+ CommitMessageImpl written =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 0,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(sorted),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+
+ List<CommitMessage> result =
+ new SortCompactCommitMessageRewriter(table, 0L,
Collections.singletonList(split))
+ .rewrite(Collections.singletonList(written));
+
+ assertThat(result).hasSize(1);
+ CommitMessageImpl compact = (CommitMessageImpl) result.get(0);
+ assertThat(compact.newFilesIncrement().isEmpty()).isTrue();
+ CompactIncrement ci = compact.compactIncrement();
+ assertThat(ci.compactBefore()).containsExactly(old0, old1);
+ assertThat(ci.compactAfter()).containsExactly(asCompactAfter(sorted));
+ assertThat(ci.changelogFiles()).isEmpty();
+ }
+
+ @Test
+ public void testDetectBatchCompactCommitSucceeded() throws Exception {
+ TestAppendFileStore store = createAppendStore(tempDir,
Collections.emptyMap());
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("data-0.orc")));
+
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ CommitMessageImpl written =
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("sorted-0.orc"));
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split));
+ long snapshotIdBeforeCommit = rewriter.latestSnapshotIdOrZero();
+ List<CommitMessage> compactMessages =
rewriter.rewrite(Collections.singletonList(written));
+
assertThat(rewriter.isBatchCompactCommitSucceeded(snapshotIdBeforeCommit,
compactMessages))
+ .isFalse();
+
+ try (BatchTableCommit commit =
table.newBatchWriteBuilder().newCommit()) {
+ commit.commit(compactMessages);
+ }
+
+
assertThat(rewriter.isBatchCompactCommitSucceeded(snapshotIdBeforeCommit,
compactMessages))
+ .isTrue();
+ assertThat(table.snapshotManager().latestSnapshot().commitKind())
+ .isEqualTo(Snapshot.CommitKind.COMPACT);
+ assertThat(table.snapshotManager().latestSnapshot().commitIdentifier())
+ .isEqualTo(BatchWriteBuilder.COMMIT_IDENTIFIER);
+ }
+
+ @Test
+ public void testDetectBatchCompactCommitSucceededWhenNewerAppendExists()
throws Exception {
+ TestAppendFileStore store = createAppendStore(tempDir,
Collections.emptyMap());
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("data-0.orc")));
+
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ CommitMessageImpl written =
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("sorted-0.orc"));
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split));
+ long snapshotIdBeforeCommit = rewriter.latestSnapshotIdOrZero();
+
+ List<CommitMessage> compactMessages =
rewriter.rewrite(Collections.singletonList(written));
+ try (BatchTableCommit commit =
table.newBatchWriteBuilder().newCommit()) {
+ commit.commit(compactMessages);
+ }
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("concurrent.orc")));
+
+ assertThat(table.snapshotManager().latestSnapshot().commitKind())
+ .isEqualTo(Snapshot.CommitKind.APPEND);
+
assertThat(rewriter.isBatchCompactCommitSucceeded(snapshotIdBeforeCommit,
compactMessages))
+ .isTrue();
+ }
+
+ @Test
+ public void
testDetectBatchCompactCommitSucceededAcrossExpiredSnapshotGap() throws
Exception {
+ TestAppendFileStore store = createAppendStore(tempDir,
Collections.emptyMap());
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("data-0.orc")));
+
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ CommitMessageImpl written =
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("sorted-0.orc"));
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split));
+ long snapshotIdBeforeCommit = rewriter.latestSnapshotIdOrZero();
+ List<CommitMessage> compactMessages =
rewriter.rewrite(Collections.singletonList(written));
+
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("filler.orc")));
+ long fillerSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ try (BatchTableCommit commit =
table.newBatchWriteBuilder().newCommit()) {
+ commit.commit(compactMessages);
+ }
+ table.snapshotManager().deleteSnapshot(fillerSnapshotId);
+
+
assertThat(rewriter.isBatchCompactCommitSucceeded(snapshotIdBeforeCommit,
compactMessages))
+ .isTrue();
+ }
+
+ @Test
+ public void
testDetectBatchCompactCommitSucceededWhenCompactSnapshotExpired() throws
Exception {
+ TestAppendFileStore store = createAppendStore(tempDir,
Collections.emptyMap());
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("data-0.orc")));
+
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ CommitMessageImpl written =
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("sorted-0.orc"));
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split));
+ long snapshotIdBeforeCommit = rewriter.latestSnapshotIdOrZero();
+ List<CommitMessage> compactMessages =
rewriter.rewrite(Collections.singletonList(written));
+
+ try (BatchTableCommit commit =
table.newBatchWriteBuilder().newCommit()) {
+ commit.commit(compactMessages);
+ }
+ long compactSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("concurrent.orc")));
+ table.snapshotManager().deleteSnapshot(compactSnapshotId);
+
+ assertThat(table.snapshotManager().latestSnapshot().commitKind())
+ .isEqualTo(Snapshot.CommitKind.APPEND);
+
assertThat(rewriter.isBatchCompactCommitSucceeded(snapshotIdBeforeCommit,
compactMessages))
+ .isTrue();
+ }
+
+ @Test
+ public void
testDetectBatchCompactCommitSucceededWhenPartialOutputReplaced() throws
Exception {
+ TestAppendFileStore store = createAppendStore(tempDir,
Collections.emptyMap());
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("data-0.orc")));
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 1,
Collections.singletonList("data-1.orc")));
+
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+ DataFileMeta oldBucket0 = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta oldBucket1 = newFile("data-1.orc", 0, 0, 100, 100);
+ CommitMessageImpl written0 =
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("sorted-0.orc"));
+ CommitMessageImpl written1 =
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 1,
Collections.singletonList("sorted-1.orc"));
+ DataSplit split0 =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(oldBucket0))
+ .build();
+ DataSplit split1 =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(1)
+ .withBucketPath("bucket-1")
+ .withDataFiles(Collections.singletonList(oldBucket1))
+ .build();
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId, Arrays.asList(split0, split1));
+ long snapshotIdBeforeCommit = rewriter.latestSnapshotIdOrZero();
+ List<CommitMessage> compactMessages =
rewriter.rewrite(Arrays.asList(written0, written1));
+
+ try (BatchTableCommit commit =
table.newBatchWriteBuilder().newCommit()) {
+ commit.commit(compactMessages);
+ }
+ long compactSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ DataFileMeta mergedBucket0 = newFile("merged-0.orc", 0, 0, 100, 100);
+ CommitMessageImpl partialCompact =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 0,
+ table.coreOptions().bucket(),
+ DataIncrement.emptyIncrement(),
+ new CompactIncrement(
+ Collections.singletonList(
+ asCompactAfter(newFile("sorted-0.orc",
0, 0, 100, 100))),
+
Collections.singletonList(asCompactAfter(mergedBucket0)),
+ Collections.emptyList(),
+ Collections.emptyList(),
+ Collections.emptyList()));
+ try (BatchTableCommit commit =
table.newBatchWriteBuilder().newCommit()) {
+ commit.commit(Collections.singletonList(partialCompact));
+ }
+ table.snapshotManager().deleteSnapshot(compactSnapshotId);
+
+
assertThat(rewriter.isBatchCompactCommitSucceeded(snapshotIdBeforeCommit,
compactMessages))
+ .isTrue();
+ }
+
+ @Test
+ public void testDetectBatchCompactCommitSucceededIgnoresUnrelatedCompact()
throws Exception {
+ TestAppendFileStore store = createAppendStore(tempDir,
Collections.emptyMap());
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0, Arrays.asList("data-0.orc",
"data-1.orc")));
+
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ CommitMessageImpl written =
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("sorted-0.orc"));
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split));
+ long snapshotIdBeforeCommit = rewriter.latestSnapshotIdOrZero();
+ List<CommitMessage> compactMessages =
rewriter.rewrite(Collections.singletonList(written));
+
+ DataFileMeta otherOld = newFile("data-1.orc", 0, 101, 200, 200);
+ CommitMessageImpl otherWritten =
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("sorted-1.orc"));
+ DataSplit otherSplit =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(otherOld))
+ .build();
+ List<CommitMessage> otherCompactMessages =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(otherSplit))
+ .rewrite(Collections.singletonList(otherWritten));
+ try (BatchTableCommit commit =
table.newBatchWriteBuilder().newCommit()) {
+ commit.commit(otherCompactMessages);
+ }
+
+
assertThat(rewriter.isBatchCompactCommitSucceeded(snapshotIdBeforeCommit,
compactMessages))
+ .isFalse();
+ }
+
+ @Test
+ public void
testDetectDeleteOnlyCompactCommitSucceededWithConcurrentCompact() throws
Exception {
+ TestAppendFileStore store =
+ createAppendStore(
+ tempDir,
+ Collections.singletonMap(
+ CoreOptions.DELETION_VECTORS_ENABLED.key(),
"true"));
+
+ // data-0 and data-1 share a single DV index file.
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0, Arrays.asList("data-0.orc",
"data-1.orc")));
+ Map<String, List<Integer>> dvs = new HashMap<>();
+ dvs.put("data-0.orc", Arrays.asList(1, 3, 5));
+ dvs.put("data-1.orc", Arrays.asList(2, 4, 6));
+ store.commit(store.writeDVIndexFiles(BinaryRow.EMPTY_ROW, 0, dvs));
+
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ DataFileMeta old0 = newFile("data-0.orc", 0, 0, 100, 100);
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old0))
+ .build();
+
+ // Our delete-only sort compact (all rows filtered out): rewrite
persists a new DV index
+ // file holding data-1's deletion vector.
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split));
+ long snapshotIdBeforeCommit = rewriter.latestSnapshotIdOrZero();
+ List<CommitMessage> compactMessages =
rewriter.rewrite(Collections.emptyList());
+ CommitMessageImpl compact = (CommitMessageImpl) compactMessages.get(0);
+ assertThat(compact.compactIncrement().compactAfter()).isEmpty();
+ assertThat(compact.compactIncrement().newIndexFiles()).isNotEmpty();
+
+ // Our commit fails; meanwhile another compaction commits and deletes
the same input file.
+ List<CommitMessage> otherMessages =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split))
+ .rewrite(Collections.emptyList());
+ try (BatchTableCommit commit =
table.newBatchWriteBuilder().newCommit()) {
+ commit.commit(otherMessages);
+ }
+
+ // Our commit never landed: the check must return false so the caller
aborts our new DV
+ // index file. The concurrent compaction having deleted the same input
files must not be
+ // mistaken for our commit.
+
assertThat(rewriter.isBatchCompactCommitSucceeded(snapshotIdBeforeCommit,
compactMessages))
+ .isFalse();
+ }
+
+ @Test
+ public void testRewriteMultipleBuckets() throws Exception {
+ FileStoreTable table = createAppendTable(Collections.emptyMap());
+
+ DataFileMeta oldBucket0 = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta oldBucket1 = newFile("data-1.orc", 0, 0, 100, 100);
+ DataFileMeta sortedBucket0 = newFile("sorted-0.orc", 0, 0, 100, 100);
+ DataFileMeta sortedBucket1 = newFile("sorted-1.orc", 0, 0, 100, 100);
+
+ DataSplit split0 =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(oldBucket0))
+ .build();
+ DataSplit split1 =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(1)
+ .withBucketPath("bucket-1")
+ .withDataFiles(Collections.singletonList(oldBucket1))
+ .build();
+
+ CommitMessageImpl written0 =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 0,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(sortedBucket0),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+ CommitMessageImpl written1 =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 1,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(sortedBucket1),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+
+ List<CommitMessage> result =
+ new SortCompactCommitMessageRewriter(table, 0L,
Arrays.asList(split0, split1))
+ .rewrite(Arrays.asList(written0, written1));
+
+ assertThat(result).hasSize(2);
+ for (CommitMessage message : result) {
+ CommitMessageImpl compact = (CommitMessageImpl) message;
+ assertThat(compact.newFilesIncrement().isEmpty()).isTrue();
+ assertThat(compact.compactIncrement().compactBefore()).hasSize(1);
+ assertThat(compact.compactIncrement().compactAfter()).hasSize(1);
+ }
+ CommitMessageImpl compact0 =
+ (CommitMessageImpl) result.stream().filter(m -> m.bucket() ==
0).findFirst().get();
+
assertThat(compact0.compactIncrement().compactBefore()).containsExactly(oldBucket0);
+ assertThat(compact0.compactIncrement().compactAfter())
+ .containsExactly(asCompactAfter(sortedBucket0));
+ CommitMessageImpl compact1 =
+ (CommitMessageImpl) result.stream().filter(m -> m.bucket() ==
1).findFirst().get();
+
assertThat(compact1.compactIncrement().compactBefore()).containsExactly(oldBucket1);
+ assertThat(compact1.compactIncrement().compactAfter())
+ .containsExactly(asCompactAfter(sortedBucket1));
+ }
+
+ @Test
+ public void testRewriteKeepsPlannedDeletesIndependentFromOutputBuckets()
throws Exception {
+ FileStoreTable table = createAppendTable(Collections.emptyMap());
+
+ DataFileMeta oldBucket0 = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta sortedBucket1 = newFile("sorted-1.orc", 0, 0, 100, 100);
+
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withTotalBuckets(2)
+ .withDataFiles(Collections.singletonList(oldBucket0))
+ .build();
+
+ CommitMessageImpl written =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 1,
+ 2,
+ new DataIncrement(
+ Collections.singletonList(sortedBucket1),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+
+ List<CommitMessage> result =
+ new SortCompactCommitMessageRewriter(table, 0L,
Collections.singletonList(split))
+ .rewrite(Collections.singletonList(written));
+
+ assertThat(result).hasSize(2);
+ CommitMessageImpl compactDelete =
+ (CommitMessageImpl) result.stream().filter(m -> m.bucket() ==
0).findFirst().get();
+ assertThat(compactDelete.totalBuckets()).isEqualTo(2);
+
assertThat(compactDelete.compactIncrement().compactBefore()).containsExactly(oldBucket0);
+ assertThat(compactDelete.compactIncrement().compactAfter()).isEmpty();
+
+ CommitMessageImpl compactAdd =
+ (CommitMessageImpl) result.stream().filter(m -> m.bucket() ==
1).findFirst().get();
+ assertThat(compactAdd.totalBuckets()).isEqualTo(2);
+ assertThat(compactAdd.compactIncrement().compactBefore()).isEmpty();
+ assertThat(compactAdd.compactIncrement().compactAfter())
+ .containsExactly(asCompactAfter(sortedBucket1));
+ }
+
+ @Test
+ public void testRewriteWithDeletionVectors() throws Exception {
+ TestAppendFileStore store =
+ createAppendStore(
+ tempDir,
+ Collections.singletonMap(
+ CoreOptions.DELETION_VECTORS_ENABLED.key(),
"true"));
+
+ // write deletion vectors for two old files
+ Map<String, List<Integer>> dvs = new HashMap<>();
+ dvs.put("data-0.orc", Arrays.asList(1, 3, 5));
+ dvs.put("data-1.orc", Arrays.asList(2, 4, 6));
+ CommitMessageImpl dvMessage =
store.writeDVIndexFiles(BinaryRow.EMPTY_ROW, 0, dvs);
+ store.commit(dvMessage);
+
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ DataFileMeta old0 = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta old1 = newFile("data-1.orc", 0, 101, 200, 200);
+ DataFileMeta sorted = newFile("sorted-0.orc", 0, 0, 200, 200);
+
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Arrays.asList(old0, old1))
+ .build();
+
+ CommitMessageImpl written =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 0,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(sorted),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+
+ List<CommitMessage> result =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split))
+ .rewrite(Collections.singletonList(written));
+
+ CommitMessageImpl compact = (CommitMessageImpl) result.get(0);
+ // all old files are removed, so their DV index entries must be
cleaned up
+
assertThat(compact.compactIncrement().deletedIndexFiles()).isNotEmpty();
+ assertThat(compact.compactIncrement().newIndexFiles()).isEmpty();
+
assertThat(compact.compactIncrement().compactBefore()).containsExactly(old0,
old1);
+ assertThat(compact.compactIncrement().compactAfter())
+ .containsExactly(asCompactAfter(sorted));
+ assertThat(compact.newFilesIncrement().isEmpty()).isTrue();
+ }
+
+ @Test
+ public void testRewriteFailsWhenConcurrentDeletionVectorAdded() throws
Exception {
+ TestAppendFileStore store =
+ createAppendStore(
+ tempDir,
+ Collections.singletonMap(
+ CoreOptions.DELETION_VECTORS_ENABLED.key(),
"true"));
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("data-0.orc")));
+
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta sorted = newFile("sorted-0.orc", 0, 0, 100, 100);
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+ CommitMessageImpl written =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 0,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(sorted),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split));
+
+ Map<String, List<Integer>> concurrentDvs = new HashMap<>();
+ concurrentDvs.put("data-0.orc", Arrays.asList(1, 3, 5));
+ store.commit(store.writeDVIndexFiles(BinaryRow.EMPTY_ROW, 0,
concurrentDvs));
+
+ assertThatThrownBy(() ->
rewriter.rewrite(Collections.singletonList(written)))
+ .isInstanceOf(IllegalStateException.class)
+ .hasMessageContaining("deletion vectors on input files
changed")
+ .hasMessageContaining("restore deleted rows")
+ .hasMessageContaining("Please retry");
+ }
+
+ @Test
+ public void testRewriteFailsWhenConcurrentDeletionVectorReplaced() throws
Exception {
+ TestAppendFileStore store =
+ createAppendStore(
+ tempDir,
+ Collections.singletonMap(
+ CoreOptions.DELETION_VECTORS_ENABLED.key(),
"true"));
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("data-0.orc")));
+
+ Map<String, List<Integer>> baseDvs = new HashMap<>();
+ baseDvs.put("data-0.orc", Arrays.asList(1, 3, 5));
+ store.commit(store.writeDVIndexFiles(BinaryRow.EMPTY_ROW, 0, baseDvs));
+
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta sorted = newFile("sorted-0.orc", 0, 0, 100, 100);
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+ CommitMessageImpl written =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 0,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(sorted),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split));
+
+ Map<String, List<Integer>> concurrentDvs = new HashMap<>();
+ concurrentDvs.put("data-0.orc", Arrays.asList(2, 4, 6));
+ commitUnawareDeletionVectors(table, concurrentDvs);
+
+ assertThatThrownBy(() ->
rewriter.rewrite(Collections.singletonList(written)))
+ .isInstanceOf(IllegalStateException.class)
+ .hasMessageContaining("deletion vectors on input files
changed")
+ .hasMessageContaining("restore deleted rows")
+ .hasMessageContaining("Please retry")
+ .hasMessageContaining("baseDv=")
+ .hasMessageContaining("latestDv=");
+ }
+
+ @Test
+ public void
testRewriteFailsWhenConcurrentDeletionVectorAddedAfterBaseExpired()
+ throws Exception {
+ TestAppendFileStore store =
+ createAppendStore(
+ tempDir,
+ Collections.singletonMap(
+ CoreOptions.DELETION_VECTORS_ENABLED.key(),
"true"));
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("data-0.orc")));
+
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta sorted = newFile("sorted-0.orc", 0, 0, 100, 100);
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+ CommitMessageImpl written =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 0,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(sorted),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+
+ // Capture known-empty base DV state before the base snapshot expires.
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split));
+
+ Map<String, List<Integer>> concurrentDvs = new HashMap<>();
+ concurrentDvs.put("data-0.orc", Arrays.asList(1, 3, 5));
+ store.commit(store.writeDVIndexFiles(BinaryRow.EMPTY_ROW, 0,
concurrentDvs));
+ table.snapshotManager().deleteSnapshot(baseSnapshotId);
+
+ assertThatThrownBy(() ->
rewriter.rewrite(Collections.singletonList(written)))
+ .isInstanceOf(IllegalStateException.class)
+ .hasMessageContaining("deletion vectors on input files
changed")
+ .hasMessageContaining("restore deleted rows")
+ .hasMessageContaining("Please retry");
+ }
+
+ @Test
+ public void
testRewriteFailsWhenConcurrentDeletionVectorAddedWithCapturedEmptyMetadata()
+ throws Exception {
+ TestAppendFileStore store =
+ createAppendStore(
+ tempDir,
+ Collections.singletonMap(
+ CoreOptions.DELETION_VECTORS_ENABLED.key(),
"true"));
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("data-0.orc")));
+
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta sorted = newFile("sorted-0.orc", 0, 0, 100, 100);
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+ CommitMessageImpl written =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 0,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(sorted),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+
+ SortCompactPlanMetadata planMetadata =
+ SortCompactPlanMetadata.capture(
+ table, baseSnapshotId,
Collections.singletonList(split));
+ assertThat(planMetadata.baseSnapshotCaptured()).isTrue();
+
+ Map<String, List<Integer>> concurrentDvs = new HashMap<>();
+ concurrentDvs.put("data-0.orc", Arrays.asList(1, 3, 5));
+ store.commit(store.writeDVIndexFiles(BinaryRow.EMPTY_ROW, 0,
concurrentDvs));
+ table.snapshotManager().deleteSnapshot(baseSnapshotId);
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split), planMetadata);
+
+ assertThatThrownBy(() ->
rewriter.rewrite(Collections.singletonList(written)))
+ .isInstanceOf(IllegalStateException.class)
+ .hasMessageContaining("deletion vectors on input files
changed")
+ .hasMessageContaining("restore deleted rows")
+ .hasMessageContaining("Please retry");
+ }
+
+ @Test
+ public void testRewriteInputOnlyGroupWithDeletionVectors() throws
Exception {
+ TestAppendFileStore store =
+ createAppendStore(
+ tempDir,
+ Collections.singletonMap(
+ CoreOptions.DELETION_VECTORS_ENABLED.key(),
"true"));
+
+ Map<String, List<Integer>> dvs = new HashMap<>();
+ dvs.put("data-0.orc", Arrays.asList(1, 3, 5));
+ CommitMessageImpl dvMessage =
store.writeDVIndexFiles(BinaryRow.EMPTY_ROW, 0, dvs);
+ store.commit(dvMessage);
+
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+
+ List<CommitMessage> result =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split))
+ .rewrite(Collections.emptyList());
+
+ assertThat(result).hasSize(1);
+ CommitMessageImpl compact = (CommitMessageImpl) result.get(0);
+ assertThat(compact.newFilesIncrement().isEmpty()).isTrue();
+
assertThat(compact.compactIncrement().compactBefore()).containsExactly(old);
+ assertThat(compact.compactIncrement().compactAfter()).isEmpty();
+
assertThat(compact.compactIncrement().deletedIndexFiles()).isNotEmpty();
+ assertThat(compact.compactIncrement().newIndexFiles()).isEmpty();
+ }
+
+ @Test
+ public void testAbortCompactMessagesCleansUpNewDeletionVectorFiles()
throws Exception {
+ TestAppendFileStore store =
+ createAppendStore(
+ tempDir,
+ Collections.singletonMap(
+ CoreOptions.DELETION_VECTORS_ENABLED.key(),
"true"));
+
+ // Two old data files sharing a single DV index file.
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0, Arrays.asList("data-0.orc",
"data-1.orc")));
+ Map<String, List<Integer>> dvs = new HashMap<>();
+ dvs.put("data-0.orc", Arrays.asList(1, 3, 5));
+ dvs.put("data-1.orc", Arrays.asList(2, 4, 6));
+ store.commit(store.writeDVIndexFiles(BinaryRow.EMPTY_ROW, 0, dvs));
+
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ // Plan sort compact for only data-0; data-1 stays, so its DV must be
rewritten to a new
+ // index file by dvMaintainer.persist() during rewrite.
+ DataFileMeta old0 = newFile("data-0.orc", 0, 0, 100, 100);
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old0))
+ .build();
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split));
+ List<CommitMessage> compactMessages =
rewriter.rewrite(Collections.emptyList());
+
+ assertThat(compactMessages).hasSize(1);
+ CommitMessageImpl compact = (CommitMessageImpl) compactMessages.get(0);
+ assertThat(compact.compactIncrement().newIndexFiles())
+ .as("new DV index file rewriting data-1's deletion vector")
+ .hasSize(1);
+ assertThat(compact.compactIncrement().deletedIndexFiles())
+ .as("old shared DV index file marked for deletion")
+ .hasSize(1);
+
+ IndexFileMeta newDvFile =
compact.compactIncrement().newIndexFiles().get(0);
+ IndexFileMeta oldSharedDvFile =
compact.compactIncrement().deletedIndexFiles().get(0);
+ IndexPathFactory indexPathFactory =
+
table.store().pathFactory().indexFileFactory(BinaryRow.EMPTY_ROW,
UNAWARE_BUCKET);
+ Path newDvPath = indexPathFactory.toPath(newDvFile);
+ Path oldSharedDvPath = indexPathFactory.toPath(oldSharedDvFile);
+ assertThat(table.fileIO().exists(newDvPath)).isTrue();
+ assertThat(table.fileIO().exists(oldSharedDvPath)).isTrue();
+
+ rewriter.abortCompactMessages(compactMessages);
+
+ // The new DV file is only referenced by the uncommitted compact
messages, so abort cleans
+ // it up to avoid orphaned index files on retry.
+ assertThat(table.fileIO().exists(newDvPath)).isFalse();
+ // The old shared DV file is still referenced by the latest snapshot
(the compact did not
+ // commit), so abort must not delete it.
+ assertThat(table.fileIO().exists(oldSharedDvPath)).isTrue();
+ }
+
+ @Test
+ public void testAbortWrittenMessagesCleansUpSortedDataFiles() throws
Exception {
+ TestAppendFileStore store = createAppendStore(tempDir,
Collections.emptyMap());
+ store.commit(
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("data-0.orc")));
+
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+
+ CommitMessageImpl written =
+ store.writeDataFiles(
+ BinaryRow.EMPTY_ROW, 0,
Collections.singletonList("sorted-0.orc"));
+ DataFileMeta sorted = written.newFilesIncrement().newFiles().get(0);
+ Path sortedPath =
+ table.store()
+ .pathFactory()
+ .createDataFilePathFactory(BinaryRow.EMPTY_ROW, 0)
+ .toPath(sorted);
+ assertThat(table.fileIO().exists(sortedPath)).isTrue();
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split));
+ rewriter.abortWrittenMessages(Collections.singletonList(written));
+
+ assertThat(table.fileIO().exists(sortedPath)).isFalse();
+ }
+
+ @Test
+ public void testRewritePartialMessagesMustBeMerged() throws Exception {
+ FileStoreTable table = createAppendTable(Collections.emptyMap());
+
+ DataFileMeta oldBucket0 = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta oldBucket1 = newFile("data-1.orc", 0, 0, 100, 100);
+ DataFileMeta sortedBucket0 = newFile("sorted-0.orc", 0, 0, 100, 100);
+ DataFileMeta sortedBucket1 = newFile("sorted-1.orc", 0, 0, 100, 100);
+
+ DataSplit split0 =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(oldBucket0))
+ .build();
+ DataSplit split1 =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(1)
+ .withBucketPath("bucket-1")
+ .withDataFiles(Collections.singletonList(oldBucket1))
+ .build();
+
+ CommitMessageImpl writtenBucket0 =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 0,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(sortedBucket0),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+ CommitMessageImpl writtenBucket1 =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 1,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(sortedBucket1),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(table, 0L,
Arrays.asList(split0, split1));
+
+ // Rewriting partial outputs separately would duplicate compactBefore
for every commit.
+ List<CommitMessage> partial0 =
rewriter.rewrite(Collections.singletonList(writtenBucket0));
+ List<CommitMessage> partial1 =
rewriter.rewrite(Collections.singletonList(writtenBucket1));
+ assertThat(partial0).hasSize(2);
+ assertThat(partial1).hasSize(2);
+
+ List<CommitMessage> merged =
+ rewriter.rewrite(Arrays.asList(writtenBucket0,
writtenBucket1));
+ assertThat(merged).hasSize(2);
+ for (CommitMessage message : merged) {
+ CommitMessageImpl compact = (CommitMessageImpl) message;
+ assertThat(compact.newFilesIncrement().isEmpty()).isTrue();
+ assertThat(compact.compactIncrement().compactBefore()).hasSize(1);
+ assertThat(compact.compactIncrement().compactAfter()).hasSize(1);
+ }
+ }
+
+ @Test
+ public void testRewriteUsesCapturedBaseSnapshotMetadata() throws Exception
{
+ TestAppendFileStore store =
+ createAppendStore(
+ tempDir,
+ Collections.singletonMap(
+ CoreOptions.DELETION_VECTORS_ENABLED.key(),
"true"));
+
+ Map<String, List<Integer>> dvs = new HashMap<>();
+ dvs.put("data-0.orc", Arrays.asList(1, 3, 5));
+ CommitMessageImpl dvMessage =
store.writeDVIndexFiles(BinaryRow.EMPTY_ROW, 0, dvs);
+ store.commit(dvMessage);
+
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta sorted = newFile("sorted-0.orc", 0, 0, 100, 100);
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+
+ CommitMessageImpl written =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 0,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(sorted),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split));
+ table.snapshotManager().deleteSnapshot(baseSnapshotId);
+
+ List<CommitMessage> result =
rewriter.rewrite(Collections.singletonList(written));
+
+ CommitMessageImpl compact = (CommitMessageImpl) result.get(0);
+
assertThat(compact.compactIncrement().deletedIndexFiles()).isNotEmpty();
+
assertThat(compact.compactIncrement().compactBefore()).containsExactly(old);
+ assertThat(compact.compactIncrement().compactAfter())
+ .containsExactly(asCompactAfter(sorted));
+ }
+
+ @Test
+ public void testPlanMetadataRoundTripSerialization() throws Exception {
+ TestAppendFileStore store =
+ createAppendStore(
+ tempDir,
+ Collections.singletonMap(
+ CoreOptions.DELETION_VECTORS_ENABLED.key(),
"true"));
+
+ Map<String, List<Integer>> dvs = new HashMap<>();
+ dvs.put("data-0.orc", Arrays.asList(1, 3, 5));
+ CommitMessageImpl dvMessage =
store.writeDVIndexFiles(BinaryRow.EMPTY_ROW, 0, dvs);
+ store.commit(dvMessage);
+
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(
+
Collections.singletonList(newFile("data-0.orc", 0, 0, 100, 100)))
+ .build();
+
+ SortCompactPlanMetadata captured =
+ SortCompactPlanMetadata.capture(
+ table, baseSnapshotId,
Collections.singletonList(split));
+ SortCompactPlanMetadata restored;
+ try (ByteArrayOutputStream baos = new ByteArrayOutputStream();
+ ObjectOutputStream oos = new ObjectOutputStream(baos)) {
+ oos.writeObject(captured);
+ try (ObjectInputStream ois =
+ new ObjectInputStream(new
ByteArrayInputStream(baos.toByteArray()))) {
+ restored = (SortCompactPlanMetadata) ois.readObject();
+ }
+ }
+
+ Map<BinaryRow, List<IndexManifestEntry>> dvEntries = new HashMap<>();
+ captured.copyInto(dvEntries);
+ Map<BinaryRow, List<IndexManifestEntry>> restoredDvEntries = new
HashMap<>();
+ restored.copyInto(restoredDvEntries);
+
+ assertThat(restoredDvEntries).isEqualTo(dvEntries);
+ assertThat(captured.baseSnapshotCaptured()).isTrue();
+ assertThat(restored.baseSnapshotCaptured()).isTrue();
+ }
+
+ @Test
+ public void
testRewriteWithoutCapturedMetadataSkipsDvCleanupWhenBaseSnapshotExpired()
+ throws Exception {
+ TestAppendFileStore store =
+ createAppendStore(
+ tempDir,
+ Collections.singletonMap(
+ CoreOptions.DELETION_VECTORS_ENABLED.key(),
"true"));
+
+ Map<String, List<Integer>> dvs = new HashMap<>();
+ dvs.put("data-0.orc", Arrays.asList(1, 3, 5));
+ CommitMessageImpl dvMessage =
store.writeDVIndexFiles(BinaryRow.EMPTY_ROW, 0, dvs);
+ store.commit(dvMessage);
+
+ FileStoreTable table =
+ FileStoreTableFactory.create(
+ store.fileIO(), store.options().path(),
store.schema());
+ long baseSnapshotId = table.snapshotManager().latestSnapshotId();
+
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta sorted = newFile("sorted-0.orc", 0, 0, 100, 100);
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+
+ CommitMessageImpl written =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 0,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(sorted),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ CompactIncrement.emptyIncrement());
+
+ SortCompactPlanMetadata planMetadata =
+ SortCompactPlanMetadata.capture(
+ table, baseSnapshotId,
Collections.singletonList(split));
+ table.snapshotManager().deleteSnapshot(baseSnapshotId);
+
+ List<CommitMessage> withoutCapturedMetadata =
+ new SortCompactCommitMessageRewriter(
+ table, baseSnapshotId,
Collections.singletonList(split))
+ .rewrite(Collections.singletonList(written));
+ assertThat(
+ ((CommitMessageImpl) withoutCapturedMetadata.get(0))
+ .compactIncrement()
+ .deletedIndexFiles())
+ .isEmpty();
+
+ List<CommitMessage> withCapturedMetadata =
+ new SortCompactCommitMessageRewriter(
+ table,
+ baseSnapshotId,
+ Collections.singletonList(split),
+ planMetadata)
+ .rewrite(Collections.singletonList(written));
+ assertThat(
+ ((CommitMessageImpl) withCapturedMetadata.get(0))
+ .compactIncrement()
+ .deletedIndexFiles())
+ .isNotEmpty();
+ }
+
+ @Test
+ public void testRewriteRejectsInlineCompactionOutput() throws Exception {
+ FileStoreTable table = createAppendTable(Collections.emptyMap());
+
+ DataFileMeta old = newFile("data-0.orc", 0, 0, 100, 100);
+ DataFileMeta l0File = newFile("l0-0.orc", 0, 0, 100, 100);
+ DataFileMeta compactedFile = newFile("compacted-0.orc", 0, 0, 100,
100);
+
+ DataSplit split =
+ DataSplit.builder()
+ .withPartition(BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(Collections.singletonList(old))
+ .build();
+
+ CommitMessageImpl written =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ 0,
+ table.coreOptions().bucket(),
+ new DataIncrement(
+ Collections.singletonList(l0File),
+ Collections.emptyList(),
+ Collections.emptyList()),
+ new CompactIncrement(
+ Collections.singletonList(old),
+ Collections.singletonList(compactedFile),
+ Collections.emptyList()));
+
+ SortCompactCommitMessageRewriter rewriter =
+ new SortCompactCommitMessageRewriter(table, 0L,
Collections.singletonList(split));
+
+ assertThatThrownBy(() ->
rewriter.rewrite(Collections.singletonList(written)))
+ .isInstanceOf(IllegalStateException.class)
+ .hasMessageContaining("inline compaction changes");
+ }
+
+ private FileStoreTable createAppendTable(Map<String, String>
dynamicOptions) throws Exception {
+ TestAppendFileStore store = createAppendStore(tempDir, dynamicOptions);
+ return FileStoreTableFactory.create(store.fileIO(),
store.options().path(), store.schema());
+ }
+
+ /**
+ * Commit additional deletion vectors for an unaware-bucket append table,
correctly deleting the
+ * previous index file when replacing an existing DV.
+ */
+ private void commitUnawareDeletionVectors(
+ FileStoreTable table, Map<String, List<Integer>>
dataFileToPositions) throws Exception {
+ BaseAppendDeleteFileMaintainer maintainer =
+ BaseAppendDeleteFileMaintainer.forUnawareAppend(
+ table.store().newIndexFileHandler(),
+ table.snapshotManager().latestSnapshot(),
+ BinaryRow.EMPTY_ROW);
+ for (Map.Entry<String, List<Integer>> entry :
dataFileToPositions.entrySet()) {
+ DeletionVector deletionVector = new BitmapDeletionVector();
+ for (Integer pos : entry.getValue()) {
+ deletionVector.delete(pos);
+ }
+ maintainer.notifyNewDeletionVector(entry.getKey(), deletionVector);
+ }
+
+ List<IndexFileMeta> newIndexFiles = new ArrayList<>();
+ List<IndexFileMeta> deletedIndexFiles = new ArrayList<>();
+ for (IndexManifestEntry entry : maintainer.persist()) {
+ if (entry.kind() == FileKind.ADD) {
+ newIndexFiles.add(entry.indexFile());
+ } else {
+ deletedIndexFiles.add(entry.indexFile());
+ }
+ }
+
+ CommitMessage message =
+ new CommitMessageImpl(
+ BinaryRow.EMPTY_ROW,
+ UNAWARE_BUCKET,
+ null,
+ new DataIncrement(
+ Collections.emptyList(),
+ Collections.emptyList(),
+ Collections.emptyList(),
+ newIndexFiles,
+ deletedIndexFiles),
+ CompactIncrement.emptyIncrement());
+ try (BatchTableCommit commit =
table.newBatchWriteBuilder().newCommit()) {
+ commit.commit(Collections.singletonList(message));
+ }
+ }
+
+ private TestAppendFileStore createAppendStore(
+ java.nio.file.Path tempDir, Map<String, String> dynamicOptions)
throws Exception {
+ String root = TraceableFileIO.SCHEME + "://" + tempDir.toString();
+ Path path = new Path(tempDir.toUri());
+ FileIO fileIO = FileIOFinder.find(new Path(root));
+ SchemaManager schemaManage = new FileSystemSchemaManager(new
LocalFileIO(), path);
+
+ Map<String, String> options = new HashMap<>(dynamicOptions);
+ options.put(CoreOptions.PATH.key(), root);
+ TableSchema tableSchema =
+ SchemaUtils.forceCommit(
+ schemaManage,
+ new Schema(
+
TestKeyValueGenerator.DEFAULT_ROW_TYPE.getFields(),
+ Collections.emptyList(),
+ Collections.emptyList(),
+ options,
+ null));
+ return new TestAppendFileStore(
+ fileIO,
+ schemaManage,
+ new CoreOptions(options),
+ tableSchema,
+ RowType.of(),
+ RowType.of(),
+ TestKeyValueGenerator.DEFAULT_ROW_TYPE,
+ (new Path(root)).getName());
+ }
+}