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(