Kushagra Tyagi created FLINK-40371:
--------------------------------------
Summary: Watermark alignment does not consistently bound
cross-subtask event-time drift after restoring from a savepoint
Key: FLINK-40371
URL: https://issues.apache.org/jira/browse/FLINK-40371
Project: Flink
Issue Type: Bug
Components: Connectors / Kafka
Affects Versions: kafka-3.2.0
Reporter: Kushagra Tyagi
Attachments: catchup_using_savepoint.png, catchup_using_timestamp.png
We configured a KafkaSource with forBoundedOutOfOrderness(100ms) and
withWatermarkAlignment(groupName, maxDrift), across roughly 48 parallel source
subtasks reading a partitioned Kafka topic.
We observed a consistent difference in behavior depending on how the job was
started when it had a large backlog to catch up on:
- Fresh start from a timestamp (OffsetsInitializer.timestamp(...), no prior
state): subtasks' per-record event-time lag (currentEmitEventTimeLag) stayed
tightly bounded across all subtasks, within roughly the configured max drift
plus a small margin, even when the backlog represented many hours of lag
(observed up to ~22 hours of lag with drift staying within a few minutes across
subtasks). (catchup_using_timestamp.png)
- Restart from a savepoint/checkpoint after the job was stopped, with a
comparable backlog to catch up: subtasks' currentEmitEventTimeLag diverged from
each other by multiple minutes, growing over roughly the first 15+ minutes
after restart, well beyond the configured max drift.
(catchup_using_savepoint.png)
Notably, the watermarkAlignmentDrift metric stayed within/near the configured
bound in both cases, while currentEmitEventTimeLag diverged significantly in
the savepoint-restore case. This suggests the coordinator's reported/enforced
alignment drift may become decoupled from the actual per-subtask processing
pace specifically during the window right after restoring from a savepoint,
while it holds correctly for a subtask population that starts from a
clean/fresh state.
*Expected behavior:* cross-subtask watermark alignment should bound actual
per-record event-time drift within roughly the configured max drift, regardless
of whether the job started fresh or resumed from a savepoint.
*Actual behavior:* alignment holds correctly on a fresh start; after a
savepoint restore, subtasks can diverge by multiple minutes for an extended
period before (if ever — not yet confirmed) converging.
*Affects:* Flink 1.19.x, flink-connector-kafka 3.2.0-1.19 (observed; not yet
tested on other versions)
*Note:* this is based on production metrics observed on a live pipeline, not an
isolated minimal reproducer. Happy to help build one if that would help triage.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)