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 67c9f47725 [core] Derive row-id conflict checks from commit messages 
(#10079)
67c9f47725 is described below

commit 67c9f477254c77f7e523d1f0d83d0e6357940c0e
Author: Jingsong Lee <[email protected]>
AuthorDate: Tue Sep 22 13:52:56 2026 +0800

    [core] Derive row-id conflict checks from commit messages (#10079)
---
 .../apache/paimon/operation/FileStoreCommit.java   |   8 +--
 .../paimon/operation/FileStoreCommitImpl.java      |  59 ++++++++++++---
 .../paimon/table/sink/BatchWriteBuilderImpl.java   |  11 +--
 .../apache/paimon/table/sink/CommitMessage.java    |   8 +++
 .../paimon/table/sink/CommitMessageImpl.java       |  40 ++++++++++-
 .../paimon/table/sink/CommitMessageSerializer.java |  29 ++++++--
 .../apache/paimon/table/sink/InnerTableCommit.java |   5 +-
 .../apache/paimon/table/sink/TableCommitImpl.java  |  11 +--
 .../apache/paimon/append/VectorStoreTableTest.java |  12 +++-
 ...festCommittableSerializerCompatibilityTest.java |  30 ++++++--
 .../table/DataEvolutionDeletionVectorTest.java     |  61 +++++++++++++++-
 .../paimon/table/DataEvolutionTableTest.java       |  79 +++++++++++++++++++--
 .../table/sink/CommitMessageSerializerTest.java    |  47 ++++++++++++
 .../compatibility/manifest-committable-v14-v5      | Bin 0 -> 3155 bytes
 .../flink/action/DataEvolutionMergeIntoAction.java |   7 +-
 .../dataevolution/DataEvolutionDeleteOperator.java |  21 +++---
 .../dataevolution/DataEvolutionDeleteSink.java     |   3 +-
 .../DataEvolutionPartialWriteOperator.java         |  17 ++---
 .../dataevolution/MergeIntoUpdateChecker.java      |  17 +++--
 .../DataEvolutionCommitPreparationOperator.java    |  29 +++++++-
 ...DataEvolutionDeletionVectorMaterializeSink.java |   5 +-
 .../flink/sink/CommittableSerializerTest.java      |   4 +-
 .../MergeIntoPaimonDataEvolutionTable.scala        |   3 -
 .../procedure/DataEvolutionRewriteExecutor.java    |   3 +
 .../MaterializeDeletionVectorsProcedure.java       |  19 ++++-
 .../DataEvolutionRowIdConflictRewriter.scala       |   7 +-
 .../MergeIntoPaimonDataEvolutionTable.scala        |   3 -
 .../paimon/spark/commands/PaimonSparkWriter.scala  |   4 --
 .../spark/procedure/CompactProcedureTestBase.scala |  20 +++---
 .../paimon/spark/sql/RowTrackingTestBase.scala     |   4 --
 30 files changed, 443 insertions(+), 123 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 b039ffb9e9..050a345a0e 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
@@ -28,8 +28,6 @@ import org.apache.paimon.stats.Statistics;
 import org.apache.paimon.table.sink.CommitMessage;
 import org.apache.paimon.utils.FileStorePathFactory;
 
-import javax.annotation.Nullable;
-
 import java.util.List;
 import java.util.Map;
 
@@ -44,10 +42,8 @@ public interface FileStoreCommit extends AutoCloseable {
 
     FileStoreCommit appendCommitCheckConflict(boolean 
appendCommitCheckConflict);
 
-    FileStoreCommit rowIdCheckConflict(@Nullable Long rowIdCheckFromSnapshot);
-
-    FileStoreCommit rowIdCheckConflictForMaterializeDvCompaction(
-            @Nullable Long rowIdCheckFromSnapshot);
+    /** Use the materialize-DV row-id conflict strategy with snapshots from 
commit messages. */
+    FileStoreCommit materializeDvRowIdCheck();
 
     FileStoreCommit withOperation(Snapshot.Operation operation);
 
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 93092e20d7..926e972cc7 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
@@ -169,6 +169,7 @@ public class FileStoreCommitImpl implements FileStoreCommit 
{
     private boolean ignoreEmptyCommit;
     private CommitMetrics commitMetrics;
     private boolean appendCommitCheckConflict = false;
+    private boolean materializeDvRowIdCheck = false;
     private long lastCommittedSnapshotId = -1L;
     @Nullable private Snapshot.Operation operation;
     @Nullable private IOManager ioManager;
@@ -263,16 +264,8 @@ public class FileStoreCommitImpl implements 
FileStoreCommit {
     }
 
     @Override
-    public FileStoreCommit rowIdCheckConflict(@Nullable Long 
rowIdCheckFromSnapshot) {
-        
this.conflictDetection.setRowIdCheckFromSnapshot(rowIdCheckFromSnapshot);
-        return this;
-    }
-
-    @Override
-    public FileStoreCommit rowIdCheckConflictForMaterializeDvCompaction(
-            @Nullable Long rowIdCheckFromSnapshot) {
-        
this.conflictDetection.setRowIdCheckFromSnapshotForMaterializeDvCompaction(
-                rowIdCheckFromSnapshot);
+    public FileStoreCommit materializeDvRowIdCheck() {
+        materializeDvRowIdCheck = true;
         return this;
     }
 
@@ -339,6 +332,7 @@ public class FileStoreCommitImpl implements FileStoreCommit 
{
         int attempts = 0;
 
         List<CommitMessage> commitMessages = committable.fileCommittables();
+        configureRowIdCheckFromMessages(commitMessages);
         ManifestEntryChanges changes = collectChanges(commitMessages);
         Set<Pair<BinaryRow, Integer>> materializedBuckets = 
materializedBuckets(commitMessages);
         try {
@@ -505,6 +499,7 @@ public class FileStoreCommitImpl implements FileStoreCommit 
{
         int generatedSnapshot = 0;
         int attempts = 0;
 
+        configureRowIdCheckFromMessages(committable.fileCommittables());
         ManifestEntryChanges changes = 
collectChanges(committable.fileCommittables());
         if (!changes.appendChangelog.isEmpty() || 
!changes.compactChangelog.isEmpty()) {
             StringBuilder warnMessage =
@@ -764,6 +759,50 @@ public class FileStoreCommitImpl implements 
FileStoreCommit {
         return changes;
     }
 
+    private void configureRowIdCheckFromMessages(List<CommitMessage> 
commitMessages) {
+        Long checkFromSnapshot = null;
+        for (CommitMessage message : commitMessages) {
+            Long snapshotId = message.checkFromSnapshot();
+            if (snapshotId == null) {
+                continue;
+            }
+            checkArgument(snapshotId >= 0, "Invalid row-id check snapshot: 
%s", snapshotId);
+            checkArgument(
+                    checkFromSnapshot == null || 
checkFromSnapshot.equals(snapshotId),
+                    "Commit messages have different row-id check snapshots: %s 
and %s",
+                    checkFromSnapshot,
+                    snapshotId);
+            checkFromSnapshot = snapshotId;
+        }
+        if (checkFromSnapshot != null) {
+            for (CommitMessage message : commitMessages) {
+                if (message.checkFromSnapshot() != null) {
+                    continue;
+                }
+                checkArgument(
+                        !materializeDvRowIdCheck,
+                        "A materialize-DV commit message is missing its 
check-from snapshot.");
+                CommitMessageImpl commitMessage = (CommitMessageImpl) message;
+                checkArgument(
+                        commitMessage.newFilesIncrement().newFiles().stream()
+                                        .noneMatch(file -> file.firstRowId() 
!= null)
+                                && 
commitMessage.newFilesIncrement().deletedFiles().stream()
+                                        .noneMatch(file -> file.firstRowId() 
!= null),
+                        "A row-id commit message is missing its check-from 
snapshot.");
+            }
+        }
+        if (materializeDvRowIdCheck) {
+            checkArgument(
+                    checkFromSnapshot != null || commitMessages.isEmpty(),
+                    "A materialize-DV commit is missing its check-from 
snapshot.");
+            
conflictDetection.setRowIdCheckFromSnapshotForMaterializeDvCompaction(
+                    checkFromSnapshot);
+        } else {
+            // A committer can be reused; an untagged commit must not inherit 
a previous baseline.
+            conflictDetection.setRowIdCheckFromSnapshot(checkFromSnapshot);
+        }
+    }
+
     private Set<Pair<BinaryRow, Integer>> 
materializedBuckets(List<CommitMessage> commitMessages) {
         if (!options.dataEvolutionEnabled() || 
!options.deletionVectorsEnabled()) {
             return Collections.emptySet();
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/BatchWriteBuilderImpl.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/BatchWriteBuilderImpl.java
index d8c97405e2..445f48bafc 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/BatchWriteBuilderImpl.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/BatchWriteBuilderImpl.java
@@ -39,7 +39,6 @@ public class BatchWriteBuilderImpl implements 
BatchWriteBuilder {
     private final String commitUser;
 
     private Map<String, String> staticPartition;
-    private @Nullable Long rowIdCheckFromSnapshot = null;
 
     public BatchWriteBuilderImpl(InnerTable table) {
         this.table = table;
@@ -74,19 +73,11 @@ public class BatchWriteBuilderImpl implements 
BatchWriteBuilder {
 
     @Override
     public BatchTableCommit newCommit() {
-        InnerTableCommit commit =
-                table.newCommit(commitUser)
-                        .withOverwrite(staticPartition)
-                        .rowIdCheckConflict(rowIdCheckFromSnapshot);
+        InnerTableCommit commit = 
table.newCommit(commitUser).withOverwrite(staticPartition);
         commit.ignoreEmptyCommit(
                 Options.fromMap(table.options())
                         .getOptional(CoreOptions.SNAPSHOT_IGNORE_EMPTY_COMMIT)
                         .orElse(true));
         return commit;
     }
-
-    public BatchWriteBuilderImpl rowIdCheckConflict(@Nullable Long 
rowIdCheckFromSnapshot) {
-        this.rowIdCheckFromSnapshot = rowIdCheckFromSnapshot;
-        return this;
-    }
 }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessage.java 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessage.java
index 47f82bcfe8..3b71a9f9b6 100644
--- a/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessage.java
+++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessage.java
@@ -42,4 +42,12 @@ public interface CommitMessage extends Serializable {
     /** Total number of buckets in this partition. */
     @Nullable
     Integer totalBuckets();
+
+    /**
+     * Snapshot used to read rows before producing this message, if row-id 
conflicts need checking.
+     */
+    @Nullable
+    default Long checkFromSnapshot() {
+        return null;
+    }
 }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageImpl.java 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageImpl.java
index 8e715462f5..8527fdcc31 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageImpl.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageImpl.java
@@ -45,6 +45,7 @@ public class CommitMessageImpl implements CommitMessage {
     private transient BinaryRow partition;
     private transient int bucket;
     private transient @Nullable Integer totalBuckets;
+    private transient @Nullable Long checkFromSnapshot;
     private transient DataIncrement dataIncrement;
     private transient CompactIncrement compactIncrement;
 
@@ -54,11 +55,22 @@ public class CommitMessageImpl implements CommitMessage {
             @Nullable Integer totalBuckets,
             DataIncrement dataIncrement,
             CompactIncrement compactIncrement) {
+        this(partition, bucket, totalBuckets, dataIncrement, compactIncrement, 
null);
+    }
+
+    public CommitMessageImpl(
+            BinaryRow partition,
+            int bucket,
+            @Nullable Integer totalBuckets,
+            DataIncrement dataIncrement,
+            CompactIncrement compactIncrement,
+            @Nullable Long checkFromSnapshot) {
         this.partition = partition;
         this.bucket = bucket;
         this.totalBuckets = totalBuckets;
         this.dataIncrement = dataIncrement;
         this.compactIncrement = compactIncrement;
+        this.checkFromSnapshot = checkFromSnapshot;
     }
 
     @Override
@@ -76,6 +88,16 @@ public class CommitMessageImpl implements CommitMessage {
         return totalBuckets;
     }
 
+    @Override
+    public @Nullable Long checkFromSnapshot() {
+        return checkFromSnapshot;
+    }
+
+    public CommitMessageImpl withCheckFromSnapshot(long snapshotId) {
+        return new CommitMessageImpl(
+                partition, bucket, totalBuckets, dataIncrement, 
compactIncrement, snapshotId);
+    }
+
     public DataIncrement newFilesIncrement() {
         return dataIncrement;
     }
@@ -103,6 +125,7 @@ public class CommitMessageImpl implements CommitMessage {
         this.partition = message.partition;
         this.bucket = message.bucket;
         this.totalBuckets = message.totalBuckets;
+        this.checkFromSnapshot = message.checkFromSnapshot;
         this.dataIncrement = message.dataIncrement;
         this.compactIncrement = message.compactIncrement;
     }
@@ -120,13 +143,20 @@ public class CommitMessageImpl implements CommitMessage {
         return bucket == that.bucket
                 && Objects.equals(partition, that.partition)
                 && Objects.equals(totalBuckets, that.totalBuckets)
+                && Objects.equals(checkFromSnapshot, that.checkFromSnapshot)
                 && Objects.equals(dataIncrement, that.dataIncrement)
                 && Objects.equals(compactIncrement, that.compactIncrement);
     }
 
     @Override
     public int hashCode() {
-        return Objects.hash(partition, bucket, totalBuckets, dataIncrement, 
compactIncrement);
+        return Objects.hash(
+                partition,
+                bucket,
+                totalBuckets,
+                checkFromSnapshot,
+                dataIncrement,
+                compactIncrement);
     }
 
     @Override
@@ -136,8 +166,14 @@ public class CommitMessageImpl implements CommitMessage {
                         + "partition = %s, "
                         + "bucket = %d, "
                         + "totalBuckets = %s, "
+                        + "checkFromSnapshot = %s, "
                         + "newFilesIncrement = %s, "
                         + "compactIncrement = %s}",
-                partition, bucket, totalBuckets, dataIncrement, 
compactIncrement);
+                partition,
+                bucket,
+                totalBuckets,
+                checkFromSnapshot,
+                dataIncrement,
+                compactIncrement);
     }
 }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageSerializer.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageSerializer.java
index 8222b07c8c..2a3e3ebd81 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageSerializer.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/sink/CommitMessageSerializer.java
@@ -53,7 +53,7 @@ import static 
org.apache.paimon.utils.SerializationUtils.serializeBinaryRow;
 /** {@link VersionedSerializer} for {@link CommitMessage}. */
 public class CommitMessageSerializer implements 
VersionedSerializer<CommitMessage> {
 
-    public static final int CURRENT_VERSION = 13;
+    public static final int CURRENT_VERSION = 14;
 
     private final DataFileMetaSerializer dataFileSerializer;
     private final IndexFileMetaSerializer indexEntrySerializer;
@@ -121,6 +121,12 @@ public class CommitMessageSerializer implements 
VersionedSerializer<CommitMessag
         
dataFileSerializer.serializeList(message.compactIncrement().changelogFiles(), 
view);
         
indexEntrySerializer.serializeList(message.compactIncrement().newIndexFiles(), 
view);
         
indexEntrySerializer.serializeList(message.compactIncrement().deletedIndexFiles(),
 view);
+
+        Long checkFromSnapshot = message.checkFromSnapshot();
+        view.writeBoolean(checkFromSnapshot != null);
+        if (checkFromSnapshot != null) {
+            view.writeLong(checkFromSnapshot);
+        }
     }
 
     @Override
@@ -143,22 +149,31 @@ public class CommitMessageSerializer implements 
VersionedSerializer<CommitMessag
         IOExceptionSupplier<List<IndexFileMeta>> indexEntryDeserializer =
                 indexEntryDeserializer(version, view);
         if (version >= 10) {
-            return new CommitMessageImpl(
-                    deserializeBinaryRow(view),
-                    view.readInt(),
-                    view.readBoolean() ? view.readInt() : null,
+            BinaryRow partition = deserializeBinaryRow(view);
+            int bucket = view.readInt();
+            Integer totalBuckets = view.readBoolean() ? view.readInt() : null;
+            DataIncrement dataIncrement =
                     new DataIncrement(
                             fileDeserializer.get(),
                             fileDeserializer.get(),
                             fileDeserializer.get(),
                             indexEntryDeserializer.get(),
-                            indexEntryDeserializer.get()),
+                            indexEntryDeserializer.get());
+            CompactIncrement compactIncrement =
                     new CompactIncrement(
                             fileDeserializer.get(),
                             fileDeserializer.get(),
                             fileDeserializer.get(),
                             indexEntryDeserializer.get(),
-                            indexEntryDeserializer.get()));
+                            indexEntryDeserializer.get());
+            Long checkFromSnapshot = version >= 14 && view.readBoolean() ? 
view.readLong() : null;
+            return new CommitMessageImpl(
+                    partition,
+                    bucket,
+                    totalBuckets,
+                    dataIncrement,
+                    compactIncrement,
+                    checkFromSnapshot);
         } else {
             BinaryRow partition = deserializeBinaryRow(view);
             int bucket = view.readInt();
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 43f98d0e79..f5c51a7473 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
@@ -56,10 +56,7 @@ public interface InnerTableCommit extends StreamTableCommit, 
BatchTableCommit {
 
     InnerTableCommit appendCommitCheckConflict(boolean 
appendCommitCheckConflict);
 
-    InnerTableCommit rowIdCheckConflict(@Nullable Long rowIdCheckFromSnapshot);
-
-    InnerTableCommit rowIdCheckConflictForMaterializeDvCompaction(
-            @Nullable Long rowIdCheckFromSnapshot);
+    InnerTableCommit materializeDvRowIdCheck();
 
     @Override
     InnerTableCommit withMetricRegistry(MetricRegistry registry);
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 014b5e64da..23f1132fd7 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
@@ -176,15 +176,8 @@ public class TableCommitImpl implements InnerTableCommit {
     }
 
     @Override
-    public TableCommitImpl rowIdCheckConflict(@Nullable Long 
rowIdCheckFromSnapshot) {
-        commit.rowIdCheckConflict(rowIdCheckFromSnapshot);
-        return this;
-    }
-
-    @Override
-    public TableCommitImpl rowIdCheckConflictForMaterializeDvCompaction(
-            @Nullable Long rowIdCheckFromSnapshot) {
-        
commit.rowIdCheckConflictForMaterializeDvCompaction(rowIdCheckFromSnapshot);
+    public TableCommitImpl materializeDvRowIdCheck() {
+        commit.materializeDvRowIdCheck();
         return this;
     }
 
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/append/VectorStoreTableTest.java 
b/paimon-core/src/test/java/org/apache/paimon/append/VectorStoreTableTest.java
index 63d6618197..3ecd11bf33 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/append/VectorStoreTableTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/append/VectorStoreTableTest.java
@@ -44,8 +44,8 @@ import org.apache.paimon.table.Table;
 import org.apache.paimon.table.sink.BatchTableCommit;
 import org.apache.paimon.table.sink.BatchTableWrite;
 import org.apache.paimon.table.sink.BatchWriteBuilder;
-import org.apache.paimon.table.sink.BatchWriteBuilderImpl;
 import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.sink.CommitMessageImpl;
 import org.apache.paimon.table.sink.StreamTableWrite;
 import org.apache.paimon.table.sink.StreamWriteBuilder;
 import org.apache.paimon.table.source.DataSplit;
@@ -320,7 +320,7 @@ public class VectorStoreTableTest extends 
DataEvolutionTestBase {
             throws Exception {
         FileStoreTable table = getTableDefault();
         BatchWriteBuilder builder = table.newBatchWriteBuilder();
-        ((BatchWriteBuilderImpl) 
builder).rowIdCheckConflict(table.latestSnapshot().get().id());
+        long readSnapshotId = table.latestSnapshot().get().id();
         try (BatchTableWrite writer =
                         
builder.newWrite().withWriteType(table.rowType().project(columns));
                 BatchTableCommit commit = builder.newCommit()) {
@@ -329,7 +329,13 @@ public class VectorStoreTableTest extends 
DataEvolutionTestBase {
             }
             List<CommitMessage> messages = writer.prepareCommit();
             setFirstRowId(messages, firstRowId);
-            commit.commit(messages);
+            commit.commit(
+                    messages.stream()
+                            .map(
+                                    message ->
+                                            ((CommitMessageImpl) message)
+                                                    
.withCheckFromSnapshot(readSnapshotId))
+                            .collect(Collectors.toList()));
         }
     }
 
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java
index c25bee603e..dde62eb42b 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestCommittableSerializerCompatibilityTest.java
@@ -50,7 +50,7 @@ public class ManifestCommittableSerializerCompatibilityTest {
             "generateManifestCommittableGoldenFiles";
 
     @Test
-    public void testCompatibilityToV5CommitV13() throws IOException {
+    public void testCompatibilityToV5CommitV13AndV14() throws IOException {
         DataFileMeta dataFile =
                 DataFileMeta.create(
                                 "column-sequence-file",
@@ -79,16 +79,24 @@ public class ManifestCommittableSerializerCompatibilityTest 
{
         IndexFileMeta indexFile =
                 new IndexFileMeta(
                         "index-type", "index-file", 100L, 10L, 
(GlobalIndexMeta) null, null);
-        ManifestCommittable committable =
+        ManifestCommittable legacyCommittable =
                 createManifestCommittable(
                         Collections.singletonList(dataFile), indexFile, 
indexFile);
+        CommitMessageImpl legacyMessage =
+                (CommitMessageImpl) 
legacyCommittable.fileCommittables().get(0);
+        ManifestCommittable committable =
+                new ManifestCommittable(
+                        legacyCommittable.identifier(),
+                        legacyCommittable.watermark(),
+                        
Collections.singletonList(legacyMessage.withCheckFromSnapshot(3L)),
+                        legacyCommittable.properties());
 
         ManifestCommittableSerializer serializer = new 
ManifestCommittableSerializer();
         byte[] current = serializer.serialize(committable);
         byte[] serialized;
         if (Boolean.parseBoolean(
                 
System.getProperties().getProperty(GENERATE_GOLDEN_FILES_PROPERTY))) {
-            
CompatibilityUtils.writeCompatibilityFile("manifest-committable-v13-v5", 
current);
+            
CompatibilityUtils.writeCompatibilityFile("manifest-committable-v14-v5", 
current);
             serialized = current;
         } else {
             serialized =
@@ -96,12 +104,24 @@ public class 
ManifestCommittableSerializerCompatibilityTest {
                             
ManifestCommittableSerializerCompatibilityTest.class
                                     .getClassLoader()
                                     .getResourceAsStream(
-                                            
"compatibility/manifest-committable-v13-v5"),
+                                            
"compatibility/manifest-committable-v14-v5"),
                             true);
         }
 
         assertThat(current).isEqualTo(serialized);
-        assertThat(serializer.deserialize(5, 
serialized)).isEqualTo(committable);
+        ManifestCommittable restored = serializer.deserialize(5, serialized);
+        assertThat(restored).isEqualTo(committable);
+        
assertThat(restored.fileCommittables().get(0).checkFromSnapshot()).isEqualTo(3L);
+
+        byte[] legacySerialized =
+                IOUtils.readFully(
+                        ManifestCommittableSerializerCompatibilityTest.class
+                                .getClassLoader()
+                                
.getResourceAsStream("compatibility/manifest-committable-v13-v5"),
+                        true);
+        ManifestCommittable restoredLegacy = serializer.deserialize(5, 
legacySerialized);
+        assertThat(restoredLegacy).isEqualTo(legacyCommittable);
+        
assertThat(restoredLegacy.fileCommittables().get(0).checkFromSnapshot()).isNull();
     }
 
     @Test
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionDeletionVectorTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionDeletionVectorTest.java
index ef7eef3196..1ec1343ab1 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionDeletionVectorTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionDeletionVectorTest.java
@@ -941,6 +941,57 @@ public class DataEvolutionDeletionVectorTest extends 
DataEvolutionTestBase {
                 .containsExactlyElementsOf(concurrentValues);
     }
 
+    @Test
+    public void testMaterializeRejectsMissingCheckFromSnapshot() throws 
Exception {
+        createTableDefault();
+        FileStoreTable table = getTableDefault();
+        writeBaseRows(table);
+        commitDeletionVectors(table, DEFAULT_DV_SPECS);
+
+        Snapshot materializeSnapshot = table.latestSnapshot().get();
+        List<CommitMessage> messages =
+                prepareMaterializeDeletionVectors(table, materializeSnapshot, 
null);
+
+        assertThatThrownBy(
+                        () -> {
+                            try (TableCommitImpl commit =
+                                    table.newCommit("test-missing-snapshot")) {
+                                
commit.materializeDvRowIdCheck().commit(messages);
+                            }
+                        })
+                .isInstanceOf(IllegalArgumentException.class)
+                .hasMessageContaining("materialize-DV commit is missing its 
check-from snapshot");
+        
assertThat(table.latestSnapshot().get().id()).isEqualTo(materializeSnapshot.id());
+    }
+
+    @Test
+    public void testMaterializeRejectsMixedTaggedAndUntaggedMessages() throws 
Exception {
+        createTableDefault();
+        FileStoreTable table = getTableDefault();
+        writeBaseRows(table);
+        commitDeletionVectors(table, DEFAULT_DV_SPECS);
+
+        Snapshot materializeSnapshot = table.latestSnapshot().get();
+        CommitMessage message =
+                prepareMaterializeDeletionVectors(table, materializeSnapshot, 
null).get(0);
+        List<CommitMessage> messages =
+                Arrays.asList(
+                        ((CommitMessageImpl) message)
+                                
.withCheckFromSnapshot(materializeSnapshot.id()),
+                        message);
+
+        assertThatThrownBy(
+                        () -> {
+                            try (TableCommitImpl commit = 
table.newCommit("test-mixed-snapshot")) {
+                                
commit.materializeDvRowIdCheck().commit(messages);
+                            }
+                        })
+                .isInstanceOf(IllegalArgumentException.class)
+                .hasMessageContaining(
+                        "materialize-DV commit message is missing its 
check-from snapshot");
+        
assertThat(table.latestSnapshot().get().id()).isEqualTo(materializeSnapshot.id());
+    }
+
     @Test
     public void testStaleMaterializeAllowsNonOverlappingConcurrentUpdate() 
throws Exception {
         FileStoreTable table =
@@ -1434,8 +1485,14 @@ public class DataEvolutionDeletionVectorTest extends 
DataEvolutionTestBase {
             String commitUser)
             throws Exception {
         try (TableCommitImpl commit = table.newCommit(commitUser)) {
-            commit.rowIdCheckConflictForMaterializeDvCompaction(snapshot.id())
-                    .commit(commitMessages);
+            List<CommitMessage> checkedMessages =
+                    commitMessages.stream()
+                            .map(
+                                    message ->
+                                            ((CommitMessageImpl) message)
+                                                    
.withCheckFromSnapshot(snapshot.id()))
+                            .collect(Collectors.toList());
+            commit.materializeDvRowIdCheck().commit(checkedMessages);
         }
     }
 
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionTableTest.java 
b/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionTableTest.java
index c10e8ad575..aa95263808 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionTableTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionTableTest.java
@@ -45,7 +45,6 @@ import org.apache.paimon.schema.SchemaChange;
 import org.apache.paimon.table.sink.BatchTableCommit;
 import org.apache.paimon.table.sink.BatchTableWrite;
 import org.apache.paimon.table.sink.BatchWriteBuilder;
-import org.apache.paimon.table.sink.BatchWriteBuilderImpl;
 import org.apache.paimon.table.sink.CommitMessage;
 import org.apache.paimon.table.sink.CommitMessageImpl;
 import org.apache.paimon.table.source.DataSplit;
@@ -1242,18 +1241,24 @@ public class DataEvolutionTableTest extends 
DataEvolutionTestBase {
         long readSnapshotId = table.latestSnapshot().get().id();
 
         RowType writeType = 
table.rowType().project(Collections.singletonList("f2"));
-        BatchWriteBuilderImpl staleBuilder = (BatchWriteBuilderImpl) 
table.newBatchWriteBuilder();
+        BatchWriteBuilder staleBuilder = table.newBatchWriteBuilder();
         List<CommitMessage> staleMessages;
         try (BatchTableWrite write = 
staleBuilder.newWrite().withWriteType(writeType)) {
             write.write(GenericRow.of(BinaryString.fromString("stale-10")));
             write.write(GenericRow.of(BinaryString.fromString("stale-11")));
-            staleMessages = write.prepareCommit();
-            setFirstRowId(staleMessages, firstRowId);
+            List<CommitMessage> prepared = write.prepareCommit();
+            setFirstRowId(prepared, firstRowId);
+            staleMessages =
+                    prepared.stream()
+                            .map(
+                                    message ->
+                                            ((CommitMessageImpl) message)
+                                                    
.withCheckFromSnapshot(readSnapshotId))
+                            .collect(Collectors.toList());
         }
 
         updateF2(table, firstRowId, 100, 101);
         long concurrentSnapshotId = table.latestSnapshot().get().id();
-        staleBuilder.rowIdCheckConflict(readSnapshotId);
 
         assertThatThrownBy(
                         () -> {
@@ -1267,6 +1272,70 @@ public class DataEvolutionTableTest extends 
DataEvolutionTestBase {
         
assertThat(readF0AndF2(table)).isEqualTo(Arrays.asList("10|updated-100", 
"11|updated-101"));
     }
 
+    @Test
+    public void testRejectDifferentRowIdCheckSnapshotsInOneCommit() throws 
Exception {
+        createTableDefault();
+        FileStoreTable table = getTableDefault();
+        long firstRowId = writeFullRows(table, 10);
+        long readSnapshotId = table.latestSnapshot().get().id();
+
+        BatchWriteBuilder builder = table.newBatchWriteBuilder();
+        List<CommitMessage> messages;
+        try (BatchTableWrite write =
+                builder.newWrite()
+                        
.withWriteType(table.rowType().project(Collections.singletonList("f2")))) {
+            write.write(GenericRow.of(BinaryString.fromString("updated")));
+            messages = write.prepareCommit();
+            setFirstRowId(messages, firstRowId);
+        }
+
+        CommitMessageImpl message = (CommitMessageImpl) messages.get(0);
+        long snapshotBeforeCommit = table.latestSnapshot().get().id();
+        try (BatchTableCommit commit = builder.newCommit()) {
+            assertThatThrownBy(
+                            () ->
+                                    commit.commit(
+                                            Arrays.asList(
+                                                    
message.withCheckFromSnapshot(readSnapshotId),
+                                                    
message.withCheckFromSnapshot(
+                                                            readSnapshotId + 
1))))
+                    .isInstanceOf(IllegalArgumentException.class)
+                    .hasMessageContaining("different row-id check snapshots");
+        }
+        
assertThat(table.latestSnapshot().get().id()).isEqualTo(snapshotBeforeCommit);
+    }
+
+    @Test
+    public void testRejectMissingRowIdCheckSnapshotInMixedCommit() throws 
Exception {
+        createTableDefault();
+        FileStoreTable table = getTableDefault();
+        long firstRowId = writeFullRows(table, 10);
+        long readSnapshotId = table.latestSnapshot().get().id();
+
+        BatchWriteBuilder builder = table.newBatchWriteBuilder();
+        List<CommitMessage> messages;
+        try (BatchTableWrite write =
+                builder.newWrite()
+                        
.withWriteType(table.rowType().project(Collections.singletonList("f2")))) {
+            write.write(GenericRow.of(BinaryString.fromString("updated")));
+            messages = write.prepareCommit();
+            setFirstRowId(messages, firstRowId);
+        }
+
+        CommitMessageImpl message = (CommitMessageImpl) messages.get(0);
+        try (BatchTableCommit commit = builder.newCommit()) {
+            assertThatThrownBy(
+                            () ->
+                                    commit.commit(
+                                            Arrays.asList(
+                                                    
message.withCheckFromSnapshot(readSnapshotId),
+                                                    message)))
+                    .isInstanceOf(IllegalArgumentException.class)
+                    .hasMessageContaining("missing its check-from snapshot");
+        }
+        
assertThat(table.latestSnapshot().get().id()).isEqualTo(readSnapshotId);
+    }
+
     @Test
     public void 
testCompactPreservesConcurrentPartialUpdateWithinCandidateRange() throws 
Exception {
         createTableDefault();
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/sink/CommitMessageSerializerTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/table/sink/CommitMessageSerializerTest.java
index 04f40f71e5..418ea005d3 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/table/sink/CommitMessageSerializerTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/sink/CommitMessageSerializerTest.java
@@ -23,7 +23,11 @@ import org.apache.paimon.io.DataIncrement;
 
 import org.junit.jupiter.api.Test;
 
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
+import java.io.ObjectInputStream;
+import java.io.ObjectOutputStream;
 import java.util.Arrays;
 
 import static 
org.apache.paimon.index.IndexFileMetaSerializerTest.randomIndexFile;
@@ -74,5 +78,48 @@ public class CommitMessageSerializerTest {
         
assertThat(newCommittable.totalBuckets()).isEqualTo(committable.totalBuckets());
         
assertThat(newCommittable.compactIncrement()).isEqualTo(committable.compactIncrement());
         
assertThat(newCommittable.newFilesIncrement()).isEqualTo(committable.newFilesIncrement());
+        assertThat(newCommittable.checkFromSnapshot()).isNull();
+
+        CommitMessageImpl checked = committable.withCheckFromSnapshot(42L);
+        CommitMessageImpl checkedRoundTrip =
+                (CommitMessageImpl)
+                        serializer.deserialize(
+                                serializer.getVersion(), 
serializer.serialize(checked));
+        assertThat(checkedRoundTrip).isEqualTo(checked);
+        assertThat(checkedRoundTrip.checkFromSnapshot()).isEqualTo(42L);
+
+        byte[] serializedWithoutSnapshot = serializer.serialize(committable);
+        CommitMessageImpl oldVersion =
+                (CommitMessageImpl)
+                        serializer.deserialize(
+                                13,
+                                Arrays.copyOf(
+                                        serializedWithoutSnapshot,
+                                        serializedWithoutSnapshot.length - 1));
+        assertThat(oldVersion.checkFromSnapshot()).isNull();
+        
assertThat(oldVersion.newFilesIncrement()).isEqualTo(committable.newFilesIncrement());
+    }
+
+    @Test
+    public void testJavaSerializationPreservesCheckFromSnapshot() throws 
Exception {
+        CommitMessageImpl message =
+                new CommitMessageImpl(
+                                row(0),
+                                1,
+                                null,
+                                randomNewFilesIncrement(),
+                                CompactIncrement.emptyIncrement())
+                        .withCheckFromSnapshot(42L);
+
+        ByteArrayOutputStream bytes = new ByteArrayOutputStream();
+        try (ObjectOutputStream output = new ObjectOutputStream(bytes)) {
+            output.writeObject(message);
+        }
+        try (ObjectInputStream input =
+                new ObjectInputStream(new 
ByteArrayInputStream(bytes.toByteArray()))) {
+            CommitMessageImpl restored = (CommitMessageImpl) 
input.readObject();
+            assertThat(restored).isEqualTo(message);
+            assertThat(restored.checkFromSnapshot()).isEqualTo(42L);
+        }
     }
 }
diff --git 
a/paimon-core/src/test/resources/compatibility/manifest-committable-v14-v5 
b/paimon-core/src/test/resources/compatibility/manifest-committable-v14-v5
new file mode 100644
index 0000000000..a13b2909e6
Binary files /dev/null and 
b/paimon-core/src/test/resources/compatibility/manifest-committable-v14-v5 
differ
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/DataEvolutionMergeIntoAction.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/DataEvolutionMergeIntoAction.java
index 5bda307641..da95b4f36e 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/DataEvolutionMergeIntoAction.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/DataEvolutionMergeIntoAction.java
@@ -717,7 +717,8 @@ public class DataEvolutionMergeIntoAction extends 
TableActionBase {
                 written.transform(
                                 "Updated Column Check",
                                 new CommittableTypeInfo(),
-                                new MergeIntoUpdateChecker(storeTable, 
updatedColumns))
+                                new MergeIntoUpdateChecker(
+                                        storeTable, updatedColumns, 
baseSnapshotId))
                         .setParallelism(1)
                         .setMaxParallelism(1);
 
@@ -729,9 +730,7 @@ public class DataEvolutionMergeIntoAction extends 
TableActionBase {
                         context ->
                                 new StoreCommitter(
                                         storeTable,
-                                        storeTable
-                                                
.newCommit(context.commitUser())
-                                                
.rowIdCheckConflict(baseSnapshotId),
+                                        
storeTable.newCommit(context.commitUser()),
                                         context),
                         new NoopCommittableStateManager());
 
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/DataEvolutionDeleteOperator.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/DataEvolutionDeleteOperator.java
index 5cf1c94dff..d030b3b4b9 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/DataEvolutionDeleteOperator.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/DataEvolutionDeleteOperator.java
@@ -256,16 +256,17 @@ public class DataEvolutionDeleteOperator
 
             CommitMessage commitMessage =
                     new CommitMessageImpl(
-                            maintainer.getPartition(),
-                            UNAWARE_BUCKET,
-                            null,
-                            new DataIncrement(
-                                    Collections.emptyList(),
-                                    Collections.emptyList(),
-                                    Collections.emptyList(),
-                                    addedIndexFiles,
-                                    deletedIndexFiles),
-                            CompactIncrement.emptyIncrement());
+                                    maintainer.getPartition(),
+                                    UNAWARE_BUCKET,
+                                    null,
+                                    new DataIncrement(
+                                            Collections.emptyList(),
+                                            Collections.emptyList(),
+                                            Collections.emptyList(),
+                                            addedIndexFiles,
+                                            deletedIndexFiles),
+                                    CompactIncrement.emptyIncrement())
+                            .withCheckFromSnapshot(baseSnapshotId);
             output.collect(
                     new StreamRecord<>(
                             new 
Committable(BatchWriteBuilder.COMMIT_IDENTIFIER, commitMessage)));
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/DataEvolutionDeleteSink.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/DataEvolutionDeleteSink.java
index eddcf3bf91..4b86cfce08 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/DataEvolutionDeleteSink.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/DataEvolutionDeleteSink.java
@@ -129,8 +129,7 @@ public class DataEvolutionDeleteSink implements 
Serializable {
                                 new StoreCommitter(
                                         table,
                                         table.newCommit(context.commitUser())
-                                                
.withOperation(Snapshot.Operation.DELETE)
-                                                
.rowIdCheckConflict(baseSnapshotId),
+                                                
.withOperation(Snapshot.Operation.DELETE),
                                         context),
                         new NoopCommittableStateManager());
 
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/DataEvolutionPartialWriteOperator.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/DataEvolutionPartialWriteOperator.java
index d7b2a2eabb..7b1e79604f 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/DataEvolutionPartialWriteOperator.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/DataEvolutionPartialWriteOperator.java
@@ -328,14 +328,15 @@ public class DataEvolutionPartialWriteOperator
 
                 CommitMessage commitMessage =
                         new CommitMessageImpl(
-                                partition,
-                                0,
-                                null,
-                                new DataIncrement(
-                                        Collections.singletonList(fileMeta),
-                                        Collections.emptyList(),
-                                        Collections.emptyList()),
-                                CompactIncrement.emptyIncrement());
+                                        partition,
+                                        0,
+                                        null,
+                                        new DataIncrement(
+                                                
Collections.singletonList(fileMeta),
+                                                Collections.emptyList(),
+                                                Collections.emptyList()),
+                                        CompactIncrement.emptyIncrement())
+                                .withCheckFromSnapshot(baseSnapshotId);
 
                 return new Committable(Long.MAX_VALUE, commitMessage);
             } finally {
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/MergeIntoUpdateChecker.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/MergeIntoUpdateChecker.java
index 0f749c18bd..39518e6147 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/MergeIntoUpdateChecker.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/dataevolution/MergeIntoUpdateChecker.java
@@ -57,12 +57,15 @@ public class MergeIntoUpdateChecker extends 
BoundedOneInputOperator<Committable,
 
     private final FileStoreTable table;
     private final Set<String> updatedColumns;
+    private final long baseSnapshotId;
 
     private transient Set<BinaryRow> affectedPartitions;
 
-    public MergeIntoUpdateChecker(FileStoreTable table, Set<String> 
updatedColumns) {
+    public MergeIntoUpdateChecker(
+            FileStoreTable table, Set<String> updatedColumns, long 
baseSnapshotId) {
         this.table = table;
         this.updatedColumns = updatedColumns;
+        this.baseSnapshotId = baseSnapshotId;
     }
 
     @Override
@@ -149,11 +152,13 @@ public class MergeIntoUpdateChecker extends 
BoundedOneInputOperator<Committable,
 
                         CommitMessage commitMessage =
                                 new CommitMessageImpl(
-                                        entry.getKey(),
-                                        0,
-                                        null,
-                                        
DataIncrement.deleteIndexIncrement(entry.getValue()),
-                                        CompactIncrement.emptyIncrement());
+                                                entry.getKey(),
+                                                0,
+                                                null,
+                                                
DataIncrement.deleteIndexIncrement(
+                                                        entry.getValue()),
+                                                
CompactIncrement.emptyIncrement())
+                                        .withCheckFromSnapshot(baseSnapshotId);
 
                         Committable committable = new 
Committable(Long.MAX_VALUE, commitMessage);
 
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/DataEvolutionCommitPreparationOperator.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/DataEvolutionCommitPreparationOperator.java
index 787fd9c2c5..f2c0aa7bff 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/DataEvolutionCommitPreparationOperator.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/DataEvolutionCommitPreparationOperator.java
@@ -23,6 +23,7 @@ import 
org.apache.paimon.append.dataevolution.DataEvolutionCompactionCommitPrepa
 import org.apache.paimon.options.Options;
 import org.apache.paimon.table.FileStoreTable;
 import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.sink.CommitMessageImpl;
 
 import org.apache.flink.streaming.api.operators.StreamOperator;
 import org.apache.flink.streaming.api.operators.StreamOperatorFactory;
@@ -39,15 +40,18 @@ public class DataEvolutionCommitPreparationOperator
 
     private final FileStoreTable table;
     private final Snapshot snapshot;
+    private final boolean materializeDvRowIdCheck;
     private final List<Committable> committables;
 
     private DataEvolutionCommitPreparationOperator(
             StreamOperatorParameters<Committable> parameters,
             FileStoreTable table,
-            Snapshot snapshot) {
+            Snapshot snapshot,
+            boolean materializeDvRowIdCheck) {
         super(parameters, Options.fromMap(table.options()));
         this.table = table;
         this.snapshot = snapshot;
+        this.materializeDvRowIdCheck = materializeDvRowIdCheck;
         this.committables = new ArrayList<>();
     }
 
@@ -73,7 +77,18 @@ public class DataEvolutionCommitPreparationOperator
                 new DataEvolutionCompactionCommitPreparation(table, 
snapshot).prepare(messages)) {
             toCommit.add(new Committable(toCommit.get(0).checkpointId(), 
message));
         }
-        return toCommit;
+        if (!materializeDvRowIdCheck) {
+            return toCommit;
+        }
+        List<Committable> checked = new ArrayList<>(toCommit.size());
+        for (Committable committable : toCommit) {
+            checked.add(
+                    new Committable(
+                            committable.checkpointId(),
+                            ((CommitMessageImpl) committable.commitMessage())
+                                    .withCheckFromSnapshot(snapshot.id())));
+        }
+        return checked;
     }
 
     /** {@link StreamOperatorFactory} of {@link 
DataEvolutionCommitPreparationOperator}. */
@@ -81,18 +96,26 @@ public class DataEvolutionCommitPreparationOperator
 
         private final FileStoreTable table;
         private final Snapshot snapshot;
+        private final boolean materializeDvRowIdCheck;
 
         public Factory(FileStoreTable table, Snapshot snapshot) {
+            this(table, snapshot, false);
+        }
+
+        public Factory(FileStoreTable table, Snapshot snapshot, boolean 
materializeDvRowIdCheck) {
             super(Options.fromMap(table.options()));
             this.table = table;
             this.snapshot = snapshot;
+            this.materializeDvRowIdCheck = materializeDvRowIdCheck;
         }
 
         @Override
         @SuppressWarnings("unchecked")
         public <T extends StreamOperator<Committable>> T createStreamOperator(
                 StreamOperatorParameters<Committable> parameters) {
-            return (T) new DataEvolutionCommitPreparationOperator(parameters, 
table, snapshot);
+            return (T)
+                    new DataEvolutionCommitPreparationOperator(
+                            parameters, table, snapshot, 
materializeDvRowIdCheck);
         }
 
         @Override
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/DataEvolutionDeletionVectorMaterializeSink.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/DataEvolutionDeletionVectorMaterializeSink.java
index 226300a3da..3ca51a69d5 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/DataEvolutionDeletionVectorMaterializeSink.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/DataEvolutionDeletionVectorMaterializeSink.java
@@ -58,7 +58,8 @@ public class DataEvolutionDeletionVectorMaterializeSink
                                 "Data Evolution Deletion Vector Materialize 
Commit Preparation : "
                                         + table.name(),
                                 new CommittableTypeInfo(),
-                                new 
DataEvolutionCommitPreparationOperator.Factory(table, snapshot))
+                                new 
DataEvolutionCommitPreparationOperator.Factory(
+                                        table, snapshot, true))
                         .forceNonParallel();
         return doCommit(written, initialCommitUser);
     }
@@ -73,7 +74,7 @@ public class DataEvolutionDeletionVectorMaterializeSink
     protected Committer.Factory<Committable, ManifestCommittable> 
createCommitterFactory() {
         return context -> {
             TableCommitImpl commit = table.newCommit(context.commitUser());
-            commit.rowIdCheckConflictForMaterializeDvCompaction(snapshot.id());
+            commit.materializeDvRowIdCheck();
             return new StoreCommitter(table, commit, context);
         };
     }
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/CommittableSerializerTest.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/CommittableSerializerTest.java
index f3cc99292c..ced80333e7 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/CommittableSerializerTest.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/CommittableSerializerTest.java
@@ -45,7 +45,8 @@ public class CommittableSerializerTest {
         DataIncrement dataIncrement = randomNewFilesIncrement();
         CompactIncrement compactIncrement = randomCompactIncrement();
         CommitMessage committable =
-                new CommitMessageImpl(row(0), 1, 2, dataIncrement, 
compactIncrement);
+                new CommitMessageImpl(row(0), 1, 2, dataIncrement, 
compactIncrement)
+                        .withCheckFromSnapshot(42L);
         CommitMessage newCommittable =
                 serializer
                         .deserialize(
@@ -53,5 +54,6 @@ public class CommittableSerializerTest {
                                 serializer.serialize(new Committable(9, 
committable)))
                         .commitMessage();
         assertThat(newCommittable).isEqualTo(committable);
+        assertThat(newCommittable.checkFromSnapshot()).isEqualTo(42L);
     }
 }
diff --git 
a/paimon-spark/paimon-spark-4.0/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
 
b/paimon-spark/paimon-spark-4.0/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
index 15a8f89429..fc887a135e 100644
--- 
a/paimon-spark/paimon-spark-4.0/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
+++ 
b/paimon-spark/paimon-spark-4.0/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
@@ -423,9 +423,6 @@ case class MergeIntoPaimonDataEvolutionTable(
           insertActionInvoke(sparkSession, touchedFileTargetRelation, 
persistSourceDss)
         else Nil
 
-      if (readSnapshot != null) {
-        writer.rowIdCheckConflict(readSnapshot.id())
-      }
       DataEvolutionRowIdConflictCommitter.commit(
         sparkSession,
         table,
diff --git 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/DataEvolutionRewriteExecutor.java
 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/DataEvolutionRewriteExecutor.java
index e964f7b61c..a6610bfff8 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/DataEvolutionRewriteExecutor.java
+++ 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/DataEvolutionRewriteExecutor.java
@@ -286,6 +286,7 @@ final class DataEvolutionRewriteExecutor {
             abortMessages.addAll(preparationArtifacts);
             try (TableCommitImpl commit = table.newCommit(commitUser)) {
                 commitConfigurer.configure(commit);
+                commitConfigurer.prepareMessages(attemptSnapshot, 
preparedMessages);
                 try {
                     commit.commit(preparedMessages);
                 } catch (RuntimeException conflict) {
@@ -398,6 +399,8 @@ final class DataEvolutionRewriteExecutor {
     interface CommitConfigurer {
 
         void configure(TableCommitImpl commit);
+
+        default void prepareMessages(Snapshot snapshot, List<CommitMessage> 
commitMessages) {}
     }
 
     @FunctionalInterface
diff --git 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/MaterializeDeletionVectorsProcedure.java
 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/MaterializeDeletionVectorsProcedure.java
index 339f31adbd..7377d8fb6e 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/MaterializeDeletionVectorsProcedure.java
+++ 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/MaterializeDeletionVectorsProcedure.java
@@ -26,6 +26,9 @@ import org.apache.paimon.partition.PartitionPredicate;
 import org.apache.paimon.spark.utils.SparkProcedureUtils;
 import org.apache.paimon.table.BucketMode;
 import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.sink.CommitMessageImpl;
+import org.apache.paimon.table.sink.TableCommitImpl;
 import org.apache.paimon.utils.ProcedureUtils;
 import org.apache.paimon.utils.StringUtils;
 
@@ -169,7 +172,21 @@ public class MaterializeDeletionVectorsProcedure extends 
BaseProcedure {
                 taskPlanner,
                 javaSparkContext,
                 sparkSession,
-                commit -> 
commit.rowIdCheckConflictForMaterializeDvCompaction(snapshot.id()));
+                new DataEvolutionRewriteExecutor.CommitConfigurer() {
+                    @Override
+                    public void configure(TableCommitImpl commit) {
+                        commit.materializeDvRowIdCheck();
+                    }
+
+                    @Override
+                    public void prepareMessages(
+                            Snapshot planningSnapshot, List<CommitMessage> 
messages) {
+                        messages.replaceAll(
+                                message ->
+                                        ((CommitMessageImpl) message)
+                                                
.withCheckFromSnapshot(planningSnapshot.id()));
+                    }
+                });
     }
 
     private boolean blank(InternalRow args, int index) {
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionRowIdConflictRewriter.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionRowIdConflictRewriter.scala
index ebe2b07108..4cb9b80079 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionRowIdConflictRewriter.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionRowIdConflictRewriter.scala
@@ -370,7 +370,12 @@ private[spark] object DataEvolutionRowIdConflictCommitter {
 
     while (true) {
       try {
-        writer.commit(currentUpdateMessages ++ otherMessages, operation)
+        val messages = currentUpdateMessages ++ otherMessages
+        writer.commit(
+          if (readSnapshotId < 0) messages
+          else
+            
messages.map(_.asInstanceOf[CommitMessageImpl].withCheckFromSnapshot(readSnapshotId)),
+          operation)
         return
       } catch {
         case conflict: RuntimeException if isRowIdExistenceConflict(conflict) 
=>
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
index 15a8f89429..fc887a135e 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
@@ -423,9 +423,6 @@ case class MergeIntoPaimonDataEvolutionTable(
           insertActionInvoke(sparkSession, touchedFileTargetRelation, 
persistSourceDss)
         else Nil
 
-      if (readSnapshot != null) {
-        writer.rowIdCheckConflict(readSnapshot.id())
-      }
       DataEvolutionRowIdConflictCommitter.commit(
         sparkSession,
         table,
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonSparkWriter.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonSparkWriter.scala
index 45f408fb78..76bcbf2505 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonSparkWriter.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonSparkWriter.scala
@@ -450,10 +450,6 @@ case class PaimonSparkWriter(
       .map(deserializeCommitMessage(serializer, _))
   }
 
-  def rowIdCheckConflict(rowIdCheckFromSnapshot: Long): Unit = {
-    
writeBuilder.asInstanceOf[BatchWriteBuilderImpl].rowIdCheckConflict(rowIdCheckFromSnapshot)
-  }
-
   def commit(commitMessages: Seq[CommitMessage]): Unit = {
     commit(commitMessages, null)
   }
diff --git 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala
 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala
index 5f57058fee..6d3997633a 100644
--- 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala
+++ 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala
@@ -1941,8 +1941,9 @@ abstract class CompactProcedureTestBase extends 
PaimonSparkTestBase with StreamT
                 .writePartialFields(updateRows, Seq("value"))
 
             val writer = PaimonSparkWriter(table)
-            writer.rowIdCheckConflict(updateSnapshot.id())
-            writer.commit(updateMessages)
+            writer.commit(
+              updateMessages.map(
+                
_.asInstanceOf[CommitMessageImpl].withCheckFromSnapshot(updateSnapshot.id())))
             assert(table.latestSnapshot().get().operation() == null)
           }
         }
@@ -2016,8 +2017,9 @@ abstract class CompactProcedureTestBase extends 
PaimonSparkTestBase with StreamT
                   .writePartialFields(updateRows, Seq("value", "extra"))
 
               val writer = PaimonSparkWriter(evolvedTable)
-              writer.rowIdCheckConflict(updateSnapshot.id())
-              writer.commit(updateMessages)
+              writer.commit(
+                updateMessages.map(
+                  
_.asInstanceOf[CommitMessageImpl].withCheckFromSnapshot(updateSnapshot.id())))
             }
           }
         )
@@ -2494,8 +2496,9 @@ abstract class CompactProcedureTestBase extends 
PaimonSparkTestBase with StreamT
                 .writePartialFields(updateRows, Seq("value"))
 
             val writer = PaimonSparkWriter(table)
-            writer.rowIdCheckConflict(updateSnapshot.id())
-            writer.commit(updateMessages)
+            writer.commit(
+              updateMessages.map(
+                
_.asInstanceOf[CommitMessageImpl].withCheckFromSnapshot(updateSnapshot.id())))
           }
         }
       )
@@ -2586,8 +2589,9 @@ abstract class CompactProcedureTestBase extends 
PaimonSparkTestBase with StreamT
                 .writePartialFields(updateRows, Seq("value"))
 
             val writer = PaimonSparkWriter(table)
-            writer.rowIdCheckConflict(updateSnapshot.id())
-            writer.commit(updateMessages)
+            writer.commit(
+              updateMessages.map(
+                
_.asInstanceOf[CommitMessageImpl].withCheckFromSnapshot(updateSnapshot.id())))
           }
         }
       )
diff --git 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/RowTrackingTestBase.scala
 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/RowTrackingTestBase.scala
index fef629d1cf..8f8b091ebb 100644
--- 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/RowTrackingTestBase.scala
+++ 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/RowTrackingTestBase.scala
@@ -176,7 +176,6 @@ abstract class RowTrackingTestBase extends 
PaimonSparkTestBase with AdaptiveSpar
       sql("CALL sys.compact(table => 't')").collect()
 
       val writer = PaimonSparkWriter(table)
-      writer.rowIdCheckConflict(readSnapshot.id())
       val targetRelation =
         
PaimonRelation.getPaimonRelation(spark.table("t").queryExecution.analyzed)
       DataEvolutionRowIdConflictCommitter.commit(
@@ -233,7 +232,6 @@ abstract class RowTrackingTestBase extends 
PaimonSparkTestBase with AdaptiveSpar
       sql("CALL sys.compact(table => 't')").collect()
 
       val writer = PaimonSparkWriter(table)
-      writer.rowIdCheckConflict(readSnapshot.id())
       val targetRelation =
         
PaimonRelation.getPaimonRelation(spark.table("t").queryExecution.analyzed)
       DataEvolutionRowIdConflictCommitter.commit(
@@ -290,7 +288,6 @@ abstract class RowTrackingTestBase extends 
PaimonSparkTestBase with AdaptiveSpar
       sql("UPDATE t SET b = 99 WHERE id = 1").collect()
 
       val writer = PaimonSparkWriter(table)
-      writer.rowIdCheckConflict(readSnapshot.id())
       val targetRelation =
         
PaimonRelation.getPaimonRelation(spark.table("t").queryExecution.analyzed)
       val exception = intercept[RuntimeException] {
@@ -481,7 +478,6 @@ abstract class RowTrackingTestBase extends 
PaimonSparkTestBase with AdaptiveSpar
       sql("CALL sys.compact(table => 't')").collect()
 
       val writer = PaimonSparkWriter(table)
-      writer.rowIdCheckConflict(readSnapshot.id())
       val targetRelation =
         
PaimonRelation.getPaimonRelation(spark.table("t").queryExecution.analyzed)
       DataEvolutionRowIdConflictCommitter.commit(

Reply via email to