[
https://issues.apache.org/jira/browse/FLINK-40371?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Kushagra Tyagi updated FLINK-40371:
-----------------------------------
Component/s: Runtime / Coordination
> 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, Runtime / Coordination
> Affects Versions: kafka-3.2.0
> Reporter: Kushagra Tyagi
> Priority: Major
> 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)