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 a70556b639 [core] Optimize data evolution conflict scans (#9334)
a70556b639 is described below

commit a70556b6399dc0941e84bc05c0a8c2c421269852
Author: YeJunHao <[email protected]>
AuthorDate: Fri Aug 21 14:14:58 2026 +0800

    [core] Optimize data evolution conflict scans (#9334)
---
 .../commit/DataEvolutionConflictDetection.java     |  72 +-----------
 .../paimon/operation/FileStoreCommitTest.java      |  23 ++++
 .../operation/commit/ConflictDetectionTest.java    | 126 +++++++++++++++++----
 3 files changed, 131 insertions(+), 90 deletions(-)

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 3ab5c67c5b..bb6e429138 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
@@ -49,6 +49,7 @@ import javax.annotation.Nullable;
 
 import java.util.ArrayList;
 import java.util.Collection;
+import java.util.Collections;
 import java.util.HashMap;
 import java.util.LinkedHashMap;
 import java.util.List;
@@ -149,66 +150,6 @@ public class DataEvolutionConflictDetection extends 
ConflictDetection {
             CommitKind commitKind,
             @Nullable CommitFailRetryResult previousAttempt,
             boolean hasOverwriteSincePreviousAttempt) {
-        if (commitKind == CommitKind.COMPACT) {
-            return scanCompactBaseDataFiles(
-                    latestSnapshot,
-                    changedPartitions,
-                    deltaFiles,
-                    indexFiles,
-                    previousAttempt,
-                    hasOverwriteSincePreviousAttempt);
-        }
-        if (commitKind == CommitKind.OVERWRITE) {
-            return scanOverwriteBaseDataFiles(
-                    latestSnapshot,
-                    changedPartitions,
-                    deltaFiles,
-                    indexFiles,
-                    previousAttempt,
-                    hasOverwriteSincePreviousAttempt);
-        }
-        return super.scanBaseDataFiles(
-                latestSnapshot,
-                changedPartitions,
-                deltaFiles,
-                indexFiles,
-                commitKind,
-                previousAttempt,
-                hasOverwriteSincePreviousAttempt);
-    }
-
-    private List<SimpleFileEntry> scanCompactBaseDataFiles(
-            Snapshot latestSnapshot,
-            List<BinaryRow> changedPartitions,
-            List<ManifestEntry> deltaFiles,
-            List<IndexManifestEntry> indexFiles,
-            @Nullable CommitFailRetryResult previousAttempt,
-            boolean hasOverwriteSincePreviousAttempt) {
-        List<Range> changedRowRanges = changedRowRanges(deltaFiles, 
indexFiles);
-        if (changedRowRanges.isEmpty()) {
-            return super.scanBaseDataFiles(
-                    latestSnapshot,
-                    changedPartitions,
-                    deltaFiles,
-                    indexFiles,
-                    CommitKind.COMPACT,
-                    previousAttempt,
-                    hasOverwriteSincePreviousAttempt);
-        }
-        return scanChangedRowRanges(
-                latestSnapshot,
-                changedPartitions,
-                changedRowRanges,
-                referencedDataFiles(deltaFiles, indexFiles));
-    }
-
-    private List<SimpleFileEntry> scanOverwriteBaseDataFiles(
-            Snapshot latestSnapshot,
-            List<BinaryRow> changedPartitions,
-            List<ManifestEntry> deltaFiles,
-            List<IndexManifestEntry> indexFiles,
-            @Nullable CommitFailRetryResult previousAttempt,
-            boolean hasOverwriteSincePreviousAttempt) {
         List<Range> changedRowRanges = changedRowRanges(deltaFiles, 
indexFiles);
         Set<String> referencedDataFiles = referencedDataFiles(deltaFiles, 
indexFiles);
         if (!changedRowRanges.isEmpty()) {
@@ -220,14 +161,7 @@ public class DataEvolutionConflictDetection extends 
ConflictDetection {
                     .readAllEntriesFromDataFiles(
                             latestSnapshot, changedPartitions, 
referencedDataFiles);
         }
-        return super.scanBaseDataFiles(
-                latestSnapshot,
-                changedPartitions,
-                deltaFiles,
-                indexFiles,
-                CommitKind.OVERWRITE,
-                previousAttempt,
-                hasOverwriteSincePreviousAttempt);
+        return Collections.emptyList();
     }
 
     private List<SimpleFileEntry> scanChangedRowRanges(
@@ -263,9 +197,9 @@ public class DataEvolutionConflictDetection extends 
ConflictDetection {
 
     private Set<String> referencedDataFiles(
             List<ManifestEntry> deltaFiles, List<IndexManifestEntry> 
indexFiles) {
+        // Include ADD files to detect replay even if their row IDs have 
changed or are unassigned.
         Set<String> referencedDataFiles =
                 deltaFiles.stream()
-                        .filter(entry -> entry.kind() == FileKind.DELETE)
                         .map(entry -> entry.file().fileName())
                         .collect(Collectors.toSet());
         for (IndexManifestEntry indexFile : indexFiles) {
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreCommitTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreCommitTest.java
index c19f318073..f6d150d023 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreCommitTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreCommitTest.java
@@ -682,6 +682,29 @@ public class FileStoreCommitTest {
         }
     }
 
+    @Test
+    public void testCommitOldSnapshotAgainForDataEvolution() throws Exception {
+        TestFileStore store = createRowTrackingDataEvolutionStore();
+        List<ManifestCommittable> committables = new ArrayList<>();
+
+        store.commitDataImpl(
+                generateDataList(10),
+                gen::getPartition,
+                kv -> 0,
+                false,
+                0L,
+                null,
+                Collections.emptyList(),
+                (commit, committable) -> {
+                    commit.commit(committable, false);
+                    committables.add(committable);
+                });
+
+        assertThatThrownBy(() -> store.newCommit().commit(committables.get(0), 
true))
+                .isInstanceOf(RuntimeException.class)
+                .hasMessageContaining("Give up committing.");
+    }
+
     @Test
     public void testCommitWatermarkWithValue() throws Exception {
         TestFileStore store = createStore(false, 2);
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 1fa8430b66..a0bb1459bd 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
@@ -62,6 +62,7 @@ import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
 import static org.mockito.Mockito.when;
 
 class ConflictDetectionTest {
@@ -96,7 +97,7 @@ class ConflictDetectionTest {
     }
 
     @Test
-    void testDataEvolutionCompactScansChangedRowRanges() {
+    void testDataEvolutionCompactScansChangedRowRangesAndAddedFile() {
         CommitScanner scanner = mock(CommitScanner.class);
         DataEvolutionConflictDetection detection =
                 (DataEvolutionConflictDetection) 
createConflictDetection(scanner, true, false);
@@ -109,11 +110,15 @@ class ConflictDetectionTest {
         DataFileMeta dataFile = mock(DataFileMeta.class);
         when(delta.kind()).thenReturn(ADD);
         when(delta.file()).thenReturn(dataFile);
+        when(dataFile.fileName()).thenReturn("added");
         when(dataFile.firstRowId()).thenReturn(10L);
         when(dataFile.nonNullRowIdRange()).thenReturn(changedRange);
         when(scanner.readAllEntriesFromChangedRowRanges(
                         snapshot, changedPartitions, 
Collections.singletonList(changedRange)))
                 .thenReturn(Collections.emptyList());
+        when(scanner.readAllEntriesFromDataFiles(
+                        snapshot, changedPartitions, 
Collections.singleton("added")))
+                .thenReturn(Collections.emptyList());
 
         assertThat(
                         detection.scanBaseDataFiles(
@@ -128,6 +133,9 @@ class ConflictDetectionTest {
         verify(scanner)
                 .readAllEntriesFromChangedRowRanges(
                         snapshot, changedPartitions, 
Collections.singletonList(changedRange));
+        verify(scanner)
+                .readAllEntriesFromDataFiles(
+                        snapshot, changedPartitions, 
Collections.singleton("added"));
     }
 
     @Test
@@ -146,7 +154,7 @@ class ConflictDetectionTest {
                 createFileEntryWithRowId("range-base", ADD, partition, 0, 10L, 
10L);
         SimpleFileEntry legacyBase = createFileEntry("legacy", ADD);
         SimpleFileEntry dvBase = createFileEntry("dv-base", ADD);
-        Set<String> referencedFiles = new HashSet<>(Arrays.asList("legacy", 
"dv-base"));
+        Set<String> referencedFiles = new HashSet<>(Arrays.asList("added", 
"legacy", "dv-base"));
         when(scanner.readAllEntriesFromChangedRowRanges(
                         snapshot, changedPartitions, 
Collections.singletonList(changedRange)))
                 .thenReturn(Collections.singletonList(rangeBase));
@@ -169,7 +177,7 @@ class ConflictDetectionTest {
     }
 
     @Test
-    void testDataEvolutionCompactWithoutRowRangesReusesRetryScan() {
+    void testDataEvolutionCompactWithoutRowRangesSkipsScan() {
         CommitScanner scanner = mock(CommitScanner.class);
         DataEvolutionConflictDetection detection =
                 (DataEvolutionConflictDetection) 
createConflictDetection(scanner, true, false);
@@ -177,12 +185,8 @@ class ConflictDetectionTest {
         Snapshot latestSnapshot = snapshot(2);
         List<BinaryRow> changedPartitions = 
Collections.singletonList(BinaryRow.singleColumn(1));
         SimpleFileEntry oldBase = createFileEntry("old", ADD);
-        SimpleFileEntry removedBase = createFileEntry("old", DELETE);
-        SimpleFileEntry newBase = createFileEntry("new", ADD);
         List<SimpleFileEntry> cachedBase = Collections.singletonList(oldBase);
         CommitFailRetryResult previousAttempt = 
commitFailRetryResult(previousSnapshot, cachedBase);
-        when(scanner.readIncrementalChanges(previousSnapshot, latestSnapshot, 
changedPartitions))
-                .thenReturn(Arrays.asList(removedBase, newBase));
 
         assertThat(
                         detection.scanBaseDataFiles(
@@ -193,10 +197,8 @@ class ConflictDetectionTest {
                                 Snapshot.CommitKind.COMPACT,
                                 previousAttempt,
                                 false))
-                .containsExactly(newBase);
-        assertThat(cachedBase).containsExactly(oldBase);
-        verify(scanner, never())
-                .readAllEntriesFromChangedPartitions(latestSnapshot, 
changedPartitions);
+                .isEmpty();
+        verifyNoInteractions(scanner);
     }
 
     @Test
@@ -326,7 +328,9 @@ class ConflictDetectionTest {
                         snapshot, changedPartitions, 
Collections.singletonList(changedRange)))
                 .thenReturn(Collections.singletonList(rangeBase));
         when(scanner.readAllEntriesFromDataFiles(
-                        snapshot, changedPartitions, 
Collections.singleton("legacy")))
+                        snapshot,
+                        changedPartitions,
+                        new HashSet<>(Arrays.asList("added", "legacy"))))
                 .thenReturn(Collections.singletonList(legacyBase));
 
         assertThat(
@@ -343,15 +347,12 @@ class ConflictDetectionTest {
     }
 
     @Test
-    void testDataEvolutionOverwriteFallsBackWithoutSelectors() {
+    void testDataEvolutionOverwriteWithoutSelectorsSkipsScan() {
         CommitScanner scanner = mock(CommitScanner.class);
         DataEvolutionConflictDetection detection =
                 (DataEvolutionConflictDetection) 
createConflictDetection(scanner, true, false);
         Snapshot snapshot = snapshot(1);
         List<BinaryRow> changedPartitions = 
Collections.singletonList(BinaryRow.singleColumn(1));
-        List<SimpleFileEntry> expected = 
Collections.singletonList(createFileEntry("base", ADD));
-        when(scanner.readAllEntriesFromChangedPartitions(snapshot, 
changedPartitions))
-                .thenReturn(expected);
 
         assertThat(
                         detection.scanBaseDataFiles(
@@ -362,30 +363,113 @@ class ConflictDetectionTest {
                                 Snapshot.CommitKind.OVERWRITE,
                                 null,
                                 false))
-                .isSameAs(expected);
+                .isEmpty();
+        verifyNoInteractions(scanner);
     }
 
     @Test
-    void testDataEvolutionAppendUsesDefaultPartitionScan() {
+    void testDataEvolutionAppendDataFileWithoutRowIdScansReferencedFile() {
         CommitScanner scanner = mock(CommitScanner.class);
         DataEvolutionConflictDetection detection =
                 (DataEvolutionConflictDetection) 
createConflictDetection(scanner, true, false);
         Snapshot snapshot = snapshot(1);
         List<BinaryRow> changedPartitions = 
Collections.singletonList(BinaryRow.singleColumn(1));
-        List<SimpleFileEntry> expected = 
Collections.singletonList(createFileEntry("base", ADD));
-        when(scanner.readAllEntriesFromChangedPartitions(snapshot, 
changedPartitions))
+        Set<String> referencedDataFiles = Collections.singleton("new");
+        List<SimpleFileEntry> expected = 
Collections.singletonList(createFileEntry("new", ADD));
+        when(scanner.readAllEntriesFromDataFiles(snapshot, changedPartitions, 
referencedDataFiles))
                 .thenReturn(expected);
 
         assertThat(
                         detection.scanBaseDataFiles(
                                 snapshot,
                                 changedPartitions,
-                                Collections.emptyList(),
+                                Collections.singletonList(manifestEntry(ADD, 
"new", null, null)),
                                 Collections.emptyList(),
                                 Snapshot.CommitKind.APPEND,
                                 null,
                                 false))
                 .isSameAs(expected);
+        verify(scanner)
+                .readAllEntriesFromDataFiles(snapshot, changedPartitions, 
referencedDataFiles);
+        verify(scanner, never()).readAllEntriesFromChangedPartitions(snapshot, 
changedPartitions);
+    }
+
+    @Test
+    void testDataEvolutionAppendDataFilesUseSelectiveScans() {
+        CommitScanner scanner = mock(CommitScanner.class);
+        DataEvolutionConflictDetection detection =
+                (DataEvolutionConflictDetection) 
createConflictDetection(scanner, true, false);
+        Snapshot snapshot = snapshot(1);
+        BinaryRow partition = BinaryRow.singleColumn(1);
+        List<BinaryRow> changedPartitions = 
Collections.singletonList(partition);
+        Range changedRange = new Range(10, 19);
+        SimpleFileEntry rangeBase =
+                createFileEntryWithRowId("updated", ADD, partition, 0, 10L, 
10L);
+        SimpleFileEntry referencedBase = createFileEntry("duplicate", ADD);
+        when(scanner.readAllEntriesFromChangedRowRanges(
+                        snapshot, changedPartitions, 
Collections.singletonList(changedRange)))
+                .thenReturn(Collections.singletonList(rangeBase));
+        when(scanner.readAllEntriesFromDataFiles(
+                        snapshot, changedPartitions, 
Collections.singleton("duplicate")))
+                .thenReturn(Collections.singletonList(referencedBase));
+
+        assertThat(
+                        detection.scanBaseDataFiles(
+                                snapshot,
+                                changedPartitions,
+                                Arrays.asList(
+                                        manifestEntry(ADD, "updated", 10L, 
changedRange),
+                                        manifestEntry(ADD, "duplicate", null, 
null)),
+                                Collections.emptyList(),
+                                Snapshot.CommitKind.APPEND,
+                                null,
+                                false))
+                .containsExactly(rangeBase, referencedBase);
+        verify(scanner)
+                .readAllEntriesFromChangedRowRanges(
+                        snapshot, changedPartitions, 
Collections.singletonList(changedRange));
+        verify(scanner)
+                .readAllEntriesFromDataFiles(
+                        snapshot, changedPartitions, 
Collections.singleton("duplicate"));
+        verify(scanner, never()).readAllEntriesFromChangedPartitions(snapshot, 
changedPartitions);
+    }
+
+    @Test
+    void testDataEvolutionIndexOnlyAppendScansChangedRowRanges() {
+        CommitScanner scanner = mock(CommitScanner.class);
+        DataEvolutionConflictDetection detection =
+                (DataEvolutionConflictDetection) 
createConflictDetection(scanner, true, false);
+        Snapshot snapshot = snapshot(1);
+        BinaryRow partition = BinaryRow.singleColumn(1);
+        List<BinaryRow> changedPartitions = 
Collections.singletonList(partition);
+        Range changedRange = new Range(10, 19);
+        List<SimpleFileEntry> expected =
+                Collections.singletonList(
+                        createFileEntryWithRowId("base", ADD, partition, 0, 
10L, 10L));
+        when(scanner.readAllEntriesFromChangedRowRanges(
+                        snapshot, changedPartitions, 
Collections.singletonList(changedRange)))
+                .thenReturn(expected);
+
+        assertThat(
+                        detection.scanBaseDataFiles(
+                                snapshot,
+                                changedPartitions,
+                                Collections.emptyList(),
+                                Collections.singletonList(
+                                        createGlobalIndexEntry(
+                                                "idx",
+                                                ADD,
+                                                partition,
+                                                changedRange.from,
+                                                changedRange.to)),
+                                Snapshot.CommitKind.APPEND,
+                                null,
+                                false))
+                .containsExactlyElementsOf(expected);
+        verify(scanner)
+                .readAllEntriesFromChangedRowRanges(
+                        snapshot, changedPartitions, 
Collections.singletonList(changedRange));
+        verify(scanner, never()).readAllEntriesFromChangedPartitions(snapshot, 
changedPartitions);
     }
 
     @Test

Reply via email to