Repository navigation
[BP-2.3][FLINK-40475][runtime] Fix watermark loss and stall in StatusWatermarkValve when subpartitions realign after idleness - #29346
Merged
Conversation
…kValve when subpartitions realign after idleness (apache#29024) [FLINK-40475][runtime] Fix watermark loss and stall in StatusWatermarkValve when subpartitions realign after idleness The emission logic of StatusWatermarkValve relies on two invariants over the set of watermark-aligned subpartitions: it is only empty when no subpartition is active, and whenever it is non-empty its min watermark equals the last output watermark. Both were violated once an unaligned subpartition - one that went idle, resumed, and only partially caught up to the last output watermark - was involved: - The FLINK-7728 all-idle flush was skipped unless the last subpartition to become idle held the current min watermark. That assumption doesn't hold when unaligned subpartitions are involved, since an unaligned subpartition's watermark is invisible to the min-derive path. This made the final watermark depend on the order in which inputs became idle - exactly the defect FLINK-7728 was meant to fix. - The idle->active branch re-added a caught-up subpartition to the aligned set without re-deriving the min. When no other aligned subpartition remained, the realigned subpartition's watermark stalled indefinitely, emitted only once an even larger watermark arrived on it. Fix both by re-deriving from the aligned set unconditionally instead of relying on the broken shortcut: - Flush the max watermark unconditionally once all subpartitions go idle. findAndOutputMaxWatermarkAcrossAllSubpartitions only emits when the max advances past the last output watermark, so this cannot regress monotonicity, and it keeps the valve consistent with CombinedWatermarkStatus, which has flushed unconditionally since FLINK-38454. - Re-derive (and possibly emit) the min watermark after a subpartition realigns, after the ACTIVE status is emitted so downstream inputs don't drop the watermark while they still consider the input idle. Documents the invariants on alignedSubpartitionStatuses. Generated-by: Claude Code (claude-fable-5) (cherry picked from commit bb773b5)
Collaborator
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Backport of #29024 (bb773b5), which went to master but not to this branch.
The FLINK-40504 backport #29331 should be merged after this one: it marks a source active again when it resumes without advancing its watermark. Without this fix, the valve then skips the all-idle flush if that input is still behind the valve's watermark when it is the last one to go idle.
This changes behaviour on a patch release, see the Release Note on FLINK-40475: once all inputs are idle, windows and timers can fire earlier than before.
Verified locally on JDK 17: the three new
StatusWatermarkValveTestcases fail with the valve of the base commit and pass with this change. The watermark and source operator tests inflink-coreandflink-runtime,OneInputStreamTaskTest,TwoInputStreamTaskTest,MultipleInputStreamTaskTestandspotless:checkpass.Was generative AI tooling used to co-author this PR?
Generated-by: Claude Code (Claude Opus 5.5)