Skip to content

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

Merged
rkhachatryan merged 5 commits into
apache:masterfrom
Admaing:FLINK-40504-reattivate-idle-output
Sep 28, 2026
Merged

rkhachatryan merged 5 commits into
apache:masterfrom
Admaing:FLINK-40504-reattivate-idle-output

Conversation

@Admaing

@Admaing Admaing commented Sep 19, 2026 •

Copy link
Copy Markdown
Contributor

What is the purpose of the change

A FLIP-27 source that has been announced as IDLE downstream stays IDLE after it resumes, unless its
next watermark strictly advances the combined watermark. 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:

  1. WatermarkOutputMultiplexer.updateCombinedWatermark() reports the combined idle state via
    underlyingOutput.markIdle(), but nothing when the status is active. An output that is active
    without advancing its watermark therefore produces neither an emitWatermark nor a markIdle
    call, and the underlying output never learns that it is active again. 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.
  2. WatermarksWithIdleness.onEvent() resets its idleness timer without reporting that the output is
    active again, so the output stays announced as idle until the wrapped generator produces an
    advancing watermark.
  3. WatermarkToDataOutput.emitWatermark() returns at the monotonicity guard
    (newWatermark <= maxWatermarkSoFar) before marking the output active, although the
    WatermarkOutput contract states that emitting a watermark implicitly marks the stream active.

Downstream, StatusWatermarkValve then excludes the channel: the watermark stalls at parallelism 1,
and at higher parallelism the records arrive behind the downstream watermark and are dropped as late.
Every sibling path re-activates eagerly (StatusWatermarkValve, legacy StreamSourceContexts, the
table WatermarkAssignerOperator, and DataStream V2); the FLIP-27 multiplexer path was the only one
without a re-activation signal.

Brief change log

  • PartialWatermark is tri-state (UNKNOWN/ACTIVE/IDLE) and starts UNKNOWN. An output that is
    registered for an assigned split but never reports anything holds the combined watermark back the
    way an active output does, but it does not count as evidence of activity.
  • WatermarkOutputMultiplexer reports the combined status whenever it is idle or active, but it only
    announces activity when some output actually reported being active
    (CombinedWatermarkStatus#hasActiveOutput()). Reporting idleness on one output therefore does not
    re-activate outputs that have not said anything.
  • WatermarksWithIdleness.onEvent() 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.
  • WatermarkToDataOutput.emitWatermark() marks the output active before the monotonicity guard.
  • Tests: 4 new WatermarkOutputMultiplexerTest cases, 1 new WatermarksWithIdlenessTest case, 1 new
    WatermarkToDataOutputTest case and 5 new end-to-end cases in
    ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest.

Verifying this change

This change added tests and can be verified as follows:

  • mvn -pl flink-core -am -Dtest='WatermarkOutputMultiplexerTest,WatermarksWithIdlenessTest' -Dsurefire.failIfNoSpecifiedTests=false test
  • mvn -pl flink-runtime -am -Dtest='WatermarkToDataOutputTest,SourceOperatorEventTimeTest,ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest,SourceOperatorAlignmentTest' -Dsurefire.failIfNoSpecifiedTests=false test
  • SourceOperatorAlignmentTest is the regression test for the activity signal: when the combined
    status was reported on every active update, a split that is registered for an assigned split but
    never reports anything kept the subtask active, and testWatermarkAlignmentWithIdleness never
    reported MAX_WATERMARK.
  • WatermarkOutputMultiplexerTest#whenRegisteredOutputReportsNothingNothingIsReported and
    ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest#subtaskStaysIdleWhenOnlyAnUnknownAndAnIdleSplitAreRegistered
    cover the same case at the multiplexer and subtask level; the latter also covers the mixed
    UNKNOWN + IDLE state.
  • With the production changes reverted to upstream/master, 10 of the new/updated cases fail.
  • Broader regression over all watermark-related test classes in flink-core and flink-runtime
    (22 classes, 163 tests) passes.

Notes for reviewers:

  • Reporting "active" is only meaningful when an output actually said so. IdlenessManager's
    IdlenessAwareWatermarkOutput starts with isIdle = true, and an output is registered for every
    assigned split even when the reader never reports through it (it emits everything through the main
    output, the strategy does not use idleness, or auto-watermark-interval is 0). Treating such an
    output as active pins the split branch active forever and makes the main output's idleness timer
    inert.
  • WatermarksWithIdleness is @Public, but only its idle -> active behaviour changes: it now sends
    the markActive() that its own markIdle() left outstanding, plus one for the first record.
  • Scope note: this does not cover the new-split immediacy of FLINK-22926. Registering a split is
    not evidence that it is producing, and it is indistinguishable from the registered-but-silent split
    above, so a subtask that is idle re-activates on the new split's first record instead. FLINK-22926
    stays open unless we want a separate mechanism for it.
  • A release note is needed for both the WatermarkToDataOutput behaviour change and
    WatermarksWithIdleness. I cannot set the JIRA Release Note field myself, so I asked a committer to
    set it on the ticket.
  • This is complementary to [FLINK-40499][runtime] Emit ACTIVE before final MAX_WATERMARK at drain #29037 (FLINK-40499), which re-activates at drain time in the stream
    tasks; this PR fixes the root cause and does not touch those files.

Does this pull request potentially affect one of the following parts:

  • Dependencies (does it add or upgrade a dependency): no
  • The public API: no (WatermarkOutputMultiplexer and WatermarkToDataOutput are @Internal;
    WatermarksWithIdleness is @Public and keeps its API, only its idleness handling is fixed)
  • The serializers: no
  • The runtime per-record code paths (performance sensitive): yes, minimally — the multiplexer now
    reports the active state on combined updates, which is one check plus a call the underlying output
    deduplicates; WatermarksWithIdleness.onEvent() reports activity only on the first record and after
    each idle period
  • Anything that affects deployment or recovery: no
  • The S3 file system connector: no

Documentation

  • Does this pull request introduce a new feature? no
  • If yes, how is the feature documented? not applicable

  • Yes (please specify the tool below)

Generated-by: DeepSeek Harness (deepseek-v4.1-flash)

@flinkbot

flinkbot commented Sep 19, 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
rkhachatryan self-requested a review September 22, 2026 13:27

@rkhachatryan rkhachatryan left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi @Admaing, thanks for picking this up. The diagnosis is right, and I agree with the WatermarkToDataOutput change: any watermark should end idleness, whether or not it advances.

While working on the same code path I found two related defects that this PR leaves open. I've put up a fix for all of them in #29270. Could you either adopt the two changes below here, or review that PR? I'm happy to rebase mine on top of yours if that's easier.

  1. A subtask can go idle while one of its splits is still producing records

ProgressiveTimestampsAndWatermarks.IdlenessManager marks the subtask idle only when both the split side and the main-output side are idle. Both IdlenessAwareWatermarkOutputs start out with isIdle = true.

This PR forwards markActive() only when the multiplexer's reported state goes from idle to active. The split side's initial idle flag is therefore only cleared by a combined watermark that actually advances. A split that produces records whose watermark never advances never clears it: the multiplexer never reported idle, so there's no transition to forward. This happens with a generator that emits nothing, or when the watermark never moves past its first value.

For a reader that only uses split outputs, the main output never sees a record, so its idleness timer always expires. At that point both sides look idle and the subtask announces IDLE while records keep flowing. At parallelism 1 the watermark stalls; at higher parallelism those records arrive behind the downstream watermark.

Fix: forward markActive() on every combined update where the status is active, not only on the transition. This needs a guard for an empty multiplexer. With no outputs, CombinedWatermarkStatus.updateCombinedWatermark() returns early and leaves idle at its old value. Without the guard, a subtask with no splits (routine when parallelism exceeds the split count) would re-activate the split side on every periodic emit and could never go idle via the main-output timer. See CombinedWatermarkStatus#hasOutputs() in the linked PR.

Test: ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest#subtaskMustNotGoIdleWhileRecordsFlowThroughItsSplits fails with this PR's production changes applied. subtaskWithoutSplitsGoesIdleViaTheMainOutputTimer checks the guard.

  1. WatermarksWithIdleness doesn't undo its own markIdle()

WatermarksWithIdleness.onEvent() resets the idleness timer and clears isIdleNow, but it never calls output.markActive(). The output stays idle until the wrapped generator emits a watermark:

  • With a periodic generator, recovery waits for the next watermark tick, and records sent in that window go out on a channel still marked idle.
  • If the wrapped generator never emits, the output stays idle for good even though records are flowing.

A record is evidence of activity on its own, so the fix is to call output.markActive() in onEvent() when isIdleNow is set. WatermarksWithIdleness is @public, but only its idle→active behaviour changes: it now sends the markActive() that its own markIdle() left outstanding.

Test: WatermarksWithIdlenessTest#testMarksActiveOnFirstEventAfterIdleness.

Two smaller notes

  • WatermarkToDataOutput now undoes an explicit markIdle() if the generator keeps re-emitting the same watermark. That matches the WatermarkOutput contract, but it's a change that source implementers can observe. It would be worth a test and a line in the release notes.
  • Your eager re-activation in registerNewOutput() (FLINK-22926) is something my PR doesn't have. I'd keep it in whichever version lands.

@Admaing
Admaing force-pushed the FLINK-40504-reattivate-idle-output branch 2 times, most recently from f26a8bc to 40e286b Compare September 23, 2026 16:10
@Admaing

Admaing commented Sep 23, 2026

Copy link
Copy Markdown
Contributor Author

@rkhachatryan thanks for the careful review — I verified both points in the code and you're right on
both.

Adopted, pushed as 40e286b:

  • WatermarkOutputMultiplexer now reports the combined status on every update (level-based, no transition tracking), and returns early while no output is registered, so a subtask without splits can still go idle through the main output's activity timer.
  • WatermarksWithIdleness.onEvent() calls output.markActive() when isIdleNow is set, and the class Javadoc no longer claims that only a watermark ends idleness.
  • WatermarkToDataOutput.emitWatermark() marks the output active before the monotonicity guard.
  • The registerNewOutput() re-activation (FLINK-22926) is kept, as you suggested.

Tests: I added the two subtask-idleness cases you described
(subtaskMustNotGoIdleWhileRecordsFlowThroughItsSplits,
subtaskWithoutSplitsGoesIdleViaTheMainOutputTimer) plus two along the same lines — all splits idle ⇒ the subtask stands down, and an idle split resuming behind the flushed watermark ⇒ active again —
as well as WatermarksWithIdlenessTest#testMarksActiveOnFirstEventAfterIdleness and a multiplexer case asserting that nothing is reported while there is no output. With the production changes reverted to master, 10 of the new/updated cases fail; a 21-class / 142-test watermark regression passes.

Two things to flag:

  • Keeping the registerNewOutput() re-activation changes the event order in the existing
    SourceOperatorEventTimeTest#testMainAndPerSplitWatermarkIdleness: registering a split now emits ACTIVE before that split's first record (FLINK-22926). I updated that expectation.
  • I added a release note line covering the WatermarkToDataOutput behaviour change and
    WatermarksWithIdleness being @Public. I cannot set the JIRA Release Note field myself, so I asked a committer to set it on the ticket.

Could you take another look at 40e286b?

@rkhachatryan rkhachatryan left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the quick turnaround — both changes are what I had in mind.

CI is red on 40e286b , and I think the failure is on me rather than on you: it falls out of the level-based reporting I asked for in point 1. SourceOperatorAlignmentTest has 6 failures across testWatermarkAlignmentWithIdleness and ...AllSubtasksIdle — expected ReportedWatermarkEvent{watermark=9223372036854775807} but was {watermark=1}, i.e. the subtask never goes idle.

The cause is that the multiplexer's "active" is a default, not an observation. A fresh PartialWatermark is {MIN_VALUE, idle = false}, so an output that has never reported anything still makes isIdle() return false, and level-based reporting turns that into a positive markActive() on IdlenessManager.splitLocalOutput. Only a combined idle can set that back, which needs every registered split to have marked itself idle — so a split that never does pins the split side active forever and the main output's idleness timer becomes inert. That test registers a split at AddSplitEvent time but runs everything, including the explicit markIdle(), through the main output. More generally it hits any silent registered split: a strategy without withIdleness, a reader not using per-split outputs, or auto-watermark-interval = 0.

Worth being clear that this isn't the FLINK-22926 hook in registerNewOutput() — that only makes it immediate, since onPeriodicEmit() calls updateCombinedWatermark() unconditionally anyway. My "report on every update where the status is active" needs qualifying: report it when it reflects something an output actually said.

One way to do that is to make PartialWatermark tri-state — UNKNOWN/ACTIVE/IDLE, starting UNKNOWN, with setIdle() moving it to ACTIVE/IDLE. UNKNOWN combines exactly as idle = false does today, so a newly assigned split still holds the combined watermark back, but it stops counting as evidence. CombinedWatermarkStatus then exposes a hasKnownActivity() (any partial not UNKNOWN) that updateCombinedWatermark() guards on instead of hasOutputs() — it returns false for an empty list, so it subsumes that guard, and it lets isIdle() go back to private, which resolves the visibility widening currently in the diff with no caller. Point 1 stays fixed, since a split emitting non-advancing watermarks goes through setWatermark() and is therefore known-active; FLINK-22926 stays fixed too, since an already-idle split has known activity.

Another option I see is a single flag on the multiplexer, set by ImmediateOutput/DeferredOutput.

Could you also add a multiplexer case for the regression itself: register an output, report nothing through it, assert the underlying output receives nothing. The existing "nothing is reported while there is no output" case passes either way.

@Admaing
Admaing force-pushed the FLINK-40504-reattivate-idle-output branch from 40e286b to 627b6ce Compare September 24, 2026 00:23
@Admaing

Admaing commented Sep 24, 2026

Copy link
Copy Markdown
Contributor Author

@rkhachatryan you're right on all counts, and thanks for digging into it.

Confirmed and fixed in 627b6ce:

  • PartialWatermark is now tri-state (UNKNOWN/ACTIVE/IDLE) and starts UNKNOWN, and
    updateCombinedWatermark() guards on CombinedWatermarkStatus#hasKnownActivity() instead of
    hasOutputs() — which also lets PartialWatermark#isIdle() go back to private, as you noted.
  • SourceOperatorAlignmentTest passes again locally (21/21).
  • Added WatermarkOutputMultiplexerTest#whenRegisteredOutputReportsNothingNothingIsReported; it is
    red against the previous commit (40e286b) and green now.

One thing your proposal did not cover, which I hit while running the tests: your
subtaskMustNotGoIdleWhileRecordsFlowThroughItsSplits uses a generator that never emits, so the
split never reaches setWatermark() and stays UNKNOWN — with the tri-state alone the subtask still
went idle while records flowed. So WatermarksWithIdleness.onEvent() now reports activity on the
first record (and again after each idle period), not only when undoing its own markIdle():

if (isIdleNow || !isActivityReported) { output.markActive(); ... }

That keeps "a generator that emits nothing" (explicitly listed in the ticket) covered without
reporting on every record.

Two consequences worth knowing:

  • testMainAndPerSplitWatermarkIdleness is back to its original expectation: registering a split no
    longer emits ACTIVE by itself. FLINK-22926 still works when at least one other split has reported
    its state (the common case), but a subtask with no splits at all and an idle main output now
    re-activates on the first record of the new split rather than at registration. I can add the
    immediacy back if you think it matters.
  • Regression over all watermark-related classes in flink-core/flink-runtime: 22 classes, 164 tests,
    green.

@Jackeyzhe

Copy link
Copy Markdown
Contributor

I found one mixed-state case that the new hasKnownActivity() guard does not cover. If a subtask is already idle, register a silent split (UNKNOWN), then register another split and call markIdle() on it. The latter makes hasKnownActivity() true, but the UNKNOWN split keeps the combined idle false, so updateCombinedWatermark() forwards markActive() even though no split has produced a record or watermark. A temporary ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest on 627b6ce expected the downstream idle updates to remain [true], but got [true, false]. Could we distinguish known ACTIVE from the mixed UNKNOWN+IDLE case for the downstream activity signal, while retaining UNKNOWN's watermark-holding behavior? A regression for the mixed state would help.

@Admaing
Admaing force-pushed the FLINK-40504-reattivate-idle-output branch from 627b6ce to be7924c Compare September 24, 2026 16:17
@Admaing

Admaing commented Sep 25, 2026

Copy link
Copy Markdown
Contributor Author

@Jackeyzhe confirmed and fixed in be7924c — I reproduced [true] → [true, false] for UNKNOWN + IDLE.

The activity signal is now guarded by hasActiveOutput(), so only outputs that reported ACTIVE count;
UNKNOWN still holds the combined watermark back, so its watermark behaviour is unchanged. Added
subtaskStaysIdleWhenOnlyAnUnknownAndAnIdleSplitAreRegistered for the mixed state — red on 627b6ce,
green now.

Note: registration is not evidence of activity either, so an idle subtask now re-activates on the new
split's first record rather than when it is assigned; FLINK-22926's registration-time immediacy is
dropped and noted as out of scope in the PR description. test_ci core is green on be7924c.

@rkhachatryan rkhachatryan left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for updating the PR!

Mostly LGTM.

Would you mind splitting the change into several commits?

E.g.

1. [FLINK-40504][core] Make PartialWatermark idleness tri-state: an unknown/active/idle enum plus hasActiveOutput(), with no behaviour change, since "unknown" combines the same way idle = false does today. 
This keeps the refactor apart from the behaviour changes.

2. [FLINK-40504][runtime] Mark WatermarkToDataOutput active on non-advancing watermarks: the change to WatermarkToDataOutput and its test.
This is the contract fix that needs a release note.

3. [FLINK-40504][core] Report activity from WatermarksWithIdleness on events: the onEvent() change and testMarksActiveOnFirstEventAfterIdleness.
It's a @Public class, so a separate commit makes its release note easy to trace.

4. [FLINK-40504][core] Propagate active combined status from WatermarkOutputMultiplexer: 
the hasActiveOutput() call in updateCombinedWatermark(), 
the 4 multiplexer tests, 
the new ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest and 
the new SourceOperatorEventTimeTest case.

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 FLINK-40504-reattivate-idle-output branch from be7924c to 32fd114 Compare September 26, 2026 17:39
@Admaing

Admaing commented Sep 26, 2026

Copy link
Copy Markdown
Contributor Author

Thanks for updating the PR!

Mostly LGTM.

Would you mind splitting the change into several commits?

E.g.

1. [FLINK-40504][core] Make PartialWatermark idleness tri-state: an unknown/active/idle enum plus hasActiveOutput(), with no behaviour change, since "unknown" combines the same way idle = false does today. 
This keeps the refactor apart from the behaviour changes.

2. [FLINK-40504][runtime] Mark WatermarkToDataOutput active on non-advancing watermarks: the change to WatermarkToDataOutput and its test.
This is the contract fix that needs a release note.

3. [FLINK-40504][core] Report activity from WatermarksWithIdleness on events: the onEvent() change and testMarksActiveOnFirstEventAfterIdleness.
It's a @Public class, so a separate commit makes its release note easy to trace.

4. [FLINK-40504][core] Propagate active combined status from WatermarkOutputMultiplexer: 
the hasActiveOutput() call in updateCombinedWatermark(), 
the 4 multiplexer tests, 
the new ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest and 
the new SourceOperatorEventTimeTest case.

Done — split into the four commits you suggested.

@Jackeyzhe Jackeyzhe left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for addressing the UNKNOWN+IDLE case.

I revalidated 72a0711 with six focused watermark test classes: all 77 tests passed. The activity guard and mixed-state regression look good to me. LGTM on the code changes.

CI still needs follow-up.

@github-actions github-actions Bot added the community-reviewed PR has been reviewed by the community. label Sep 28, 2026

@rkhachatryan rkhachatryan left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@rkhachatryan

Copy link
Copy Markdown
Contributor

@flinkbot run azure

@rkhachatryan

Copy link
Copy Markdown
Contributor

CI failure seems to be unrelated (FLINK-40069)
I'm going to merge the PR.

Thank you @Admaing and @Jackeyzhe for fixing and reviewing!

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

community-reviewed PR has been reviewed by the community.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants