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 fe4934979d [flink] Fix postpone bucket not handled for
IncrementalSplit (#8698)
fe4934979d is described below
commit fe4934979df99ff4e3f6c9d08c41222f65c18c18
Author: Vova Kolmakov <[email protected]>
AuthorDate: Thu Jul 16 21:04:52 2026 +0700
[flink] Fix postpone bucket not handled for IncrementalSplit (#8698)
---
.../source/ContinuousFileSplitEnumerator.java | 33 ++++++++--
.../source/ContinuousFileSplitEnumeratorTest.java | 76 ++++++++++++++++++++++
2 files changed, 104 insertions(+), 5 deletions(-)
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/ContinuousFileSplitEnumerator.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/ContinuousFileSplitEnumerator.java
index a8c77aab7f..75c75cfddb 100644
---
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/ContinuousFileSplitEnumerator.java
+++
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/ContinuousFileSplitEnumerator.java
@@ -22,6 +22,7 @@ import org.apache.paimon.annotation.VisibleForTesting;
import org.apache.paimon.flink.source.assigners.FIFOSplitAssigner;
import org.apache.paimon.flink.source.assigners.PreAssignSplitAssigner;
import org.apache.paimon.flink.source.assigners.SplitAssigner;
+import org.apache.paimon.io.DataFileMeta;
import org.apache.paimon.postpone.PostponeBucketFileStoreWrite;
import org.apache.paimon.table.BucketMode;
import org.apache.paimon.table.sink.ChannelComputer;
@@ -327,13 +328,26 @@ public class ContinuousFileSplitEnumerator
}
protected int assignSuggestedTask(FileStoreSourceSplit split) {
+ int task;
if (split.split() instanceof DataSplit) {
- return assignSuggestedTask((DataSplit) split.split());
+ task = assignSuggestedTask((DataSplit) split.split());
} else if (split.split() instanceof ChainSplit) {
- return assignSuggestedTask((ChainSplit) split.split());
+ task = assignSuggestedTask((ChainSplit) split.split());
} else {
- return assignSuggestedTask((IncrementalSplit) split.split());
+ task = assignSuggestedTask((IncrementalSplit) split.split());
}
+
+ // Split assigners keep splits in a map keyed by task, but only ever
hand out splits for
+ // tasks within the parallelism, so a task outside of it silently
strands the split instead
+ // of failing.
+ int parallelism = context.currentParallelism();
+ checkArgument(
+ task >= 0 && task < parallelism,
+ "Split %s is suggested to task %s, which is out of the
parallelism [0, %s). This is unexpected.",
+ split.splitId(),
+ task,
+ parallelism);
+ return task;
}
protected int assignSuggestedTask(DataSplit split) {
@@ -358,8 +372,17 @@ public class ContinuousFileSplitEnumerator
protected int assignSuggestedTask(IncrementalSplit split) {
int parallelism = context.currentParallelism();
- // TODO how to deal with postpone bucket?
- int bucketId = split.bucket();
+ int bucketId;
+ if (split.bucket() == BucketMode.POSTPONE_BUCKET) {
+ // A diff which only removes files has no after files, so fall
back to the before files.
+ List<DataFileMeta> files =
+ split.afterFiles().isEmpty() ? split.beforeFiles() :
split.afterFiles();
+ bucketId =
+
PostponeBucketFileStoreWrite.getWriteId(files.get(0).fileName()) % parallelism;
+ } else {
+ bucketId = split.bucket();
+ }
+
if (shuffleBucketWithPartition) {
return ChannelComputer.select(split.partition(), bucketId,
parallelism);
} else {
diff --git
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/ContinuousFileSplitEnumeratorTest.java
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/ContinuousFileSplitEnumeratorTest.java
index b1ec92251e..0a50038bb4 100644
---
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/ContinuousFileSplitEnumeratorTest.java
+++
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/ContinuousFileSplitEnumeratorTest.java
@@ -19,8 +19,10 @@
package org.apache.paimon.flink.source;
import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.table.BucketMode;
import org.apache.paimon.table.source.DataFilePlan;
import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.table.source.IncrementalSplit;
import org.apache.paimon.table.source.StreamTableScan;
import org.apache.paimon.table.source.TableScan;
@@ -41,6 +43,7 @@ import java.util.UUID;
import java.util.concurrent.atomic.AtomicReference;
import static
org.apache.flink.connector.testutils.source.reader.TestingSplitEnumeratorContext.SplitAssignmentState;
+import static org.apache.paimon.io.DataFileTestUtils.fromMinMax;
import static org.apache.paimon.io.DataFileTestUtils.row;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatCode;
@@ -855,6 +858,79 @@ public class ContinuousFileSplitEnumeratorTest
Assertions.assertThat(enumerator.nextSnapshotId).isEqualTo(3);
}
+ @Test
+ public void testPostponeBucketIncrementalSplitAssignedByWriteId() {
+ int parallelism = 3;
+ ContinuousFileSplitEnumerator enumerator =
buildEnumerator(parallelism);
+
+ IncrementalSplit split =
+ createPostponeIncrementalSplit(
+ Collections.emptyList(),
Collections.singletonList(postponeFile(7)));
+
+ assertThat(enumerator.assignSuggestedTask(split)).isEqualTo(7 %
parallelism);
+ }
+
+ @Test
+ public void testPostponeBucketIncrementalSplitWithoutAfterFiles() {
+ int parallelism = 3;
+ ContinuousFileSplitEnumerator enumerator =
buildEnumerator(parallelism);
+
+ // a diff which only deletes files has no after files, so the write id
must be taken from
+ // the before files
+ IncrementalSplit split =
+ createPostponeIncrementalSplit(
+ Collections.singletonList(postponeFile(7)),
Collections.emptyList());
+
+ assertThat(enumerator.assignSuggestedTask(split)).isEqualTo(7 %
parallelism);
+ }
+
+ @Test
+ public void testIncrementalSplitWithRealBucketAssignedByBucket() {
+ int parallelism = 3;
+ ContinuousFileSplitEnumerator enumerator =
buildEnumerator(parallelism);
+
+ IncrementalSplit split =
+ new IncrementalSplit(
+ 1L,
+ row(1),
+ 4,
+ parallelism,
+ Collections.emptyList(),
+ null,
+ Collections.singletonList(postponeFile(7)),
+ null,
+ true);
+
+ assertThat(enumerator.assignSuggestedTask(split)).isEqualTo(4 %
parallelism);
+ }
+
+ private ContinuousFileSplitEnumerator buildEnumerator(int parallelism) {
+ return new Builder()
+
.setSplitEnumeratorContext(getSplitEnumeratorContext(parallelism))
+ .setScan(new MockScan(new TreeMap<>()))
+ .build();
+ }
+
+ private static IncrementalSplit createPostponeIncrementalSplit(
+ List<DataFileMeta> beforeFiles, List<DataFileMeta> afterFiles) {
+ return new IncrementalSplit(
+ 1L,
+ row(1),
+ BucketMode.POSTPONE_BUCKET,
+ BucketMode.POSTPONE_BUCKET,
+ beforeFiles,
+ null,
+ afterFiles,
+ null,
+ true);
+ }
+
+ /** Builds a file name matching the prefix written by {@code
PostponeBucketFileStoreWrite}. */
+ private static DataFileMeta postponeFile(int writeId) {
+ return fromMinMax(
+ String.format("data-u-%s-s-%d-w-0-0.parquet",
UUID.randomUUID(), writeId), 0, 0);
+ }
+
private void triggerCheckpointAndComplete(
ContinuousFileSplitEnumerator enumerator, long checkpointId)
throws Exception {
enumerator.snapshotState(checkpointId);