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 5321706ee0 [core] Fix scan from tag on branch table. (#9044)
5321706ee0 is described below

commit 5321706ee050b472bd23653034103ddd924084a3
Author: Wenchao Wu <[email protected]>
AuthorDate: Wed Aug 5 23:18:55 2026 +0800

    [core] Fix scan from tag on branch table. (#9044)
---
 .../snapshot/StaticFromTagStartingScanner.java     |  5 +-
 .../snapshot/StaticFromTagStartingScannerTest.java | 68 ++++++++++++++++++++++
 2 files changed, 72 insertions(+), 1 deletion(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/StaticFromTagStartingScanner.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/StaticFromTagStartingScanner.java
index 3008c970e2..b6501e5801 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/StaticFromTagStartingScanner.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/StaticFromTagStartingScanner.java
@@ -36,7 +36,10 @@ public class StaticFromTagStartingScanner extends 
ReadPlanStartingScanner {
 
     public Snapshot getSnapshot() {
         TagManager tagManager =
-                new TagManager(snapshotManager.fileIO(), 
snapshotManager.tablePath());
+                new TagManager(
+                        snapshotManager.fileIO(),
+                        snapshotManager.tablePath(),
+                        snapshotManager.branch());
         return tagManager.getOrThrow(tagName).trimToSnapshot();
     }
 
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/source/snapshot/StaticFromTagStartingScannerTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/table/source/snapshot/StaticFromTagStartingScannerTest.java
index c595359ab9..8134860c5f 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/table/source/snapshot/StaticFromTagStartingScannerTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/source/snapshot/StaticFromTagStartingScannerTest.java
@@ -18,6 +18,7 @@
 
 package org.apache.paimon.table.source.snapshot;
 
+import org.apache.paimon.table.FileStoreTable;
 import org.apache.paimon.table.sink.StreamTableCommit;
 import org.apache.paimon.table.sink.StreamTableWrite;
 import org.apache.paimon.utils.SnapshotManager;
@@ -66,6 +67,73 @@ public class StaticFromTagStartingScannerTest extends 
ScannerTestBase {
         commit.close();
     }
 
+    @Test
+    public void testGetSnapshotFromBranchTag() throws Exception {
+        StreamTableWrite write = table.newWrite(commitUser);
+        StreamTableCommit commit = table.newCommit(commitUser);
+
+        write.write(rowData(1, 10, 100L));
+        commit.commit(0, write.prepareCommit(true, 0));
+
+        table.createTag("tag1", 1);
+        table.createBranch("delta", "tag1");
+
+        FileStoreTable deltaTable = table.switchToBranch("delta");
+        String branchOnlyTag = "branch_only_tag";
+        deltaTable.createTag(branchOnlyTag, 1);
+
+        assertThat(deltaTable.tagManager().tagExists(branchOnlyTag)).isTrue();
+        assertThat(table.tagManager().tagExists(branchOnlyTag)).isFalse();
+
+        StaticFromTagStartingScanner scanner =
+                new StaticFromTagStartingScanner(deltaTable.snapshotManager(), 
branchOnlyTag);
+        assertThat(scanner.getSnapshot().id()).isEqualTo(1);
+
+        write.close();
+        commit.close();
+    }
+
+    @Test
+    public void testScanFromBranchTag() throws Exception {
+        StreamTableWrite write = table.newWrite(commitUser);
+        StreamTableCommit commit = table.newCommit(commitUser);
+
+        write.write(rowData(1, 10, 100L));
+        commit.commit(0, write.prepareCommit(true, 0));
+
+        table.createBranch("delta");
+        FileStoreTable deltaTable = table.switchToBranch("delta");
+
+        StreamTableWrite deltaWrite = deltaTable.newWrite(commitUser);
+        StreamTableCommit deltaCommit = deltaTable.newCommit(commitUser);
+
+        deltaWrite.write(rowData(1, 10, 100L));
+        deltaCommit.commit(0, deltaWrite.prepareCommit(true, 0));
+
+        deltaWrite.write(rowData(2, 30, 101L));
+        deltaCommit.commit(1, deltaWrite.prepareCommit(true, 1));
+
+        String tagName = "branch_only_tag";
+        deltaTable.createTag(tagName, 2);
+
+        assertThat(deltaTable.tagManager().tagExists(tagName)).isTrue();
+        assertThat(table.tagManager().tagExists(tagName)).isFalse();
+
+        SnapshotManager deltaSnapshotManager = deltaTable.snapshotManager();
+        StaticFromTagStartingScanner scanner =
+                new StaticFromTagStartingScanner(deltaSnapshotManager, 
tagName);
+        StartingScanner.ScannedResult result =
+                (StartingScanner.ScannedResult) 
scanner.scan(deltaTable.newSnapshotReader());
+        assertThat(result.currentSnapshotId()).isEqualTo(2);
+        assertThat(getResult(deltaTable.newRead(), result.splits()))
+                .hasSameElementsAs(Arrays.asList("+I 1|10|100", "+I 
2|30|101"));
+
+        write.close();
+        commit.close();
+        deltaWrite.close();
+        deltaCommit.close();
+    }
+
     @Test
     public void testNonExistingTag() {
         SnapshotManager snapshotManager = table.snapshotManager();

Reply via email to