YordanPavlov opened a new pull request, #29349:
URL: https://github.com/apache/flink/pull/29349

   
   ## What is the purpose of the change
   
   `AsyncRequestBuffer` schedules a periodic timeout check on 
`AsyncRequestBuffer.DELAYER`, a static scheduler shared by the whole 
TaskManager. Only `AsyncExecutionController#close()` cancels it, and nothing 
calls that. So every ended attempt of an async state (State V2) operator 
(failover, restart, rescale, cancel) leaves its task queued, keeping the 
controller and the disposed keyed state backend reachable until the TaskManager 
exits. Details, a MiniCluster repro and the heap dump evidence are in 
FLINK-40872.
   
   ## Brief change log
   
   - `AbstractAsyncKeyOrderedStreamOperator#close()` and 
`AbstractAsyncStateStreamOperatorV2#close()` close the 
`AsyncExecutionController`, in a `finally` block so it also runs when the drain 
or `super.close()` throws. These are the only two operator bases that create a 
controller.
   - `AsyncExecutionController#isBufferTimeoutScheduled()` 
(`@VisibleForTesting`) lets tests observe the scheduled timeout task.
   
   ## Verifying this change
   
   This change added tests and can be verified as follows:
   
   - Added `testCloseCancelsBufferTimeout` to 
`AbstractAsyncStateStreamOperatorTest` and 
`AbstractAsyncStateStreamOperatorV2Test`: the timeout task is scheduled after 
`open()` and cancelled after `close()`. Both fail without the fix.
   - The async processing tests in `flink-runtime` and `flink-streaming-java` 
pass locally.
   - We run this change patched into Flink 2.3.0 on a production ForSt job. A 
heap dump before the fix showed 10 disposed `ForStKeyedStateBackend`s reachable 
only through the DELAYER thread; with the fix, the number of backends on each 
TaskManager matches its running tasks.
   
   ## Does this pull request potentially affect one of the following parts:
   
   - Dependencies (does it add or upgrade a dependency): no
   - The public API, i.e., is any changed class annotated with 
`@Public(Evolving)`: no
   - The serializers: no
   - The runtime per-record code paths (performance sensitive): no
   - Anything that affects deployment or recovery: JobManager (and its 
components), Checkpointing, Kubernetes/Yarn, ZooKeeper: no
   - The S3 file system connector: no
   
   ## Documentation
   
   - Does this pull request introduce a new feature? no
   - If yes, how is the feature documented? not applicable
   
   ---
   
   ##### Was generative AI tooling used to co-author this PR?
   
   - [X] Yes (please specify the tool below)
   
   Generated-by: Claude Opus 5.5
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to