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)