diff --git a/paimon-core/src/main/java/org/apache/paimon/append/AppendCompactTask.java b/paimon-core/src/main/java/org/apache/paimon/append/AppendCompactTask.java index 77d52b7c7d70..aad83a1d024c 100644 --- a/paimon-core/src/main/java/org/apache/paimon/append/AppendCompactTask.java +++ b/paimon-core/src/main/java/org/apache/paimon/append/AppendCompactTask.java @@ -84,10 +84,11 @@ public CommitMessage doCompact(FileStoreTable table, BaseAppendFileStoreWrite wr partition); compactAfter.addAll( write.compactRewrite( - partition, - UNAWARE_BUCKET, - dvIndexFileMaintainer::getDeletionVector, - compactBefore)); + partition, + UNAWARE_BUCKET, + dvIndexFileMaintainer::getDeletionVector, + compactBefore) + .after()); compactBefore.forEach( f -> dvIndexFileMaintainer.notifyRemovedDeletionVector(f.fileName())); @@ -101,7 +102,7 @@ public CommitMessage doCompact(FileStoreTable table, BaseAppendFileStoreWrite wr } } else { compactAfter.addAll( - write.compactRewrite(partition, UNAWARE_BUCKET, null, compactBefore)); + write.compactRewrite(partition, UNAWARE_BUCKET, null, compactBefore).after()); } CompactIncrement compactIncrement = diff --git a/paimon-core/src/main/java/org/apache/paimon/append/BucketedAppendCompactManager.java b/paimon-core/src/main/java/org/apache/paimon/append/BucketedAppendCompactManager.java index 71cebe7f7a8d..d81fbf348a79 100644 --- a/paimon-core/src/main/java/org/apache/paimon/append/BucketedAppendCompactManager.java +++ b/paimon-core/src/main/java/org/apache/paimon/append/BucketedAppendCompactManager.java @@ -46,8 +46,6 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; -import static java.util.Collections.emptyList; - /** Compact manager for {@link AppendOnlyFileStore}. */ public class BucketedAppendCompactManager extends CompactFutureManager { @@ -115,15 +113,15 @@ private void triggerFullCompaction() { LOG.debug("Submit full compaction with these files {}", toCompact); } - taskFuture = - executor.submit( - new FullCompactTask( - dvMaintainer, - toCompact, - compactionFileSize, - forceRewriteAllFiles, - rewriter, - metricsReporter)); + submitTask( + executor, + new FullCompactTask( + dvMaintainer, + toCompact, + compactionFileSize, + forceRewriteAllFiles, + rewriter, + metricsReporter)); recordCompactionsQueuedRequest(); compacting = new ArrayList<>(toCompact); toCompact.clear(); @@ -147,10 +145,9 @@ private void triggerCompactionWithBestEffort() { LOG.debug("Submit normal compaction with these files {}", compacting); } - taskFuture = - executor.submit( - new AutoCompactTask( - dvMaintainer, compacting, rewriter, metricsReporter)); + submitTask( + executor, + new AutoCompactTask(dvMaintainer, compacting, rewriter, metricsReporter)); recordCompactionsQueuedRequest(); } } @@ -281,7 +278,7 @@ protected CompactResult doCompact() throws Exception { // do compaction if (dvMaintainer != null) { // if deletion vector enables, always trigger compaction. - return compact(dvMaintainer, toCompact, rewriter); + return compact(this, dvMaintainer, toCompact, rewriter); } else { // compute small files int big = 0; @@ -295,9 +292,9 @@ protected CompactResult doCompact() throws Exception { } if (forceRewriteAllFiles || (small > big && toCompact.size() >= FULL_COMPACT_MIN_FILE)) { - return compact(null, toCompact, rewriter); + return compact(this, null, toCompact, rewriter); } else { - return result(emptyList(), emptyList()); + return new CompactResult(); } } } @@ -334,17 +331,20 @@ public AutoCompactTask( @Override protected CompactResult doCompact() throws Exception { - return compact(dvMaintainer, toCompact, rewriter); + return compact(this, dvMaintainer, toCompact, rewriter); } } private static CompactResult compact( + CompactTask task, @Nullable BucketedDvMaintainer dvMaintainer, List toCompact, CompactRewriter rewriter) throws Exception { - List rewrite = rewriter.rewrite(toCompact); - CompactResult result = result(toCompact, rewrite); + CompactResult result = rewriter.rewrite(toCompact); + // The files exist from here on, so hand them to the task before doing anything that can + // still fail: from now on the task is what keeps them reachable. + task.trackNewFiles(result); if (dvMaintainer != null) { toCompact.forEach(f -> dvMaintainer.removeDeletionVectorOf(f.fileName())); result.setDeletionFile(CompactDeletionFile.generateFiles(dvMaintainer)); @@ -352,12 +352,16 @@ private static CompactResult compact( return result; } - private static CompactResult result(List before, List after) { - return new CompactResult(before, after); - } - /** Compact rewriter for append-only table. */ public interface CompactRewriter { - List rewrite(List compactBefore) throws Exception; + + /** + * Rewrites the given files. + * + *

The returned result carries the handles to delete the files it has written, so that + * they can still be cleaned up if the compaction is cancelled before its result is + * committed. + */ + CompactResult rewrite(List compactBefore) throws Exception; } } diff --git a/paimon-core/src/main/java/org/apache/paimon/append/DedicatedFormatRollingFileWriter.java b/paimon-core/src/main/java/org/apache/paimon/append/DedicatedFormatRollingFileWriter.java index aface9879a11..c0a4dfd8b502 100644 --- a/paimon-core/src/main/java/org/apache/paimon/append/DedicatedFormatRollingFileWriter.java +++ b/paimon-core/src/main/java/org/apache/paimon/append/DedicatedFormatRollingFileWriter.java @@ -444,6 +444,15 @@ public void abort() { } } + /** Transfers ownership of abort executors for closed files to the caller. */ + @Override + public List drainAbortExecutors() { + Preconditions.checkState(closed, "Cannot drain abort executors unless close all writers."); + List abortExecutors = new ArrayList<>(closedWriters); + closedWriters.clear(); + return abortExecutors; + } + /** * Checks if the current file should be rolled. The row cap applies even when there is no main * writer (all fields dedicated), so blob/vector writers roll together. diff --git a/paimon-core/src/main/java/org/apache/paimon/append/cluster/BucketedAppendClusterManager.java b/paimon-core/src/main/java/org/apache/paimon/append/cluster/BucketedAppendClusterManager.java index c223084d3184..5b13ebc6bcde 100644 --- a/paimon-core/src/main/java/org/apache/paimon/append/cluster/BucketedAppendClusterManager.java +++ b/paimon-core/src/main/java/org/apache/paimon/append/cluster/BucketedAppendClusterManager.java @@ -151,7 +151,7 @@ private void submitCompaction(CompactUnit unit) { file.fileName(), file.level(), file.fileSize())) .collect(Collectors.joining(", "))); } - taskFuture = executor.submit(task); + submitTask(executor, task); } @Override diff --git a/paimon-core/src/main/java/org/apache/paimon/compact/CompactFutureManager.java b/paimon-core/src/main/java/org/apache/paimon/compact/CompactFutureManager.java index e43bec01630d..e1cb691d6c94 100644 --- a/paimon-core/src/main/java/org/apache/paimon/compact/CompactFutureManager.java +++ b/paimon-core/src/main/java/org/apache/paimon/compact/CompactFutureManager.java @@ -20,9 +20,12 @@ import org.apache.paimon.annotation.VisibleForTesting; +import javax.annotation.Nullable; + import java.util.Optional; import java.util.concurrent.CancellationException; import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; /** Base implementation of {@link CompactManager} which runs compaction in a separate thread. */ @@ -30,12 +33,29 @@ public abstract class CompactFutureManager implements CompactManager { protected Future taskFuture; + /** + * The task behind {@link #taskFuture}, kept so that its files can be deleted if it is + * cancelled. + */ + @Nullable private CompactTask task; + + /** Submits a compaction task and remembers it as the current one. */ + protected void submitTask(ExecutorService executor, CompactTask task) { + this.task = task; + this.taskFuture = executor.submit(task); + } + @Override public void cancelCompaction() { - // TODO this method may leave behind orphan files if compaction is actually finished - // but some CPU work still needs to be done if (taskFuture != null && !taskFuture.isCancelled()) { - taskFuture.cancel(true); + boolean cancelled = taskFuture.cancel(true); + if (cancelled && task != null) { + // A cancelled future throws its result away, so the files the task has already + // written become unreachable: nothing else knows their names. The task deletes + // them itself once the interrupt reaches it, but the interrupt can just as well + // land after the task is done, and then this is the only cleanup left. + task.abortNewFiles(); + } } } @@ -55,6 +75,7 @@ protected final Optional innerGetCompactionResult(boolean blockin return Optional.empty(); } finally { taskFuture = null; + task = null; } return Optional.of(result); } diff --git a/paimon-core/src/main/java/org/apache/paimon/compact/CompactResult.java b/paimon-core/src/main/java/org/apache/paimon/compact/CompactResult.java index 08d7de5dab7f..dbb44a6c4a9e 100644 --- a/paimon-core/src/main/java/org/apache/paimon/compact/CompactResult.java +++ b/paimon-core/src/main/java/org/apache/paimon/compact/CompactResult.java @@ -19,6 +19,7 @@ package org.apache.paimon.compact; import org.apache.paimon.io.DataFileMeta; +import org.apache.paimon.io.FileWriterAbortExecutor; import javax.annotation.Nullable; @@ -33,6 +34,13 @@ public class CompactResult { private final List after; private final List changelog; + /** + * Handles to delete the files this result is made of. They are only meaningful while the result + * has not been handed over to the writer: whoever throws the result away is responsible for + * aborting the files it describes. + */ + private final List abortExecutors; + @Nullable private CompactDeletionFile deletionFile; public CompactResult() { @@ -52,6 +60,7 @@ public CompactResult( this.before = new ArrayList<>(before); this.after = new ArrayList<>(after); this.changelog = new ArrayList<>(changelog); + this.abortExecutors = new ArrayList<>(); } public List before() { @@ -66,6 +75,14 @@ public List changelog() { return changelog; } + public void addAbortExecutors(List executors) { + abortExecutors.addAll(executors); + } + + public List abortExecutors() { + return abortExecutors; + } + public void setDeletionFile(@Nullable CompactDeletionFile deletionFile) { this.deletionFile = deletionFile; } @@ -79,6 +96,7 @@ public void merge(CompactResult that) { before.addAll(that.before); after.addAll(that.after); changelog.addAll(that.changelog); + abortExecutors.addAll(that.abortExecutors); if (deletionFile != null || that.deletionFile != null) { throw new UnsupportedOperationException( diff --git a/paimon-core/src/main/java/org/apache/paimon/compact/CompactTask.java b/paimon-core/src/main/java/org/apache/paimon/compact/CompactTask.java index 229b25324b2b..b8888427db6d 100644 --- a/paimon-core/src/main/java/org/apache/paimon/compact/CompactTask.java +++ b/paimon-core/src/main/java/org/apache/paimon/compact/CompactTask.java @@ -19,6 +19,7 @@ package org.apache.paimon.compact; import org.apache.paimon.io.DataFileMeta; +import org.apache.paimon.io.FileWriterAbortExecutor; import org.apache.paimon.operation.metrics.CompactionMetrics; import org.apache.paimon.operation.metrics.MetricUtils; @@ -27,6 +28,7 @@ import javax.annotation.Nullable; +import java.util.ArrayList; import java.util.List; import java.util.concurrent.Callable; @@ -38,6 +40,13 @@ public abstract class CompactTask implements Callable { @Nullable private final CompactionMetrics.Reporter metricsReporter; private final String bucketInfo; + /** + * Files this task has already written, with the handles needed to delete them again. Written by + * the compaction thread as parts of the task finish, read by whoever cancels the task, hence + * the lock. + */ + private final List newFiles = new ArrayList<>(); + public CompactTask(@Nullable CompactionMetrics.Reporter metricsReporter, String bucketInfo) { this.metricsReporter = metricsReporter; this.bucketInfo = bucketInfo; @@ -88,6 +97,15 @@ public CompactResult call() throws Exception { if (LOG.isDebugEnabled()) { LOG.debug(logMetric(startMillis, result.before(), result.after())); } + + // Keep this check last. If we were cancelled the future drops the result on the + // floor, so this task is the only one that still knows about the files it wrote. + if (Thread.currentThread().isInterrupted()) { + trackNewFiles(result); + cleanDeletionFile(result); + throw new InterruptedException( + "Compact task was cancelled after it had written its files."); + } return result; } catch (Exception e) { LOG.warn( @@ -95,6 +113,10 @@ public CompactResult call() throws Exception { bucketInfo, getClass().getSimpleName(), e); + // Parts of this task that already finished have closed their writers, so nothing but + // this task can delete their output any more. The result is never returned now, so + // those files would be left behind. + abortNewFiles(); throw e; } finally { MetricUtils.safeCall(this::stopTimer, LOG); @@ -102,6 +124,60 @@ public CompactResult call() throws Exception { } } + /** + * Takes over the files a finished part of this task has written, so that {@link + * #abortNewFiles()} can still delete them once the sub-result has been merged away. + * + *

Ownership is transferred: the handles are removed from {@code partialResult}, so a file is + * never tracked twice. + */ + public void trackNewFiles(CompactResult partialResult) { + List executors = partialResult.abortExecutors(); + if (executors.isEmpty()) { + return; + } + synchronized (newFiles) { + newFiles.addAll(executors); + } + executors.clear(); + } + + /** + * Deletes every file written by this task so far. Only call this once it is certain that the + * result of this task will not be committed, otherwise it deletes live data. + */ + public void abortNewFiles() { + List toAbort; + synchronized (newFiles) { + if (newFiles.isEmpty()) { + return; + } + toAbort = new ArrayList<>(newFiles); + newFiles.clear(); + } + + LOG.info( + "Deleting {} file(s) written by a compact task whose result is discarded: {}, taskType={}", + toAbort.size(), + bucketInfo, + getClass().getSimpleName()); + for (FileWriterAbortExecutor abortExecutor : toAbort) { + abortExecutor.abort(); + } + } + + private void cleanDeletionFile(CompactResult result) { + CompactDeletionFile deletionFile = result.deletionFile(); + if (deletionFile == null) { + return; + } + try { + deletionFile.clean(); + } catch (Throwable t) { + LOG.warn("Failed to clean the deletion file of a discarded compact task.", t); + } + } + private void decreaseCompactionsQueuedCount() { if (metricsReporter != null) { metricsReporter.decreaseCompactionsQueuedCount(); diff --git a/paimon-core/src/main/java/org/apache/paimon/io/RollingFileWriter.java b/paimon-core/src/main/java/org/apache/paimon/io/RollingFileWriter.java index 18846bcc084f..ca33188f37a6 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/RollingFileWriter.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/RollingFileWriter.java @@ -44,6 +44,17 @@ public interface RollingFileWriter extends FileWriter> { void writeBundle(BundleRecords records) throws IOException; + /** + * Transfers ownership of the abort executors for the files this writer has already closed and + * rolled over. + * + *

Once a file is rolled over its writer is dropped, so nothing but this handle can delete it + * any more. A caller which is going to throw away {@link #result()} - a compaction task whose + * result will never be committed, for instance - must drain these and abort them, otherwise the + * files are left behind with nothing referencing them. + */ + List drainAbortExecutors(); + @VisibleForTesting static FileWriterContext createFileWriterContext( FileFormat fileFormat, diff --git a/paimon-core/src/main/java/org/apache/paimon/io/RollingFileWriterImpl.java b/paimon-core/src/main/java/org/apache/paimon/io/RollingFileWriterImpl.java index 11c332b90fa5..70ff9f84424c 100644 --- a/paimon-core/src/main/java/org/apache/paimon/io/RollingFileWriterImpl.java +++ b/paimon-core/src/main/java/org/apache/paimon/io/RollingFileWriterImpl.java @@ -192,6 +192,7 @@ public List result() { } /** Transfers ownership of abort executors for closed files to the caller. */ + @Override public List drainAbortExecutors() { Preconditions.checkState(closed, "Cannot drain abort executors unless close all writers."); List abortExecutors = new ArrayList<>(closedWriters); diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/ChangelogMergeTreeRewriter.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/ChangelogMergeTreeRewriter.java index a76c590d1191..b48ba251ab5a 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/ChangelogMergeTreeRewriter.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/ChangelogMergeTreeRewriter.java @@ -206,7 +206,14 @@ private CompactResult rewriteOrProduceChangelog( changelogFileWriter != null ? changelogFileWriter.result() : Collections.emptyList(); - return new CompactResult(before, after, changelogFiles); + CompactResult result = new CompactResult(before, after, changelogFiles); + if (compactFileWriter != null) { + result.addAbortExecutors(compactFileWriter.drainAbortExecutors()); + } + if (changelogFileWriter != null) { + result.addAbortExecutors(changelogFileWriter.drainAbortExecutors()); + } + return result; } @Override diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FileRewriteCompactTask.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FileRewriteCompactTask.java index 0e52dbf14c9a..c0448af7fa61 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FileRewriteCompactTask.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FileRewriteCompactTask.java @@ -69,6 +69,8 @@ protected CompactResult doCompact() throws Exception { private void rewriteFile(DataFileMeta file, CompactResult toUpdate) throws Exception { List> candidate = singletonList(singletonList(SortedRun.fromSingle(file))); - toUpdate.merge(rewriter.rewrite(outputLevel, dropDelete, candidate)); + CompactResult rewriteResult = rewriter.rewrite(outputLevel, dropDelete, candidate); + trackNewFiles(rewriteResult); + toUpdate.merge(rewriteResult); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManager.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManager.java index b1bc38533bcb..40c724229d46 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManager.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManager.java @@ -255,7 +255,7 @@ private void submitCompaction(CompactUnit unit, boolean dropDelete) { file.fileName(), file.level(), file.fileSize())) .collect(Collectors.joining(", "))); } - taskFuture = executor.submit(task); + submitTask(executor, task); if (metricsReporter != null) { metricsReporter.increaseCompactionsQueuedCount(); metricsReporter.increaseCompactionsTotalCount(); diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactRewriter.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactRewriter.java index 9eecce5248c4..29b756577d2b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactRewriter.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactRewriter.java @@ -112,7 +112,9 @@ protected CompactResult rewriteCompaction( metricsReporter.reportSortBufferMetrics( mergeSorter.sortBufferUsedBytes(), mergeSorter.sortBufferTotalBytes()); } - return new CompactResult(before, after); + CompactResult result = new CompactResult(before, after); + result.addAbortExecutors(writer.drainAbortExecutors()); + return result; } protected RecordReader readerForMergeTree( diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactTask.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactTask.java index db6d8e23e831..06740884470c 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactTask.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactTask.java @@ -133,6 +133,7 @@ private void upgrade(DataFileMeta file, CompactResult toUpdate) throws Exception if (file.level() != outputLevel) { CompactResult upgradeResult = rewriter.upgrade(outputLevel, file); + trackNewFiles(upgradeResult); toUpdate.merge(upgradeResult); upgradeFilesNum++; } @@ -160,6 +161,9 @@ private void rewrite(List> candidate, CompactResult toUpdate) th private void rewriteImpl(List> candidate, CompactResult toUpdate) throws Exception { CompactResult rewriteResult = rewriter.rewrite(outputLevel, dropDelete, candidate); + // Take the files over before merging: this section is finished, its writer is closed, and + // if a later section fails or the task is cancelled nothing else can delete them. + trackNewFiles(rewriteResult); toUpdate.merge(rewriteResult); candidate.clear(); } diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/clustering/ClusteringCompactManager.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/clustering/ClusteringCompactManager.java index 8d094209577f..50f2948ffe04 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/clustering/ClusteringCompactManager.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/clustering/ClusteringCompactManager.java @@ -184,14 +184,14 @@ public void triggerCompaction(boolean fullCompaction) { if (taskFuture != null) { return; } - taskFuture = - executor.submit( - new CompactTask(metricsReporter, "") { - @Override - protected CompactResult doCompact() throws Exception { - return compact(fullCompaction); - } - }); + submitTask( + executor, + new CompactTask(metricsReporter, "") { + @Override + protected CompactResult doCompact() throws Exception { + return compact(fullCompaction); + } + }); } private CompactResult compact(boolean fullCompaction) throws Exception { diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java b/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java index 2e734552bd71..64952500e648 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java @@ -23,6 +23,7 @@ import org.apache.paimon.append.AppendOnlyWriter; import org.apache.paimon.append.cluster.Sorter; import org.apache.paimon.compact.CompactManager; +import org.apache.paimon.compact.CompactResult; import org.apache.paimon.data.BinaryRow; import org.apache.paimon.data.BlobConsumer; import org.apache.paimon.data.InternalRow; @@ -58,7 +59,6 @@ import javax.annotation.Nullable; import java.io.IOException; -import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.List; @@ -269,14 +269,21 @@ protected abstract CompactManager getCompactManager( ExecutorService compactExecutor, @Nullable BucketedDvMaintainer dvMaintainer); - public List compactRewrite( + /** + * Rewrites the given files into new ones. + * + *

The result also carries the handles needed to delete the files it has written: a caller + * whose result may end up being thrown away - an asynchronous compaction task that can be + * cancelled - needs those to clean up after itself. + */ + public CompactResult compactRewrite( BinaryRow partition, int bucket, @Nullable Function dvFactory, List toCompact) throws Exception { if (toCompact.isEmpty()) { - return Collections.emptyList(); + return new CompactResult(); } Exception collectedExceptions = null; RowDataRollingFileWriter rewriter = @@ -305,7 +312,9 @@ public List compactRewrite( if (collectedExceptions != null) { throw collectedExceptions; } - return rewriter.result(); + CompactResult result = new CompactResult(toCompact, rewriter.result()); + result.addAbortExecutors(rewriter.drainAbortExecutors()); + return result; } public List clusterRewrite( diff --git a/paimon-core/src/test/java/org/apache/paimon/append/AppendCompactTaskTest.java b/paimon-core/src/test/java/org/apache/paimon/append/AppendCompactTaskTest.java index 19fc3875931d..12fc6b245c32 100644 --- a/paimon-core/src/test/java/org/apache/paimon/append/AppendCompactTaskTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/append/AppendCompactTaskTest.java @@ -22,6 +22,7 @@ import org.apache.paimon.TestAppendFileStore; import org.apache.paimon.TestKeyValueGenerator; import org.apache.paimon.compact.CompactManager; +import org.apache.paimon.compact.CompactResult; import org.apache.paimon.data.BinaryRow; import org.apache.paimon.data.InternalRow; import org.apache.paimon.deletionvectors.BucketedDvMaintainer; @@ -157,13 +158,13 @@ private NoopAppendWrite( } @Override - public List compactRewrite( + public CompactResult compactRewrite( BinaryRow partition, int bucket, @Nullable Function dvFactory, List toCompact) throws Exception { - return Collections.emptyList(); + return new CompactResult(); } @Override diff --git a/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java b/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java index 3b2464d311d3..6a1af78ceb15 100644 --- a/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java @@ -19,6 +19,7 @@ package org.apache.paimon.append; import org.apache.paimon.CoreOptions; +import org.apache.paimon.compact.CompactResult; import org.apache.paimon.compact.NoopCompactManager; import org.apache.paimon.compression.CompressOptions; import org.apache.paimon.data.BinaryRow; @@ -192,7 +193,7 @@ public void testBinaryColumnStatsRoundTrip() throws Exception { true, true, Collections.emptyList(), - compactBefore -> Collections.emptyList(), + compactBefore -> new CompactResult(), options) .getKey(); @@ -1228,8 +1229,10 @@ private Pair> createWriter( compactBefore -> { latch.await(); return compactBefore.isEmpty() - ? Collections.emptyList() - : Collections.singletonList(generateCompactAfter(compactBefore)); + ? new CompactResult() + : new CompactResult( + compactBefore, + Collections.singletonList(generateCompactAfter(compactBefore))); }, options); } @@ -1248,7 +1251,7 @@ private AppendOnlyWriter createVectorStoreWriter( false, true, Collections.emptyList(), - compactBefore -> Collections.emptyList(), + compactBefore -> new CompactResult(), options) .getKey(); } diff --git a/paimon-core/src/test/java/org/apache/paimon/append/FullCompactTaskTest.java b/paimon-core/src/test/java/org/apache/paimon/append/FullCompactTaskTest.java index 2877247145dc..784c320a5774 100644 --- a/paimon-core/src/test/java/org/apache/paimon/append/FullCompactTaskTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/append/FullCompactTaskTest.java @@ -152,7 +152,7 @@ private BucketedAppendCompactManager.CompactRewriter rewriter() { compactAfter.add(newFile(minSeq, file.maxSequenceNumber())); } } - return compactAfter; + return new CompactResult(compactBefore, compactAfter); }; } } diff --git a/paimon-core/src/test/java/org/apache/paimon/compact/CompactTaskCancelTest.java b/paimon-core/src/test/java/org/apache/paimon/compact/CompactTaskCancelTest.java new file mode 100644 index 000000000000..e0518503ec60 --- /dev/null +++ b/paimon-core/src/test/java/org/apache/paimon/compact/CompactTaskCancelTest.java @@ -0,0 +1,239 @@ +/* + * 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.compact; + +import org.apache.paimon.fs.Path; +import org.apache.paimon.fs.local.LocalFileIO; +import org.apache.paimon.io.DataFileMeta; +import org.apache.paimon.io.FileWriterAbortExecutor; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.api.io.TempDir; + +import java.io.IOException; +import java.util.Collection; +import java.util.Collections; +import java.util.Optional; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** + * Tests that a {@link CompactTask} whose result is thrown away does not leave the files it has + * already written behind. + */ +public class CompactTaskCancelTest { + + @TempDir java.nio.file.Path tempDir; + + private LocalFileIO fileIO; + private ExecutorService executor; + + @BeforeEach + public void before() { + fileIO = LocalFileIO.create(); + executor = Executors.newSingleThreadExecutor(); + } + + @AfterEach + public void after() { + executor.shutdownNow(); + // a test may leave the flag set on the main thread + Thread.interrupted(); + } + + /** + * A section that has been rewritten has closed its writer, so its files can only be reached + * through the result of the whole task. When a later section fails that result is never + * returned, and the earlier files used to be leaked. + */ + @Test + public void testFilesOfFinishedPartsAreDeletedWhenALaterPartFails() throws Exception { + Path first = writeFile("first"); + Path second = writeFile("second"); + + CompactTask task = + new CompactTask(null, "") { + @Override + protected CompactResult doCompact() { + CompactResult result = new CompactResult(); + for (Path finished : new Path[] {first, second}) { + CompactResult part = finishPart(finished); + trackNewFiles(part); + result.merge(part); + } + throw new RuntimeException("rewriting the third section failed"); + } + }; + + assertThatThrownBy(task::call).hasMessageContaining("third section"); + + assertThat(fileIO.exists(first)).isFalse(); + assertThat(fileIO.exists(second)).isFalse(); + } + + /** + * The interrupt of a cancellation can arrive after the task has written everything. The future + * then drops the result, so the task itself has to clean up. + */ + @Test + public void testFilesAreDeletedWhenCancelledAfterWriting() throws Exception { + Path written = writeFile("written"); + + CompactTask task = + new CompactTask(null, "") { + @Override + protected CompactResult doCompact() { + CompactResult result = new CompactResult(); + result.merge(finishPart(written)); + // the cancellation lands here, once all the files are on disk + Thread.currentThread().interrupt(); + return result; + } + }; + + assertThatThrownBy(task::call).isInstanceOf(InterruptedException.class); + + assertThat(fileIO.exists(written)).isFalse(); + } + + /** A task which completes normally keeps its files - they are about to be committed. */ + @Test + public void testFilesAreKeptWhenTaskSucceeds() throws Exception { + Path written = writeFile("written"); + + CompactTask task = + new CompactTask(null, "") { + @Override + protected CompactResult doCompact() { + CompactResult result = new CompactResult(); + result.merge(finishPart(written)); + return result; + } + }; + + assertThat(task.call()).isNotNull(); + assertThat(fileIO.exists(written)).isTrue(); + } + + /** + * Covers the case the {@code cancelCompaction} TODO described: the task is done writing but + * still busy, so the interrupt never reaches a point where the task can react to it. Only the + * manager can clean up then. + */ + @Test + @Timeout(30) + public void testCancelCompactionDeletesFilesOfAnUnresponsiveTask() throws Exception { + Path written = writeFile("written"); + + CountDownLatch tracked = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + + CompactTask task = + new CompactTask(null, "") { + @Override + protected CompactResult doCompact() { + CompactResult result = new CompactResult(); + CompactResult part = finishPart(written); + trackNewFiles(part); + result.merge(part); + tracked.countDown(); + // busy with work that does not observe interrupts + boolean released = false; + while (!released) { + try { + released = release.await(1, TimeUnit.SECONDS); + } catch (InterruptedException ignored) { + // deliberately swallowed + } + } + return result; + } + }; + + TestCompactManager manager = new TestCompactManager(); + manager.submit(executor, task); + assertThat(tracked.await(30, TimeUnit.SECONDS)).isTrue(); + + manager.cancelCompaction(); + + assertThat(fileIO.exists(written)).isFalse(); + release.countDown(); + } + + /** Builds the result of one finished part of a task, and hands its files to the task. */ + private CompactResult finishPart(Path file) { + CompactResult part = new CompactResult(); + part.addAbortExecutors( + Collections.singletonList(new FileWriterAbortExecutor(fileIO, file))); + return part; + } + + private Path writeFile(String name) throws IOException { + Path path = new Path(tempDir.toUri().toString(), name); + fileIO.tryToWriteAtomic(path, "some compacted data"); + assertThat(fileIO.exists(path)).isTrue(); + return path; + } + + private static class TestCompactManager extends CompactFutureManager { + + void submit(ExecutorService executor, CompactTask task) { + submitTask(executor, task); + } + + @Override + public boolean shouldWaitForLatestCompaction() { + return false; + } + + @Override + public boolean shouldWaitForPreparingCheckpoint() { + return false; + } + + @Override + public void addNewFile(DataFileMeta file) {} + + @Override + public Collection allFiles() { + return Collections.emptyList(); + } + + @Override + public void triggerCompaction(boolean fullCompaction) {} + + @Override + public Optional getCompactionResult(boolean blocking) + throws ExecutionException, InterruptedException { + return innerGetCompactionResult(blocking); + } + + @Override + public void close() {} + } +} diff --git a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/CompactOrphanFileTest.java b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/CompactOrphanFileTest.java new file mode 100644 index 000000000000..d8e528ab2b77 --- /dev/null +++ b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/CompactOrphanFileTest.java @@ -0,0 +1,390 @@ +/* + * 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.mergetree.compact; + +import org.apache.paimon.CoreOptions; +import org.apache.paimon.KeyValue; +import org.apache.paimon.compact.CompactResult; +import org.apache.paimon.data.GenericRow; +import org.apache.paimon.data.InternalRow; +import org.apache.paimon.deletionvectors.DeletionVector; +import org.apache.paimon.format.FileFormat; +import org.apache.paimon.fs.FileStatus; +import org.apache.paimon.fs.Path; +import org.apache.paimon.fs.local.LocalFileIO; +import org.apache.paimon.io.DataFileMeta; +import org.apache.paimon.io.KeyValueFileReaderFactory; +import org.apache.paimon.io.KeyValueFileWriterFactory; +import org.apache.paimon.io.RollingFileWriter; +import org.apache.paimon.manifest.FileSource; +import org.apache.paimon.mergetree.Levels; +import org.apache.paimon.mergetree.MergeSorter; +import org.apache.paimon.mergetree.SortedRun; +import org.apache.paimon.options.MemorySize; +import org.apache.paimon.options.Options; +import org.apache.paimon.schema.KeyValueFieldsExtractor; +import org.apache.paimon.schema.SchemaManager; +import org.apache.paimon.schema.TableSchema; +import org.apache.paimon.table.SchemaEvolutionTableTestBase; +import org.apache.paimon.types.DataField; +import org.apache.paimon.types.IntType; +import org.apache.paimon.types.RowKind; +import org.apache.paimon.types.RowType; +import org.apache.paimon.utils.FileStorePathFactory; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.api.io.TempDir; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Comparator; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.function.Function; +import java.util.stream.Collectors; + +import static java.util.Collections.singletonList; +import static org.apache.paimon.utils.FileStorePathFactoryTest.createNonPartFactory; +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** + * A compaction task writes real files before its result reaches the writer. These tests run a real + * merge tree compaction, break it half way through, and check that the files the finished part has + * already written to the bucket directory do not survive as orphans. + */ +public class CompactOrphanFileTest { + + @TempDir java.nio.file.Path tempDir; + + private final LocalFileIO fileIO = LocalFileIO.create(); + private final RowType keyType = + new RowType(singletonList(new DataField(0, "k", new IntType()))); + private final RowType valueType = + new RowType(singletonList(new DataField(1, "v", new IntType()))); + + private CoreOptions options; + private Comparator comparator; + private KeyValueFileReaderFactory readerFactory; + private KeyValueFileWriterFactory writerFactory; + private Path bucketDir; + private ExecutorService executor; + private long sequenceNumber = 0; + + @BeforeEach + public void before() throws IOException { + Path root = new Path(tempDir.toString()); + FileStorePathFactory pathFactory = createNonPartFactory(root); + comparator = Comparator.comparingInt(o -> o.getInt(0)); + executor = Executors.newSingleThreadExecutor(); + + Options conf = new Options(); + conf.set(CoreOptions.TARGET_FILE_SIZE, MemorySize.ofMebiBytes(1)); + options = new CoreOptions(conf); + + FileFormat avro = FileFormat.fromIdentifier("avro", new Options()); + SchemaManager schemaManager = testingSchemaManager(root); + KeyValueFileReaderFactory.Builder readerBuilder = + KeyValueFileReaderFactory.builder( + fileIO, + schemaManager, + schemaManager.schema(0), + keyType, + valueType, + ignore -> avro, + pathFactory, + new TestKeyValueFieldsExtractor(), + new CoreOptions(new HashMap<>())); + readerFactory = + readerBuilder.build( + org.apache.paimon.data.BinaryRow.EMPTY_ROW, + 0, + DeletionVector.emptyFactory()); + + Function pathFactoryMap = k -> pathFactory; + writerFactory = + KeyValueFileWriterFactory.builder( + fileIO, + 0, + keyType, + valueType, + avro, + pathFactoryMap, + options.targetFileSize(true)) + .build(org.apache.paimon.data.BinaryRow.EMPTY_ROW, 0, options); + + bucketDir = writerFactory.pathFactory(0).newPath().getParent(); + fileIO.mkdirs(bucketDir); + } + + @AfterEach + public void after() { + executor.shutdownNow(); + } + + /** + * The compaction is cancelled while it is still running. Everything the sections it already + * rewrote wrote to the bucket directory is only reachable through the result of the task, which + * the cancelled future drops. + */ + @Test + @Timeout(60) + public void testCancelledCompactionLeavesNoOrphanFile() throws Exception { + List inputs = writeInputFiles(); + + CountDownLatch reachedSecondRewrite = new CountDownLatch(1); + CountDownLatch neverReleased = new CountDownLatch(1); + BreakingRewriter rewriter = + new BreakingRewriter( + realRewriter(), + 2, + () -> { + reachedSecondRewrite.countDown(); + neverReleased.await(); + }); + + MergeTreeCompactManager manager = createCompactManager(inputs, rewriter); + manager.triggerCompaction(true); + assertThat(reachedSecondRewrite.await(60, TimeUnit.SECONDS)).isTrue(); + + // the first section has been rewritten by now, its files are on disk + assertThat(orphanFiles(inputs)).isNotEmpty(); + + manager.cancelCompaction(); + + assertThat(orphanFiles(inputs)).isEmpty(); + } + + /** + * The same files are left behind when a later section of the same task fails: the task never + * returns the result that names them. + */ + @Test + @Timeout(60) + public void testFailedCompactionLeavesNoOrphanFile() throws Exception { + List inputs = writeInputFiles(); + + BreakingRewriter rewriter = + new BreakingRewriter( + realRewriter(), + 2, + () -> { + throw new IOException("rewriting the second section failed"); + }); + + MergeTreeCompactManager manager = createCompactManager(inputs, rewriter); + manager.triggerCompaction(true); + + assertThatThrownBy(() -> manager.getCompactionResult(true)) + .hasRootCauseMessage("rewriting the second section failed"); + + assertThat(orphanFiles(inputs)).isEmpty(); + } + + /** A compaction that runs to the end keeps its files - they are about to be committed. */ + @Test + @Timeout(60) + public void testSuccessfulCompactionKeepsItsFiles() throws Exception { + List inputs = writeInputFiles(); + + MergeTreeCompactManager manager = createCompactManager(inputs, realRewriter()); + manager.triggerCompaction(true); + + CompactResult result = manager.getCompactionResult(true).orElseThrow(AssertionError::new); + assertThat(result.after()).isNotEmpty(); + + Set onDisk = filesOnDisk(); + for (DataFileMeta file : result.after()) { + assertThat(onDisk).contains(file.fileName()); + } + } + + /** + * Five level 0 files: two overlapping small ones, one large one on its own, then two more + * overlapping small ones. {@link MergeTreeCompactTask} rewrites the first pair, upgrades the + * large file in between, and rewrites the last pair - two separate rewrites in one task. + */ + private List writeInputFiles() throws Exception { + List files = new ArrayList<>(); + files.add(writeFile(1, 10)); + files.add(writeFile(5, 15)); + files.add(writeFile(100, 700)); + files.add(writeFile(2000, 2010)); + files.add(writeFile(2005, 2015)); + return files; + } + + /** Between the small files and the large one, so the large one is upgraded, not rewritten. */ + private long compactionFileSize(List inputs) { + long large = inputs.get(2).fileSize(); + long smallest = + inputs.stream() + .mapToLong(DataFileMeta::fileSize) + .min() + .orElseThrow(AssertionError::new); + assertThat(large).isGreaterThan(smallest * 2); + return (large + smallest) / 2; + } + + private MergeTreeCompactManager createCompactManager( + List inputs, CompactRewriter rewriter) { + return new MergeTreeCompactManager( + executor, + new Levels(comparator, inputs, options.numLevels()), + new UniversalCompaction( + options.maxSizeAmplificationPercent(), + options.sortedRunSizeRatio(), + options.numSortedRunCompactionTrigger(), + null, + null), + comparator, + compactionFileSize(inputs), + options.numSortedRunStopTrigger(), + rewriter, + null, + null, + false, + false, + null, + false, + false, + ""); + } + + private MergeTreeCompactRewriter realRewriter() { + return new MergeTreeCompactRewriter( + readerFactory, + writerFactory, + comparator, + null, + DeduplicateMergeFunction.factory(), + new MergeSorter(options, keyType, valueType, null)); + } + + private DataFileMeta writeFile(int minKey, int maxKey) throws Exception { + RollingFileWriter writer = + writerFactory.createRollingMergeTreeFileWriter(0, FileSource.APPEND); + for (int k = minKey; k <= maxKey; k++) { + writer.write( + new KeyValue() + .replace( + GenericRow.of(k), + sequenceNumber++, + RowKind.INSERT, + GenericRow.of(k))); + } + writer.close(); + List result = writer.result(); + assertThat(result).hasSize(1); + return result.get(0); + } + + private Set filesOnDisk() throws IOException { + FileStatus[] statuses = fileIO.listStatus(bucketDir); + return Arrays.stream(statuses).map(s -> s.getPath().getName()).collect(Collectors.toSet()); + } + + /** Files present in the bucket directory that nothing refers to any more. */ + private Set orphanFiles(List inputs) throws IOException { + Set orphans = new HashSet<>(filesOnDisk()); + inputs.forEach(f -> orphans.remove(f.fileName())); + return orphans; + } + + private SchemaManager testingSchemaManager(Path path) { + TableSchema schema = + new TableSchema( + 0, + new ArrayList<>(), + -1, + new ArrayList<>(), + new ArrayList<>(), + new HashMap<>(), + ""); + Map schemas = new HashMap<>(); + schemas.put(schema.id(), schema); + return new SchemaEvolutionTableTestBase.TestingSchemaManager(path, schemas); + } + + private class TestKeyValueFieldsExtractor implements KeyValueFieldsExtractor { + + private static final long serialVersionUID = 1L; + + @Override + public List keyFields(TableSchema schema) { + return keyType.getFields(); + } + + @Override + public List valueFields(TableSchema schema) { + return valueType.getFields(); + } + } + + /** Runs a real rewriter, but breaks on the n-th call to {@link #rewrite}. */ + private static class BreakingRewriter implements CompactRewriter { + + private final CompactRewriter delegate; + private final int breakOnCall; + private final Break onBreak; + + private int calls = 0; + + private BreakingRewriter(CompactRewriter delegate, int breakOnCall, Break onBreak) { + this.delegate = delegate; + this.breakOnCall = breakOnCall; + this.onBreak = onBreak; + } + + @Override + public CompactResult rewrite( + int outputLevel, boolean dropDelete, List> sections) + throws Exception { + if (++calls == breakOnCall) { + onBreak.run(); + } + return delegate.rewrite(outputLevel, dropDelete, sections); + } + + @Override + public CompactResult upgrade(int outputLevel, DataFileMeta file) throws Exception { + return delegate.upgrade(outputLevel, file); + } + + @Override + public void close() throws IOException { + delegate.close(); + } + + interface Break { + void run() throws Exception; + } + } +}