Skip to content

[FLINK-40780][runtime] Make FORWARD-to-REBALANCE downgrade on parallelism mismatch configurable - #29272

Open
1996fanrui wants to merge 4 commits into
apache:masterfrom
1996fanrui:FLINK-40780
Open

1996fanrui wants to merge 4 commits into
apache:masterfrom
1996fanrui:FLINK-40780

Conversation

@1996fanrui

@1996fanrui 1996fanrui commented Sep 23, 2026 •

Copy link
Copy Markdown
Member

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), StreamTask silently replaces the ForwardPartitioner with a RebalancePartitioner (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

  • Add 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#replaceForwardPartitionerIfConsumerParallelismDoesNotMatch applies the configured mode.

Verifying this change

  • StreamTaskTest covers the three modes.
  • ForwardEdgeParallelismMismatchOverridesITCase and ForwardEdgeParallelismMismatchRestRescaleITCase cover every combination of scheduler (Default, Adaptive, AdaptiveBatch), parallelism change and mode, changing the parallelism at submission via pipeline.jobvertex-parallelism-overrides or at runtime via the resource requirements API.

Does this pull request potentially affect one of the following parts

  • Dependencies: no
  • The public API: yes (new @PublicEvolving config option)
  • The serializers: no
  • The runtime per-record code paths: no (the check runs once when record writers are created)
  • Anything that affects deployment or recovery: no
  • The S3 file system connector: no

Documentation

  • New feature: yes
  • Docs: the config docs were regenerated (pipeline_configuration.html)

@flinkbot

flinkbot commented Sep 23, 2026 •

Copy link
Copy Markdown
Collaborator

CI report:

Bot commands The @flinkbot bot supports the following commands:
  • @flinkbot run azure re-run the last Azure build

Comment on lines +2051 to +2053
final int producerParallelism = environment.getTaskInfo().getNumberOfParallelSubtasks();
final int consumerParallelism =
environment.getWriter(outputIndex).getNumberOfSubpartitions();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.
…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.
…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.
@1996fanrui
1996fanrui marked this pull request as ready for review October 3, 2026 14:53
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants