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 bbecc0fc44 [core][flink] Support anchor-based merge in starting phase 
for chain table streaming read (#8723)
bbecc0fc44 is described below

commit bbecc0fc44e74a4558936fe25629a271d3b77d9a
Author: Juntao Zhang <[email protected]>
AuthorDate: Mon Jul 20 11:40:37 2026 +0800

    [core][flink] Support anchor-based merge in starting phase for chain table 
streaming read (#8723)
---
 docs/docs/primary-key-table/chain-table.mdx        |  36 +-
 docs/generated/core_configuration.html             |   6 +
 .../main/java/org/apache/paimon/CoreOptions.java   |  20 +
 .../apache/paimon/table/ChainGroupReadTable.java   |  77 +--
 .../apache/paimon/table/ChainTableStreamScan.java  | 170 ++++++-
 .../org/apache/paimon/utils/ChainTableUtils.java   |  77 +++
 .../apache/paimon/flink/FlinkChainTableITCase.java | 562 +++++++++------------
 7 files changed, 546 insertions(+), 402 deletions(-)

diff --git a/docs/docs/primary-key-table/chain-table.mdx 
b/docs/docs/primary-key-table/chain-table.mdx
index 9556434692..031c0a03e4 100644
--- a/docs/docs/primary-key-table/chain-table.mdx
+++ b/docs/docs/primary-key-table/chain-table.mdx
@@ -222,9 +222,11 @@ you will get the following result:
 
 Chain tables support Flink streaming read. A streaming read job operates in 
two phases:
 
-1. **Full load phase**: Produces a full result by reading the latest snapshot 
partition (per group)
-   and delta partitions that come after it. For each partition group, only the 
most recent snapshot
-   partition is included — older snapshot partitions are considered outdated 
and excluded.
+1. **Full load phase**: By default it produces a lightweight result by reading 
the latest
+   snapshot partition (per group) and delta partitions that come after it. 
Older snapshot
+   partitions are excluded. You can enable 
`chain-table.streaming.merge-snapshot` to perform
+   anchor-based chain merging in this phase, allowing cross-branch `DELETE` 
records to be
+   resolved together with the snapshot data.
 2. **Incremental phase**: Continuously reads new commits from the delta branch 
as they arrive.
 
 ### Write-Side Requirements
@@ -251,6 +253,28 @@ SET 'execution.runtime-mode' = 'streaming';
 INSERT INTO downstream_sink SELECT * FROM default.t;
 ```
 
+### Merge Snapshot in Full Load Phase
+
+By default, the full-load phase is lightweight: for each group it reads the 
latest snapshot and later
+delta partitions as separate splits. This is fast but cross-branch deletes are 
invisible — the `DELETE`
+records in the delta branch cannot be deleted in the snapshot branch.
+
+If you need a fully reconciled starting snapshot, enable merge mode:
+
+```sql
+ALTER TABLE default.t SET (
+  'chain-table.streaming.merge-snapshot' = 'true'
+);
+```
+
+With merge mode enabled, the full-load phase merges the latest snapshot 
partition per group with
+delta partitions whose chain key is strictly greater than the snapshot's, so 
cross-branch
+deletes are correctly resolved. The trade-off is a heavier startup scan.
+
+To reduce the overhead, run `CALL sys.compact_chain_table(...)` periodically.
+After compaction, only the delta changes that arrived after compaction need to 
be merged.
+
+
 ### Limitations
 
 - The incremental phase only monitors the **delta branch**. Writes to the 
snapshot branch are
@@ -268,6 +292,12 @@ INSERT INTO downstream_sink SELECT * FROM default.t;
   specific partition, use batch mode instead.
 - The delta branch must use the `DEDUPLICATE` merge engine (default). Other 
merge engine
   types are not supported.
+- The `changelog-producer` option must be `none` (default) or `input`; 
`lookup` and `full-compaction`
+  are not supported for chain tables.
+  - When `changelog-producer` is `none`, Flink's operator normalizes records 
by the full
+    primary key including the chain partition. The records of `-D`/`-U` in a 
different chain partition
+    than the original records of `+I` will be dropped. Use `input` if 
downstream must receive cross-partition
+    changelog records.
 
 ## Lookup Join
 
diff --git a/docs/generated/core_configuration.html 
b/docs/generated/core_configuration.html
index 143c244451..29947a5e44 100644
--- a/docs/generated/core_configuration.html
+++ b/docs/generated/core_configuration.html
@@ -158,6 +158,12 @@ under the License.
             <td>Boolean</td>
             <td>Whether enabled chain table.</td>
         </tr>
+        <tr>
+            <td><h5>chain-table.streaming.merge-snapshot</h5></td>
+            <td style="word-wrap: break-word;">false</td>
+            <td>Boolean</td>
+            <td>If true, the starting phase of chain table streaming read 
performs anchor-based chain merging: for each group it merges the latest 
snapshot partition with delta partitions whose chain key is strictly greater 
than the snapshot chain key. This allows streaming readers to see cross-branch 
deletions and updates at the cost of a heavier startup scan. When false 
(default), the starting phase only reads the latest snapshot partition per 
group and later delta partitions as separa [...]
+        </tr>
         <tr>
             <td><h5>changelog-file.compression</h5></td>
             <td style="word-wrap: break-word;">(none)</td>
diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java 
b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
index 4f1cb31111..459f4a1545 100644
--- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
+++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
@@ -295,6 +295,22 @@ public class CoreOptions implements Serializable {
                                     + "suffix of the table's partition keys. 
Comma-separated. "
                                     + "If not set, all partition keys 
participate in chain.");
 
+    public static final ConfigOption<Boolean> 
CHAIN_TABLE_STREAMING_MERGE_SNAPSHOT =
+            key("chain-table.streaming.merge-snapshot")
+                    .booleanType()
+                    .defaultValue(false)
+                    .withDescription(
+                            "If true, the starting phase of chain table 
streaming read performs "
+                                    + "anchor-based chain merging: for each 
group it merges the "
+                                    + "latest snapshot partition with delta 
partitions whose chain "
+                                    + "key is strictly greater than the 
snapshot chain key. This "
+                                    + "allows streaming readers to see 
cross-branch deletions and "
+                                    + "updates at the cost of a heavier 
startup scan. When false "
+                                    + "(default), the starting phase only 
reads the latest snapshot "
+                                    + "partition per group and later delta 
partitions as separate "
+                                    + "splits, which is lightweight but may 
not reflect cross-branch "
+                                    + "deletes.");
+
     public static final String FILE_FORMAT_ORC = "orc";
     public static final String FILE_FORMAT_AVRO = "avro";
     public static final String FILE_FORMAT_PARQUET = "parquet";
@@ -4185,6 +4201,10 @@ public class CoreOptions implements Serializable {
         return 
Arrays.stream(value.split(",")).map(String::trim).collect(Collectors.toList());
     }
 
+    public boolean chainTableStreamingMergeSnapshot() {
+        return options.get(CHAIN_TABLE_STREAMING_MERGE_SNAPSHOT);
+    }
+
     public boolean formatTableImplementationIsPaimon() {
         return options.get(FORMAT_TABLE_IMPLEMENTATION) == 
FormatTableImplementation.PAIMON;
     }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/ChainGroupReadTable.java 
b/paimon-core/src/main/java/org/apache/paimon/table/ChainGroupReadTable.java
index 535b4a575b..c1e67e8a95 100644
--- a/paimon-core/src/main/java/org/apache/paimon/table/ChainGroupReadTable.java
+++ b/paimon-core/src/main/java/org/apache/paimon/table/ChainGroupReadTable.java
@@ -60,7 +60,6 @@ import java.util.function.Function;
 import java.util.stream.Collectors;
 
 import static org.apache.paimon.utils.Preconditions.checkArgument;
-import static org.apache.paimon.utils.Preconditions.checkNotNull;
 
 /**
  * Chain table which mainly read from the snapshot branch. However, if the 
snapshot branch does not
@@ -360,7 +359,6 @@ public class ChainGroupReadTable extends 
FallbackReadFileStoreTable {
                 for (List<BinaryRow> deltaPartitionsInGroup : 
groupedDeltaPartitions.values()) {
 
                     // Sort delta by chain dimension ascending.
-                    // chainPartitionForCompare avoids copying BinaryRow in 
the comparator hot path.
                     deltaPartitionsInGroup.sort(
                             (a, b) ->
                                     chainPartitionComparator.compare(
@@ -432,69 +430,26 @@ public class ChainGroupReadTable extends 
FallbackReadFileStoreTable {
                             
deltaScan.withPartitionFilter(selectedDeltaPartitions);
                         }
 
-                        List<Split> subSplits = deltaScan.plan().splits();
-                        Set<String> snapshotFileNames = new HashSet<>();
+                        List<DataSplit> deltaSubSplits =
+                                deltaScan.plan().splits().stream()
+                                        .map(s -> (DataSplit) s)
+                                        .collect(Collectors.toList());
+                        List<DataSplit> snapshotSubSplits = new ArrayList<>();
                         if (partitionPairs.getValue() != null) {
                             snapshotScan.withPartitionFilter(
                                     
Collections.singletonList(partitionPairs.getValue()));
-                            List<Split> mainSubSplits = 
snapshotScan.plan().splits();
-                            snapshotFileNames =
-                                    mainSubSplits.stream()
-                                            .flatMap(
-                                                    s ->
-                                                            ((DataSplit) s)
-                                                                    
.dataFiles().stream()
-                                                                            
.map(
-                                                                               
     DataFileMeta
-                                                                               
             ::fileName))
-                                            .collect(Collectors.toSet());
-                            subSplits.addAll(mainSubSplits);
-                        }
-                        Map<Integer, List<DataSplit>> bucketSplits = new 
LinkedHashMap<>();
-                        Integer bucketInAll = null;
-                        for (Split split : subSplits) {
-                            DataSplit dataSplit = (DataSplit) split;
-                            Integer totalBuckets = dataSplit.totalBuckets();
-                            checkNotNull(totalBuckets);
-                            if (bucketInAll == null) {
-                                bucketInAll = totalBuckets;
-                            } else {
-                                checkArgument(
-                                        totalBuckets.equals(bucketInAll),
-                                        "Inconsistent bucket num " + 
dataSplit.bucket());
-                            }
-
-                            bucketSplits
-                                    .computeIfAbsent(dataSplit.bucket(), k -> 
new ArrayList<>())
-                                    .add(dataSplit);
-                        }
-                        for (Map.Entry<Integer, List<DataSplit>> entry : 
bucketSplits.entrySet()) {
-                            HashMap<String, String> fileBucketPathMapping = 
new HashMap<>();
-                            HashMap<String, String> fileBranchMapping = new 
HashMap<>();
-                            List<DataSplit> splitList = entry.getValue();
-                            for (DataSplit dataSplit : splitList) {
-                                for (DataFileMeta file : 
dataSplit.dataFiles()) {
-                                    fileBucketPathMapping.put(
-                                            file.fileName(), 
dataSplit.bucketPath());
-                                    String branch =
-                                            
snapshotFileNames.contains(file.fileName())
-                                                    ? 
options.scanFallbackSnapshotBranch()
-                                                    : 
options.scanFallbackDeltaBranch();
-                                    fileBranchMapping.put(file.fileName(), 
branch);
-                                }
-                            }
-                            ChainSplit split =
-                                    new ChainSplit(
-                                            partitionPairs.getKey(),
-                                            entry.getValue().stream()
-                                                    .flatMap(
-                                                            dataSplit ->
-                                                                    
dataSplit.dataFiles().stream())
-                                                    
.collect(Collectors.toList()),
-                                            fileBranchMapping,
-                                            fileBucketPathMapping);
-                            splits.add(split);
+                            snapshotSubSplits =
+                                    snapshotScan.plan().splits().stream()
+                                            .map(s -> (DataSplit) s)
+                                            .collect(Collectors.toList());
                         }
+                        splits.addAll(
+                                ChainTableUtils.buildChainSplits(
+                                        partitionPairs.getKey(),
+                                        snapshotSubSplits,
+                                        deltaSubSplits,
+                                        options.scanFallbackSnapshotBranch(),
+                                        options.scanFallbackDeltaBranch()));
                     }
                 }
             }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/ChainTableStreamScan.java 
b/paimon-core/src/main/java/org/apache/paimon/table/ChainTableStreamScan.java
index fcd07cb26b..5812a2611b 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/ChainTableStreamScan.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/ChainTableStreamScan.java
@@ -118,8 +118,16 @@ public class ChainTableStreamScan implements 
StreamDataTableScan {
     /** Maximum number of retries when race condition is detected during 
position capture. */
     private static final int MAX_RACE_RETRIES = 3;
 
+    /**
+     * If true, the starting phase uses the same anchor-based chain merging 
plan as batch mode,
+     * allowing streaming readers to see deletions/updates that require 
merging historical snapshot
+     * partitions with delta partitions.
+     */
+    private final boolean mergeSnapshot;
+
     public ChainTableStreamScan(ChainGroupReadTable chainGroupReadTable) {
         this.chainGroupReadTable = chainGroupReadTable;
+        this.mergeSnapshot = 
chainGroupReadTable.coreOptions().chainTableStreamingMergeSnapshot();
         this.batchScan =
                 new ChainGroupReadTable.ChainTableBatchScan(
                         chainGroupReadTable.schema(), chainGroupReadTable);
@@ -176,8 +184,10 @@ public class ChainTableStreamScan implements 
StreamDataTableScan {
      * come after it. Older snapshot partitions are excluded. Each primary key 
appears exactly once
      * under its natural partition.
      *
-     * <p>Unlike batch full scan, anchor-based chain merging is not performed. 
This keeps Phase 1
-     * lightweight for long-running jobs.
+     * <p>By default anchor-based chain merging is skipped to keep Phase 1 
lightweight. When {@code
+     * chain-table.streaming.merge-snapshot} is true, the latest snapshot 
partition per group is
+     * merged with delta partitions whose chain key is strictly greater than 
the snapshot chain key,
+     * allowing streaming readers to see cross-branch deletions and updates.
      */
     private TableScan.Plan planStarting() {
         FileStoreTable deltaTable = chainGroupReadTable.other();
@@ -274,9 +284,52 @@ public class ChainTableStreamScan implements 
StreamDataTableScan {
         }
 
         // 4. Build ChainSplits:
-        //    - Snapshot partitions are already filtered to latest per group 
at the pinned snapshot.
-        //    - Delta partitions: include partitions with chain key > latest 
snapshot chain key for
-        //      that group, or all partitions if no snapshot exists for that 
group.
+        //    - Lightweight mode: snapshot partitions are read directly; delta 
partitions are
+        //      included only if their chain key is greater than the latest 
snapshot chain key.
+        //    - Merge mode: for each group, merge the latest snapshot 
partition with delta
+        //      partitions whose chain key is strictly greater than the 
snapshot chain key.
+        //      This allows streaming readers to see deletions/updates that 
span both branches.
+        List<Split> allSplits =
+                mergeSnapshot
+                        ? buildMergedStartingSplits(
+                                snapshotBranch,
+                                deltaBranch,
+                                snapshotSplitsByPartition,
+                                deltaSplitsByPartition,
+                                latestChainPartitionPerGroup)
+                        : buildLightweightStartingSplits(
+                                snapshotBranch,
+                                deltaBranch,
+                                snapshotSplitsByPartition,
+                                deltaSplitsByPartition,
+                                latestChainPartitionPerGroup);
+
+        LOG.info(
+                "ChainTableStreamScan.planStarting [snapshot={}, delta={}]: "
+                        + "{} delta partitions, {} snapshot partitions, "
+                        + "{} latest snapshot groups, {} total splits",
+                snapshotBranch,
+                deltaBranch,
+                deltaSplitsByPartition.size(),
+                snapshotSplitsByPartition.size(),
+                latestChainPartitionPerGroup.size(),
+                allSplits.size());
+
+        startingDone = true;
+        return new DataFilePlan<>(allSplits);
+    }
+
+    /**
+     * Lightweight starting splits: read the latest snapshot partition per 
group directly, and only
+     * include delta partitions whose chain key is strictly greater than the 
latest snapshot chain
+     * key for that group.
+     */
+    private List<Split> buildLightweightStartingSplits(
+            String snapshotBranch,
+            String deltaBranch,
+            Map<BinaryRow, List<DataSplit>> snapshotSplitsByPartition,
+            Map<BinaryRow, List<DataSplit>> deltaSplitsByPartition,
+            Map<Object, BinaryRow> latestChainPartitionPerGroup) {
         List<Split> allSplits = new ArrayList<>();
 
         for (Map.Entry<BinaryRow, List<DataSplit>> entry : 
snapshotSplitsByPartition.entrySet()) {
@@ -303,19 +356,102 @@ public class ChainTableStreamScan implements 
StreamDataTableScan {
             }
         }
 
-        LOG.info(
-                "ChainTableStreamScan.planStarting [snapshot={}, delta={}]: "
-                        + "{} delta partitions, {} snapshot partitions, "
-                        + "{} latest snapshot groups, {} total splits",
-                snapshotBranch,
-                deltaBranch,
-                deltaSplitsByPartition.size(),
-                snapshotSplitsByPartition.size(),
-                latestChainPartitionPerGroup.size(),
-                allSplits.size());
+        return allSplits;
+    }
 
-        startingDone = true;
-        return new DataFilePlan<>(allSplits);
+    /**
+     * Merge-mode starting splits: for each group, merge the latest snapshot 
partition (if any) with
+     * all delta partitions whose chain key is strictly greater than the 
snapshot chain key. Groups
+     * without a snapshot merge all their delta partitions into the latest 
delta partition. This
+     * makes cross-branch deletions and updates visible in the streaming 
starting phase.
+     */
+    private List<Split> buildMergedStartingSplits(
+            String snapshotBranch,
+            String deltaBranch,
+            Map<BinaryRow, List<DataSplit>> snapshotSplitsByPartition,
+            Map<BinaryRow, List<DataSplit>> deltaSplitsByPartition,
+            Map<Object, BinaryRow> latestChainPartitionPerGroup) {
+        List<Split> allSplits = new ArrayList<>();
+
+        // Pre-group delta splits and find the latest delta partition per 
group.
+        Map<Object, List<DataSplit>> deltaSplitsByGroup = new HashMap<>();
+        Map<Object, BinaryRow> latestDeltaPartitionPerGroup = new HashMap<>();
+        for (Map.Entry<BinaryRow, List<DataSplit>> e : 
deltaSplitsByPartition.entrySet()) {
+            BinaryRow deltaPartition = e.getKey();
+            Object groupKey = toGroupKey(deltaPartition);
+            deltaSplitsByGroup
+                    .computeIfAbsent(groupKey, k -> new ArrayList<>())
+                    .addAll(e.getValue());
+
+            BinaryRow currentLatest = 
latestDeltaPartitionPerGroup.get(groupKey);
+            if (currentLatest == null
+                    || chainPartitionComparator.compare(
+                                    
partitionProjector.extractChainPartition(deltaPartition),
+                                    
partitionProjector.extractChainPartition(currentLatest))
+                            > 0) {
+                latestDeltaPartitionPerGroup.put(groupKey, deltaPartition);
+            }
+        }
+
+        // Groups that have a snapshot anchor.
+        for (Map.Entry<Object, BinaryRow> entry : 
latestChainPartitionPerGroup.entrySet()) {
+            Object groupKey = entry.getKey();
+            BinaryRow snapshotPartition = entry.getValue();
+            List<DataSplit> snapshotSplits =
+                    snapshotSplitsByPartition.getOrDefault(
+                            snapshotPartition, Collections.emptyList());
+
+            BinaryRow latestDeltaPartition = 
latestDeltaPartitionPerGroup.get(groupKey);
+            boolean hasDeltaAfterSnapshot =
+                    latestDeltaPartition != null
+                            && chainPartitionComparator.compare(
+                                            
partitionProjector.extractChainPartition(
+                                                    latestDeltaPartition),
+                                            
partitionProjector.extractChainPartition(
+                                                    snapshotPartition))
+                                    > 0;
+
+            List<DataSplit> selectedDeltaSplits = new ArrayList<>();
+            if (hasDeltaAfterSnapshot) {
+                for (DataSplit dataSplit : deltaSplitsByGroup.get(groupKey)) {
+                    BinaryRow deltaPartition = dataSplit.partition();
+                    if (chainPartitionComparator.compare(
+                                    
partitionProjector.extractChainPartition(deltaPartition),
+                                    
partitionProjector.extractChainPartition(snapshotPartition))
+                            > 0) {
+                        selectedDeltaSplits.add(dataSplit);
+                    }
+                }
+            }
+
+            BinaryRow logicalPartition =
+                    hasDeltaAfterSnapshot ? latestDeltaPartition : 
snapshotPartition;
+            allSplits.addAll(
+                    ChainTableUtils.buildChainSplits(
+                            logicalPartition,
+                            snapshotSplits,
+                            selectedDeltaSplits,
+                            snapshotBranch,
+                            deltaBranch));
+        }
+
+        // Delta-only groups: there is no snapshot anchor, so merge all delta 
partitions in the
+        // group into the latest delta partition.
+        for (Map.Entry<Object, List<DataSplit>> entry : 
deltaSplitsByGroup.entrySet()) {
+            Object groupKey = entry.getKey();
+            if (!latestChainPartitionPerGroup.containsKey(groupKey)) {
+                BinaryRow logicalPartition = 
latestDeltaPartitionPerGroup.get(groupKey);
+                allSplits.addAll(
+                        ChainTableUtils.buildChainSplits(
+                                logicalPartition,
+                                Collections.emptyList(),
+                                entry.getValue(),
+                                snapshotBranch,
+                                deltaBranch));
+            }
+        }
+
+        return allSplits;
     }
 
     /**
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/utils/ChainTableUtils.java 
b/paimon-core/src/main/java/org/apache/paimon/utils/ChainTableUtils.java
index d338df5191..498cc4dacc 100644
--- a/paimon-core/src/main/java/org/apache/paimon/utils/ChainTableUtils.java
+++ b/paimon-core/src/main/java/org/apache/paimon/utils/ChainTableUtils.java
@@ -22,12 +22,15 @@ import org.apache.paimon.CoreOptions;
 import org.apache.paimon.codegen.RecordComparator;
 import org.apache.paimon.data.BinaryRow;
 import org.apache.paimon.data.serializer.InternalRowSerializer;
+import org.apache.paimon.io.DataFileMeta;
 import org.apache.paimon.partition.PartitionTimeExtractor;
 import org.apache.paimon.predicate.Predicate;
 import org.apache.paimon.predicate.PredicateBuilder;
 import org.apache.paimon.table.ChainGroupReadTable;
 import org.apache.paimon.table.FallbackReadFileStoreTable;
 import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.source.ChainSplit;
+import org.apache.paimon.table.source.DataSplit;
 import org.apache.paimon.types.RowType;
 
 import java.time.LocalDateTime;
@@ -37,11 +40,15 @@ import java.util.HashMap;
 import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.Set;
 import java.util.function.BiFunction;
 import java.util.regex.Matcher;
 import java.util.regex.Pattern;
 import java.util.stream.Collectors;
 
+import static org.apache.paimon.utils.Preconditions.checkArgument;
+import static org.apache.paimon.utils.Preconditions.checkNotNull;
+
 /** Utils for chain table. */
 public class ChainTableUtils {
 
@@ -365,6 +372,76 @@ public class ChainTableUtils {
         return PredicateBuilder.and(conditions);
     }
 
+    /**
+     * Builds per-bucket {@link ChainSplit}s from the given snapshot and delta 
splits. Files that
+     * originate from the snapshot splits are tagged with {@code 
snapshotBranch}; all other files
+     * are tagged with {@code deltaBranch}.
+     *
+     * @param logicalPartition the logical partition for the resulting 
ChainSplits
+     * @param snapshotSplits splits from the snapshot branch
+     * @param deltaSplits splits from the delta branch
+     * @param snapshotBranch name of the snapshot branch
+     * @param deltaBranch name of the delta branch
+     * @return one ChainSplit per bucket
+     */
+    public static List<ChainSplit> buildChainSplits(
+            BinaryRow logicalPartition,
+            List<DataSplit> snapshotSplits,
+            List<DataSplit> deltaSplits,
+            String snapshotBranch,
+            String deltaBranch) {
+        Set<String> snapshotFileNames =
+                snapshotSplits.stream()
+                        .flatMap(s -> 
s.dataFiles().stream().map(DataFileMeta::fileName))
+                        .collect(Collectors.toSet());
+
+        Map<Integer, List<DataSplit>> bucketSplits = new LinkedHashMap<>();
+        Integer bucketInAll = null;
+        for (DataSplit ds : snapshotSplits) {
+            bucketInAll = addToBucketMap(ds, bucketSplits, bucketInAll);
+        }
+        for (DataSplit ds : deltaSplits) {
+            bucketInAll = addToBucketMap(ds, bucketSplits, bucketInAll);
+        }
+
+        List<ChainSplit> result = new ArrayList<>();
+        for (Map.Entry<Integer, List<DataSplit>> entry : 
bucketSplits.entrySet()) {
+            Map<String, String> fileBranchMapping = new HashMap<>();
+            Map<String, String> fileBucketPathMapping = new HashMap<>();
+            for (DataSplit ds : entry.getValue()) {
+                for (DataFileMeta file : ds.dataFiles()) {
+                    fileBucketPathMapping.put(file.fileName(), 
ds.bucketPath());
+                    String branch =
+                            snapshotFileNames.contains(file.fileName())
+                                    ? snapshotBranch
+                                    : deltaBranch;
+                    fileBranchMapping.put(file.fileName(), branch);
+                }
+            }
+            result.add(
+                    new ChainSplit(
+                            logicalPartition,
+                            entry.getValue().stream()
+                                    .flatMap(ds -> ds.dataFiles().stream())
+                                    .collect(Collectors.toList()),
+                            fileBranchMapping,
+                            fileBucketPathMapping));
+        }
+        return result;
+    }
+
+    private static Integer addToBucketMap(
+            DataSplit ds, Map<Integer, List<DataSplit>> bucketSplits, Integer 
bucketInAll) {
+        Integer totalBuckets = ds.totalBuckets();
+        checkNotNull(totalBuckets, "totalBuckets should not be null");
+        if (bucketInAll != null) {
+            checkArgument(
+                    totalBuckets.equals(bucketInAll), "Inconsistent bucket num 
" + ds.bucket());
+        }
+        bucketSplits.computeIfAbsent(ds.bucket(), k -> new 
ArrayList<>()).add(ds);
+        return totalBuckets;
+    }
+
     /**
      * Validates that the chain table configuration is compatible with 
incremental read paths
      * (streaming read and lookup join). All validation rules for incremental 
reads should be
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FlinkChainTableITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FlinkChainTableITCase.java
index 6a9b2c8058..a89a103ab1 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FlinkChainTableITCase.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FlinkChainTableITCase.java
@@ -68,6 +68,7 @@ import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.TimeUnit;
 import java.util.stream.Collectors;
 
+import static java.lang.String.format;
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
 
@@ -286,21 +287,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'partition.timestamp-formatter' = 'yyyyMMdd 
HH:mm:ss'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_test_hourly', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_test_hourly', 'delta')", db);
-        sql(
-                "ALTER TABLE chain_test_hourly SET ("
-                        + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                        + "  'scan.fallback-delta-branch' = 'delta')");
-        sql(
-                "ALTER TABLE `chain_test_hourly$branch_snapshot` SET ("
-                        + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                        + "  'scan.fallback-delta-branch' = 'delta')");
-        sql(
-                "ALTER TABLE `chain_test_hourly$branch_delta` SET ("
-                        + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                        + "  'scan.fallback-delta-branch' = 'delta')");
+        setupChainTableBranches("chain_test_hourly");
 
         // Write main branch
         sql(
@@ -434,21 +421,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'partition.timestamp-formatter' = 'yyyyMMdd'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_test_partial', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_test_partial', 'delta')", db);
-        sql(
-                "ALTER TABLE chain_test_partial SET ("
-                        + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                        + "  'scan.fallback-delta-branch' = 'delta')");
-        sql(
-                "ALTER TABLE `chain_test_partial$branch_snapshot` SET ("
-                        + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                        + "  'scan.fallback-delta-branch' = 'delta')");
-        sql(
-                "ALTER TABLE `chain_test_partial$branch_delta` SET ("
-                        + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                        + "  'scan.fallback-delta-branch' = 'delta')");
+        setupChainTableBranches("chain_test_partial");
 
         // Write main branch
         sql(
@@ -553,21 +526,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'chain-table.chain-partition-keys' = 'dt'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_test_group', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_test_group', 'delta')", db);
-        sql(
-                "ALTER TABLE chain_test_group SET ("
-                        + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                        + "  'scan.fallback-delta-branch' = 'delta')");
-        sql(
-                "ALTER TABLE `chain_test_group$branch_snapshot` SET ("
-                        + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                        + "  'scan.fallback-delta-branch' = 'delta')");
-        sql(
-                "ALTER TABLE `chain_test_group$branch_delta` SET ("
-                        + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                        + "  'scan.fallback-delta-branch' = 'delta')");
+        setupChainTableBranches("chain_test_group");
 
         // Write main branch
         sql(
@@ -764,6 +723,36 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
         env.execute();
     }
 
+    /**
+     * Write Row data (with RowKind and a group partition column) to a 
specific branch using
+     * DataStream API.
+     */
+    private void writeChangelogToBranchWithRegion(
+            String db, String tableName, String branch, Row... rows) throws 
Exception {
+        FileStoreTable table = paimonTable(tableName + "$branch_" + branch);
+
+        StreamExecutionEnvironment env =
+                streamExecutionEnvironmentBuilder()
+                        .streamingMode()
+                        .checkpointIntervalMs(100)
+                        .parallelism(1)
+                        .build();
+
+        DataStream<Row> stream = env.fromCollection(Arrays.asList(rows));
+
+        new FlinkSinkBuilder(table)
+                .forRow(
+                        stream,
+                        DataTypes.ROW(
+                                DataTypes.FIELD("k", DataTypes.BIGINT()),
+                                DataTypes.FIELD("seq", DataTypes.BIGINT()),
+                                DataTypes.FIELD("v", DataTypes.STRING()),
+                                DataTypes.FIELD("region", DataTypes.STRING()),
+                                DataTypes.FIELD("dt", DataTypes.STRING())))
+                .build();
+        env.execute();
+    }
+
     /**
      * Collect n rows from a streaming iterator with a timeout. If no data 
arrives within
      * timeoutSeconds, the iterator is closed and an AssertionError is thrown. 
This is necessary
@@ -827,18 +816,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + ")");
 
         String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_life_cl', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_life_cl', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_life_cl", "chain_life_cl$branch_snapshot", 
"chain_life_cl$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_life_cl");
 
         // === Phase 1: Delta-only initial data (all inserts) ===
         sql(
@@ -995,19 +973,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'sequence.field' = 'seq'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_restart', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_restart', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_restart", "chain_restart$branch_snapshot", 
"chain_restart$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_restart");
 
         // Configure checkpoint for stateful restart
         org.apache.flink.configuration.Configuration config = 
sEnv.getConfig().getConfiguration();
@@ -1159,18 +1125,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + ")");
 
         String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_overlap', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_overlap', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_overlap", "chain_overlap$branch_snapshot", 
"chain_overlap$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_overlap");
 
         // Write snapshot data: dt=20250807 (snapshot-only) and dt=20250808 
(overlapping)
         sql(
@@ -1222,6 +1177,192 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
         it.close();
     }
 
+    /**
+     * Tests streaming read with {@code 
chain-table.streaming.merge-snapshot=true}. Verifies that
+     * the starting phase merges the latest snapshot partition with later 
delta partitions, so
+     * cross-branch deletes and updates are visible in the initial snapshot.
+     */
+    @ParameterizedTest
+    @ValueSource(strings = {"input", "none"})
+    @Timeout(120)
+    public void testStreamingReadWithMergeSnapshot(String changelogProducer) 
throws Exception {
+        String tableName = "chain_merge_stream_" + changelogProducer;
+        sql(
+                format(
+                        "CREATE TABLE %s ("
+                                + "  k BIGINT, seq BIGINT, v STRING, dt STRING"
+                                + ") PARTITIONED BY (dt) WITH ("
+                                + "  'primary-key' = 'dt,k',"
+                                + "  'bucket-key' = 'k',"
+                                + "  'bucket' = '2',"
+                                + "  'sequence.field' = 'seq',"
+                                + "  'merge-engine' = 'deduplicate',"
+                                + "  'changelog-producer' = '%s',"
+                                + "  'chain-table.enabled' = 'true',"
+                                + "  'chain-table.streaming.merge-snapshot' = 
'true',"
+                                + "  'partition.timestamp-pattern' = '$dt',"
+                                + "  'partition.timestamp-formatter' = 
'yyyyMMdd',"
+                                + "  'continuous.discovery-interval' = '1ms'"
+                                + ")",
+                        tableName, changelogProducer));
+
+        String db = tEnv.getCurrentDatabase();
+        setupChainTableBranches(tableName);
+
+        // Write snapshot data at dt=20250808
+        sql(
+                "INSERT INTO `"
+                        + tableName
+                        + "$branch_snapshot` PARTITION (dt = '20250808')"
+                        + " VALUES (1, 1, 'snap_1'), (2, 1, 'snap_2')");
+
+        // Write delta data spanning dt=20250809 and dt=20250810:
+        // - delete k=1 at dt=20250809
+        // - update k=2: -U old snapshot value at dt=20250809, +U new delta 
value at dt=20250810
+        // - insert k=3 at dt=20250810
+        writeChangelogToBranch(
+                db,
+                tableName,
+                "delta",
+                Row.ofKind(RowKind.DELETE, 1L, 2L, "snap_1", "20250809"),
+                Row.ofKind(RowKind.UPDATE_BEFORE, 2L, 2L, "snap_2", 
"20250809"),
+                Row.ofKind(RowKind.UPDATE_AFTER, 2L, 3L, "delta_2", 
"20250810"),
+                Row.ofKind(RowKind.INSERT, 3L, 1L, "delta_3", "20250810"));
+
+        CloseableIterator<Row> it = sEnv.executeSql("SELECT * FROM " + 
tableName).collect();
+
+        // Starting (merge mode): snapshot@20250808 is anchored to the latest 
delta partition
+        // dt=20250810.
+        // k=1 is deleted; k=2 is updated from snapshot value to delta value; 
k=3 is newly inserted.
+        // The logical partition of the merged ChainSplit is the latest delta 
partition 20250810.
+        // With changelog-producer=input the update is emitted as +U; with 
changelog-producer=none
+        // the upsert result is emitted as +I.
+        String updatedRowKind = "input".equals(changelogProducer) ? "+U" : 
"+I";
+        List<String> startingRows = collectRows(it, 2);
+        assertThat(startingRows)
+                .as(
+                        "Starting with merge-snapshot: cross-branch 
delete/update should be applied, "
+                                + "updated/inserted rows should use the latest 
delta partition")
+                .containsExactlyInAnyOrder(
+                        updatedRowKind + "[2, 3, delta_2, 20250810]",
+                        "+I[3, 1, delta_3, 20250810]");
+
+        // Incremental: write new delta and verify it streams through
+        writeChangelogToBranch(
+                db, tableName, "delta", Row.ofKind(RowKind.INSERT, 4L, 1L, 
"delta_4", "20250811"));
+
+        List<String> incr = collectRows(it, 1);
+        assertThat(incr)
+                .as("Incremental: new delta data should stream through")
+                .containsExactlyInAnyOrder("+I[4, 1, delta_4, 20250811]");
+
+        it.close();
+    }
+
+    /**
+     * Tests streaming read with {@code 
chain-table.streaming.merge-snapshot=true} and a group
+     * partition (region). Verifies that each group is handled independently:
+     *
+     * <ul>
+     *   <li>CN: snapshot anchor at 20250809 + delta at 20250810 (cross-branch 
delete/insert).
+     *   <li>UK: snapshot anchor at 20250808 + delta at 20250809 (cross-branch 
delete).
+     *   <li>US: delta-only group (inserts at 20250811, delete at 20250812, 
later insert at
+     *       20250813).
+     * </ul>
+     */
+    @ParameterizedTest
+    @ValueSource(strings = {"input", "none"})
+    @Timeout(120)
+    public void testStreamingReadWithMergeSnapshotAndGroup(String 
changelogProducer)
+            throws Exception {
+        String tableName = "chain_merge_stream_group_" + changelogProducer;
+        sql(
+                format(
+                        "CREATE TABLE %s ("
+                                + "  k BIGINT, seq BIGINT, v STRING, region 
STRING, dt STRING"
+                                + ") PARTITIONED BY (region, dt) WITH ("
+                                + "  'primary-key' = 'region,dt,k',"
+                                + "  'bucket-key' = 'k',"
+                                + "  'bucket' = '2',"
+                                + "  'sequence.field' = 'seq',"
+                                + "  'merge-engine' = 'deduplicate',"
+                                + "  'changelog-producer' = '%s',"
+                                + "  'chain-table.enabled' = 'true',"
+                                + "  'chain-table.streaming.merge-snapshot' = 
'true',"
+                                + "  'partition.timestamp-pattern' = '$dt',"
+                                + "  'partition.timestamp-formatter' = 
'yyyyMMdd',"
+                                + "  'chain-table.chain-partition-keys' = 
'dt',"
+                                + "  'continuous.discovery-interval' = '1ms'"
+                                + ")",
+                        tableName, changelogProducer));
+
+        String db = tEnv.getCurrentDatabase();
+        setupChainTableBranches(tableName);
+
+        // Snapshot branch: CN and UK have anchors; US has no snapshot.
+        sql(
+                "INSERT INTO `"
+                        + tableName
+                        + "$branch_snapshot`"
+                        + " PARTITION (region = 'CN', dt = '20250809')"
+                        + " VALUES (1, 1, 'cn_snap_1'), (2, 1, 'cn_snap_2')");
+        sql(
+                "INSERT INTO `"
+                        + tableName
+                        + "$branch_snapshot`"
+                        + " PARTITION (region = 'UK', dt = '20250808')"
+                        + " VALUES (21, 1, 'uk_snap_21'), (22, 1, 
'uk_snap_22')");
+
+        // First delta commit:
+        // - CN: at dt=20250810, delete k=1 and insert k=3 (delta > snapshot 
anchor 20250809).
+        // - UK: at dt=20250809, delete k=21 (delta > snapshot anchor 
20250808).
+        // - US: delta-only group, insert k=11 and k=12 at dt=20250811.
+        writeChangelogToBranchWithRegion(
+                db,
+                tableName,
+                "delta",
+                Row.ofKind(RowKind.DELETE, 1L, 2L, "cn_snap_1", "CN", 
"20250810"),
+                Row.ofKind(RowKind.INSERT, 3L, 1L, "cn_delta_3", "CN", 
"20250810"),
+                Row.ofKind(RowKind.INSERT, 11L, 1L, "us_delta_11", "US", 
"20250811"),
+                Row.ofKind(RowKind.INSERT, 12L, 1L, "us_delta_12", "US", 
"20250811"),
+                Row.ofKind(RowKind.DELETE, 21L, 2L, "uk_snap_21", "UK", 
"20250809"));
+        writeChangelogToBranchWithRegion(
+                db,
+                tableName,
+                "delta",
+                Row.ofKind(RowKind.DELETE, 11L, 2L, "us_delta_11", "US", 
"20250812"));
+
+        // Start streaming read (pinned at the second delta commit)
+        CloseableIterator<Row> it = sEnv.executeSql("SELECT * FROM " + 
tableName).collect();
+
+        // Starting (merge mode):
+        List<String> startingRows = collectRows(it, 4);
+        assertThat(startingRows)
+                .as(
+                        "Starting with merge-snapshot and group: each group 
should be handled "
+                                + "independently")
+                .containsExactlyInAnyOrder(
+                        "+I[2, 1, cn_snap_2, CN, 20250810]",
+                        "+I[3, 1, cn_delta_3, CN, 20250810]",
+                        "+I[12, 1, us_delta_12, US, 20250812]",
+                        "+I[22, 1, uk_snap_22, UK, 20250809]");
+
+        // Third delta commit (Phase 2 incremental): US delta-only group 
inserts k=11 at
+        // dt=20250813.
+        writeChangelogToBranchWithRegion(
+                db,
+                tableName,
+                "delta",
+                Row.ofKind(RowKind.INSERT, 11L, 3L, "us_delta_11", "US", 
"20250813"));
+
+        List<String> incr = collectRows(it, 1);
+        assertThat(incr)
+                .as("Incremental: new insert in delta-only group should stream 
through")
+                .containsExactlyInAnyOrder("+I[11, 3, us_delta_11, US, 
20250813]");
+
+        it.close();
+    }
+
     /**
      * T2: Tests that non-default startup modes throw an error for chain table 
streaming read. When
      * scan.mode=latest is specified, an {@link UnsupportedOperationException} 
is thrown with a
@@ -1245,20 +1386,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'partition.timestamp-formatter' = 'yyyyMMdd',"
                         + "  'continuous.discovery-interval' = '1ms'"
                         + ")");
-
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_bypass', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_bypass', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_bypass", "chain_bypass$branch_snapshot", 
"chain_bypass$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_bypass");
 
         // Write data to main table (so snapshots exist for copy() to resolve)
         sql(
@@ -1304,25 +1432,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'continuous.discovery-interval' = '1ms'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_consumer', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_consumer', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_consumer",
-                    "chain_consumer$branch_snapshot",
-                    "chain_consumer$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
-
-        sql(
-                "INSERT INTO `chain_consumer$branch_delta` PARTITION (dt = 
'20250808')"
-                        + " VALUES (1, 1, 'v1'), (2, 1, 'v2')");
+        setupChainTableBranches("chain_consumer");
 
         FileStoreTable table = paimonTable("chain_consumer");
 
@@ -1372,19 +1482,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'continuous.discovery-interval' = '1ms'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_no_cl', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_no_cl', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_no_cl", "chain_no_cl$branch_snapshot", 
"chain_no_cl$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_no_cl");
 
         // Phase 1: Insert initial data into delta branch
         sql(
@@ -1433,21 +1531,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'continuous.discovery-interval' = '1ms'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_stream_group', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_stream_group', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_stream_group",
-                    "chain_stream_group$branch_snapshot",
-                    "chain_stream_group$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_stream_group");
 
         // Write initial delta data for two regions
         sql(
@@ -1520,21 +1604,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'partition.timestamp-formatter' = 'yyyyMMdd'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_restore_all', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_restore_all', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_restore_all",
-                    "chain_restore_all$branch_snapshot",
-                    "chain_restore_all$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_restore_all");
 
         sql(
                 "INSERT INTO `chain_restore_all$branch_delta` PARTITION (dt = 
'20250808')"
@@ -1588,21 +1658,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'partition.timestamp-formatter' = 'yyyyMMdd'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_restore_null', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_restore_null', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_restore_null",
-                    "chain_restore_null$branch_snapshot",
-                    "chain_restore_null$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_restore_null");
 
         sql(
                 "INSERT INTO `chain_restore_null$branch_delta` PARTITION (dt = 
'20250808')"
@@ -1646,22 +1702,8 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'partition.timestamp-formatter' = 'yyyyMMdd',"
                         + "  'continuous.discovery-interval' = '1ms'"
                         + ")");
-
         String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_empty_delta', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_empty_delta', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_empty_delta",
-                    "chain_empty_delta$branch_snapshot",
-                    "chain_empty_delta$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_empty_delta");
 
         // Write ONLY to snapshot branch, delta stays empty
         sql(
@@ -1711,21 +1753,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'continuous.discovery-interval' = '1ms'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_empty_snap', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_empty_snap', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_empty_snap",
-                    "chain_empty_snap$branch_snapshot",
-                    "chain_empty_snap$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_empty_snap");
 
         // Write ONLY to delta branch, snapshot stays empty
         sql(
@@ -1767,19 +1795,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'partition.timestamp-formatter' = 'yyyyMMdd'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_shard', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_shard', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_shard", "chain_shard$branch_snapshot", 
"chain_shard$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_shard");
 
         sql(
                 "INSERT INTO `chain_shard$branch_delta` PARTITION (dt = 
'20250808')"
@@ -1823,21 +1839,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'continuous.discovery-interval' = '1ms'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_both_empty', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_both_empty', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_both_empty",
-                    "chain_both_empty$branch_snapshot",
-                    "chain_both_empty$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_both_empty");
 
         // Both branches are empty — Phase 1 should produce no splits
         FileStoreTable table = paimonTable("chain_both_empty");
@@ -1874,21 +1876,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'continuous.discovery-interval' = '1ms'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_overwrite_p2', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_overwrite_p2', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_overwrite_p2",
-                    "chain_overwrite_p2$branch_snapshot",
-                    "chain_overwrite_p2$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_overwrite_p2");
 
         // Initial delta data
         sql(
@@ -1935,21 +1923,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'partition.timestamp-formatter' = 'yyyyMMdd'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_restore_newdata', 'snapshot')", 
db);
-        sql("CALL sys.create_branch('%s.chain_restore_newdata', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_restore_newdata",
-                    "chain_restore_newdata$branch_snapshot",
-                    "chain_restore_newdata$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_restore_newdata");
 
         // Write initial snapshot + delta data
         sql(
@@ -2078,22 +2052,8 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'partition.timestamp-formatter' = 'yyyyMMdd',"
                         + "  'continuous.discovery-interval' = '1ms'"
                         + ")");
-
         String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_data_filter', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_data_filter', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_data_filter",
-                    "chain_data_filter$branch_snapshot",
-                    "chain_data_filter$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_data_filter");
 
         // Write initial delta data with mixed values of v
         sql(
@@ -2279,21 +2239,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'partition.timestamp-formatter' = 'yyyyMMdd',"
                         + "  'continuous.discovery-interval' = '1ms'"
                         + ")");
-
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_race', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_race', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_race", "chain_race$branch_snapshot", 
"chain_race$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta'"
-                            + ")",
-                    tbl);
-        }
+        setupChainTableBranches("chain_race");
 
         // Step 1: Write delta data at dt=20250808
         sql(
@@ -2357,20 +2303,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + "  'continuous.discovery-interval' = '1ms'"
                         + ")");
 
-        String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_phase2', 'snapshot')", db);
-        sql("CALL sys.create_branch('%s.chain_phase2', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_phase2", "chain_phase2$branch_snapshot", 
"chain_phase2$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta'"
-                            + ")",
-                    tbl);
-        }
+        setupChainTableBranches("chain_phase2");
 
         // Write snapshot data at dt=20250808
         sql(
@@ -3387,20 +3320,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
                         + ")");
 
         String db = tEnv.getCurrentDatabase();
-        sql("CALL sys.create_branch('%s.chain_bucket_filter', 'snapshot')", 
db);
-        sql("CALL sys.create_branch('%s.chain_bucket_filter', 'delta')", db);
-        for (String tbl :
-                new String[] {
-                    "chain_bucket_filter",
-                    "chain_bucket_filter$branch_snapshot",
-                    "chain_bucket_filter$branch_delta"
-                }) {
-            sql(
-                    "ALTER TABLE `%s` SET ("
-                            + "  'scan.fallback-snapshot-branch' = 'snapshot',"
-                            + "  'scan.fallback-delta-branch' = 'delta')",
-                    tbl);
-        }
+        setupChainTableBranches("chain_bucket_filter");
 
         // Write main branch
         sql(
@@ -3410,7 +3330,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
         // Write delta data across many keys to guarantee both buckets are 
populated
         for (int i = 1; i <= 50; i++) {
             sql(
-                    String.format(
+                    format(
                             "INSERT INTO `chain_bucket_filter$branch_delta`"
                                     + " PARTITION (dt = '%d') VALUES (%d, 1, 
'v%d')",
                             20250809 + (i % 5), i, i));
@@ -3609,7 +3529,7 @@ public class FlinkChainTableITCase extends 
CatalogITCaseBase {
 
         // Submit lookup join job BEFORE inserting source data
         String query =
-                String.format(
+                format(
                         "INSERT INTO sink_refresh "
                                 + "SELECT S.id, D.k, D.v "
                                 + "FROM source_refresh AS S "

Reply via email to