From f18ad022a0ebea476978c84db3856a7294a7945d Mon Sep 17 00:00:00 2001 From: Xiangyi Zhu <82511136+zhuxiangyi@users.noreply.github.com> Date: Sun, 6 Sep 2026 23:32:39 +0800 Subject: [PATCH] [core] Delete the files of a compaction task whose result is discarded A compaction task writes its output files before its CompactResult reaches the writer, and that result is the only thing that knows their names. Two paths throw it away and leak the files: - The task is cancelled. CompactFutureManager#cancelCompaction interrupts the thread and the future drops the result, so nothing deletes what was already written. This is the case the TODO there described. - A section of a merge tree compaction fails. Sections that finished have closed their writers, so their files can only be reached through the result of the whole task, which is never returned now. The files stay behind until an orphan file clean runs. CompactResult now carries the abort executors of the files it describes, and CompactTask takes them over as each part of it finishes. A task deletes them itself when it fails or when it notices it was cancelled, and cancelCompaction deletes them for a task that never gets far enough to observe the interrupt. --- .../paimon/append/AppendCompactTask.java | 11 +- .../append/BucketedAppendCompactManager.java | 56 +-- .../DedicatedFormatRollingFileWriter.java | 9 + .../cluster/BucketedAppendClusterManager.java | 2 +- .../paimon/compact/CompactFutureManager.java | 27 +- .../apache/paimon/compact/CompactResult.java | 18 + .../apache/paimon/compact/CompactTask.java | 76 ++++ .../apache/paimon/io/RollingFileWriter.java | 11 + .../paimon/io/RollingFileWriterImpl.java | 1 + .../compact/ChangelogMergeTreeRewriter.java | 9 +- .../compact/FileRewriteCompactTask.java | 4 +- .../compact/MergeTreeCompactManager.java | 2 +- .../compact/MergeTreeCompactRewriter.java | 4 +- .../compact/MergeTreeCompactTask.java | 4 + .../clustering/ClusteringCompactManager.java | 16 +- .../operation/BaseAppendFileStoreWrite.java | 17 +- .../paimon/append/AppendCompactTaskTest.java | 5 +- .../paimon/append/AppendOnlyWriterTest.java | 11 +- .../paimon/append/FullCompactTaskTest.java | 2 +- .../paimon/compact/CompactTaskCancelTest.java | 239 +++++++++++ .../compact/CompactOrphanFileTest.java | 390 ++++++++++++++++++ 21 files changed, 856 insertions(+), 58 deletions(-) create mode 100644 paimon-core/src/test/java/org/apache/paimon/compact/CompactTaskCancelTest.java create mode 100644 paimon-core/src/test/java/org/apache/paimon/mergetree/compact/CompactOrphanFileTest.java 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; + } + } +}