1996fanrui commented on code in PR #29272:
URL: https://github.com/apache/flink/pull/29272#discussion_r4173748231


##########
flink-runtime/src/main/java/org/apache/flink/streaming/runtime/tasks/StreamTask.java:
##########
@@ -2046,13 +2048,52 @@ 
List<RecordWriter<SerializationDelegate<StreamRecord<OUT>>>> createRecordWriters
 
     private static void 
replaceForwardPartitionerIfConsumerParallelismDoesNotMatch(
             Environment environment, NonChainedOutput streamOutput, int 
outputIndex) {
-        if (streamOutput.getPartitioner() instanceof ForwardPartitioner
-                && 
environment.getWriter(outputIndex).getNumberOfSubpartitions()
-                        != 
environment.getTaskInfo().getNumberOfParallelSubtasks()) {
-            LOG.debug(
-                    "Replacing forward partitioner with rebalance for {}",
-                    environment.getTaskInfo().getTaskNameWithSubtasks());
-            streamOutput.setPartitioner(new RebalancePartitioner<>());
+        final int producerParallelism = 
environment.getTaskInfo().getNumberOfParallelSubtasks();
+        final int consumerParallelism =
+                environment.getWriter(outputIndex).getNumberOfSubpartitions();

Review Comment:
   Good catch, I created a new [pr](https://github.com/apache/flink/pull/29374) 
to fix it since it is an existing bug.
   
   Regarding tests, added two ITCases, each covering 11 parallelism changes × 3 
modes:
   
   - `ForwardEdgeParallelismMismatchOverridesITCase`: changes the parallelism 
via pipeline.jobvertex-parallelism-overrides under the Default, Adaptive and 
AdaptiveBatch schedulers (99 cases).
   - `ForwardEdgeParallelismMismatchRestRescaleITCase`: changes it at runtime 
via the resource requirements API under the AdaptiveScheduler (30 cases).
   
   All 129 pass. 
   
   Reading the consumer parallelism in `StreamTask` from 
`getWriter(outputIndex).getNumberOfSubpartitions()` again (a one-line revert) 
makes 117 fail; only the 1 -> 2 cases pass, since there the two values happen 
to match.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to