[ 
https://issues.apache.org/jira/browse/FLINK-40371?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Kushagra Tyagi updated FLINK-40371:
-----------------------------------
    Description: 
*Setup*

We configured a KafkaSource with forBoundedOutOfOrderness(100ms) and 
withWatermarkAlignment(groupName, maxDrift), across roughly 48 parallel source 
subtasks reading a partitioned Kafka topic.

*What we observed*

The behavior differs depending on how the job is started when it has a large 
backlog to catch up on.
 * Fresh start from a timestamp (OffsetsInitializer.timestamp(...), no prior 
state):
Subtasks stayed tightly aligned. Per record event time lag 
(currentEmitEventTimeLag) stayed within roughly the configured max drift plus a 
small margin across all subtasks, even with a backlog representing many hours 
of lag. We observed this holding at up to about 22 hours of lag, with drift 
staying within a few minutes across subtasks. See attachment: 
catchup_using_timestamp.png.
 * Restart from a savepoint or checkpoint after the job was stopped, with a 
comparable backlog:
Subtasks diverged from each other by multiple minutes. currentEmitEventTimeLag 
grew apart over roughly the first 15 minutes after restart, well beyond the 
configured max drift, before beginning to converge. See attachment: 
catchup_using_savepoint.png.

*Key finding*

watermarkAlignmentDrift stayed within or near the configured bound in both 
cases. currentEmitEventTimeLag diverged significantly, but only in the 
savepoint restore case.

This suggests the coordinator's reported alignment drift can become decoupled 
from the actual per subtask processing pace, specifically in the window right 
after restoring from a savepoint. It holds correctly when subtasks start from a 
clean state.

*Metrics referenced*
- flink.operator.currentEmitEventTimeLag: per subtask, per record event time 
lag at emission. This is the metric that showed the real divergence between 
subtasks after a savepoint resume.
- flink.operator.watermarkAlignmentDrift: the coordinator's own reported 
alignment drift. This stayed within bound in both cases, including the case 
where real lag diverged.

*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 roughly 10 to 15 minutes before converging.

*Affects*

Flink 1.19.x, flink connector kafka 3.2.0 for 1.19. 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 with 
triage.

  was:
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.

*Metric:* flink.operator.currentEmitEventTimeLag — per subtask, per record 
event time lag at emission. This is the one that showed the real divergence 
between subtasks after a savepoint resume.

*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.


> 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
>
>
> *Setup*
> We configured a KafkaSource with forBoundedOutOfOrderness(100ms) and 
> withWatermarkAlignment(groupName, maxDrift), across roughly 48 parallel 
> source subtasks reading a partitioned Kafka topic.
> *What we observed*
> The behavior differs depending on how the job is started when it has a large 
> backlog to catch up on.
>  * Fresh start from a timestamp (OffsetsInitializer.timestamp(...), no prior 
> state):
> Subtasks stayed tightly aligned. Per record event time lag 
> (currentEmitEventTimeLag) stayed within roughly the configured max drift plus 
> a small margin across all subtasks, even with a backlog representing many 
> hours of lag. We observed this holding at up to about 22 hours of lag, with 
> drift staying within a few minutes across subtasks. See attachment: 
> catchup_using_timestamp.png.
>  * Restart from a savepoint or checkpoint after the job was stopped, with a 
> comparable backlog:
> Subtasks diverged from each other by multiple minutes. 
> currentEmitEventTimeLag grew apart over roughly the first 15 minutes after 
> restart, well beyond the configured max drift, before beginning to converge. 
> See attachment: catchup_using_savepoint.png.
> *Key finding*
> watermarkAlignmentDrift stayed within or near the configured bound in both 
> cases. currentEmitEventTimeLag diverged significantly, but only in the 
> savepoint restore case.
> This suggests the coordinator's reported alignment drift can become decoupled 
> from the actual per subtask processing pace, specifically in the window right 
> after restoring from a savepoint. It holds correctly when subtasks start from 
> a clean state.
> *Metrics referenced*
> - flink.operator.currentEmitEventTimeLag: per subtask, per record event time 
> lag at emission. This is the metric that showed the real divergence between 
> subtasks after a savepoint resume.
> - flink.operator.watermarkAlignmentDrift: the coordinator's own reported 
> alignment drift. This stayed within bound in both cases, including the case 
> where real lag diverged.
> *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 roughly 10 to 15 minutes before 
> converging.
> *Affects*
> Flink 1.19.x, flink connector kafka 3.2.0 for 1.19. 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 with 
> triage.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to