Martijn Visser created FLINK-40475:
--------------------------------------

             Summary: StatusWatermarkValve loses or stalls watermarks when 
subpartitions realign after idleness
                 Key: FLINK-40475
                 URL: https://issues.apache.org/jira/browse/FLINK-40475
             Project: Flink
          Issue Type: Bug
          Components: Runtime / Task
            Reporter: Martijn Visser


When a subpartition of {{StatusWatermarkValve}} reactivates after idleness and 
rejoins the aligned set, the min watermark across aligned subpartitions is not 
re-derived. This breaks the invariant that the aligned min equals the last 
output watermark, which two guards ({{subpartitionStatus.watermark == 
lastOutputWatermark}}) silently rely on. Two symptoms:

The all-idle flush-to-max from FLINK-7728 is skipped when the last subpartition 
to become idle is unaligned (went idle, resumed, only partially caught up). The 
final watermark then depends on the order in which inputs became idle — the 
exact defect FLINK-7728 was created to fix. Example: inputs at 100/60/70, last 
output watermark 70; if the input at 60 idles last, the valve emits only IDLE 
and the watermark stays at 70 instead of flushing 100. A different idle order 
emits 100. This also diverges from {{CombinedWatermarkStatus}}, which flushes 
unconditionally since FLINK-38454 despite the "keep in sync" contract.

A subpartition that reactivates as the only aligned subpartition has its 
watermark stalled indefinitely — it is only emitted once an even larger 
watermark arrives on that subpartition. No all-idle transition is involved.

The pattern is reachable with standard sources: {{WatermarkToDataOutput}} 
enforces watermark monotonicity only per subtask and emits ACTIVE before the 
resuming watermark, so e.g. a quiet Kafka partition with 
{{WatermarksWithIdleness}} resuming with old-timestamped backlog reactivates 
below the downstream valve's last output watermark.

Fix: re-derive the min watermark after realignment on reactivation, and drop 
both guards — {{findAndOutputMaxWatermarkAcrossAllSubpartitions}} and 
{{findAndOutputNewMinWatermarkAcrossAlignedSubpartitions}} already only emit on 
advancement, so monotonicity is unaffected. Unit tests reproducing both 
symptoms (failing on current master) are included in the PR.



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

Reply via email to