multiPartUploadStore(FileIO fileIO) throws IOException {
if (fileIO instanceof RESTTokenFileIO) {
RESTTokenFileIO restTokenFileIO = (RESTTokenFileIO) fileIO;
diff --git a/paimon-common/src/main/java/org/apache/paimon/fs/TwoPhaseOutputStream.java b/paimon-common/src/main/java/org/apache/paimon/fs/TwoPhaseOutputStream.java
index 931969ec68cb..1fa0d7cfd371 100644
--- a/paimon-common/src/main/java/org/apache/paimon/fs/TwoPhaseOutputStream.java
+++ b/paimon-common/src/main/java/org/apache/paimon/fs/TwoPhaseOutputStream.java
@@ -53,6 +53,19 @@ public interface Committer extends Serializable {
*/
void discard(FileIO fileIO) throws IOException;
+ /**
+ * Discards staged resources without deleting {@link #targetPath()}.
+ *
+ * This is used when a commit may have taken effect and its target must therefore be
+ * preserved. The default delegates to {@link #clean}. Override this method if a failed or
+ * uncertain commit can leave staged resources that {@code clean} does not release.
+ *
+ * @throws IOException if an I/O error occurs during cleanup
+ */
+ default void discardStaging(FileIO fileIO) throws IOException {
+ clean(fileIO);
+ }
+
Path targetPath();
/**
diff --git a/paimon-core/src/main/java/org/apache/paimon/table/format/FormatBatchWriteBuilder.java b/paimon-core/src/main/java/org/apache/paimon/table/format/FormatBatchWriteBuilder.java
index 73c11d774166..f9d08a5c1bfb 100644
--- a/paimon-core/src/main/java/org/apache/paimon/table/format/FormatBatchWriteBuilder.java
+++ b/paimon-core/src/main/java/org/apache/paimon/table/format/FormatBatchWriteBuilder.java
@@ -78,6 +78,14 @@ public BatchTableCommit newCommit() {
CoreOptions options = new CoreOptions(table.options());
boolean formatTablePartitionOnlyValueInPath = options.formatTablePartitionOnlyValueInPath();
String syncHiveUri = options.formatTableCommitSyncPartitionHiveUri();
+ int cleanupThreadNum =
+ table.partitionManager() != null && !table.partitionKeys().isEmpty()
+ ? options.formatTableCommitCleanupThreadNum()
+ : 1;
+ int publishThreadNum =
+ table.partitionManager() != null && !table.partitionKeys().isEmpty()
+ ? options.formatTableCommitPublishThreadNum()
+ : 1;
return new FormatTableCommit(
table.location(),
table.partitionKeys(),
@@ -90,7 +98,9 @@ public BatchTableCommit newCommit() {
syncHiveUri,
table.catalogContext(),
table.partitionManager(),
- options.dynamicPartitionOverwrite());
+ options.dynamicPartitionOverwrite(),
+ cleanupThreadNum,
+ publishThreadNum);
}
@Override
diff --git a/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableCommit.java b/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableCommit.java
index be2105930a57..76de97dcfb39 100644
--- a/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableCommit.java
+++ b/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableCommit.java
@@ -38,6 +38,9 @@
import org.apache.paimon.table.sink.TableCommit;
import org.apache.paimon.utils.Pair;
import org.apache.paimon.utils.PartitionPathUtils;
+import org.apache.paimon.utils.ThreadPoolUtils;
+
+import org.apache.paimon.shade.guava30.com.google.common.collect.Iterators;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -46,24 +49,42 @@
import java.io.FileNotFoundException;
import java.io.IOException;
+import java.io.UncheckedIOException;
import java.lang.reflect.Method;
+import java.security.AccessControlContext;
+import java.security.AccessController;
+import java.security.PrivilegedAction;
+import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashSet;
+import java.util.Iterator;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
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.atomic.AtomicBoolean;
+import java.util.function.Consumer;
+import java.util.function.Function;
import java.util.stream.Collectors;
import static org.apache.paimon.table.format.FormatBatchWriteBuilder.validateStaticPartition;
+import static org.apache.paimon.utils.ExceptionUtils.firstOrSuppressed;
/** Commit for Format Table. */
public class FormatTableCommit implements BatchTableCommit {
private static final Logger LOG = LoggerFactory.getLogger(FormatTableCommit.class);
+ private static final int MAX_COMMIT_THREAD_NUM = 64;
+
+ private static final ExecutorService COMMIT_EXECUTOR =
+ ThreadPoolUtils.createCachedThreadPool(
+ MAX_COMMIT_THREAD_NUM, "FORMAT-TABLE-COMMIT-THREAD-POOL");
+
private String location;
private final boolean formatTablePartitionOnlyValueInPath;
private final String defaultPartName;
@@ -75,6 +96,8 @@ public class FormatTableCommit implements BatchTableCommit {
private Identifier tableIdentifier;
@Nullable private final FormatTablePartitionManager partitionManager;
private final boolean dynamicPartitionOverwrite;
+ private final int cleanupThreadNum;
+ private final int publishThreadNum;
public FormatTableCommit(
String location,
@@ -89,6 +112,50 @@ public FormatTableCommit(
CatalogContext catalogContext,
@Nullable FormatTablePartitionManager partitionManager,
boolean dynamicPartitionOverwrite) {
+ this(
+ location,
+ partitionKeys,
+ fileIO,
+ formatTablePartitionOnlyValueInPath,
+ defaultPartName,
+ overwrite,
+ tableIdentifier,
+ staticPartitions,
+ syncHiveUri,
+ catalogContext,
+ partitionManager,
+ dynamicPartitionOverwrite,
+ 1,
+ 1);
+ }
+
+ FormatTableCommit(
+ String location,
+ List partitionKeys,
+ FileIO fileIO,
+ boolean formatTablePartitionOnlyValueInPath,
+ String defaultPartName,
+ boolean overwrite,
+ Identifier tableIdentifier,
+ @Nullable Map staticPartitions,
+ @Nullable String syncHiveUri,
+ CatalogContext catalogContext,
+ @Nullable FormatTablePartitionManager partitionManager,
+ boolean dynamicPartitionOverwrite,
+ int cleanupThreadNum,
+ int publishThreadNum) {
+ if (cleanupThreadNum < 1 || cleanupThreadNum > MAX_COMMIT_THREAD_NUM) {
+ throw new IllegalArgumentException(
+ String.format(
+ "Format Table cleanup thread number must be between 1 and %s, but was %s.",
+ MAX_COMMIT_THREAD_NUM, cleanupThreadNum));
+ }
+ if (publishThreadNum < 1 || publishThreadNum > MAX_COMMIT_THREAD_NUM) {
+ throw new IllegalArgumentException(
+ String.format(
+ "Format Table publish thread number must be between 1 and %s, but was %s.",
+ MAX_COMMIT_THREAD_NUM, publishThreadNum));
+ }
this.location = location;
this.fileIO = fileIO;
this.formatTablePartitionOnlyValueInPath = formatTablePartitionOnlyValueInPath;
@@ -100,6 +167,8 @@ public FormatTableCommit(
this.tableIdentifier = tableIdentifier;
this.partitionManager = partitionManager;
this.dynamicPartitionOverwrite = dynamicPartitionOverwrite;
+ this.cleanupThreadNum = cleanupThreadNum;
+ this.publishThreadNum = publishThreadNum;
if (syncHiveUri != null) {
try {
Options options = new Options();
@@ -136,6 +205,7 @@ public void commit(List commitMessages) {
Set