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(

Reply via email to