This is an automated email from the ASF dual-hosted git repository.

JingsongLi 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 2d42fb8da0 [core] Support overwriting multiple specified partitions 
(#8789)
2d42fb8da0 is described below

commit 2d42fb8da07a71c3a7b01b517e7af1669bf2cb1e
Author: tsreaper <[email protected]>
AuthorDate: Thu Jul 23 18:36:38 2026 +0800

    [core] Support overwriting multiple specified partitions (#8789)
---
 .../apache/paimon/operation/FileStoreCommit.java   |  9 ++--
 .../paimon/operation/FileStoreCommitImpl.java      | 58 ++++++++++++++--------
 .../apache/paimon/table/sink/InnerTableCommit.java | 10 +++-
 .../apache/paimon/table/sink/TableCommitImpl.java  | 27 +++++++---
 .../test/java/org/apache/paimon/TestFileStore.java |  3 +-
 .../apache/paimon/operation/FileDeletionTest.java  |  6 +--
 .../apache/paimon/operation/TestCommitThread.java  |  3 +-
 .../apache/paimon/table/OverwriteTableTest.java    | 58 ++++++++++++++++++++++
 8 files changed, 132 insertions(+), 42 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommit.java 
b/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommit.java
index 42211d56ca..a9936833a9 100644
--- a/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommit.java
+++ b/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommit.java
@@ -19,6 +19,7 @@
 package org.apache.paimon.operation;
 
 import org.apache.paimon.Snapshot;
+import org.apache.paimon.data.BinaryRow;
 import org.apache.paimon.disk.IOManager;
 import org.apache.paimon.fs.FileIO;
 import org.apache.paimon.manifest.ManifestCommittable;
@@ -61,10 +62,10 @@ public interface FileStoreCommit extends AutoCloseable {
      *     note that this partition does not necessarily equal to the 
partitions of the newly added
      *     key-values. This is just the partition to be cleaned up.
      */
-    int overwritePartition(
-            Map<String, String> partition,
-            ManifestCommittable committable,
-            Map<String, String> properties);
+    int overwritePartition(Map<String, String> partition, ManifestCommittable 
committable);
+
+    /** Overwrite from manifest committable and specified partitions. */
+    int overwriteStaticPartitions(List<BinaryRow> partitions, 
ManifestCommittable committable);
 
     /**
      * Drop multiple partitions. The {@link Snapshot.CommitKind} of generated 
snapshot is {@link
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommitImpl.java
 
b/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommitImpl.java
index b440320212..1eacfa71f7 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommitImpl.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommitImpl.java
@@ -451,21 +451,44 @@ public class FileStoreCommitImpl implements 
FileStoreCommit {
     }
 
     @Override
-    public int overwritePartition(
-            Map<String, String> partition,
-            ManifestCommittable committable,
-            Map<String, String> properties) {
-        LOG.info(
-                "Ready to overwrite to table {}, number of commit messages: 
{}",
-                tableName,
-                committable.fileCommittables().size());
+    public int overwritePartition(Map<String, String> partition, 
ManifestCommittable committable) {
         if (LOG.isDebugEnabled()) {
             LOG.debug(
-                    "Ready to overwrite partition {}\nManifestCommittable: 
{}\nProperties: {}",
+                    "Ready to overwrite partition {}\nManifestCommittable: {}",
                     partition,
-                    committable,
-                    properties);
+                    committable);
         }
+        return overwritePartition(
+                () -> {
+                    Predicate partitionPredicate =
+                            createPartitionPredicate(
+                                    partition, partitionType, 
options.partitionDefaultName());
+                    return PartitionPredicate.fromPredicate(partitionType, 
partitionPredicate);
+                },
+                committable);
+    }
+
+    @Override
+    public int overwriteStaticPartitions(
+            List<BinaryRow> staticPartitions, ManifestCommittable committable) 
{
+        checkArgument(!staticPartitions.isEmpty(), "Partitions list cannot be 
empty.");
+        if (LOG.isDebugEnabled()) {
+            LOG.debug(
+                    "Ready to overwrite partitions {}\nManifestCommittable: 
{}",
+                    staticPartitions,
+                    committable);
+        }
+        return overwritePartition(
+                () -> PartitionPredicate.fromMultiple(partitionType, 
staticPartitions),
+                committable);
+    }
+
+    private int overwritePartition(
+            Supplier<PartitionPredicate> staticPartitionFilter, 
ManifestCommittable committable) {
+        LOG.info(
+                "Ready to overwrite to table {}, number of commit messages: 
{}",
+                tableName,
+                committable.fileCommittables().size());
 
         long started = System.nanoTime();
         int generatedSnapshot = 0;
@@ -490,7 +513,7 @@ public class FileStoreCommitImpl implements FileStoreCommit 
{
 
         try {
             boolean skipOverwrite = false;
-            // partition filter is built from static or dynamic partition 
according to properties
+            // partition filter is built from static or dynamic partitions
             PartitionPredicate partitionFilter = null;
             if (partitionType.getFieldCount() > 0 && 
options.dynamicPartitionOverwrite()) {
                 if (changes.appendTableFiles.isEmpty()) {
@@ -504,20 +527,13 @@ public class FileStoreCommitImpl implements 
FileStoreCommit {
                     partitionFilter = 
PartitionPredicate.fromMultiple(partitionType, partitions);
                 }
             } else {
-                // partition may be partial partition fields, so here must use 
predicate way.
-                Predicate partitionPredicate =
-                        createPartitionPredicate(
-                                partition, partitionType, 
options.partitionDefaultName());
-                partitionFilter =
-                        PartitionPredicate.fromPredicate(partitionType, 
partitionPredicate);
+                partitionFilter = staticPartitionFilter.get();
                 // sanity check, all changes must be done within the given 
partition
                 if (partitionFilter != null) {
                     for (ManifestEntry entry : changes.appendTableFiles) {
                         if (!partitionFilter.test(entry.partition())) {
                             throw new IllegalArgumentException(
-                                    "Trying to overwrite partition "
-                                            + partition
-                                            + ", but the changes in "
+                                    "The changes in "
                                             + 
pathFactory.getPartitionString(entry.partition())
                                             + " does not belong to this 
partition");
                         }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/InnerTableCommit.java 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/InnerTableCommit.java
index a73771f218..f5bb81d667 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/InnerTableCommit.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/InnerTableCommit.java
@@ -18,17 +18,25 @@
 
 package org.apache.paimon.table.sink;
 
+import org.apache.paimon.data.BinaryRow;
 import org.apache.paimon.metrics.MetricRegistry;
 
 import javax.annotation.Nullable;
 
+import java.util.List;
 import java.util.Map;
 
 /** Inner {@link TableCommit} contains overwrite setter. */
 public interface InnerTableCommit extends StreamTableCommit, BatchTableCommit {
 
     /** Overwrite writing, same as the 'INSERT OVERWRITE T PARTITION (...)' 
semantics of SQL. */
-    InnerTableCommit withOverwrite(@Nullable Map<String, String> 
staticPartition);
+    InnerTableCommit withOverwrite(@Nullable Map<String, String> spec);
+
+    /**
+     * Overwrite specified partitions. Unlike {@link 
InnerTableCommit#withOverwrite}, the given
+     * partitions must have all the partition keys.
+     */
+    InnerTableCommit withOverwriteStaticPartitions(List<BinaryRow> 
staticPartitions);
 
     /**
      * If this is set to true, when there is no new data, no snapshot will be 
generated. By default,
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/TableCommitImpl.java 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/TableCommitImpl.java
index f5b8c7892b..b66cca67e3 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/TableCommitImpl.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/TableCommitImpl.java
@@ -21,6 +21,7 @@ package org.apache.paimon.table.sink;
 import org.apache.paimon.Snapshot;
 import org.apache.paimon.annotation.VisibleForTesting;
 import org.apache.paimon.consumer.ConsumerManager;
+import org.apache.paimon.data.BinaryRow;
 import org.apache.paimon.disk.IOManager;
 import org.apache.paimon.fs.Path;
 import org.apache.paimon.index.IndexPathFactory;
@@ -55,7 +56,6 @@ import java.io.UncheckedIOException;
 import java.time.Duration;
 import java.time.LocalDateTime;
 import java.util.ArrayList;
-import java.util.Collections;
 import java.util.Comparator;
 import java.util.List;
 import java.util.Map;
@@ -91,7 +91,8 @@ public class TableCommitImpl implements InnerTableCommit {
     private final boolean forceCreatingSnapshot;
     private final ThreadPoolExecutor fileCheckExecutor;
 
-    @Nullable private Map<String, String> overwritePartition = null;
+    @Nullable private Map<String, String> overwritePartitionSpec = null;
+    @Nullable private List<BinaryRow> overwriteStaticPartitions = null;
     private boolean batchCommitted = false;
     private boolean expireForEmptyCommit = true;
 
@@ -135,7 +136,7 @@ public class TableCommitImpl implements InnerTableCommit {
         if (this.forceCreatingSnapshot) {
             return true;
         }
-        if (overwritePartition != null) {
+        if (overwritePartitionSpec != null || overwriteStaticPartitions != 
null) {
             return true;
         }
         return tagAutoManager != null
@@ -144,8 +145,16 @@ public class TableCommitImpl implements InnerTableCommit {
     }
 
     @Override
-    public TableCommitImpl withOverwrite(@Nullable Map<String, String> 
overwritePartitions) {
-        this.overwritePartition = overwritePartitions;
+    public TableCommitImpl withOverwrite(@Nullable Map<String, String> spec) {
+        this.overwritePartitionSpec = spec;
+        this.overwriteStaticPartitions = null;
+        return this;
+    }
+
+    @Override
+    public TableCommitImpl withOverwriteStaticPartitions(List<BinaryRow> 
overwritePartitions) {
+        this.overwritePartitionSpec = null;
+        this.overwriteStaticPartitions = overwritePartitions;
         return this;
     }
 
@@ -267,7 +276,7 @@ public class TableCommitImpl implements InnerTableCommit {
     }
 
     public void commitMultiple(List<ManifestCommittable> committables, boolean 
checkAppendFiles) {
-        if (overwritePartition == null) {
+        if (overwritePartitionSpec == null && overwriteStaticPartitions == 
null) {
             int newSnapshots = 0;
             for (ManifestCommittable committable : committables) {
                 newSnapshots += commit.commit(committable, checkAppendFiles);
@@ -293,8 +302,10 @@ public class TableCommitImpl implements InnerTableCommit {
                 committable = new ManifestCommittable(Long.MAX_VALUE);
             }
             int newSnapshots =
-                    commit.overwritePartition(
-                            overwritePartition, committable, 
Collections.emptyMap());
+                    overwriteStaticPartitions == null
+                            ? 
commit.overwritePartition(overwritePartitionSpec, committable)
+                            : commit.overwriteStaticPartitions(
+                                    overwriteStaticPartitions, committable);
             maintain(
                     committable.identifier(),
                     maintainExecutor,
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 09595d9188..6f4a9aa89d 100644
--- a/paimon-core/src/test/java/org/apache/paimon/TestFileStore.java
+++ b/paimon-core/src/test/java/org/apache/paimon/TestFileStore.java
@@ -248,8 +248,7 @@ public class TestFileStore extends KeyValueFileStore {
                 null,
                 null,
                 Collections.emptyList(),
-                (commit, committable) ->
-                        commit.overwritePartition(partition, committable, 
Collections.emptyMap()));
+                (commit, committable) -> commit.overwritePartition(partition, 
committable));
     }
 
     public Snapshot dropPartitions(List<Map<String, String>> partitions) {
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 5806626079..ee4eccd1d5 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
@@ -153,14 +153,12 @@ public class FileDeletionTest {
         FileStoreCommitImpl commit = store.newCommit();
         Map<String, String> partitionSpec = new HashMap<>();
         partitionSpec.put("dt", "0401");
-        commit.overwritePartition(
-                partitionSpec, new ManifestCommittable(commitIdentifier++), 
Collections.emptyMap());
+        commit.overwritePartition(partitionSpec, new 
ManifestCommittable(commitIdentifier++));
 
         // step 3: generate snapshot 3 by cleaning partition dt=0402/hr=10
         partitionSpec.put("dt", "0402");
         partitionSpec.put("hr", "8");
-        commit.overwritePartition(
-                partitionSpec, new ManifestCommittable(commitIdentifier++), 
Collections.emptyMap());
+        commit.overwritePartition(partitionSpec, new 
ManifestCommittable(commitIdentifier++));
         commit.close();
 
         // step 4: generate snapshot 4 by cleaning dt=0402/hr=12/bucket-0
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/operation/TestCommitThread.java 
b/paimon-core/src/test/java/org/apache/paimon/operation/TestCommitThread.java
index cdd2d8fba1..40b8c6f192 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/operation/TestCommitThread.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/operation/TestCommitThread.java
@@ -171,8 +171,7 @@ public class TestCommitThread extends Thread {
                 () ->
                         commit.overwritePartition(
                                 
TestKeyValueGenerator.toPartitionMap(partition, MULTI_PARTITIONED),
-                                committable,
-                                Collections.emptyMap()));
+                                committable));
     }
 
     private void doFinalCompact() {
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/OverwriteTableTest.java 
b/paimon-core/src/test/java/org/apache/paimon/table/OverwriteTableTest.java
index eccadd38b2..2d1271d165 100644
--- a/paimon-core/src/test/java/org/apache/paimon/table/OverwriteTableTest.java
+++ b/paimon-core/src/test/java/org/apache/paimon/table/OverwriteTableTest.java
@@ -20,17 +20,20 @@ package org.apache.paimon.table;
 
 import org.apache.paimon.CoreOptions;
 import org.apache.paimon.catalog.Identifier;
+import org.apache.paimon.data.BinaryRow;
 import org.apache.paimon.data.BinaryString;
 import org.apache.paimon.data.DataFormatTestUtil;
 import org.apache.paimon.data.GenericRow;
 import org.apache.paimon.data.InternalRow;
 import org.apache.paimon.schema.Schema;
+import org.apache.paimon.table.sink.CommitMessage;
 import org.apache.paimon.table.sink.InnerTableCommit;
 import org.apache.paimon.table.sink.StreamTableWrite;
 import org.apache.paimon.table.source.Split;
 import org.apache.paimon.table.source.TableRead;
 import org.apache.paimon.types.DataTypes;
 
+import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.Arguments;
 import org.junit.jupiter.params.provider.MethodSource;
@@ -40,6 +43,8 @@ import java.util.Arrays;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
+import java.util.function.Function;
+import java.util.stream.Collectors;
 
 import static org.apache.paimon.CoreOptions.BUCKET;
 import static org.apache.paimon.table.SimpleTableTestBase.getResult;
@@ -73,6 +78,59 @@ public class OverwriteTableTest extends TableTestBase {
                 true, dynamicPartitionOverwrite, overwriteData, 
overwritePartition, expected);
     }
 
+    @Test
+    public void testOverwriteMultiplePartitions() throws Exception {
+        Schema schema =
+                Schema.newBuilder()
+                        .column("pk", DataTypes.INT())
+                        .column("pt0", DataTypes.INT())
+                        .column("pt1", DataTypes.STRING())
+                        .column("v", DataTypes.STRING())
+                        .partitionKeys("pt0", "pt1")
+                        .option(CoreOptions.DYNAMIC_PARTITION_OVERWRITE.key(), 
"false")
+                        .build();
+        catalog.createTable(identifier(), schema, false);
+        FileStoreTable table = (FileStoreTable) catalog.getTable(identifier());
+        write(
+                table,
+                overwriteRow(1, 1, "A", "old-1A"),
+                overwriteRow(2, 1, "B", "old-1B"),
+                overwriteRow(3, 2, "A", "old-2A"));
+
+        Function<InternalRow, String> rowToString =
+                row -> DataFormatTestUtil.toStringNoRowKind(row, 
table.rowType());
+        try (StreamTableWrite write = 
table.newWrite(commitUser).withIgnorePreviousFiles(true);
+                InnerTableCommit commit = table.newCommit(commitUser)) {
+            write.write(overwriteRow(4, 1, "A", "list-1A"));
+            write.write(overwriteRow(5, 2, "A", "list-2A"));
+            List<CommitMessage> messages = write.prepareCommit(true, 1);
+            List<BinaryRow> overwritePartitions =
+                    messages.stream()
+                            .map(message -> message.partition().copy())
+                            .distinct()
+                            .collect(Collectors.toList());
+            
commit.withOverwriteStaticPartitions(overwritePartitions).commit(1, messages);
+            assertThat(read(table))
+                    .extracting(rowToString)
+                    .containsExactlyInAnyOrder(
+                            "2, 1, B, old-1B", "4, 1, A, list-1A", "5, 2, A, 
list-2A");
+
+            write.write(overwriteRow(6, 1, "A", "map-1A"));
+            messages = write.prepareCommit(true, 2);
+            commit.withOverwrite(Collections.singletonMap("pt0", 
"1")).commit(2, messages);
+            assertThat(read(table))
+                    .extracting(rowToString)
+                    .containsExactlyInAnyOrder("5, 2, A, list-2A", "6, 1, A, 
map-1A");
+
+            write.write(overwriteRow(7, 2, "A", "list-again-2A"));
+            messages = write.prepareCommit(true, 3);
+            
commit.withOverwriteStaticPartitions(overwritePartitions).commit(3, messages);
+            assertThat(read(table))
+                    .extracting(rowToString)
+                    .containsExactly("7, 2, A, list-again-2A");
+        }
+    }
+
     private void innerTestOverwrite(
             boolean withPrimaryKey,
             boolean dynamicPartitionOverwrite,

Reply via email to