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 796b698b6c [core] Fix incremental diff scan missing first snapshot
(#9234)
796b698b6c is described below
commit 796b698b6c2d37544599efdb31c80d6d771b416c
Author: Arnav Balyan <[email protected]>
AuthorDate: Sat Aug 15 17:41:38 2026 +0530
[core] Fix incremental diff scan missing first snapshot (#9234)
---
.../snapshot/IncrementalDiffStartingScanner.java | 14 ++++++----
.../table/IncrementalTimeStampTableTest.java | 32 ++++++++++++++++++++++
2 files changed, 41 insertions(+), 5 deletions(-)
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/IncrementalDiffStartingScanner.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/IncrementalDiffStartingScanner.java
index f7bb09e112..6beb550e26 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/IncrementalDiffStartingScanner.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/IncrementalDiffStartingScanner.java
@@ -101,18 +101,22 @@ public class IncrementalDiffStartingScanner extends
AbstractStartingScanner {
return new IncrementalDiffStartingScanner(snapshotManager, start, end);
}
- public static IncrementalDiffStartingScanner betweenTimestamps(
+ public static StartingScanner betweenTimestamps(
long startTimestamp, long endTimestamp, SnapshotManager
snapshotManager) {
Snapshot startSnapshot =
snapshotManager.earlierOrEqualTimeMills(startTimestamp);
- if (startSnapshot == null) {
- startSnapshot = snapshotManager.earliestSnapshot();
- }
-
Snapshot endSnapshot =
snapshotManager.earlierOrEqualTimeMills(endTimestamp);
if (endSnapshot == null) {
endSnapshot = snapshotManager.latestSnapshot();
}
+ if (startSnapshot == null) {
+ Snapshot earliestSnapshot = snapshotManager.earliestSnapshot();
+ if (earliestSnapshot.id() == Snapshot.FIRST_SNAPSHOT_ID) {
+ return new StaticFromSnapshotStartingScanner(snapshotManager,
endSnapshot.id());
+ }
+ startSnapshot = earliestSnapshot;
+ }
+
return new IncrementalDiffStartingScanner(snapshotManager,
startSnapshot, endSnapshot);
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/IncrementalTimeStampTableTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/IncrementalTimeStampTableTest.java
index c6a8e2abf2..aed2b352bc 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/IncrementalTimeStampTableTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/IncrementalTimeStampTableTest.java
@@ -35,6 +35,7 @@ import org.junit.jupiter.api.Test;
import java.time.LocalDateTime;
import java.util.List;
+import static org.apache.paimon.CoreOptions.INCREMENTAL_BETWEEN_SCAN_MODE;
import static org.apache.paimon.CoreOptions.INCREMENTAL_BETWEEN_TIMESTAMP;
import static org.apache.paimon.SnapshotTest.newSnapshotManager;
import static org.apache.paimon.io.DataFileTestUtils.row;
@@ -45,6 +46,37 @@ import static org.assertj.core.api.Assertions.assertThat;
/** Test for {@link CoreOptions#INCREMENTAL_BETWEEN_TIMESTAMP}. */
public class IncrementalTimeStampTableTest extends TableTestBase {
+ @Test
+ public void testDiffFromBeforeFirstSnapshot() throws Exception {
+ Identifier identifier = identifier("T");
+ Schema schema =
+ Schema.newBuilder()
+ .column("k", DataTypes.INT())
+ .column("v", DataTypes.INT())
+ .build();
+ catalog.createTable(identifier, schema, true);
+ Table table = catalog.getTable(identifier);
+
+ write(table, GenericRow.of(1, 1));
+ write(table, GenericRow.of(2, 2));
+
+ SnapshotManager snapshotManager =
+ newSnapshotManager(
+ LocalFileIO.create(),
+ new Path(String.format("%s/%s.db/%s", warehouse,
database, "T")));
+ long startTimestamp = snapshotManager.snapshot(1).timeMillis() - 1;
+ long endTimestamp = snapshotManager.snapshot(2).timeMillis();
+
+ assertThat(
+ read(
+ table,
+ Pair.of(
+ INCREMENTAL_BETWEEN_TIMESTAMP,
+ String.format("%s,%s", startTimestamp,
endTimestamp)),
+ Pair.of(INCREMENTAL_BETWEEN_SCAN_MODE,
"diff")))
+ .containsExactlyInAnyOrder(GenericRow.of(1, 1),
GenericRow.of(2, 2));
+ }
+
@Test
public void testPrimaryKeyTable() throws Exception {
Identifier identifier = identifier("T");