[
https://issues.apache.org/jira/browse/FLINK-40234?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Keith Lee updated FLINK-40234:
------------------------------
Description:
Consider the following:
# WatermarksWithIdleness#onEvent() runs on subtask thread. It increments split
activity whenever it is called.
# WatermarksWithIdleness#checkIfIdle() also runs on subtask thread. It checks
for idleness based on relative clock's elapsed time since inactivity was first
detected[
https://github.com/apache/flink/blob/f9485ff173b940c5ddf606d8ca1a92a18bb23ff3/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarksWithIdleness.java#L140|https://github.com/apache/flink/blob/f9485ff173b940c5ddf606d8ca1a92a18bb23ff3/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarksWithIdleness.java#L140]
If records processing (also subtask thread) takes a long time, records flow
through slower and onEvent() for a split with records may not have been called
yet because SourceOperator has not reach the same split. In this scenario,
relative clock's elapsed time can exceed threshold and the split is marked as
idle, even though the split has records.
See PR where unit test is added to demonstrate this issue
[https://github.com/apache/flink/pull/28816]
was:
See PR where unit test is added to demonstrate this issue
[https://github.com/apache/flink/pull/28816]
> Busy subtask thread may cause splits to be marked as idle erroneously and
> cause dropped records
> -----------------------------------------------------------------------------------------------
>
> Key: FLINK-40234
> URL: https://issues.apache.org/jira/browse/FLINK-40234
> Project: Flink
> Issue Type: Bug
> Components: Runtime / Task
> Affects Versions: 2.3.0, 2.2.1
> Reporter: Keith Lee
> Priority: Major
> Labels: pull-request-available
>
> Consider the following:
> # WatermarksWithIdleness#onEvent() runs on subtask thread. It increments
> split activity whenever it is called.
> # WatermarksWithIdleness#checkIfIdle() also runs on subtask thread. It
> checks for idleness based on relative clock's elapsed time since inactivity
> was first detected[
> https://github.com/apache/flink/blob/f9485ff173b940c5ddf606d8ca1a92a18bb23ff3/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarksWithIdleness.java#L140|https://github.com/apache/flink/blob/f9485ff173b940c5ddf606d8ca1a92a18bb23ff3/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarksWithIdleness.java#L140]
> If records processing (also subtask thread) takes a long time, records flow
> through slower and onEvent() for a split with records may not have been
> called yet because SourceOperator has not reach the same split. In this
> scenario, relative clock's elapsed time can exceed threshold and the split is
> marked as idle, even though the split has records.
> See PR where unit test is added to demonstrate this issue
> [https://github.com/apache/flink/pull/28816]
--
This message was sent by Atlassian Jira
(v8.20.10#820010)