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);
+ }
}