Skip to content

[BP-2.2][FLINK-40504][runtime] Re-activate idle FLIP-27 outputs when they resume - #29332

Merged
MartijnVisser merged 4 commits into
apache:release-2.2from
Admaing:bp-2.2-FLINK-40504
Oct 2, 2026
Merged

MartijnVisser merged 4 commits into
apache:release-2.2from
Admaing:bp-2.2-FLINK-40504

Conversation

@Admaing

@Admaing Admaing commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

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: WatermarkOutputMultiplexer never reported the active state of the combined status, WatermarksWithIdleness never undid the markIdle() it had emitted, and WatermarkToDataOutput marked 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:check passes.
  • *Watermark*Test, SourceOperatorEventTimeTest and SourceOperatorAlignmentTest pass (21 test classes, 139 tests), including the new ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest and the SourceOperatorAlignmentTest case that an earlier version of the fix broke.
  • The full unit test suites of flink-core and flink-runtime were 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-sensitive TaskExecutor/ClusterEntrypoint tests) that need the Linux CI container, and no watermark- or source-operator-related test fails.

@flinkbot

flinkbot commented Sep 29, 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

@MartijnVisser

Copy link
Copy Markdown
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
Admaing force-pushed the bp-2.2-FLINK-40504 branch from c026d66 to 5683e18 Compare October 1, 2026 09:19
@Admaing

Admaing commented Oct 1, 2026

Copy link
Copy Markdown
Contributor Author

@Admaing Can you please rebase?

@MartijnVisser
Done

@MartijnVisser
MartijnVisser merged commit fc918d4 into apache:release-2.2 Oct 2, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants