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