Rui Fan created FLINK-40904:
-------------------------------

             Summary: FORWARD-to-REBALANCE downgrade compares producer 
parallelism with the number of subpartitions instead of the consumer parallelism
                 Key: FLINK-40904
                 URL: https://issues.apache.org/jira/browse/FLINK-40904
             Project: Flink
          Issue Type: Bug
          Components: Runtime / Task
            Reporter: Rui Fan
            Assignee: Rui Fan


FLINK-30213 replaces a ForwardPartitioner with a RebalancePartitioner in 
StreamTask when the producer and consumer parallelism of a FORWARD edge differ. 
It detects the mismatch by comparing the producer parallelism with 
ResultPartitionWriter#getNumberOfSubpartitions(), but on a POINTWISE edge that 
is the number of consumer subtasks connected to one producer subtask, not the 
consumer parallelism.

This was found in 
https://github.com/apache/flink/pull/29272#discussion_r4088147465.

The check is wrong in both directions:
- False positive: equal parallelism > 1 (e.g. 4 -> 4) gives 1 subpartition per 
producer, so it is treated as a mismatch. This was harmless so far because a 
rebalance over a single channel behaves like forward.
- False negative: with the DefaultScheduler and 
pipeline.jobvertex-parallelism-overrides, 2 -> 4 gives 2 subpartitions per 
producer, so the mismatch is not detected. The forward partitioner is kept, and 
half of the consumer subtasks receive no data.

Fix: pass the parallelism of the consuming job vertices to the task via 
ResultPartitionDeploymentDescriptor, and use it for the check.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to