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