[ 
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)

Reply via email to