This is an automated email from the ASF dual-hosted git repository.
yuzelin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new b5c03b2d1c [core] Snapshot expiration deletes data and changelog files
in batches (#9409)
b5c03b2d1c is described below
commit b5c03b2d1ca79e15ef8598453f121141e857258b
Author: yuzelin <[email protected]>
AuthorDate: Wed Aug 26 20:14:56 2026 +0800
[core] Snapshot expiration deletes data and changelog files in batches
(#9409)
---
.../apache/paimon/operation/FileDeletionBase.java | 5 ---
.../paimon/table/AbstractFileStoreTable.java | 3 +-
.../apache/paimon/table/ExpireSnapshotsImpl.java | 41 +++++++++++++++++-----
.../test/java/org/apache/paimon/TestFileStore.java | 6 ++--
.../paimon/operation/ExpireSnapshotsTest.java | 12 +++++--
.../apache/paimon/operation/FileDeletionTest.java | 18 ++++++++--
6 files changed, 62 insertions(+), 23 deletions(-)
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/FileDeletionBase.java
b/paimon-core/src/main/java/org/apache/paimon/operation/FileDeletionBase.java
index 724c0354cd..a0545c87e4 100644
---
a/paimon-core/src/main/java/org/apache/paimon/operation/FileDeletionBase.java
+++
b/paimon-core/src/main/java/org/apache/paimon/operation/FileDeletionBase.java
@@ -82,7 +82,6 @@ public abstract class FileDeletionBase<T extends Snapshot> {
protected final Map<BinaryRow, Set<Integer>> deletionBuckets;
private final Executor fileExecutor;
- private final int fileOperationParallelism;
@Nullable private final Integer manifestReadParallelism;
protected boolean changelogDecoupled;
@@ -112,10 +111,6 @@ public abstract class FileDeletionBase<T extends Snapshot>
{
this.cleanEmptyDirectories = cleanEmptyDirectories;
this.deletionBuckets = new ConcurrentHashMap<>();
this.fileExecutor =
FileOperationThreadPool.getExecutorService(fileOperationThreadNum);
- this.fileOperationParallelism =
- fileOperationThreadNum > 0
- ? fileOperationThreadNum
- : Runtime.getRuntime().availableProcessors();
this.manifestReadParallelism = manifestReadParallelism;
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/AbstractFileStoreTable.java
b/paimon-core/src/main/java/org/apache/paimon/table/AbstractFileStoreTable.java
index 88ee3f6a4b..97058f2b8e 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/AbstractFileStoreTable.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/AbstractFileStoreTable.java
@@ -450,7 +450,8 @@ abstract class AbstractFileStoreTable implements
FileStoreTable {
snapshotManager(),
changelogManager(),
store().newSnapshotDeletion(),
- store().newTagManager());
+ store().newTagManager(),
+ store().options().scanManifestParallelism());
}
@Override
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/ExpireSnapshotsImpl.java
b/paimon-core/src/main/java/org/apache/paimon/table/ExpireSnapshotsImpl.java
index 32124bfc2a..b7cabb4d32 100644
--- a/paimon-core/src/main/java/org/apache/paimon/table/ExpireSnapshotsImpl.java
+++ b/paimon-core/src/main/java/org/apache/paimon/table/ExpireSnapshotsImpl.java
@@ -35,6 +35,8 @@ import org.apache.paimon.utils.TagManager;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import javax.annotation.Nullable;
+
import java.io.FileNotFoundException;
import java.io.IOException;
import java.io.UncheckedIOException;
@@ -68,6 +70,7 @@ public class ExpireSnapshotsImpl implements ExpireSnapshots {
private final SnapshotDeletion snapshotDeletion;
private final Executor fileExecutor;
private final TagManager tagManager;
+ private final int snapshotExpireBatchSize;
private ExpireConfig expireConfig;
private Supplier<Long> currentTimeMillis = System::currentTimeMillis;
@@ -76,7 +79,8 @@ public class ExpireSnapshotsImpl implements ExpireSnapshots {
SnapshotManager snapshotManager,
ChangelogManager changelogManager,
SnapshotDeletion snapshotDeletion,
- TagManager tagManager) {
+ TagManager tagManager,
+ @Nullable Integer scanManifestParallelism) {
this.snapshotManager = snapshotManager;
this.changelogManager = changelogManager;
this.consumerManager =
@@ -88,6 +92,10 @@ public class ExpireSnapshotsImpl implements ExpireSnapshots {
this.tagManager = tagManager;
this.expireConfig = ExpireConfig.builder().build();
this.fileExecutor = snapshotDeletion.fileExecutor();
+ this.snapshotExpireBatchSize =
+ scanManifestParallelism == null || scanManifestParallelism <= 0
+ ? Runtime.getRuntime().availableProcessors()
+ : scanManifestParallelism;
}
@VisibleForTesting
@@ -207,12 +215,11 @@ public class ExpireSnapshotsImpl implements
ExpireSnapshots {
// delete merge tree files
// deleted merge tree files in a snapshot are not used by the next
snapshot, so the range of
// id should be (beginInclusiveId, endExclusiveId]
- snapshotDeletion.cleanDataFiles(
- collectDataFilesToDelete(snapshotsIncludingEnd,
taggedSnapshots, beginInclusiveId));
+ cleanDataFiles(snapshotsIncludingEnd, taggedSnapshots,
beginInclusiveId);
// delete changelog files
if (!expireConfig.isChangelogDecoupled()) {
-
snapshotDeletion.cleanDataFiles(collectChangelogFilesToDelete(snapshotsExcludingEnd));
+ cleanChangelogFiles(snapshotsExcludingEnd);
}
// data files and changelog files in bucket directories has been
deleted
@@ -262,7 +269,7 @@ public class ExpireSnapshotsImpl implements ExpireSnapshots
{
return snapshotsExcludingEnd.size();
}
- private Collection<Path> collectDataFilesToDelete(
+ private void cleanDataFiles(
List<Snapshot> snapshotsIncludingEnd,
List<Snapshot> taggedSnapshots,
long beginInclusiveId)
@@ -288,6 +295,7 @@ public class ExpireSnapshotsImpl implements ExpireSnapshots
{
collectTagSkippers(tags.values());
Predicate<ExpireFileEntry> deleteAll = entry -> false;
List<CompletableFuture<List<Path>>> futures = new ArrayList<>();
+ int plannedSnapshots = 0;
for (Snapshot snapshot : snapshotsIncludingEnd) {
long id = snapshot.id();
if (id == beginInclusiveId) {
@@ -308,15 +316,18 @@ public class ExpireSnapshotsImpl implements
ExpireSnapshots {
id);
continue;
}
-
futures.add(
CompletableFuture.supplyAsync(
() ->
snapshotDeletion.planDeletedInDeltaManifest(
snapshot, skipper.get()),
fileExecutor));
+ if (++plannedSnapshots >= snapshotExpireBatchSize) {
+ cleanBatch(futures);
+ plannedSnapshots = 0;
+ }
}
- return flatten(getAll(futures));
+ cleanBatch(futures);
}
private Map<Long, Optional<Predicate<ExpireFileEntry>>> collectTagSkippers(
@@ -350,9 +361,10 @@ public class ExpireSnapshotsImpl implements
ExpireSnapshots {
return skippers;
}
- private Collection<Path> collectChangelogFilesToDelete(List<Snapshot>
snapshots)
+ private void cleanChangelogFiles(List<Snapshot> snapshots)
throws ExecutionException, InterruptedException {
List<CompletableFuture<List<Path>>> futures = new ArrayList<>();
+ int plannedSnapshots = 0;
for (Snapshot snapshot : snapshots) {
if (LOG.isDebugEnabled()) {
LOG.debug("Ready to delete changelog files from snapshot #{}",
snapshot.id());
@@ -362,9 +374,20 @@ public class ExpireSnapshotsImpl implements
ExpireSnapshots {
CompletableFuture.supplyAsync(
() ->
snapshotDeletion.planAddedInChangelogManifest(snapshot),
fileExecutor));
+ if (++plannedSnapshots >= snapshotExpireBatchSize) {
+ cleanBatch(futures);
+ plannedSnapshots = 0;
+ }
}
}
- return flatten(getAll(futures));
+ cleanBatch(futures);
+ }
+
+ private void cleanBatch(List<CompletableFuture<List<Path>>> futures)
+ throws ExecutionException, InterruptedException {
+ Collection<Path> paths = flatten(getAll(futures));
+ futures.clear();
+ snapshotDeletion.cleanDataFiles(paths);
}
private Collection<Runnable> collectManifestDeletionTasks(
diff --git a/paimon-core/src/test/java/org/apache/paimon/TestFileStore.java
b/paimon-core/src/test/java/org/apache/paimon/TestFileStore.java
index b670ffa48f..178905917c 100644
--- a/paimon-core/src/test/java/org/apache/paimon/TestFileStore.java
+++ b/paimon-core/src/test/java/org/apache/paimon/TestFileStore.java
@@ -175,7 +175,8 @@ public class TestFileStore extends KeyValueFileStore {
snapshotManager(),
changelogManager(),
newSnapshotDeletion(),
- new TagManager(fileIO, options.path()))
+ new TagManager(fileIO, options.path()),
+ null)
.config(
ExpireConfig.builder()
.snapshotRetainMax(numRetainedMax)
@@ -189,7 +190,8 @@ public class TestFileStore extends KeyValueFileStore {
snapshotManager(),
changelogManager(),
newSnapshotDeletion(),
- new TagManager(fileIO, options.path()))
+ new TagManager(fileIO, options.path()),
+ null)
.config(expireConfig);
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/operation/ExpireSnapshotsTest.java
b/paimon-core/src/test/java/org/apache/paimon/operation/ExpireSnapshotsTest.java
index 946144af58..0d2538e826 100644
---
a/paimon-core/src/test/java/org/apache/paimon/operation/ExpireSnapshotsTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/operation/ExpireSnapshotsTest.java
@@ -635,7 +635,8 @@ public class ExpireSnapshotsTest {
blockingSnapshotManager,
changelogManager,
store.newSnapshotDeletion(),
- store.newTagManager());
+ store.newTagManager(),
+ store.options().scanManifestParallelism());
expire.expireUntil(1, latestSnapshotId);
@@ -961,7 +962,8 @@ public class ExpireSnapshotsTest {
failingSnapshotManager,
changelogManager,
store.newSnapshotDeletion(),
- store.newTagManager());
+ store.newTagManager(),
+ store.options().scanManifestParallelism());
expire.config(config);
expire.setCurrentTimeMillis(() -> 6000L);
@@ -1217,7 +1219,11 @@ public class ExpireSnapshotsTest {
SnapshotManager snapshotManager,
SnapshotDeletion snapshotDeletion) {
return new ExpireSnapshotsImpl(
- snapshotManager, store.changelogManager(), snapshotDeletion,
store.newTagManager());
+ snapshotManager,
+ store.changelogManager(),
+ snapshotDeletion,
+ store.newTagManager(),
+ store.options().scanManifestParallelism());
}
private void rewriteSnapshotTime(long snapshotId, long newTimeMillis)
throws IOException {
diff --git
a/paimon-core/src/test/java/org/apache/paimon/operation/FileDeletionTest.java
b/paimon-core/src/test/java/org/apache/paimon/operation/FileDeletionTest.java
index 49345b0e98..31ba29a2b9 100644
---
a/paimon-core/src/test/java/org/apache/paimon/operation/FileDeletionTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/operation/FileDeletionTest.java
@@ -682,7 +682,11 @@ public class FileDeletionTest {
// result: exist A & B (because of tag2)
ExpireSnapshots expireSnapshots =
new ExpireSnapshotsImpl(
- snapshotManager, changelogManager,
store.newSnapshotDeletion(), tagManager);
+ snapshotManager,
+ changelogManager,
+ store.newSnapshotDeletion(),
+ tagManager,
+ store.options().scanManifestParallelism());
expireSnapshots
.config(
ExpireConfig.builder()
@@ -755,7 +759,11 @@ public class FileDeletionTest {
ExpireSnapshots expireSnapshots =
new ExpireSnapshotsImpl(
- snapshotManager, changelogManager, snapshotDeletion,
tagManager);
+ snapshotManager,
+ changelogManager,
+ snapshotDeletion,
+ tagManager,
+ store.options().scanManifestParallelism());
snapshotDeletion.readMergedDataFilesThrowException = true;
expireSnapshots
.config(
@@ -820,7 +828,11 @@ public class FileDeletionTest {
store.options().scanManifestParallelism());
ExpireSnapshots expireSnapshots =
new ExpireSnapshotsImpl(
- snapshotManager, changelogManager, snapshotDeletion,
tagManager);
+ snapshotManager,
+ changelogManager,
+ snapshotDeletion,
+ tagManager,
+ store.options().scanManifestParallelism());
snapshotDeletion.manifestSkippingSetThrowException = true;
expireSnapshots
.config(