[
https://issues.apache.org/jira/browse/FLINK-40456?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated FLINK-40456:
-----------------------------------
Labels: pull-request-available (was: )
> AsyncWaitOperator can permanently drop an element's result when a timeout
> races with a retry
> --------------------------------------------------------------------------------------------
>
> Key: FLINK-40456
> URL: https://issues.apache.org/jira/browse/FLINK-40456
> Project: Flink
> Issue Type: Sub-task
> Components: API / DataStream, Build System / CI
> Reporter: Sergey Nuyanzin
> Assignee: Sergey Nuyanzin
> Priority: Major
> Labels: pull-request-available
>
> When async retries are enabled and an element also has a timeout configured,
> the operator can lose that element's result. The element is then never
> emitted and never fails — it stays "in flight" forever, so the job (or a
> test) hangs
> indefinitely.
> When an element times out, the operator is supposed to emit the timeout's
> default result and stop. But that timeout result is passed through the retry
> code path, which begins with a guard meaning "a retry is already in progress
> for this
> element — ignore this completion." If a retry for the same element was just
> scheduled or is being processed at that moment, the timeout result hits that
> guard and is silently discarded. By then the timeout handler has already
> cancelled the
> element's retry timer, and the element was never marked as completed — so
> nothing remains that could ever complete it. The element is stuck permanently.
> Steps to reproduce
> {noformat}
> docker run --rm --cpus=2 -v {Path-to-flink-local-repo}:/flink -v
> {path-to-m2}:/root/.m2 -w /flink -e JAVA_HOME=/usr/lib/jvm/jdk-25.0.3+9
> apache/flink-ci-docker:java_8_11_17_21_25_maven_386_jammy ./mvnw -o -q
> -pl flink-streaming-java -Phadoop3-tests,hive3 -Djdk25 -Pjava25-target
> -Dsurefire.failIfNoSpecifiedTests=false
> -Dtest='AsyncWaitOperatorTest#testProcessingTimeWithTimeoutFunctionUnorderedWithRetry'
> surefire:test
> {noformat}
> it fails as
> {noformat}
> [ERROR]
> org.apache.flink.streaming.api.operators.async.AsyncWaitOperatorTest.testProcessingTimeWithTimeoutFunctionUnorderedWithRetry
> -- Time elapsed: 100.1 s <<< ERROR!
> java.util.concurrent.TimeoutException:
> testProcessingTimeWithTimeoutFunctionUnorderedWithRetry() timed out after 100
> seconds
> at
> java.base/java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:511)
> at
> java.base/java.util.concurrent.ForkJoinPool$WorkQueue.topLevelExec(ForkJoinPool.java:1450)
> at
> java.base/java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:2019)
> at
> java.base/java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:187)
> Suppressed: java.lang.InterruptedException: sleep interrupted
> at java.base/java.lang.Thread.sleepNanos0(Native Method)
> at java.base/java.lang.Thread.sleepNanos(Thread.java:509)
> at java.base/java.lang.Thread.sleep(Thread.java:540)
> at
> org.apache.flink.streaming.api.operators.async.AsyncWaitOperatorTest.testProcessingTimeAlwaysTimeoutFunctionWithRetry(AsyncWaitOperatorTest.java:1444)
> at
> org.apache.flink.streaming.api.operators.async.AsyncWaitOperatorTest.testProcessingTimeWithTimeoutFunctionUnorderedWithRetry(AsyncWaitOperatorTest.java:1399)
> at java.base/java.lang.reflect.Method.invoke(Method.java:565)
> ... 4 more
> {noformat}
> it is very likely the reason it appeared in 100% of cases for java25 (however
> might be reproduced with the line above for other jdks) is that
> {quote}
> - Before JDK 25 (11/17/21): CompletableFuture.runAsync(...) with no explicit
> executor, when the common pool has parallelism < 2 (e.g. a 2‑CPU box:
> parallelism 1), spawned a fresh thread per async task. So the two elements'
> blocking callbacks
> ran concurrently (~545ms)=> the timeout-vs-retry race almost never
> happened.
> - JDK 25: those async methods now run on the common pool itself. At
> parallelism 1 there's a single worker, so the callbacks serialize (~1065ms)
> => which is the timing that triggers the operator's dropped-completion race.
> {quote}
> probably this change https://bugs.openjdk.org/browse/JDK-8362881
--
This message was sent by Atlassian Jira
(v8.20.10#820010)