deltaPartitions =
@@ -449,7 +471,10 @@ public Plan plan() {
snapshotSubSplits,
deltaSubSplits,
options.scanFallbackSnapshotBranch(),
- options.scanFallbackDeltaBranch()));
+ options.scanFallbackDeltaBranch(),
+ keyComparator,
+ targetSplitSize,
+ openFileCost));
}
}
}
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 498cc4dacc2c..35659f06d8b9 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
@@ -21,8 +21,11 @@
import org.apache.paimon.CoreOptions;
import org.apache.paimon.codegen.RecordComparator;
import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.serializer.InternalRowSerializer;
import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.mergetree.SortedRun;
+import org.apache.paimon.mergetree.compact.IntervalPartition;
import org.apache.paimon.partition.PartitionTimeExtractor;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
@@ -33,15 +36,20 @@
import org.apache.paimon.table.source.DataSplit;
import org.apache.paimon.types.RowType;
+import javax.annotation.Nullable;
+
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.ArrayList;
+import java.util.Collections;
+import java.util.Comparator;
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.function.Function;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
@@ -377,6 +385,9 @@ public static Predicate createGroupChainPredicate(
* originate from the snapshot splits are tagged with {@code snapshotBranch}; all other files
* are tagged with {@code deltaBranch}.
*
+ * This overload produces one {@link ChainSplit} per bucket (no key-range splitting). It is
+ * used by streaming read, where splitting a bucket by key range is not desired.
+ *
* @param logicalPartition the logical partition for the resulting ChainSplits
* @param snapshotSplits splits from the snapshot branch
* @param deltaSplits splits from the delta branch
@@ -390,6 +401,49 @@ public static List buildChainSplits(
List deltaSplits,
String snapshotBranch,
String deltaBranch) {
+ return buildChainSplits(
+ logicalPartition,
+ snapshotSplits,
+ deltaSplits,
+ snapshotBranch,
+ deltaBranch,
+ null,
+ 0L,
+ 0L);
+ }
+
+ /**
+ * Builds {@link ChainSplit}s from the given snapshot and delta splits, optionally splitting
+ * each bucket's files into multiple splits by key range to improve batch read parallelism.
+ *
+ * When {@code keyComparator} is non-null, each bucket's snapshot and delta files are split
+ * via interval partition: files with intersecting key ranges always go to the same split, so
+ * that all versions of a key across the snapshot and delta branches are merged together within
+ * one split. Key-disjoint sections are then packed into splits up to {@code targetSplitSize}.
+ * This is the same invariant the ordinary primary-key table relies on (see {@link
+ * org.apache.paimon.table.source.MergeTreeSplitGenerator}), except the chain path always goes
+ * through the interval partition because snapshot and delta files from different branches may
+ * have intersecting key ranges.
+ *
+ *
When {@code keyComparator} is null, each bucket produces a single {@link ChainSplit} (the
+ * original behavior, used by streaming).
+ *
+ * @param keyComparator primary-key comparator used to order files by min/max key; null disables
+ * key-range splitting
+ * @param targetSplitSize target size (in bytes) of each split; ignored when {@code
+ * keyComparator} is null
+ * @param openFileCost per-file open cost accounted for in packing; ignored when {@code
+ * keyComparator} is null
+ */
+ public static List buildChainSplits(
+ BinaryRow logicalPartition,
+ List snapshotSplits,
+ List deltaSplits,
+ String snapshotBranch,
+ String deltaBranch,
+ @Nullable Comparator keyComparator,
+ long targetSplitSize,
+ long openFileCost) {
Set snapshotFileNames =
snapshotSplits.stream()
.flatMap(s -> s.dataFiles().stream().map(DataFileMeta::fileName))
@@ -408,6 +462,7 @@ public static List buildChainSplits(
for (Map.Entry> entry : bucketSplits.entrySet()) {
Map fileBranchMapping = new HashMap<>();
Map fileBucketPathMapping = new HashMap<>();
+ List bucketFiles = new ArrayList<>();
for (DataSplit ds : entry.getValue()) {
for (DataFileMeta file : ds.dataFiles()) {
fileBucketPathMapping.put(file.fileName(), ds.bucketPath());
@@ -416,20 +471,82 @@ public static List buildChainSplits(
? snapshotBranch
: deltaBranch;
fileBranchMapping.put(file.fileName(), branch);
+ bucketFiles.add(file);
}
}
- result.add(
- new ChainSplit(
- logicalPartition,
- entry.getValue().stream()
- .flatMap(ds -> ds.dataFiles().stream())
- .collect(Collectors.toList()),
- fileBranchMapping,
- fileBucketPathMapping));
+
+ List> fileGroups;
+ if (keyComparator == null) {
+ fileGroups = Collections.singletonList(bucketFiles);
+ } else {
+ fileGroups =
+ splitBucketByKeyRange(
+ bucketFiles, keyComparator, targetSplitSize, openFileCost);
+ }
+
+ for (List groupFiles : fileGroups) {
+ result.add(
+ new ChainSplit(
+ logicalPartition,
+ groupFiles,
+ subMapping(fileBranchMapping, groupFiles),
+ subMapping(fileBucketPathMapping, groupFiles)));
+ }
}
return result;
}
+ /**
+ * Splits a bucket's files (snapshot + delta combined) into key-range groups. The interval
+ * partition first sorts files by (minKey, maxKey) and groups key-overlapping files into the
+ * same section, guaranteeing all versions of a key stay in one split. Sections are key-disjoint
+ * and key-ordered, so packing consecutive sections up to {@code targetSplitSize} preserves that
+ * invariant.
+ */
+ private static List> splitBucketByKeyRange(
+ List files,
+ Comparator keyComparator,
+ long targetSplitSize,
+ long openFileCost) {
+ List> sections = new ArrayList<>();
+ for (List section : new IntervalPartition(files, keyComparator).partition()) {
+ List sectionFiles = new ArrayList<>();
+ for (SortedRun run : section) {
+ sectionFiles.addAll(run.files());
+ }
+ sections.add(sectionFiles);
+ }
+
+ Function, Long> weightFunc =
+ section -> Math.max(totalSize(section), openFileCost);
+ return BinPacking.packForOrdered(sections, weightFunc, targetSplitSize).stream()
+ .map(ChainTableUtils::flatFiles)
+ .collect(Collectors.toList());
+ }
+
+ private static long totalSize(List files) {
+ long size = 0L;
+ for (DataFileMeta file : files) {
+ size += file.fileSize();
+ }
+ return size;
+ }
+
+ private static List flatFiles(List> packedSections) {
+ List files = new ArrayList<>();
+ packedSections.forEach(files::addAll);
+ return files;
+ }
+
+ private static Map subMapping(
+ Map mapping, List files) {
+ Map sub = new HashMap<>();
+ for (DataFileMeta file : files) {
+ sub.put(file.fileName(), mapping.get(file.fileName()));
+ }
+ return sub;
+ }
+
private static Integer addToBucketMap(
DataSplit ds, Map> bucketSplits, Integer bucketInAll) {
Integer totalBuckets = ds.totalBuckets();
diff --git a/paimon-core/src/test/java/org/apache/paimon/utils/ChainTableUtilsTest.java b/paimon-core/src/test/java/org/apache/paimon/utils/ChainTableUtilsTest.java
index 971407b92815..50e74d4d4e7f 100644
--- a/paimon-core/src/test/java/org/apache/paimon/utils/ChainTableUtilsTest.java
+++ b/paimon-core/src/test/java/org/apache/paimon/utils/ChainTableUtilsTest.java
@@ -24,9 +24,15 @@
import org.apache.paimon.data.BinaryRowWriter;
import org.apache.paimon.data.BinaryString;
import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.data.Timestamp;
+import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.manifest.FileSource;
import org.apache.paimon.options.Options;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
+import org.apache.paimon.stats.StatsTestUtils;
+import org.apache.paimon.table.source.ChainSplit;
+import org.apache.paimon.table.source.DataSplit;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowType;
@@ -37,9 +43,12 @@
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
+import java.util.Set;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@@ -767,4 +776,145 @@ public void testFindFirstLatestPartitionsWithMultiChainKeys() {
assertThat(getString(matched2, 0)).isEqualTo("20250809");
assertThat(getString(matched2, 1)).isEqualTo("02");
}
+
+ @Test
+ public void testBuildChainSplitsWithKeyRangeSplitting() {
+ // Single bucket. Snapshot files S1=[a,b], S2=[m,n]; delta files D1=[a,b] (overlaps S1),
+ // D2=[m,n] (overlaps S2). The two sections are key-disjoint, so they form two splits;
+ // within each section the overlapping snapshot and delta files must stay together so that
+ // all versions of a key are merged in the same split.
+ BinaryRow partition = row(Lists.newArrayList("p0"));
+ String snapBranch = "snap";
+ String deltaBranch = "delta";
+
+ DataFileMeta s1 = makeFile("S1", "a", "b", 100);
+ DataFileMeta s2 = makeFile("S2", "m", "n", 100);
+ DataFileMeta d1 = makeFile("D1", "a", "b", 100);
+ DataFileMeta d2 = makeFile("D2", "m", "n", 100);
+
+ DataSplit snapshotSplit = dataSplit(partition, 0, Lists.newArrayList(s1, s2));
+ DataSplit deltaSplit = dataSplit(partition, 0, Lists.newArrayList(d1, d2));
+
+ // Small targetSplitSize so each section forms its own split.
+ List splits =
+ ChainTableUtils.buildChainSplits(
+ partition,
+ Collections.singletonList(snapshotSplit),
+ Collections.singletonList(deltaSplit),
+ snapBranch,
+ deltaBranch,
+ KEY_COMPARATOR,
+ 1L,
+ 1L);
+
+ assertThat(splits).hasSize(2);
+
+ List> groups =
+ splits.stream().map(ChainTableUtilsTest::fileNames).collect(Collectors.toList());
+ // Overlapping files are kept together: {S1,D1} and {S2,D2}, never mixed.
+ assertThat(groups)
+ .containsExactlyInAnyOrder(
+ new HashSet<>(Arrays.asList("S1", "D1")),
+ new HashSet<>(Arrays.asList("S2", "D2")));
+
+ // Branch tagging is preserved per file across the resulting splits.
+ ChainSplit s1Split =
+ splits.stream().filter(s -> fileNames(s).contains("S1")).findFirst().get();
+ assertThat(s1Split.fileBranchMapping().get("S1")).isEqualTo(snapBranch);
+ assertThat(s1Split.fileBranchMapping().get("D1")).isEqualTo(deltaBranch);
+ ChainSplit s2Split =
+ splits.stream().filter(s -> fileNames(s).contains("S2")).findFirst().get();
+ assertThat(s2Split.fileBranchMapping().get("S2")).isEqualTo(snapBranch);
+ assertThat(s2Split.fileBranchMapping().get("D2")).isEqualTo(deltaBranch);
+ }
+
+ @Test
+ public void testBuildChainSplitsKeyRangeSplittingLargeTargetKeepsOneSplit() {
+ // With a large targetSplitSize all key-disjoint sections are packed into a single split.
+ BinaryRow partition = row(Lists.newArrayList("p0"));
+ DataFileMeta s1 = makeFile("S1", "a", "b", 100);
+ DataFileMeta s2 = makeFile("S2", "m", "n", 100);
+ DataFileMeta d1 = makeFile("D1", "a", "b", 100);
+ DataFileMeta d2 = makeFile("D2", "m", "n", 100);
+
+ DataSplit snapshotSplit = dataSplit(partition, 0, Lists.newArrayList(s1, s2));
+ DataSplit deltaSplit = dataSplit(partition, 0, Lists.newArrayList(d1, d2));
+
+ List splits =
+ ChainTableUtils.buildChainSplits(
+ partition,
+ Collections.singletonList(snapshotSplit),
+ Collections.singletonList(deltaSplit),
+ "snap",
+ "delta",
+ KEY_COMPARATOR,
+ Long.MAX_VALUE,
+ 1L);
+
+ assertThat(splits).hasSize(1);
+ assertThat(fileNames(splits.get(0))).containsExactlyInAnyOrder("S1", "S2", "D1", "D2");
+ }
+
+ @Test
+ public void testBuildChainSplitsWithoutKeyRangeSplitting() {
+ // A null keyComparator keeps the original one-split-per-bucket behavior (used by
+ // streaming).
+ BinaryRow partition = row(Lists.newArrayList("p0"));
+ DataFileMeta s1 = makeFile("S1", "a", "b", 100);
+ DataFileMeta d1 = makeFile("D1", "a", "b", 100);
+ DataSplit snapshotSplit = dataSplit(partition, 0, Lists.newArrayList(s1));
+ DataSplit deltaSplit = dataSplit(partition, 0, Lists.newArrayList(d1));
+
+ List splits =
+ ChainTableUtils.buildChainSplits(
+ partition,
+ Collections.singletonList(snapshotSplit),
+ Collections.singletonList(deltaSplit),
+ "snap",
+ "delta",
+ null,
+ 0L,
+ 0L);
+
+ assertThat(splits).hasSize(1);
+ assertThat(fileNames(splits.get(0))).containsExactlyInAnyOrder("S1", "D1");
+ }
+
+ private static DataFileMeta makeFile(String name, String min, String max, long size) {
+ return DataFileMeta.create(
+ name,
+ size,
+ 10L,
+ row(Lists.newArrayList(min)),
+ row(Lists.newArrayList(max)),
+ StatsTestUtils.newEmptySimpleStats(),
+ StatsTestUtils.newEmptySimpleStats(),
+ 0L,
+ 9L,
+ 0L,
+ 0,
+ Collections.emptyList(),
+ Timestamp.fromEpochMillis(1000L),
+ 0L,
+ null,
+ FileSource.APPEND,
+ null,
+ null,
+ null,
+ null);
+ }
+
+ private static DataSplit dataSplit(BinaryRow partition, int bucket, List files) {
+ return DataSplit.builder()
+ .withPartition(partition)
+ .withBucket(bucket)
+ .withBucketPath("bucket_" + bucket)
+ .withTotalBuckets(1)
+ .withDataFiles(files)
+ .build();
+ }
+
+ private static Set fileNames(ChainSplit split) {
+ return split.dataFiles().stream().map(DataFileMeta::fileName).collect(Collectors.toSet());
+ }
}