From bcfd450608c82e13f351c340eeff08c2de6f41e3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BC=A0=E4=B8=87=E4=B9=89?= Date: Sun, 23 Aug 2026 14:34:29 +0800 Subject: [PATCH] [Bug] Fix data evolution self-merge ABA across rollback snapshot lineage Add snapshot UUID-based lineage validation to prevent staged row-ID partial updates from being applied to wrong snapshots after a rollback reuses the same numeric snapshot ID. Three-layer validation in checkForRowIdFromSnapshot: 1. Fail closed when latest snapshot ID < base snapshot ID (rollback deleted the update's base snapshot) 2. Detect missing base snapshot (race with concurrent cleanup) 3. ABA detection: compare base snapshot UUID with current snapshot UUID at the same ID (different lineage) The baseSnapshotUuid field is nullable for backward compatibility. Callers that don't pass UUID get existing behavior without the ABA protection. Closes #9352 --- .../apache/paimon/errors/ErrorMessages.java | 4 ++ .../paimon/operation/FileStoreCommit.java | 3 + .../paimon/operation/FileStoreCommitImpl.java | 7 ++ .../operation/commit/ConflictDetection.java | 5 ++ .../DataEvolutionConflictDetection.java | 44 +++++++++++- .../table/sink/BatchWriteBuilderImpl.java | 9 ++- .../paimon/table/sink/InnerTableCommit.java | 3 + .../paimon/table/sink/TableCommitImpl.java | 7 ++ .../commit/ConflictDetectionTest.java | 70 +++++++++++++++++++ .../action/DataEvolutionMergeIntoAction.java | 6 +- .../DataEvolutionDeleteSink.java | 5 +- 11 files changed, 157 insertions(+), 6 deletions(-) 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());