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 "