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 3130869e62 [flink] Avoid slot deadlock when sorting normalized index 
keys (#9394)
3130869e62 is described below

commit 3130869e6248983d89a1c2e49d9142dc4420c869
Author: YeJunHao <[email protected]>
AuthorDate: Wed Aug 26 14:51:56 2026 +0800

    [flink] Avoid slot deadlock when sorting normalized index keys (#9394)
---
 .../flink/globalindex/SortedIndexTopoBuilder.java  | 24 ++++++++++++++-
 .../globalindex/SortedIndexTopoBuilderTest.java    | 34 ++++++++++++++++++++++
 2 files changed, 57 insertions(+), 1 deletion(-)

diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilder.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilder.java
index cda2275903..c5bcb89b0a 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilder.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilder.java
@@ -66,10 +66,14 @@ import org.apache.paimon.types.DataTypes;
 import org.apache.paimon.types.RowType;
 import org.apache.paimon.utils.Range;
 
+import org.apache.flink.api.common.functions.Partitioner;
 import org.apache.flink.api.java.functions.KeySelector;
 import org.apache.flink.streaming.api.datastream.DataStream;
 import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
 import org.apache.flink.streaming.api.operators.OneInputStreamOperatorFactory;
+import org.apache.flink.streaming.api.transformations.PartitionTransformation;
+import org.apache.flink.streaming.api.transformations.StreamExchangeMode;
+import org.apache.flink.streaming.runtime.partitioner.CustomPartitionerWrapper;
 import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
 import org.apache.flink.table.data.RowData;
 
@@ -96,6 +100,8 @@ public class SortedIndexTopoBuilder {
 
     private static final String BUILD_TASK_ID_FIELD = 
"_SORTED_INDEX_BUILD_TASK_ID";
     private static final int BUILD_TASK_ID_FIELD_ID = -1;
+    static final Partitioner<Integer> BUILD_TASK_PARTITIONER =
+            (taskId, numPartitions) -> Math.floorMod(taskId, numPartitions);
     private static final HashSet<String> SUPPORTED_INDEX_TYPES =
             new HashSet<>(Arrays.asList("btree", "bitmap", "multivalue"));
 
@@ -384,7 +390,12 @@ public class SortedIndexTopoBuilder {
             RowType readType,
             int parallelism) {
         if (!identity) {
-            return input.keyBy((KeySelector<InternalRow, Integer>) row -> 
row.getInt(taskIdPos))
+            // Only co-location by build task is required here. Using keyBy in 
batch mode makes
+            // Flink insert another managed-memory sorter before Paimon's 
heap-and-spill sorter.
+            // Materialize the repartition result so batch scheduling can 
reuse slots between the
+            // reader and sorter stages. A pipelined exchange can deadlock 
when maxSlot is smaller
+            // than their combined parallelism.
+            return partitionByBuildTask(input, taskIdPos)
                     .transform(
                             "Sort Normalized Index Keys",
                             InternalTypeInfo.fromRowType(readType),
@@ -424,6 +435,17 @@ public class SortedIndexTopoBuilder {
                 .setParallelism(parallelism);
     }
 
+    static DataStream<InternalRow> partitionByBuildTask(
+            DataStream<InternalRow> input, int taskIdPos) {
+        KeySelector<InternalRow, Integer> keySelector = row -> 
row.getInt(taskIdPos);
+        return new DataStream<>(
+                input.getExecutionEnvironment(),
+                new PartitionTransformation<>(
+                        input.getTransformation(),
+                        new CustomPartitionerWrapper<>(BUILD_TASK_PARTITIONER, 
keySelector),
+                        StreamExchangeMode.BATCH));
+    }
+
     static int calculateParallelism(
             List<SortedBuildTask> buildTasks, long recordsPerRange, int 
maxParallelism) {
         long totalRecords = 0;
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilderTest.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilderTest.java
index a29eb9e7be..38f1b59abf 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilderTest.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilderTest.java
@@ -18,7 +18,11 @@
 
 package org.apache.paimon.flink.globalindex;
 
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.BinaryRowWriter;
+import org.apache.paimon.data.InternalRow;
 import 
org.apache.paimon.flink.globalindex.SortedIndexTopoBuilder.SortedBuildTask;
+import org.apache.paimon.flink.utils.InternalTypeInfo;
 import org.apache.paimon.globalindex.GlobalIndexSingleColumnWriter;
 import org.apache.paimon.globalindex.sorted.SortedGlobalIndexScanner;
 import org.apache.paimon.globalindex.sorted.SortedGlobalIndexWriter;
@@ -27,9 +31,13 @@ import org.apache.paimon.options.Options;
 import org.apache.paimon.table.FileStoreTable;
 import org.apache.paimon.table.SpecialFields;
 import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.types.RowType;
 import org.apache.paimon.utils.Range;
 
+import org.apache.flink.streaming.api.datastream.DataStream;
 import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.transformations.PartitionTransformation;
+import org.apache.flink.streaming.api.transformations.StreamExchangeMode;
 import org.junit.jupiter.api.Test;
 
 import java.io.Closeable;
@@ -180,4 +188,30 @@ public class SortedIndexTopoBuilderTest {
         assertThat(SortedIndexTopoBuilder.createSortColumns("task-id", 
"index-key"))
                 .containsExactly("task-id", "index-key", 
SpecialFields.ROW_ID.name());
     }
+
+    @Test
+    public void testBuildTaskPartitioner() {
+        assertThat(SortedIndexTopoBuilder.BUILD_TASK_PARTITIONER.partition(0, 
4)).isEqualTo(0);
+        assertThat(SortedIndexTopoBuilder.BUILD_TASK_PARTITIONER.partition(5, 
4)).isEqualTo(1);
+        assertThat(SortedIndexTopoBuilder.BUILD_TASK_PARTITIONER.partition(9, 
4)).isEqualTo(1);
+    }
+
+    @Test
+    public void testBuildTaskPartitionUsesBatchExchange() {
+        StreamExecutionEnvironment env = 
StreamExecutionEnvironment.getExecutionEnvironment();
+        BinaryRow row = new BinaryRow(1);
+        BinaryRowWriter writer = new BinaryRowWriter(row);
+        writer.writeInt(0, 0);
+        writer.complete();
+        DataStream<InternalRow> input =
+                env.fromData(
+                        Collections.<InternalRow>singletonList(row),
+                        
InternalTypeInfo.fromRowType(RowType.of(DataTypes.INT())));
+
+        DataStream<InternalRow> partitioned = 
SortedIndexTopoBuilder.partitionByBuildTask(input, 0);
+
+        
assertThat(partitioned.getTransformation()).isInstanceOf(PartitionTransformation.class);
+        assertThat(((PartitionTransformation<?>) 
partitioned.getTransformation()).getExchangeMode())
+                .isEqualTo(StreamExchangeMode.BATCH);
+    }
 }

Reply via email to