Skip to content

[FLINK-40504][core] Re-activate the source output on any sign of activity - #29270

Closed
rkhachatryan wants to merge 1 commit into
apache:masterfrom
rkhachatryan:f40504
Closed

rkhachatryan wants to merge 1 commit into
apache:masterfrom
rkhachatryan:f40504

Conversation

@rkhachatryan

Copy link
Copy Markdown
Contributor

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.

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

@rkhachatryan

Copy link
Copy Markdown
Contributor Author

Superseded by #29239

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