From 26390120836ca6358e1ce554d5bd6ac3d29f9f05 Mon Sep 17 00:00:00 2001
From: Fantasy-Jay <13631435453@163.com>
Date: Sun, 16 Aug 2026 14:04:10 +0800
Subject: [PATCH 1/2] [flink] Support table batching for orphan file cleanup
---
docs/content/flink/procedures.md | 8 +-
.../procedure/RemoveOrphanFilesProcedure.java | 15 ++-
.../flink/orphan/FlinkOrphanFilesClean.java | 97 ++++++++++++++-----
.../procedure/RemoveOrphanFilesProcedure.java | 12 ++-
.../RemoveOrphanFilesActionITCaseBase.java | 12 ++-
5 files changed, 111 insertions(+), 33 deletions(-)
diff --git a/docs/content/flink/procedures.md b/docs/content/flink/procedures.md
index 3757f0e66122..d5a70684a670 100644
--- a/docs/content/flink/procedures.md
+++ b/docs/content/flink/procedures.md
@@ -332,13 +332,13 @@ All available procedures are listed below.
remove_orphan_files |
-- Use named argument
- CALL [catalog.]sys.remove_orphan_files(`table` => 'identifier', older_than => 'olderThan', dry_run => 'dryRun', mode => 'mode')
+ CALL [catalog.]sys.remove_orphan_files(`table` => 'identifier', older_than => 'olderThan', dry_run => 'dryRun', mode => 'mode', table_batch_size => 'tableBatchSize')
-- Use indexed argument
CALL [catalog.]sys.remove_orphan_files('identifier')
CALL [catalog.]sys.remove_orphan_files('identifier', 'olderThan')
CALL [catalog.]sys.remove_orphan_files('identifier', 'olderThan', 'dryRun')
CALL [catalog.]sys.remove_orphan_files('identifier', 'olderThan', 'dryRun','parallelism')
- CALL [catalog.]sys.remove_orphan_files('identifier', 'olderThan', 'dryRun','parallelism','mode')
+ CALL [catalog.]sys.remove_orphan_files('identifier', 'olderThan', 'dryRun','parallelism','mode','tableBatchSize')
|
To remove the orphan data files and metadata files. Arguments:
@@ -349,12 +349,14 @@ All available procedures are listed below.
dryRun: when true, view only orphan files, don't actually remove files. Default is false.
parallelism: The maximum number of concurrent deleting files. By default is the number of processors available to the Java virtual machine.
mode: The mode of remove orphan clean procedure (local or distributed) . By default is distributed.
+ tableBatchSize: The maximum number of tables cleaned by each distributed Flink job. Default is 10.
|
CALL sys.remove_orphan_files(`table` => 'default.T', older_than => '2023-10-31 12:00:00')
CALL sys.remove_orphan_files(`table` => 'default.*', older_than => '2023-10-31 12:00:00')
CALL sys.remove_orphan_files(`table` => 'default.T', older_than => '2023-10-31 12:00:00', dry_run => true)
CALL sys.remove_orphan_files(`table` => 'default.T', older_than => '2023-10-31 12:00:00', dry_run => false, parallelism => '5')
- CALL sys.remove_orphan_files(`table` => 'default.T', older_than => '2023-10-31 12:00:00', dry_run => false, parallelism => '5', mode => 'local')
+ CALL sys.remove_orphan_files(`table` => 'default.T', older_than => '2023-10-31 12:00:00', dry_run => false, parallelism => '5', mode => 'local')
+ CALL sys.remove_orphan_files(`table` => 'default.*', older_than => '2023-10-31 12:00:00', table_batch_size => 5)
|
diff --git a/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/RemoveOrphanFilesProcedure.java b/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/RemoveOrphanFilesProcedure.java
index 70f797b03d1c..a9c503aeec45 100644
--- a/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/RemoveOrphanFilesProcedure.java
+++ b/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/RemoveOrphanFilesProcedure.java
@@ -80,6 +80,18 @@ public String[] call(
Integer parallelism,
String mode)
throws Exception {
+ return call(procedureContext, tableId, olderThan, dryRun, parallelism, mode, null);
+ }
+
+ public String[] call(
+ ProcedureContext procedureContext,
+ String tableId,
+ String olderThan,
+ boolean dryRun,
+ Integer parallelism,
+ String mode,
+ Integer tableBatchSize)
+ throws Exception {
Identifier identifier = Identifier.fromString(tableId);
String databaseName = identifier.getDatabaseName();
String tableName = identifier.getObjectName();
@@ -99,7 +111,8 @@ public String[] call(
dryRun,
parallelism,
databaseName,
- tableName);
+ tableName,
+ tableBatchSize);
break;
case "LOCAL":
cleanOrphanFilesResult =
diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/orphan/FlinkOrphanFilesClean.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/orphan/FlinkOrphanFilesClean.java
index a98282f45d02..b43147beade2 100644
--- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/orphan/FlinkOrphanFilesClean.java
+++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/orphan/FlinkOrphanFilesClean.java
@@ -70,9 +70,11 @@
import static org.apache.flink.util.Preconditions.checkState;
import static org.apache.paimon.utils.Preconditions.checkArgument;
-/** Flink {@link OrphanFilesClean}, it will submit a job for a table. */
+/** Flink {@link OrphanFilesClean}, it will submit jobs in table batches. */
public class FlinkOrphanFilesClean extends OrphanFilesClean {
+ public static final int DEFAULT_TABLE_BATCH_SIZE = 10;
+
protected static final Logger LOG = LoggerFactory.getLogger(FlinkOrphanFilesClean.class);
@Nullable protected final Integer parallelism;
@@ -376,40 +378,85 @@ public static CleanOrphanFilesResult executeDatabaseOrphanFiles(
String databaseName,
@Nullable String tableName)
throws Catalog.DatabaseNotExistException, Catalog.TableNotExistException {
+ return executeDatabaseOrphanFiles(
+ env,
+ catalog,
+ olderThanMillis,
+ dryRun,
+ parallelism,
+ databaseName,
+ tableName,
+ DEFAULT_TABLE_BATCH_SIZE);
+ }
+
+ public static CleanOrphanFilesResult executeDatabaseOrphanFiles(
+ StreamExecutionEnvironment env,
+ Catalog catalog,
+ long olderThanMillis,
+ boolean dryRun,
+ @Nullable Integer parallelism,
+ String databaseName,
+ @Nullable String tableName,
+ @Nullable Integer tableBatchSize)
+ throws Catalog.DatabaseNotExistException, Catalog.TableNotExistException {
List tableNames = Collections.singletonList(tableName);
if (tableName == null || "*".equals(tableName)) {
tableNames = catalog.listTables(databaseName);
}
- List> orphanFilesCleans =
- new ArrayList<>(tableNames.size());
- for (String t : tableNames) {
- Identifier identifier = new Identifier(databaseName, t);
- Table table = catalog.getTable(identifier);
- checkArgument(
- table instanceof FileStoreTable,
- "Only FileStoreTable supports remove-orphan-files action. The table type is '%s'.",
- table.getClass().getName());
-
- DataStream clean =
- new FlinkOrphanFilesClean(
- (FileStoreTable) table, olderThanMillis, dryRun, parallelism)
- .doOrphanClean(env);
- if (clean != null) {
- orphanFilesCleans.add(clean);
+ int batchSize = tableBatchSize == null ? DEFAULT_TABLE_BATCH_SIZE : tableBatchSize;
+ checkArgument(batchSize > 0, "Table batch size must be greater than 0.");
+
+ long deletedFilesCount = 0;
+ long deletedFilesLenInBytes = 0;
+ for (int start = 0; start < tableNames.size(); start += batchSize) {
+ int end = Math.min(start + batchSize, tableNames.size());
+ int batchNumber = start / batchSize + 1;
+ int tableCount = end - start;
+ long batchStart = System.currentTimeMillis();
+ LOG.info(
+ "Starting orphan files clean batch #{} with {} tables.",
+ batchNumber,
+ tableCount);
+
+ List> orphanFilesCleans =
+ new ArrayList<>(tableCount);
+ for (String t : tableNames.subList(start, end)) {
+ Identifier identifier = new Identifier(databaseName, t);
+ Table table = catalog.getTable(identifier);
+ checkArgument(
+ table instanceof FileStoreTable,
+ "Only FileStoreTable supports remove-orphan-files action. The table type is '%s'.",
+ table.getClass().getName());
+
+ DataStream clean =
+ new FlinkOrphanFilesClean(
+ (FileStoreTable) table,
+ olderThanMillis,
+ dryRun,
+ parallelism)
+ .doOrphanClean(env);
+ if (clean != null) {
+ orphanFilesCleans.add(clean);
+ }
}
- }
- DataStream result = null;
- for (DataStream clean : orphanFilesCleans) {
- if (result == null) {
- result = clean;
- } else {
- result = result.union(clean);
+ DataStream result = null;
+ for (DataStream clean : orphanFilesCleans) {
+ result = result == null ? clean : result.union(clean);
}
+
+ CleanOrphanFilesResult batchResult = sum(result);
+ deletedFilesCount += batchResult.getDeletedFileCount();
+ deletedFilesLenInBytes += batchResult.getDeletedFileTotalLenInBytes();
+ LOG.info(
+ "Finished orphan files clean batch #{} with {} tables in {} ms.",
+ batchNumber,
+ tableCount,
+ System.currentTimeMillis() - batchStart);
}
- return sum(result);
+ return new CleanOrphanFilesResult(deletedFilesCount, deletedFilesLenInBytes);
}
private static CleanOrphanFilesResult sum(DataStream deleted) {
diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/RemoveOrphanFilesProcedure.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/RemoveOrphanFilesProcedure.java
index c3983c8aa528..ff7735be7127 100644
--- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/RemoveOrphanFilesProcedure.java
+++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/RemoveOrphanFilesProcedure.java
@@ -59,7 +59,11 @@ public class RemoveOrphanFilesProcedure extends ProcedureBase {
isOptional = true),
@ArgumentHint(name = "dry_run", type = @DataTypeHint("BOOLEAN"), isOptional = true),
@ArgumentHint(name = "parallelism", type = @DataTypeHint("INT"), isOptional = true),
- @ArgumentHint(name = "mode", type = @DataTypeHint("STRING"), isOptional = true)
+ @ArgumentHint(name = "mode", type = @DataTypeHint("STRING"), isOptional = true),
+ @ArgumentHint(
+ name = "table_batch_size",
+ type = @DataTypeHint("INT"),
+ isOptional = true)
})
public String[] call(
ProcedureContext procedureContext,
@@ -67,7 +71,8 @@ public String[] call(
String olderThan,
Boolean dryRun,
Integer parallelism,
- String mode)
+ String mode,
+ Integer tableBatchSize)
throws Exception {
Identifier identifier = Identifier.fromString(tableId);
String databaseName = identifier.getDatabaseName();
@@ -87,7 +92,8 @@ public String[] call(
dryRun != null && dryRun,
parallelism,
databaseName,
- tableName);
+ tableName,
+ tableBatchSize);
break;
case "LOCAL":
cleanOrphanFilesResult =
diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/RemoveOrphanFilesActionITCaseBase.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/RemoveOrphanFilesActionITCaseBase.java
index 50fbd7dac10b..42d4cd5bcabf 100644
--- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/RemoveOrphanFilesActionITCaseBase.java
+++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/RemoveOrphanFilesActionITCaseBase.java
@@ -241,6 +241,17 @@ public void testRemoveDatabaseOrphanFilesITCase(boolean isNamedArgument) throws
ImmutableList actualDryRunDeleteFile = ImmutableList.copyOf(executeSQL(withDryRun));
assertThat(actualDryRunDeleteFile).containsOnly(Row.of("4"));
+ String withBatchSize =
+ String.format(
+ isNamedArgument
+ ? "CALL sys.remove_orphan_files(`table` => '%s.%s', older_than => '%s', dry_run => true, table_batch_size => 1)"
+ : "CALL sys.remove_orphan_files('%s.%s', '%s', true, 5, 'distributed', 1)",
+ database,
+ "*",
+ olderThan);
+ ImmutableList actualBatchDeleteFile = ImmutableList.copyOf(executeSQL(withBatchSize));
+ assertThat(actualBatchDeleteFile).containsOnly(Row.of("4"));
+
String withOlderThan =
String.format(
isNamedArgument
@@ -250,7 +261,6 @@ public void testRemoveDatabaseOrphanFilesITCase(boolean isNamedArgument) throws
"*",
olderThan);
ImmutableList actualDeleteFile = ImmutableList.copyOf(executeSQL(withOlderThan));
-
assertThat(actualDeleteFile).containsOnly(Row.of("4"));
}
From 026cf0551f97b0564e51283ff506d0cd759683a7 Mon Sep 17 00:00:00 2001
From: Fantasy-Jay <13631435453@163.com>
Date: Fri, 21 Aug 2026 01:14:53 +0800
Subject: [PATCH 2/2] [flink] Distinguish orphan file cleanup batch jobs
---
.../flink/orphan/FlinkOrphanFilesClean.java | 12 ++++++++----
.../RemoveOrphanFilesActionITCaseBase.java | 16 ++++++++++++++++
2 files changed, 24 insertions(+), 4 deletions(-)
diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/orphan/FlinkOrphanFilesClean.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/orphan/FlinkOrphanFilesClean.java
index 64442dd1142f..3bfaf03773b9 100644
--- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/orphan/FlinkOrphanFilesClean.java
+++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/orphan/FlinkOrphanFilesClean.java
@@ -506,7 +506,7 @@ public static CleanOrphanFilesResult executeDatabaseOrphanFiles(
result = result == null ? clean : result.union(clean);
}
- CleanOrphanFilesResult batchResult = sum(result);
+ CleanOrphanFilesResult batchResult = executeAndAggregateResults(result, batchNumber);
deletedFilesCount += batchResult.getDeletedFileCount();
deletedFilesLenInBytes += batchResult.getDeletedFileTotalLenInBytes();
LOG.info(
@@ -519,13 +519,17 @@ public static CleanOrphanFilesResult executeDatabaseOrphanFiles(
return new CleanOrphanFilesResult(deletedFilesCount, deletedFilesLenInBytes);
}
- private static CleanOrphanFilesResult sum(DataStream deleted) {
+ private static CleanOrphanFilesResult executeAndAggregateResults(
+ DataStream cleanResults, int batchNumber) {
long deletedFilesCount = 0;
long deletedFilesLenInBytes = 0;
- if (deleted != null) {
+ if (cleanResults != null) {
try {
CloseableIterator iterator =
- deleted.global().executeAndCollect("OrphanFilesClean");
+ cleanResults
+ .global()
+ .executeAndCollect(
+ String.format("OrphanFilesClean-Batch-%d", batchNumber));
while (iterator.hasNext()) {
CleanOrphanFilesResult cleanOrphanFilesResult = iterator.next();
deletedFilesCount += cleanOrphanFilesResult.getDeletedFileCount();
diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/RemoveOrphanFilesActionITCaseBase.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/RemoveOrphanFilesActionITCaseBase.java
index c6ecfca3dd42..b90c44af65a8 100644
--- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/RemoveOrphanFilesActionITCaseBase.java
+++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/RemoveOrphanFilesActionITCaseBase.java
@@ -40,6 +40,8 @@
import org.apache.paimon.shade.guava30.com.google.common.collect.ImmutableList;
+import org.apache.flink.api.common.JobStatus;
+import org.apache.flink.client.program.ClusterClient;
import org.apache.flink.types.Row;
import org.apache.flink.util.CloseableIterator;
import org.junit.jupiter.api.Test;
@@ -69,6 +71,15 @@ public abstract class RemoveOrphanFilesActionITCaseBase extends ActionITCaseBase
private static final String ORPHAN_FILE_1 = "bucket-0/orphan_file1";
private static final String ORPHAN_FILE_2 = "bucket-0/orphan_file2";
+ private long countFinishedOrphanFilesCleanJobs() throws Exception {
+ try (ClusterClient> client = MINI_CLUSTER_EXTENSION.createRestClusterClient()) {
+ return client.listJobs().get().stream()
+ .filter(job -> job.getJobName().startsWith("OrphanFilesClean-Batch-"))
+ .filter(job -> job.getJobState() == JobStatus.FINISHED)
+ .count();
+ }
+ }
+
private FileStoreTable createTableAndWriteData(String tableName) throws Exception {
RowType rowType =
RowType.of(
@@ -248,8 +259,10 @@ public void testRemoveDatabaseOrphanFilesITCase(boolean isNamedArgument) throws
database,
"*",
olderThan);
+ long defaultBatchJobCountBefore = countFinishedOrphanFilesCleanJobs();
ImmutableList actualDryRunDeleteFile = ImmutableList.copyOf(executeSQL(withDryRun));
assertThat(actualDryRunDeleteFile).containsOnly(Row.of("4"));
+ assertThat(countFinishedOrphanFilesCleanJobs() - defaultBatchJobCountBefore).isEqualTo(1);
String withBatchSize =
String.format(
@@ -259,8 +272,11 @@ public void testRemoveDatabaseOrphanFilesITCase(boolean isNamedArgument) throws
database,
"*",
olderThan);
+ long configuredBatchJobCountBefore = countFinishedOrphanFilesCleanJobs();
ImmutableList actualBatchDeleteFile = ImmutableList.copyOf(executeSQL(withBatchSize));
assertThat(actualBatchDeleteFile).containsOnly(Row.of("4"));
+ assertThat(countFinishedOrphanFilesCleanJobs() - configuredBatchJobCountBefore)
+ .isEqualTo(2);
String withOlderThan =
String.format(