From f14b6b4c42b12bc8cdf4bc36185768b8723157fe Mon Sep 17 00:00:00 2001 From: Rui Fan <1996fanrui@gmail.com> Date: Fri, 2 Oct 2026 23:43:42 +0200 Subject: [PATCH 1/2] [FLINK-40904][runtime] Pass consumer parallelism to tasks via ResultPartitionDeploymentDescriptor A task can only see the number of subpartitions of its outputs, which is not the consumer parallelism on a POINTWISE edge. Pass the parallelism of the consuming job vertices through ResultPartitionDeploymentDescriptor and expose it via Environment#getWriterConsumerParallelism. No behavior change. --- .../api/runtime/SavepointEnvironment.java | 5 ++++ .../ResultPartitionDeploymentDescriptor.java | 30 ++++++++++++++++++- .../flink/runtime/execution/Environment.java | 8 +++++ .../runtime/executiongraph/Execution.java | 5 +++- .../taskmanager/RuntimeEnvironment.java | 10 +++++++ .../flink/runtime/taskmanager/Task.java | 8 +++++ ...sultPartitionDeploymentDescriptorTest.java | 8 ++++- .../deployment/ShuffleDescriptorTest.java | 3 +- .../AbstractPartitionTrackerTest.java | 1 + .../network/partition/PartitionTestUtils.java | 3 +- .../partition/ResultPartitionFactoryTest.java | 1 + .../operators/testutils/DummyEnvironment.java | 5 ++++ .../operators/testutils/MockEnvironment.java | 6 ++++ .../shuffle/NettyShuffleUtilsTest.java | 3 +- .../TaskExecutorSubmissionTest.java | 5 ++-- .../flink/runtime/taskmanager/TaskTest.java | 3 +- .../runtime/tasks/StreamMockEnvironment.java | 13 ++++++++ .../runtime/tasks/StreamTaskITCase.java | 1 + 18 files changed, 109 insertions(+), 9 deletions(-) diff --git a/flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/runtime/SavepointEnvironment.java b/flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/runtime/SavepointEnvironment.java index a15f6a5ab0081..9b963a344bc8c 100644 --- a/flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/runtime/SavepointEnvironment.java +++ b/flink-libraries/flink-state-processing-api/src/main/java/org/apache/flink/state/api/runtime/SavepointEnvironment.java @@ -322,6 +322,11 @@ public ResultPartitionWriter getWriter(int index) { throw new UnsupportedOperationException(ERROR_MSG); } + @Override + public int getWriterConsumerParallelism(int index) { + throw new UnsupportedOperationException(ERROR_MSG); + } + @Override public ResultPartitionWriter[] getAllWriters() { return new ResultPartitionWriter[0]; diff --git a/flink-runtime/src/main/java/org/apache/flink/runtime/deployment/ResultPartitionDeploymentDescriptor.java b/flink-runtime/src/main/java/org/apache/flink/runtime/deployment/ResultPartitionDeploymentDescriptor.java index 6cca01c788f69..79b3910f007a9 100644 --- a/flink-runtime/src/main/java/org/apache/flink/runtime/deployment/ResultPartitionDeploymentDescriptor.java +++ b/flink-runtime/src/main/java/org/apache/flink/runtime/deployment/ResultPartitionDeploymentDescriptor.java @@ -28,6 +28,7 @@ import java.io.Serializable; +import static org.apache.flink.util.Preconditions.checkArgument; import static org.apache.flink.util.Preconditions.checkNotNull; /** @@ -45,14 +46,31 @@ public class ResultPartitionDeploymentDescriptor implements Serializable { private final int maxParallelism; + /** Parallelism of the consuming job vertices, or {@link #UNKNOWN_CONSUMER_PARALLELISM}. */ + private final int consumerParallelism; + + /** + * The consumer parallelism is not decided yet when the producer is deployed. This only happens + * in dynamic graphs (AdaptiveBatchScheduler), where the parallelism of a consumer vertex may be + * decided after its producers are deployed. Every other partition carries the actual consumer + * parallelism. + */ + public static final int UNKNOWN_CONSUMER_PARALLELISM = -1; + public ResultPartitionDeploymentDescriptor( PartitionDescriptor partitionDescriptor, ShuffleDescriptor shuffleDescriptor, - int maxParallelism) { + int maxParallelism, + int consumerParallelism) { this.partitionDescriptor = checkNotNull(partitionDescriptor); this.shuffleDescriptor = checkNotNull(shuffleDescriptor); KeyGroupRangeAssignment.checkParallelismPreconditions(maxParallelism); this.maxParallelism = maxParallelism; + checkArgument( + consumerParallelism > 0 || consumerParallelism == UNKNOWN_CONSUMER_PARALLELISM, + "Invalid consumer parallelism %s.", + consumerParallelism); + this.consumerParallelism = consumerParallelism; } public IntermediateDataSetID getResultId() { @@ -88,6 +106,16 @@ public int getMaxParallelism() { return maxParallelism; } + /** + * Returns the parallelism of the job vertices consuming this partition, or {@link + * #UNKNOWN_CONSUMER_PARALLELISM} if it is not decided yet. Unlike {@link + * #getNumberOfSubpartitions()}, this is the actual consumer parallelism for every distribution + * pattern. + */ + public int getConsumerParallelism() { + return consumerParallelism; + } + public ShuffleDescriptor getShuffleDescriptor() { return shuffleDescriptor; } diff --git a/flink-runtime/src/main/java/org/apache/flink/runtime/execution/Environment.java b/flink-runtime/src/main/java/org/apache/flink/runtime/execution/Environment.java index 4b449a4d12950..d989f67477433 100644 --- a/flink-runtime/src/main/java/org/apache/flink/runtime/execution/Environment.java +++ b/flink-runtime/src/main/java/org/apache/flink/runtime/execution/Environment.java @@ -32,6 +32,7 @@ import org.apache.flink.runtime.checkpoint.TaskStateSnapshot; import org.apache.flink.runtime.checkpoint.channel.ChannelStateWriteRequestExecutorFactory; import org.apache.flink.runtime.checkpoint.channel.ChannelStateWriter; +import org.apache.flink.runtime.deployment.ResultPartitionDeploymentDescriptor; import org.apache.flink.runtime.executiongraph.ExecutionAttemptID; import org.apache.flink.runtime.externalresource.ExternalResourceInfoProvider; import org.apache.flink.runtime.io.disk.iomanager.IOManager; @@ -248,6 +249,13 @@ void acknowledgeCheckpoint( ResultPartitionWriter getWriter(int index); + /** + * Returns the parallelism of the job vertices consuming the partition written by {@link + * #getWriter(int)}, or {@link ResultPartitionDeploymentDescriptor#UNKNOWN_CONSUMER_PARALLELISM} + * if it is not decided yet. + */ + int getWriterConsumerParallelism(int index); + ResultPartitionWriter[] getAllWriters(); IndexedInputGate getInputGate(int index); diff --git a/flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/Execution.java b/flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/Execution.java index 46d5a8d60b022..95d9610e3fbdf 100644 --- a/flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/Execution.java +++ b/flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/Execution.java @@ -561,7 +561,10 @@ private static ResultPartitionDeploymentDescriptor createResultPartitionDeployme IntermediateResultPartition partition, ShuffleDescriptor shuffleDescriptor) { return new ResultPartitionDeploymentDescriptor( - partitionDescriptor, shuffleDescriptor, getPartitionMaxParallelism(partition)); + partitionDescriptor, + shuffleDescriptor, + getPartitionMaxParallelism(partition), + partition.getIntermediateResult().getConsumersParallelism()); } /** diff --git a/flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/RuntimeEnvironment.java b/flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/RuntimeEnvironment.java index 69df0116802cd..131af5bc3351c 100644 --- a/flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/RuntimeEnvironment.java +++ b/flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/RuntimeEnvironment.java @@ -58,6 +58,7 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; +import static org.apache.flink.util.Preconditions.checkArgument; import static org.apache.flink.util.Preconditions.checkNotNull; import static org.apache.flink.util.Preconditions.checkState; @@ -93,6 +94,7 @@ public class RuntimeEnvironment implements Environment { private final Map> distCacheEntries; private final ResultPartitionWriter[] writers; + private final int[] writerConsumerParallelisms; private final IndexedInputGate[] inputGates; private final TaskEventDispatcher taskEventDispatcher; @@ -145,6 +147,7 @@ public RuntimeEnvironment( InputSplitProvider splitProvider, Map> distCacheEntries, ResultPartitionWriter[] writers, + int[] writerConsumerParallelisms, IndexedInputGate[] inputGates, TaskEventDispatcher taskEventDispatcher, CheckpointResponder checkpointResponder, @@ -177,6 +180,8 @@ public RuntimeEnvironment( this.splitProvider = checkNotNull(splitProvider); this.distCacheEntries = checkNotNull(distCacheEntries); this.writers = checkNotNull(writers); + this.writerConsumerParallelisms = checkNotNull(writerConsumerParallelisms); + checkArgument(writers.length == writerConsumerParallelisms.length); this.inputGates = checkNotNull(inputGates); this.taskEventDispatcher = checkNotNull(taskEventDispatcher); this.checkpointResponder = checkNotNull(checkpointResponder); @@ -306,6 +311,11 @@ public ResultPartitionWriter getWriter(int index) { return writers[index]; } + @Override + public int getWriterConsumerParallelism(int index) { + return writerConsumerParallelisms[index]; + } + @Override public ResultPartitionWriter[] getAllWriters() { return writers; diff --git a/flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java b/flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java index 3ef341c137297..52e01e59f9675 100644 --- a/flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java +++ b/flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java @@ -230,6 +230,9 @@ public class Task private final ResultPartitionWriter[] partitionWriters; + /** Consumer parallelism of each partition in {@link #partitionWriters}, in the same order. */ + private final int[] partitionConsumerParallelisms; + private final IndexedInputGate[] inputGates; /** Connection to the task manager. */ @@ -423,6 +426,10 @@ public Task( .toArray(new ResultPartitionWriter[] {}); this.partitionWriters = resultPartitionWriters; + this.partitionConsumerParallelisms = + resultPartitionDeploymentDescriptors.stream() + .mapToInt(ResultPartitionDeploymentDescriptor::getConsumerParallelism) + .toArray(); // consumed intermediate result partitions final IndexedInputGate[] gates = @@ -731,6 +738,7 @@ private void doRun() { inputSplitProvider, distributedCacheEntries, partitionWriters, + partitionConsumerParallelisms, inputGates, taskEventDispatcher, checkpointResponder, diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/deployment/ResultPartitionDeploymentDescriptorTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/deployment/ResultPartitionDeploymentDescriptorTest.java index eddb41d6c7848..181f807fff80b 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/deployment/ResultPartitionDeploymentDescriptorTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/deployment/ResultPartitionDeploymentDescriptorTest.java @@ -51,6 +51,8 @@ class ResultPartitionDeploymentDescriptorTest { private static final ResultPartitionType partitionType = ResultPartitionType.PIPELINED; private static final int numberOfSubpartitions = 24; + + private static final int consumerParallelism = 12; private static final int connectionIndex = 10; private static final boolean isBroadcast = false; private static final boolean isAllToAllDistribution = true; @@ -112,7 +114,10 @@ void testSerializationWithNettyShuffleDescriptor() throws IOException { ShuffleDescriptor shuffleDescriptor) throws IOException { ResultPartitionDeploymentDescriptor orig = new ResultPartitionDeploymentDescriptor( - partitionDescriptor, shuffleDescriptor, numberOfSubpartitions); + partitionDescriptor, + shuffleDescriptor, + numberOfSubpartitions, + consumerParallelism); ResultPartitionDeploymentDescriptor copy = CommonTestUtils.createCopySerializable(orig); verifyResultPartitionDeploymentDescriptorCopy(copy); return copy; @@ -125,5 +130,6 @@ private static void verifyResultPartitionDeploymentDescriptorCopy( assertThat(partitionId).isEqualTo(copy.getPartitionId()); assertThat(partitionType).isEqualTo(copy.getPartitionType()); assertThat(numberOfSubpartitions).isEqualTo(copy.getNumberOfSubpartitions()); + assertThat(consumerParallelism).isEqualTo(copy.getConsumerParallelism()); } } diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/deployment/ShuffleDescriptorTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/deployment/ShuffleDescriptorTest.java index df4291170ca2c..79cae3fb573b9 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/deployment/ShuffleDescriptorTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/deployment/ShuffleDescriptorTest.java @@ -274,6 +274,7 @@ private static ResultPartitionDeploymentDescriptor createResultPartitionDeployme .registerPartitionWithProducer( jobID, partitionDescriptor, producerDescriptor) .get(); - return new ResultPartitionDeploymentDescriptor(partitionDescriptor, shuffleDescriptor, 1); + return new ResultPartitionDeploymentDescriptor( + partitionDescriptor, shuffleDescriptor, 1, 1); } } diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/AbstractPartitionTrackerTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/AbstractPartitionTrackerTest.java index 6f90313cea846..a49075b76dbdc 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/AbstractPartitionTrackerTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/AbstractPartitionTrackerTest.java @@ -87,6 +87,7 @@ public Optional storesLocalResourcesOn() { : Optional.empty(); } }, + 1, 1); } diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/PartitionTestUtils.java b/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/PartitionTestUtils.java index 20547adc2c77c..1e921f5b3bf91 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/PartitionTestUtils.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/PartitionTestUtils.java @@ -137,7 +137,8 @@ public static ResultPartitionDeploymentDescriptor createPartitionDeploymentDescr .setPartitionId(shuffleDescriptor.getResultPartitionID().getPartitionId()) .setPartitionType(partitionType) .build(); - return new ResultPartitionDeploymentDescriptor(partitionDescriptor, shuffleDescriptor, 1); + return new ResultPartitionDeploymentDescriptor( + partitionDescriptor, shuffleDescriptor, 1, 1); } public static PartitionedFile createPartitionedFile( diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/ResultPartitionFactoryTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/ResultPartitionFactoryTest.java index bd8156634df1c..95907f5ab54db 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/ResultPartitionFactoryTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/ResultPartitionFactoryTest.java @@ -193,6 +193,7 @@ private static ResultPartition createResultPartition( .setIsBroadcast(isBroadcast) .build(), NettyShuffleDescriptorBuilder.newBuilder().buildLocal(), + 1, 1); // guard our test assumptions diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/operators/testutils/DummyEnvironment.java b/flink-runtime/src/test/java/org/apache/flink/runtime/operators/testutils/DummyEnvironment.java index fe291c1b4e5f3..367b28d3d30cf 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/operators/testutils/DummyEnvironment.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/operators/testutils/DummyEnvironment.java @@ -264,6 +264,11 @@ public ResultPartitionWriter getWriter(int index) { return null; } + @Override + public int getWriterConsumerParallelism(int index) { + return taskInfo.getNumberOfParallelSubtasks(); + } + @Override public ResultPartitionWriter[] getAllWriters() { return new ResultPartitionWriter[0]; diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/operators/testutils/MockEnvironment.java b/flink-runtime/src/test/java/org/apache/flink/runtime/operators/testutils/MockEnvironment.java index 101ff4386cb29..84114ee9ee13b 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/operators/testutils/MockEnvironment.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/operators/testutils/MockEnvironment.java @@ -332,6 +332,12 @@ public ResultPartitionWriter getWriter(int index) { return outputs.get(index); } + /** Outputs are consumed with the same parallelism as this task. */ + @Override + public int getWriterConsumerParallelism(int index) { + return taskInfo.getNumberOfParallelSubtasks(); + } + @Override public ResultPartitionWriter[] getAllWriters() { return outputs.toArray(new ResultPartitionWriter[outputs.size()]); diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/shuffle/NettyShuffleUtilsTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/shuffle/NettyShuffleUtilsTest.java index 76f3bf38341c9..575d7addf238c 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/shuffle/NettyShuffleUtilsTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/shuffle/NettyShuffleUtilsTest.java @@ -194,7 +194,8 @@ private ResultPartition createResultPartition( true, false); ResultPartitionDeploymentDescriptor resultPartitionDeploymentDescriptor = - new ResultPartitionDeploymentDescriptor(partitionDescriptor, shuffleDescriptor, 1); + new ResultPartitionDeploymentDescriptor( + partitionDescriptor, shuffleDescriptor, 1, 1); ExecutionAttemptID consumerID = createExecutionAttemptId(); Collection resultPartitions = diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorSubmissionTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorSubmissionTest.java index 8e87b86fcadd6..65c37643f86cc 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorSubmissionTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorSubmissionTest.java @@ -503,7 +503,7 @@ void testGetPartitionWithMetrics() throws Exception { .build(); ResultPartitionDeploymentDescriptor resultPartitionDeploymentDescriptor = - new ResultPartitionDeploymentDescriptor(partitionDescriptor, sdd, 1); + new ResultPartitionDeploymentDescriptor(partitionDescriptor, sdd, 1, 1); TaskDeploymentDescriptor tdd = createTestTaskDeploymentDescriptor( "task", @@ -798,7 +798,8 @@ private TaskDeploymentDescriptor createSender( .setPartitionId(shuffleDescriptor.getResultPartitionID().getPartitionId()) .build(); ResultPartitionDeploymentDescriptor resultPartitionDeploymentDescriptor = - new ResultPartitionDeploymentDescriptor(partitionDescriptor, shuffleDescriptor, 1); + new ResultPartitionDeploymentDescriptor( + partitionDescriptor, shuffleDescriptor, 1, 1); return createTestTaskDeploymentDescriptor( "Sender", shuffleDescriptor.getResultPartitionID().getProducerId(), diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/taskmanager/TaskTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/taskmanager/TaskTest.java index 36a248934e14b..eaa99ef366424 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/taskmanager/TaskTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/taskmanager/TaskTest.java @@ -316,7 +316,8 @@ public void testExecutionFailsInNetworkRegistrationForPartitions() throws Except final ShuffleDescriptor shuffleDescriptor = NettyShuffleDescriptorBuilder.newBuilder().buildLocal(); final ResultPartitionDeploymentDescriptor dummyPartition = - new ResultPartitionDeploymentDescriptor(partitionDescriptor, shuffleDescriptor, 1); + new ResultPartitionDeploymentDescriptor( + partitionDescriptor, shuffleDescriptor, 1, 1); testExecutionFailsInNetworkRegistration( Collections.singletonList(dummyPartition), Collections.emptyList()); } diff --git a/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamMockEnvironment.java b/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamMockEnvironment.java index dd4ef8994c2e5..27f37f24cf7f8 100644 --- a/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamMockEnvironment.java +++ b/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamMockEnvironment.java @@ -108,6 +108,9 @@ public class StreamMockEnvironment implements Environment { private List outputs; + /** Consumer parallelism of all outputs, equal to the parallelism of this task by default. */ + private int writerConsumerParallelism; + private final ExecutionAttemptID executionAttemptID; private final BroadcastVariableManager bcVarManager = new BroadcastVariableManager(); @@ -193,6 +196,7 @@ public StreamMockEnvironment( this.taskConfiguration = taskConfig; this.inputs = new LinkedList<>(); this.outputs = new LinkedList(); + this.writerConsumerParallelism = taskInfo.getNumberOfParallelSubtasks(); this.memManager = MemoryManagerBuilder.newBuilder().setMemorySize(offHeapMemorySize).build(); this.sharedResources = new SharedResources(); @@ -324,6 +328,15 @@ public ResultPartitionWriter getWriter(int index) { return outputs.get(index); } + @Override + public int getWriterConsumerParallelism(int index) { + return writerConsumerParallelism; + } + + public void setWriterConsumerParallelism(int writerConsumerParallelism) { + this.writerConsumerParallelism = writerConsumerParallelism; + } + @Override public ResultPartitionWriter[] getAllWriters() { return outputs.toArray(new ResultPartitionWriter[outputs.size()]); diff --git a/flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamTaskITCase.java b/flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamTaskITCase.java index 1dac2e2d8018b..74d1e024a7886 100644 --- a/flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamTaskITCase.java +++ b/flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamTaskITCase.java @@ -132,6 +132,7 @@ private TestTaskBuilder taskBuilderWithConfiguredRecordWriter( new ResultPartitionDeploymentDescriptor( PartitionDescriptorBuilder.newBuilder().build(), NettyShuffleDescriptorBuilder.newBuilder().buildLocal(), + 1, 1); return new TestTaskBuilder(shuffleEnvironment) .setInvokable(NoOpStreamTask.class) From 7c587b2fafda4c98574c554d346c1e56c3dfd2d2 Mon Sep 17 00:00:00 2001 From: Rui Fan <1996fanrui@gmail.com> Date: Fri, 2 Oct 2026 23:44:13 +0200 Subject: [PATCH 2/2] [FLINK-40904][runtime] Detect FORWARD parallelism mismatch with the actual consumer parallelism The FORWARD-to-REBALANCE downgrade compared the producer parallelism with the number of subpartitions, which is not the consumer parallelism on a POINTWISE edge. Equal parallelism > 1 was taken as a mismatch, and N -> M with N > 1 could go undetected. Compare with the actual consumer parallelism instead and include both values in the log. --- .../streaming/runtime/tasks/StreamTask.java | 18 ++++++++++++++---- .../runtime/tasks/StreamTaskTest.java | 1 + 2 files changed, 15 insertions(+), 4 deletions(-) diff --git a/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/tasks/StreamTask.java b/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/tasks/StreamTask.java index 7bccacd79ddc8..1b83957632675 100644 --- a/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/tasks/StreamTask.java +++ b/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/tasks/StreamTask.java @@ -49,6 +49,7 @@ import org.apache.flink.runtime.checkpoint.channel.RecoveryCheckpointTrigger; import org.apache.flink.runtime.checkpoint.channel.SequentialChannelStateReader; import org.apache.flink.runtime.checkpoint.filemerging.FileMergingSnapshotManager; +import org.apache.flink.runtime.deployment.ResultPartitionDeploymentDescriptor; import org.apache.flink.runtime.execution.CancelTaskException; import org.apache.flink.runtime.execution.Environment; import org.apache.flink.runtime.io.AvailabilityProvider; @@ -2046,12 +2047,21 @@ List>>> createRecordWriters private static void replaceForwardPartitionerIfConsumerParallelismDoesNotMatch( Environment environment, NonChainedOutput streamOutput, int outputIndex) { + final int producerParallelism = environment.getTaskInfo().getNumberOfParallelSubtasks(); + // The number of subpartitions isn't the consumer parallelism on a POINTWISE edge, so the + // actual consumer parallelism is passed through the deployment descriptor. + final int consumerParallelism = environment.getWriterConsumerParallelism(outputIndex); if (streamOutput.getPartitioner() instanceof ForwardPartitioner - && environment.getWriter(outputIndex).getNumberOfSubpartitions() - != environment.getTaskInfo().getNumberOfParallelSubtasks()) { + // An undecided consumer parallelism is not a mismatch. + && consumerParallelism + != ResultPartitionDeploymentDescriptor.UNKNOWN_CONSUMER_PARALLELISM + && consumerParallelism != producerParallelism) { LOG.debug( - "Replacing forward partitioner with rebalance for {}", - environment.getTaskInfo().getTaskNameWithSubtasks()); + "Replacing forward partitioner with rebalance for {} " + + "(producer parallelism {} != consumer parallelism {}).", + environment.getTaskInfo().getTaskNameWithSubtasks(), + producerParallelism, + consumerParallelism); streamOutput.setPartitioner(new RebalancePartitioner<>()); } } diff --git a/flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamTaskTest.java b/flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamTaskTest.java index 1b42449f539ee..2884383863290 100644 --- a/flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamTaskTest.java +++ b/flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamTaskTest.java @@ -1871,6 +1871,7 @@ public int getNumberOfSubpartitions() { } }); harness.streamMockEnvironment.setOutputs(newOutputs); + harness.streamMockEnvironment.setWriterConsumerParallelism(2); // Re-create outputs recordWriterDelegate =