diff --git a/docs/docs/primary-key-table/chain-table.mdx b/docs/docs/primary-key-table/chain-table.mdx index 955643469277..031c0a03e406 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 401d451b8693..0ca452170231 100644 --- a/docs/generated/core_configuration.html +++ b/docs/generated/core_configuration.html @@ -158,6 +158,12 @@ Boolean Whether enabled chain table. + +
chain-table.streaming.merge-snapshot
+ false + Boolean + 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. +
changelog-file.compression
(none) 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 b353b5632790..61b2e87c5e1d 100644 --- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java +++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java @@ -291,6 +291,22 @@ public InlineElement getDescription() { + "suffix of the table's partition keys. Comma-separated. " + "If not set, all partition keys participate in chain."); + public static final ConfigOption 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"; @@ -4148,6 +4164,10 @@ public List chainTableChainPartitionKeys() { 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 1428d26e9d5f..477d7f81f545 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 @@ -59,7 +59,6 @@ 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 @@ -343,7 +342,6 @@ public Plan plan() { for (List deltaPartitionsInGroup : groupedDeltaPartitions.values()) { // Sort delta by chain dimension ascending. - // chainPartitionForCompare avoids copying BinaryRow in the comparator hot path. deltaPartitionsInGroup.sort( (a, b) -> chainPartitionComparator.compare( @@ -415,69 +413,26 @@ public Plan plan() { deltaScan.withPartitionFilter(selectedDeltaPartitions); } - List subSplits = deltaScan.plan().splits(); - Set snapshotFileNames = new HashSet<>(); + List deltaSubSplits = + deltaScan.plan().splits().stream() + .map(s -> (DataSplit) s) + .collect(Collectors.toList()); + List snapshotSubSplits = new ArrayList<>(); if (partitionPairs.getValue() != null) { snapshotScan.withPartitionFilter( Collections.singletonList(partitionPairs.getValue())); - List mainSubSplits = snapshotScan.plan().splits(); - snapshotFileNames = - mainSubSplits.stream() - .flatMap( - s -> - ((DataSplit) s) - .dataFiles().stream() - .map( - DataFileMeta - ::fileName)) - .collect(Collectors.toSet()); - subSplits.addAll(mainSubSplits); - } - Map> 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> entry : bucketSplits.entrySet()) { - HashMap fileBucketPathMapping = new HashMap<>(); - HashMap fileBranchMapping = new HashMap<>(); - List 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 fcd07cb26bd0..5812a2611bc7 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 TableScan.Plan plan() { * come after it. Older snapshot partitions are excluded. Each primary key appears exactly once * under its natural partition. * - *

Unlike batch full scan, anchor-based chain merging is not performed. This keeps Phase 1 - * lightweight for long-running jobs. + *

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 @@ private TableScan.Plan planStarting() { } // 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 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 buildLightweightStartingSplits( + String snapshotBranch, + String deltaBranch, + Map> snapshotSplitsByPartition, + Map> deltaSplitsByPartition, + Map latestChainPartitionPerGroup) { List allSplits = new ArrayList<>(); for (Map.Entry> entry : snapshotSplitsByPartition.entrySet()) { @@ -303,19 +356,102 @@ private TableScan.Plan planStarting() { } } - 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 buildMergedStartingSplits( + String snapshotBranch, + String deltaBranch, + Map> snapshotSplitsByPartition, + Map> deltaSplitsByPartition, + Map latestChainPartitionPerGroup) { + List allSplits = new ArrayList<>(); + + // Pre-group delta splits and find the latest delta partition per group. + Map> deltaSplitsByGroup = new HashMap<>(); + Map latestDeltaPartitionPerGroup = new HashMap<>(); + for (Map.Entry> 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 entry : latestChainPartitionPerGroup.entrySet()) { + Object groupKey = entry.getKey(); + BinaryRow snapshotPartition = entry.getValue(); + List 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 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> 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 d338df5191e4..498cc4dacc2c 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.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.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 static Predicate createGroupChainPredicate( 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 buildChainSplits( + BinaryRow logicalPartition, + List snapshotSplits, + List deltaSplits, + String snapshotBranch, + String deltaBranch) { + Set snapshotFileNames = + snapshotSplits.stream() + .flatMap(s -> s.dataFiles().stream().map(DataFileMeta::fileName)) + .collect(Collectors.toSet()); + + Map> 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 result = new ArrayList<>(); + for (Map.Entry> entry : bucketSplits.entrySet()) { + Map fileBranchMapping = new HashMap<>(); + Map 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> 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 6a9b2c8058d9..a89a103ab1ce 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.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 void testHourlyChainTable() throws Exception { + " '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 void testChainTableWithPartialUpdate() throws Exception { + " '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 void testChainTableWithGroupPartition() throws Exception { + " '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 @@ private void writeChangelogToBranch(String db, String tableName, String branch, 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 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 void testStreamingReadChainTableLifecycleWithInputChangelog() throws Exce + ")"); 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 void testStreamingReadChainTableStatefulRestart() throws Exception { + " '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 void testStreamingReadWithSnapshotDeltaOverlap() throws Exception { + ")"); 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 void testStreamingReadWithSnapshotDeltaOverlap() throws Exception { 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 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 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 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: + * + *

    + *
  • CN: snapshot anchor at 20250809 + delta at 20250810 (cross-branch delete/insert). + *
  • UK: snapshot anchor at 20250808 + delta at 20250809 (cross-branch delete). + *
  • US: delta-only group (inserts at 20250811, delete at 20250812, later insert at + * 20250813). + *
+ */ + @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 it = sEnv.executeSql("SELECT * FROM " + tableName).collect(); + + // Starting (merge mode): + List 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 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 void testStreamingReadRejectsNonDefaultStartup() throws Exception { + " '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 void testStreamingReadRejectsConsumerMode() throws Exception { + " '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 void testStreamingReadWithNoChangelogProducer() throws Exception { + " '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 void testStreamingReadWithGroupPartition() throws Exception { + " '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 void testRestoreScanAll() throws Exception { + " '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 void testRestoreNullScanAll() throws Exception { + " '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 void testStreamingReadEmptyDelta() throws Exception { + " '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 void testStreamingReadEmptySnapshot() throws Exception { + " '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 void testWithShardForwarding() throws Exception { + " '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 void testStreamingReadBothBranchesEmpty() throws Exception { + " '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 void testStreamingReadDeltaOverwriteInPhase2() throws Exception { + " '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 void testStreamingReadRestoreAfterNewData() throws Exception { + " '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 void testStreamingReadWithNonPartitionFilter() throws Exception { + " '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 void testStreamingReadSnapshotBranchRaceCondition() throws Exception { + " '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 void testStreamingReadPhase2BranchAwareRead() throws Exception { + " '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 void testStreamingReadBucketFilter() throws Exception { + ")"); 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 void testStreamingReadBucketFilter() throws Exception { // 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 void testLookupJoinRefresh(boolean asyncRefresh) throws Exception { // 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 "