[FLINK-40780][runtime] Make FORWARD-to-REBALANCE downgrade on parallelism mismatch configurable - #29272
[FLINK-40780][runtime] Make FORWARD-to-REBALANCE downgrade on parallelism mismatch configurable#292721996fanrui wants to merge 4 commits into
Conversation
| final int producerParallelism = environment.getTaskInfo().getNumberOfParallelSubtasks(); | ||
| final int consumerParallelism = | ||
| environment.getWriter(outputIndex).getNumberOfSubpartitions(); |
There was a problem hiding this comment.
getNumberOfSubpartitions() isn't the consumer parallelism on a POINTWISE (forward) edge. In a static graph it's the size of the consumer vertex group for this producer subtask (IntermediateResultPartition#getNumberOfSubpartitions), so it's 1 when producer and consumer parallelism are equal. In a dynamic graph (adaptive batch) a forward result always has exactly 1 subpartition (computeNumberOfSubpartitionsForDynamicGraph). So the != producerParallelism check is wrong both ways.
False positive: any non-chained forward edge with parallelism > 1 and no rescale counts as a mismatch. Before this PR that was harmless, since the result was a RebalancePartitioner over a single channel. Now FAIL rejects such jobs and KEEP_FORWARD logs a WARN for every subtask.
False negative: producer 2 → consumer 4 gives 2 subpartitions per producer, which equals the producer parallelism of 2, so nothing is detected. FAIL doesn't fail and REBALANCE doesn't rebalance, and half the consumers get no data.
I checked this with a MiniCluster ITCase on this branch, with mode FAIL, source → forward() → map (startNewChain):
| Case | Result |
|---|---|
| streaming 4 → 4, no rescale | fails: producer parallelism 4 != consumer parallelism 1 |
| batch 4 → 4, no rescale | fails, same error |
| streaming 2 → 4 (via pipeline.jobvertex-parallelism-overrides) | succeeds, records per consumer subtask {0=500, 1=0, 2=500, 3=0} |
| streaming 1 → 2 (the case covered by StreamTaskTest) | fails, as expected |
The existing tests only cover 1 → 2, where the two quantities happen to coincide.
The error and WARN messages have the same mix-up: they report the subpartition count as "consumer parallelism", e.g. "consumer parallelism 1" for a consumer running at 4.
I think the check needs the actual consumer vertex parallelism, which I don't believe is available in the task today, so it may have to be passed through the deployment descriptor or NonChainedOutput. Alternatively, the mismatch could be decided at scheduling time. Could you also add tests for equal parallelism > 1 (streaming and batch) and for N → M with N > 1?
There was a problem hiding this comment.
Good catch, I created a new pr 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.
…artitionDeploymentDescriptor 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.
60a3487 to
1b897e7
Compare
…ctual 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.
…lism mismatch configurable When a FORWARD edge connects a producer and consumer whose parallelism no longer match (e.g. after a rescale via the AdaptiveScheduler or pipeline.jobvertex-parallelism-overrides), the runtime silently replaces the ForwardPartitioner with a RebalancePartitioner (FLINK-30213). This is safe for append-only streams but silently reorders order-sensitive (changelog) streams, producing incorrect results. Add pipeline.forward-edge.parallelism-mismatch-mode to control this behavior: - REBALANCE (default): current behavior, debug log. - FAIL: reject the job with an exception. - KEEP_FORWARD: keep the forward partitioner (order-preserving), warn log.
1b897e7 to
6c77218
Compare
6c77218 to
e6ec054
Compare
…h modes Cover every combination of scheduler (Default, Adaptive, AdaptiveBatch), parallelism change and mismatch mode, with the new parallelism applied either at submission via pipeline.jobvertex-parallelism-overrides or at runtime via the resource requirements API of the AdaptiveScheduler.
e6ec054 to
3686690
Compare
This PR is stacked on #29374 (FLINK-40904), which fixes the mismatch detection. Only the last two commits belong to this PR.
What is the purpose of the change
When the producer and consumer of a FORWARD edge end up with different parallelism (e.g. after rescaling via the AdaptiveScheduler or
pipeline.jobvertex-parallelism-overrides),StreamTasksilently replaces theForwardPartitionerwith aRebalancePartitioner(FLINK-30213). That's fine for append-only streams, but it reorders records in changelog streams and produces wrong results. This PR makes the behavior configurable.Brief change log
pipeline.forward-edge.parallelism-mismatch-mode(enum):REBALANCE(default): current behavior.FAIL: fail the job with an exception.KEEP_FORWARD: keep the forward partitioner (order is preserved, and records go to a single consumer subtask).StreamTask#replaceForwardPartitionerIfConsumerParallelismDoesNotMatchapplies the configured mode.Verifying this change
StreamTaskTestcovers the three modes.ForwardEdgeParallelismMismatchOverridesITCaseandForwardEdgeParallelismMismatchRestRescaleITCasecover every combination of scheduler (Default, Adaptive, AdaptiveBatch), parallelism change and mode, changing the parallelism at submission viapipeline.jobvertex-parallelism-overridesor at runtime via the resource requirements API.Does this pull request potentially affect one of the following parts
@PublicEvolvingconfig option)Documentation
pipeline_configuration.html)