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)