[ 
https://issues.apache.org/jira/browse/FLINK-40371?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18109207#comment-18109207
 ] 

Martijn Visser commented on FLINK-40371:
----------------------------------------

A possible mechanism for this, found while auditing watermark alignment + 
idleness:
the SourceCoordinator never removes failed/reset subtasks from the alignment
aggregation - {{executionAttemptFailed}}/{{subtaskReset}} clear readers and 
gateways
but the {{WatermarkAggregator}} has no removal API, so a stale
{{ReportedWatermarkEvent}} from before a restart keeps constraining
{{maxAllowedWatermark}} for the whole group until the restarted attempt 
re-reports
(and after downscale, forever). Filed with a reproducing unit test as 
FLINK-40505.

If your restore scenario involved failover or rescaling, this would match: 
could you
check whether the coordinator-side {{maxAllowedWatermark}} (or the announced 
alignment
events) stays pinned at a pre-restore value while the diverging subtasks run 
ahead?
That would confirm FLINK-40505. as the root cause here.

> 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.
> *Ruled out*
> We checked whether the savepoint itself held stale or misaligned state. It 
> did not. The job was well aligned at the moment it was stopped, so the 
> savepoint captured an already aligned position for every subtask. The 
> divergence appears after resume despite that, which suggests the difference 
> is in how the alignment coordinator initializes or reconnects with subtasks 
> on restore, not in the correctness of the saved state itself.
> *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