Skip to content

[FLINK-40904][runtime] Detect FORWARD parallelism mismatch with the actual consumer parallelism - #29374

Open
1996fanrui wants to merge 2 commits into
apache:masterfrom
1996fanrui:FLINK-40904
Open

1996fanrui wants to merge 2 commits into
apache:masterfrom
1996fanrui:FLINK-40904

Conversation

@1996fanrui

Copy link
Copy Markdown
Member

What is the purpose of the change

StreamTask replaces a ForwardPartitioner with a RebalancePartitioner when the parallelism of a FORWARD edge doesn't match (FLINK-30213). The check compares the producer parallelism with the number of subpartitions, which is not the consumer parallelism on a POINTWISE edge. Equal parallelism > 1 is taken as a mismatch, and N -> M with N > 1 can go undetected (e.g. 2 -> 4 with the DefaultScheduler, where half of the consumers get no data). Found in #29272 (comment).

Brief change log

  • Pass the parallelism of the consuming job vertices through ResultPartitionDeploymentDescriptor and expose it via Environment#getWriterConsumerParallelism.
  • StreamTask compares the producer parallelism with it instead of the number of subpartitions, and logs both values.

Verifying this change

Does this pull request potentially affect one of the following parts

  • Dependencies: no
  • The public API: no
  • 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: yes (ResultPartitionDeploymentDescriptor carries one more int)
  • The S3 file system connector: no

Documentation

  • New feature: no

…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.
@flinkbot

flinkbot commented Oct 2, 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

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.

2 participants