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(