[FLINK-40504][core] Re-activate the source output on any sign of activity - #29270
Closed
rkhachatryan wants to merge 1 commit into
Closed
rkhachatryan wants to merge 1 commit into
rkhachatryan wants to merge 1 commit into
Conversation
…vity
A source subtask can announce WatermarkStatus.IDLE downstream and never
take it back while it keeps emitting records. Downstream StatusWatermarkValve
then excludes the channel, stalling the watermark at parallelism 1 and
dropping the records as late at higher parallelism.
Three places on the path from a split's watermark generator to the data
output treat an advancing watermark as the only evidence of activity:
- WatermarkOutputMultiplexer.updateCombinedWatermark() forwards markIdle()
when the combined status is idle but nothing when it is active, so an
output that is active without advancing its watermark produces neither
an emitWatermark nor a markIdle call. Because the split branch of
ProgressiveTimestampsAndWatermarks.IdlenessManager starts out idle and
is only ever cleared by one of those two calls, an active split is then
ANDed into IDLE by the main output's idleness timer .
- WatermarksWithIdleness.onEvent() resets the idleness timer but does not
undo the markIdle() it emitted earlier .
- WatermarkToDataOutput.emitWatermark() gates markActiveInternally() on the
monotonicity check, so a non-advancing watermark is swallowed together
with the activation .
They compound: once all splits fall idle the combined watermark is flushed
to the maximum over all outputs, so a split that resumes behind that maximum
produces no advancing watermark by construction.
Propagate the activity in all three, guarding the multiplexer on the new
CombinedWatermarkStatus.hasOutputs(). That guard is load-bearing:
updateCombinedWatermark() returns early on an empty output list and leaves
the idle flag uncomputed, and a subtask that was never assigned a split
has to be able to stand down via its main-output idleness timer alone.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
1 task done
Collaborator
Contributor
Author
|
Superseded by #29239 |
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.
A source subtask can announce WatermarkStatus.IDLE downstream and never take it back while it keeps emitting records. Downstream StatusWatermarkValve then excludes the channel, stalling the watermark at parallelism 1 and dropping the records as late at higher parallelism.
Three places on the path from a split's watermark generator to the data output treat an advancing watermark as the only evidence of activity:
They compound: once all splits fall idle the combined watermark is flushed to the maximum over all outputs, so a split that resumes behind that maximum produces no advancing watermark by construction.
Propagate the activity in all three, guarding the multiplexer on the new CombinedWatermarkStatus.hasOutputs(). That guard is load-bearing: updateCombinedWatermark() returns early on an empty output list and leaves the idle flag uncomputed, and a subtask that was never assigned a split has to be able to stand down via its main-output idleness timer alone.