Repository navigation
[BP-2.2][FLINK-40504][runtime] Re-activate idle FLIP-27 outputs when they resume - #29332
Merged
Merged
Conversation
Collaborator
1 task done
Contributor
|
@Admaing Can you please rebase? |
An output's idleness is now UNKNOWN, ACTIVE or IDLE instead of a boolean that defaults to "not idle", and CombinedWatermarkStatus exposes hasActiveOutput() for callers that need to know whether an output actually reported being active. This is behaviour-neutral on its own: UNKNOWN combines exactly like the previous idle = false did, so it still holds the combined watermark back, and no caller uses hasActiveOutput() yet. Generated-by: DeepSeek Harness (deepseek-v4.1-flash)
…cing watermarks emitWatermark() returned at the monotonicity guard before marking the output active, although the WatermarkOutput contract states that emitting a watermark implicitly marks the stream active. A source that resumes with watermarks that are not larger than the max watermark so far therefore stayed announced as idle downstream, and its records could be dropped as late. The activation is now reported before the guard. Source implementers can observe this change, so it needs a release note. Generated-by: DeepSeek Harness (deepseek-v4.1-flash)
…ents onEvent() reset the idleness timer without undoing the markIdle() it had emitted earlier, so the output stayed announced as idle until the wrapped generator produced an advancing watermark. It now reports activity on the first record and after each idle period, so an output whose generator never emits a watermark is still known to be active. WatermarksWithIdleness is @public, so this is a separate commit to keep its release note easy to trace. Generated-by: DeepSeek Harness (deepseek-v4.1-flash)
…tputMultiplexer updateCombinedWatermark() reported the combined idle state to the underlying output, but nothing when it was active. An output that is active without advancing its watermark therefore produced neither an emitWatermark nor a markIdle call, and the underlying output never learned that it was active again. Because the split branch of IdlenessManager starts out idle and is only cleared by one of those two calls, an active split was then ANDed into IDLE by the main output's idleness timer. The combined status is now reported whenever it is idle or active, but activity is only announced when some output actually reported being active (hasActiveOutput()), so an output that is registered for an assigned split but never reports anything holds the combined watermark back without keeping the downstream output active. Generated-by: DeepSeek Harness (deepseek-v4.1-flash)
Admaing
force-pushed
the
bp-2.2-FLINK-40504
branch
from
October 1, 2026 09:19
c026d66 to
5683e18
Compare
Contributor
Author
@MartijnVisser |
MartijnVisser
approved these changes
Oct 2, 2026
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 #29239 (91e1e5f .. 0897586), which went to master but not to this branch.
A FLIP-27 source that was announced as IDLE downstream stayed IDLE after it resumed, unless its next watermark strictly advanced the combined watermark:
WatermarkOutputMultiplexernever reported the active state of the combined status,WatermarksWithIdlenessnever undid themarkIdle()it had emitted, andWatermarkToDataOutputmarked the output active only after the monotonicity guard. Records of a resumed source could therefore be dropped as late events.The four commits are cherry-picked without conflicts; the touched files are identical on
release-2.2, and the new test file is the only addition.Verified locally on JDK 17:
spotless:checkpasses.*Watermark*Test,SourceOperatorEventTimeTestandSourceOperatorAlignmentTestpass (21 test classes, 139 tests), including the newProgressiveTimestampsAndWatermarksSubtaskIdlenessTestand theSourceOperatorAlignmentTestcase that an earlier version of the fix broke.flink-coreandflink-runtimewere run as well, together with the same run on the base commit for comparison. Both report exactly the same failures (113 failing cases), all of them environment-dependent tests (RMI, SSL blob server, timing-sensitiveTaskExecutor/ClusterEntrypointtests) that need the Linux CI container, and no watermark- or source-operator-related test fails.