[FLINK-40504][runtime] Re-activate idle FLIP-27 outputs when they resume - #29239
Conversation
rkhachatryan
left a comment
There was a problem hiding this comment.
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.
- 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.
- 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.
f26a8bc to
40e286b
Compare
|
@rkhachatryan thanks for the careful review — I verified both points in the code and you're right on Adopted, pushed as 40e286b:
Tests: I added the two subtask-idleness cases you described Two things to flag:
Could you take another look at 40e286b? |
rkhachatryan
left a comment
There was a problem hiding this comment.
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.
40e286b to
627b6ce
Compare
|
@rkhachatryan you're right on all counts, and thanks for digging into it. Confirmed and fixed in 627b6ce:
One thing your proposal did not cover, which I hit while running the tests: your That keeps "a generator that emits nothing" (explicitly listed in the ticket) covered without Two consequences worth knowing:
|
|
I found one mixed-state case that the new |
627b6ce to
be7924c
Compare
|
@Jackeyzhe confirmed and fixed in be7924c — I reproduced The activity signal is now guarded by Note: registration is not evidence of activity either, so an idle subtask now re-activates on the new |
rkhachatryan
left a comment
There was a problem hiding this comment.
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)
be7924c to
32fd114
Compare
Done — split into the four commits you suggested. |
Jackeyzhe
left a comment
There was a problem hiding this comment.
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.
|
@flinkbot run azure |
|
CI failure seems to be unrelated (FLINK-40069) Thank you @Admaing and @Jackeyzhe for fixing and reviewing! |
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:
WatermarkOutputMultiplexer.updateCombinedWatermark()reports the combined idle state viaunderlyingOutput.markIdle(), but nothing when the status is active. An output that is activewithout advancing its watermark therefore produces neither an
emitWatermarknor amarkIdlecall, and the underlying output never learns that it is active again. Because the split branch of
ProgressiveTimestampsAndWatermarks.IdlenessManagerstarts out idle and is only ever cleared byone of those two calls, an active split is then ANDed into IDLE by the main output's idleness
timer.
WatermarksWithIdleness.onEvent()resets its idleness timer without reporting that the output isactive again, so the output stays announced as idle until the wrapped generator produces an
advancing watermark.
WatermarkToDataOutput.emitWatermark()returns at the monotonicity guard(
newWatermark <= maxWatermarkSoFar) before marking the output active, although theWatermarkOutputcontract states that emitting a watermark implicitly marks the stream active.Downstream,
StatusWatermarkValvethen 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, legacyStreamSourceContexts, thetable
WatermarkAssignerOperator, and DataStream V2); the FLIP-27 multiplexer path was the only onewithout a re-activation signal.
Brief change log
PartialWatermarkis tri-state (UNKNOWN/ACTIVE/IDLE) and startsUNKNOWN. An output that isregistered 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.
WatermarkOutputMultiplexerreports the combined status whenever it is idle or active, but it onlyannounces activity when some output actually reported being active
(
CombinedWatermarkStatus#hasActiveOutput()). Reporting idleness on one output therefore does notre-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.WatermarkOutputMultiplexerTestcases, 1 newWatermarksWithIdlenessTestcase, 1 newWatermarkToDataOutputTestcase and 5 new end-to-end cases inProgressiveTimestampsAndWatermarksSubtaskIdlenessTest.Verifying this change
This change added tests and can be verified as follows:
mvn -pl flink-core -am -Dtest='WatermarkOutputMultiplexerTest,WatermarksWithIdlenessTest' -Dsurefire.failIfNoSpecifiedTests=false testmvn -pl flink-runtime -am -Dtest='WatermarkToDataOutputTest,SourceOperatorEventTimeTest,ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest,SourceOperatorAlignmentTest' -Dsurefire.failIfNoSpecifiedTests=false testSourceOperatorAlignmentTestis the regression test for the activity signal: when the combinedstatus was reported on every active update, a split that is registered for an assigned split but
never reports anything kept the subtask active, and
testWatermarkAlignmentWithIdlenessneverreported
MAX_WATERMARK.WatermarkOutputMultiplexerTest#whenRegisteredOutputReportsNothingNothingIsReportedandProgressiveTimestampsAndWatermarksSubtaskIdlenessTest#subtaskStaysIdleWhenOnlyAnUnknownAndAnIdleSplitAreRegisteredcover the same case at the multiplexer and subtask level; the latter also covers the mixed
UNKNOWN+IDLEstate.upstream/master, 10 of the new/updated cases fail.flink-coreandflink-runtime(22 classes, 163 tests) passes.
Notes for reviewers:
IdlenessManager'sIdlenessAwareWatermarkOutputstarts withisIdle = true, and an output is registered for everyassigned 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-intervalis 0). Treating such anoutput as active pins the split branch active forever and makes the main output's idleness timer
inert.
WatermarksWithIdlenessis@Public, but only its idle -> active behaviour changes: it now sendsthe
markActive()that its ownmarkIdle()left outstanding, plus one for the first record.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.
WatermarkToDataOutputbehaviour change andWatermarksWithIdleness. I cannot set the JIRA Release Note field myself, so I asked a committer toset it on the ticket.
tasks; this PR fixes the root cause and does not touch those files.
Does this pull request potentially affect one of the following parts:
WatermarkOutputMultiplexerandWatermarkToDataOutputare@Internal;WatermarksWithIdlenessis@Publicand keeps its API, only its idleness handling is fixed)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 aftereach idle period
Documentation
Generated-by: DeepSeek Harness (deepseek-v4.1-flash)