[
https://issues.apache.org/jira/browse/FLINK-40371?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Kushagra Tyagi updated FLINK-40371:
-----------------------------------
Description:
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.
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.
*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
>
>
> 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.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)