[ 
https://issues.apache.org/jira/browse/FLINK-40504?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18117790#comment-18117790
 ] 

yongfu.gao edited comment on FLINK-40504 at 9/23/26 4:11 PM:
-------------------------------------------------------------

Update: the PR now covers all three places on the path that treated an 
advancing watermark as the only evidence of activity.

https://github.com/apache/flink/pull/29239

* WatermarkOutputMultiplexer reports the combined status whenever it is idle or 
active, so an output that is active without advancing its watermark is no 
longer left announced as idle. While no output is registered the combined 
status is meaningless, so nothing is reported.
* WatermarksWithIdleness.onEvent() marks the output active on the first record 
after idleness, instead of waiting for the wrapped generator to emit an 
advancing watermark.
* WatermarkToDataOutput.emitWatermark() marks the output active before the 
monotonicity guard, as required by the WatermarkOutput contract.
* registerNewOutput() re-activates the combined status right away, which covers 
the new-split case of the same defect (FLINK-22926).

Verification: 4 new WatermarkOutputMultiplexerTest cases, 1 new 
WatermarksWithIdlenessTest case, 1 new WatermarkToDataOutputTest case and 4 new 
end-to-end cases in ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest. With 
the production changes reverted to master, 10 of the new/updated cases fail; a 
21-class / 142-test watermark regression passes.


was (Author: JIRAUSER314604):
[~martijnvisser] The PR for this issue is up: 
https://github.com/apache/flink/pull/29239

It addresses the root cause of the two mechanisms described in the ticket:

* WatermarkOutputMultiplexer now reports the idle -> active transition of the 
combined status to
  the underlying output (previously it only reported markIdle()), and a 
watermark on a previously
  idle multiplexed output triggers the combined update, so the re-activation is 
no longer delayed
  until the next periodic emit.
* registerNewOutput() reports the combined status immediately, which covers the 
new-split case of
  the same defect described in FLINK-22926.
* WatermarkToDataOutput.emitWatermark() marks the output active before the 
monotonicity guard, as
  required by the WatermarkOutput contract.

Verification: 7 new tests (4 in WatermarkOutputMultiplexerTest, 1 in 
WatermarkToDataOutputTest and
2 end-to-end in SourceOperatorEventTimeTest). All 7 were verified to fail 
without the corresponding
production change and to pass with it; a broader regression over 20 
watermark-related test classes
(137 tests) passes on a clean build. Azure CI on ca8d2501 is green.

This is complementary to #29037 (FLINK-40499) and does not touch the files 
changed there.

Could a committer please take a look, or point me to someone who can review it?

> FLIP-27 source never re-emits WatermarkStatus.ACTIVE after idleness unless 
> its watermark strictly advances
> ----------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40504
>                 URL: https://issues.apache.org/jira/browse/FLINK-40504
>             Project: Flink
>          Issue Type: Bug
>          Components: API / Core, Runtime / Task
>            Reporter: Martijn Visser
>            Assignee: yongfu.gao
>            Priority: Major
>              Labels: pull-request-available
>
> Two compounding mechanisms leave a resumed FLIP-27 source announced-IDLE 
> downstream:
> # {{WatermarkOutputMultiplexer.updateCombinedWatermark()}} calls
>   {{underlyingOutput.markIdle()}} when combined-idle but has no
>   {{underlyingOutput.markActive()}} counterpart; the per-output {{markActive}}
>   implementations only mutate internal partial state.
> # {{WatermarkToDataOutput.emitWatermark}} returns at the monotonicity guard
>   ({{newWatermark <= maxWatermarkSoFar}}) *before* {{markActiveInternally()}}
>   ({{WatermarkToDataOutput.java:77-85}}).
> The only path back to ACTIVE is therefore a strictly larger watermark. 
> Scenario: all
> splits of a subtask go idle; the combined watermark is flushed to the max M 
> across
> splits (FLINK-38454). A split then resumes with backlog whose watermarks are 
> <= M
> (e.g. a recovered Kafka partition catching up). Result — proven by red unit 
> tests on
> both layers: *zero* calls reach the underlying output (no watermark, no 
> ACTIVE), even
> via {{onPeriodicEmit}}; the edge stays IDLE downstream indefinitely while 
> records flow.
> Downstream min-watermark is then driven by the other inputs, past the resumed 
> backlog,
> and its records are dropped as late.
> Every sibling path re-activates eagerly on resume: {{StatusWatermarkValve}} 
> (explicitly,
> since FLINK-40475), the legacy {{StreamSourceContexts}} (ACTIVE on every 
> collect), the
> table {{WatermarkAssignerOperator}} (ACTIVE on first record), and DataStream 
> V2
> (explicit idle=false watermark on next record). The FLIP-27 multiplexer path 
> is the only
> one without a re-activation signal.
> FLINK-22926 (open since 2021) describes the new-split registration special 
> case of this
> same defect; it should be closed into this ticket.
> Note: the fix is a semantics decision (emit ACTIVE eagerly on 
> record/markActive, aligning
> with all other paths) rather than a local patch - it also composes with 
> FLINK-40499: a
> resumed-but-not-advanced source is exactly the still-IDLE state in which its 
> final
> MAX_WATERMARK is then dropped at drain.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to