Martijn Visser created FLINK-40505:
--------------------------------------
Summary: SourceCoordinator never removes failed or reset subtasks
from watermark-alignment aggregation
Key: FLINK-40505
URL: https://issues.apache.org/jira/browse/FLINK-40505
Project: Flink
Issue Type: Bug
Components: Runtime / Coordination
Reporter: Martijn Visser
The per-subtask watermark aggregation used for watermark alignment retains a
subtask's
last {{ReportedWatermarkEvent}} across failures and resets.
{{WatermarkAggregator}}
({{SourceCoordinator.java:835}}) exposes only {{aggregate}}/{{keySet}} — there
is no
removal API — and {{executionAttemptFailed}} ({{:352-369}}) / {{subtaskReset}}
({{:371-394}}) clear registered readers and gateways but never touch the
aggregator
(neither the per-coordinator one nor the group-level one in
{{coordinatorStore}}).
Consequence: after a failover, the dead attempt's stale watermark keeps
constraining
{{maxAllowedWatermark}} for the entire alignment group until the restarted
attempt
happens to re-report. The retained value is not a conservative bound — the
restarted
attempt can resume from an earlier checkpoint position than its last report —
and on
downscale or permanent removal the entry is never refreshed, freezing the
group's
{{maxAllowedWatermark}} permanently.
Red test on master: two subtasks report watermarks 100 and 200 (announced
maxAllowed=1100 with drift 1000 — sane); after {{executionAttemptFailed}} +
{{subtaskReset}} of the subtask that reported 100, a fresh announcement still
sends
maxAllowed=1100 from the stale entry (expected: 1200).
Related facet (characterization test): a subtask that has never sent a
{{ReportedWatermarkEvent}} receives no {{WatermarkAlignmentEvent}} at all (the
announce
loop iterates {{combinedWatermark.keySet()}}), overlapping FLINK-36914.
This may be the mechanism behind FLINK-40371 (alignment drift bound not honored
after
savepoint restore); FLINK-33109 reported the same symptom and was closed as an
HA
artifact; FLINK-32414 fixed the splitless-subtask MIN_VALUE variant of the same
aggregator-poisoning shape.
Proposed fix: remove the subtask's entry in {{subtaskReset}} (and on attempt
failure
once no attempt of the subtask is running); the entry is re-established by the
restarted
attempt's next report.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)