diff --git a/paimon-api/src/main/java/org/apache/paimon/errors/ErrorMessages.java b/paimon-api/src/main/java/org/apache/paimon/errors/ErrorMessages.java index 8fd9012b1bb9..9d33437b0d46 100644 --- a/paimon-api/src/main/java/org/apache/paimon/errors/ErrorMessages.java +++ b/paimon-api/src/main/java/org/apache/paimon/errors/ErrorMessages.java @@ -25,5 +25,9 @@ public class ErrorMessages { "For Data Evolution table, multiple 'MERGE INTO' operations have encountered conflicts," + " updating the same file, which can render some updates ineffective."; + public static final String DATA_EVOLUTION_SNAPSHOT_LINEAGE_CONFLICT_MESSAGE = + "For Data Evolution table, the base snapshot lineage has changed, possibly due to a" + + " rollback. Staged updates from the old snapshot lineage cannot be committed."; + private ErrorMessages() {} } 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 b039ffb9e9fc..f78cae84dfbc 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 @@ -46,6 +46,9 @@ public interface FileStoreCommit extends AutoCloseable { FileStoreCommit rowIdCheckConflict(@Nullable Long rowIdCheckFromSnapshot); + FileStoreCommit rowIdCheckConflict( + @Nullable Long rowIdCheckFromSnapshot, @Nullable String baseSnapshotUuid); + FileStoreCommit rowIdCheckConflictForMaterializeDvCompaction( @Nullable Long rowIdCheckFromSnapshot); 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 02ce5dc05d09..7fb39e7cd7ba 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 @@ -267,6 +267,13 @@ public FileStoreCommit rowIdCheckConflict(@Nullable Long rowIdCheckFromSnapshot) return this; } + @Override + public FileStoreCommit rowIdCheckConflict( + @Nullable Long rowIdCheckFromSnapshot, @Nullable String baseSnapshotUuid) { + this.conflictDetection.setRowIdCheckFromSnapshot(rowIdCheckFromSnapshot, baseSnapshotUuid); + return this; + } + @Override public FileStoreCommit rowIdCheckConflictForMaterializeDvCompaction( @Nullable Long rowIdCheckFromSnapshot) { diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/commit/ConflictDetection.java b/paimon-core/src/main/java/org/apache/paimon/operation/commit/ConflictDetection.java index b660d50ad481..f019b73dd922 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/commit/ConflictDetection.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/commit/ConflictDetection.java @@ -162,6 +162,11 @@ public void setRowIdCheckFromSnapshot(@Nullable Long rowIdCheckFromSnapshot) { // Only Data Evolution tables support Row ID conflict detection. } + public void setRowIdCheckFromSnapshot( + @Nullable Long rowIdCheckFromSnapshot, @Nullable String baseSnapshotUuid) { + // Only Data Evolution tables support Row ID conflict detection. + } + public void setRowIdCheckFromSnapshotForMaterializeDvCompaction( @Nullable Long rowIdCheckFromSnapshot) { // Only Data Evolution tables support Row ID conflict detection. diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/commit/DataEvolutionConflictDetection.java b/paimon-core/src/main/java/org/apache/paimon/operation/commit/DataEvolutionConflictDetection.java index bb6e429138c9..4ce1a373283b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/commit/DataEvolutionConflictDetection.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/commit/DataEvolutionConflictDetection.java @@ -74,6 +74,7 @@ public class DataEvolutionConflictDetection extends ConflictDetection { private final SnapshotManager snapshotManager; private @Nullable Long rowIdCheckFromSnapshot; + private @Nullable String baseSnapshotUuid; private @Nullable RowIdConflictCheckStrategy rowIdConflictCheckStrategy; public DataEvolutionConflictDetection( @@ -103,19 +104,31 @@ public DataEvolutionConflictDetection( @Override public void setRowIdCheckFromSnapshot(@Nullable Long rowIdCheckFromSnapshot) { setRowIdCheckFromSnapshot( - rowIdCheckFromSnapshot, DataEvolutionDmlRowIdConflictCheck.INSTANCE); + rowIdCheckFromSnapshot, null, DataEvolutionDmlRowIdConflictCheck.INSTANCE); + } + + @Override + public void setRowIdCheckFromSnapshot( + @Nullable Long rowIdCheckFromSnapshot, @Nullable String baseSnapshotUuid) { + setRowIdCheckFromSnapshot( + rowIdCheckFromSnapshot, + baseSnapshotUuid, + DataEvolutionDmlRowIdConflictCheck.INSTANCE); } @Override public void setRowIdCheckFromSnapshotForMaterializeDvCompaction( @Nullable Long rowIdCheckFromSnapshot) { - setRowIdCheckFromSnapshot(rowIdCheckFromSnapshot, MaterializeDvRowIdConflictCheck.INSTANCE); + setRowIdCheckFromSnapshot( + rowIdCheckFromSnapshot, null, MaterializeDvRowIdConflictCheck.INSTANCE); } private void setRowIdCheckFromSnapshot( @Nullable Long rowIdCheckFromSnapshot, + @Nullable String baseSnapshotUuid, RowIdConflictCheckStrategy conflictCheckStrategy) { this.rowIdCheckFromSnapshot = rowIdCheckFromSnapshot; + this.baseSnapshotUuid = baseSnapshotUuid; this.rowIdConflictCheckStrategy = rowIdCheckFromSnapshot == null ? null : conflictCheckStrategy; } @@ -356,8 +369,33 @@ private Optional checkForRowIdFromSnapshot( return Optional.empty(); } + // Fail closed when the latest snapshot ID is less than the base snapshot ID. + // This indicates a rollback has deleted newer snapshots, and the staged update + // is based on a snapshot lineage that no longer exists. + if (latestSnapshot.id() < rowIdCheckFromSnapshot) { + return Optional.of( + new RuntimeException( + ErrorMessages.DATA_EVOLUTION_SNAPSHOT_LINEAGE_CONFLICT_MESSAGE)); + } + + // Detect equal snapshot IDs with different snapshot UUIDs (ABA problem). + // A rollback can delete a snapshot and a new commit can reuse the same numeric ID. + // If the base snapshot UUID differs from the current snapshot UUID at that ID, + // the staged update is based on a different snapshot lineage. + Snapshot baseSnapshot = snapshotManager.snapshot(rowIdCheckFromSnapshot); + if (baseSnapshot == null) { + return Optional.of( + new RuntimeException( + ErrorMessages.DATA_EVOLUTION_SNAPSHOT_LINEAGE_CONFLICT_MESSAGE)); + } + if (baseSnapshotUuid != null && !baseSnapshotUuid.equals(baseSnapshot.uuid())) { + return Optional.of( + new RuntimeException( + ErrorMessages.DATA_EVOLUTION_SNAPSHOT_LINEAGE_CONFLICT_MESSAGE)); + } + List changedPartitions = changedPartitions(deltaEntries, deltaIndexEntries); - Long checkNextRowId = snapshotManager.snapshot(rowIdCheckFromSnapshot).nextRowId(); + Long checkNextRowId = baseSnapshot.nextRowId(); checkState( checkNextRowId != null, "Next row id cannot be null for snapshot %s.", 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 d8c97405e2b0..da709eb366ef 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 @@ -40,6 +40,7 @@ public class BatchWriteBuilderImpl implements BatchWriteBuilder { private Map staticPartition; private @Nullable Long rowIdCheckFromSnapshot = null; + private @Nullable String baseSnapshotUuid = null; public BatchWriteBuilderImpl(InnerTable table) { this.table = table; @@ -77,7 +78,7 @@ public BatchTableCommit newCommit() { InnerTableCommit commit = table.newCommit(commitUser) .withOverwrite(staticPartition) - .rowIdCheckConflict(rowIdCheckFromSnapshot); + .rowIdCheckConflict(rowIdCheckFromSnapshot, baseSnapshotUuid); commit.ignoreEmptyCommit( Options.fromMap(table.options()) .getOptional(CoreOptions.SNAPSHOT_IGNORE_EMPTY_COMMIT) @@ -86,7 +87,13 @@ public BatchTableCommit newCommit() { } public BatchWriteBuilderImpl rowIdCheckConflict(@Nullable Long rowIdCheckFromSnapshot) { + return rowIdCheckConflict(rowIdCheckFromSnapshot, null); + } + + public BatchWriteBuilderImpl rowIdCheckConflict( + @Nullable Long rowIdCheckFromSnapshot, @Nullable String baseSnapshotUuid) { this.rowIdCheckFromSnapshot = rowIdCheckFromSnapshot; + this.baseSnapshotUuid = baseSnapshotUuid; return this; } } 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 43f98d0e7933..b7918974833d 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 @@ -58,6 +58,9 @@ public interface InnerTableCommit extends StreamTableCommit, BatchTableCommit { InnerTableCommit rowIdCheckConflict(@Nullable Long rowIdCheckFromSnapshot); + InnerTableCommit rowIdCheckConflict( + @Nullable Long rowIdCheckFromSnapshot, @Nullable String baseSnapshotUuid); + InnerTableCommit rowIdCheckConflictForMaterializeDvCompaction( @Nullable Long rowIdCheckFromSnapshot); 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 014b5e64daa1..847e609c6caf 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 @@ -181,6 +181,13 @@ public TableCommitImpl rowIdCheckConflict(@Nullable Long rowIdCheckFromSnapshot) return this; } + @Override + public TableCommitImpl rowIdCheckConflict( + @Nullable Long rowIdCheckFromSnapshot, @Nullable String baseSnapshotUuid) { + commit.rowIdCheckConflict(rowIdCheckFromSnapshot, baseSnapshotUuid); + return this; + } + @Override public TableCommitImpl rowIdCheckConflictForMaterializeDvCompaction( @Nullable Long rowIdCheckFromSnapshot) { diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/commit/ConflictDetectionTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/commit/ConflictDetectionTest.java index a0bb1459bd00..d285c96c1ba0 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/commit/ConflictDetectionTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/commit/ConflictDetectionTest.java @@ -20,6 +20,7 @@ import org.apache.paimon.Snapshot; import org.apache.paimon.data.BinaryRow; +import org.apache.paimon.errors.ErrorMessages; import org.apache.paimon.index.DeletionVectorMeta; import org.apache.paimon.index.GlobalIndexMeta; import org.apache.paimon.index.IndexFileMeta; @@ -1637,4 +1638,73 @@ private Snapshot snapshot(long id) { null, null); } + + @Test + void testRowIdCheckConflictAbaDetectsRollback() { + CommitScanner scanner = mock(CommitScanner.class); + SnapshotManager snapshotManager = mock(SnapshotManager.class); + DataEvolutionConflictDetection detection = + (DataEvolutionConflictDetection) + createConflictDetection(scanner, true, false, false, snapshotManager); + + String baseUuid = "uuid-v1"; + detection.setRowIdCheckFromSnapshot(1L, baseUuid); + + Snapshot baseSnapshot = mock(Snapshot.class); + Snapshot latestSnapshot = mock(Snapshot.class); + when(baseSnapshot.uuid()).thenReturn("uuid-v2"); + when(baseSnapshot.nextRowId()).thenReturn(100L); + when(latestSnapshot.id()).thenReturn(2L); + when(latestSnapshot.commitUser()).thenReturn("test-user"); + when(snapshotManager.snapshot(1L)).thenReturn(baseSnapshot); + + RowIdConflictChecker checker = mock(RowIdConflictChecker.class); + when(checker.isEmpty()).thenReturn(false); + + assertThat( + detection.checkConflicts( + latestSnapshot, + Collections.emptyList(), + Collections.emptyList(), + Collections.emptyList(), + checker, + Snapshot.CommitKind.APPEND)) + .isPresent() + .get() + .hasMessageContaining( + ErrorMessages.DATA_EVOLUTION_SNAPSHOT_LINEAGE_CONFLICT_MESSAGE); + } + + @Test + void testRowIdCheckConflictNoAbaWhenUuidMatches() { + CommitScanner scanner = mock(CommitScanner.class); + SnapshotManager snapshotManager = mock(SnapshotManager.class); + DataEvolutionConflictDetection detection = + (DataEvolutionConflictDetection) + createConflictDetection(scanner, true, false, false, snapshotManager); + + String baseUuid = "uuid-v1"; + detection.setRowIdCheckFromSnapshot(1L, baseUuid); + + Snapshot baseSnapshot = mock(Snapshot.class); + Snapshot latestSnapshot = mock(Snapshot.class); + when(baseSnapshot.uuid()).thenReturn("uuid-v1"); + when(baseSnapshot.nextRowId()).thenReturn(100L); + when(latestSnapshot.id()).thenReturn(2L); + when(latestSnapshot.commitUser()).thenReturn("test-user"); + when(snapshotManager.snapshot(1L)).thenReturn(baseSnapshot); + + RowIdConflictChecker checker = mock(RowIdConflictChecker.class); + when(checker.isEmpty()).thenReturn(false); + + assertThat( + detection.checkConflicts( + latestSnapshot, + Collections.emptyList(), + Collections.emptyList(), + Collections.emptyList(), + checker, + Snapshot.CommitKind.APPEND)) + .isEmpty(); + } } 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 d0babc5b6f0a..665923db40ed 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 @@ -19,6 +19,7 @@ package org.apache.paimon.flink.action; import org.apache.paimon.CoreOptions; +import org.apache.paimon.Snapshot; import org.apache.paimon.annotation.VisibleForTesting; import org.apache.paimon.data.InternalRow; import org.apache.paimon.flink.FlinkRowWrapper; @@ -405,6 +406,8 @@ public DataStream commit( FileStoreTable storeTable = (FileStoreTable) table; // copy to avoid serialization issue long baseSnapshotId = this.baseSnapshotId; + Snapshot baseSnapshot = ((FileStoreTable) table).snapshotManager().snapshot(baseSnapshotId); + String baseSnapshotUuid = baseSnapshot != null ? baseSnapshot.uuid() : null; // Check if some global-indexed columns are updated DataStream checked = @@ -425,7 +428,8 @@ public DataStream commit( storeTable, storeTable .newCommit(context.commitUser()) - .rowIdCheckConflict(baseSnapshotId), + .rowIdCheckConflict( + baseSnapshotId, baseSnapshotUuid), context), new NoopCommittableStateManager()); 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 55d6814c739c..3e9cf126f9d8 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 @@ -121,6 +121,8 @@ public DataStreamSink sinkFrom(DataStream rowIds) { .setParallelism(sinkParallelism); String commitUser = CoreOptions.createCommitUser(table.coreOptions().toConfiguration()); + Snapshot baseSnapshot = table.snapshotManager().snapshot(baseSnapshotId); + String baseSnapshotUuid = baseSnapshot != null ? baseSnapshot.uuid() : null; CommitterOperatorFactory committerOperator = new CommitterOperatorFactory<>( false, @@ -131,7 +133,8 @@ public DataStreamSink sinkFrom(DataStream rowIds) { table, table.newCommit(context.commitUser()) .withOperation(Snapshot.Operation.DELETE) - .rowIdCheckConflict(baseSnapshotId), + .rowIdCheckConflict( + baseSnapshotId, baseSnapshotUuid), context), new NoopCommittableStateManager());