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");

Reply via email to