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

Yordan Pavlov edited comment on FLINK-40872 at 10/1/26 11:57 AM:
-----------------------------------------------------------------

I've opened  [a 
PR|[https://github.com/apache/flink/pull/29349]|http://example.com]  with the 
fix and regression tests (they fail without the fix).


was (Author: yordanpavlov):
I've opened  [a 
PR|[https://github.com/apache/flink/pull/29349]|http://example.com]]  with the 
fix and regression tests (they fail without the fix).

> AsyncExecutionController is never closed. Leak the state backend through the 
> static AsyncRequestBuffer.DELAYER
> --------------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40872
>                 URL: https://issues.apache.org/jira/browse/FLINK-40872
>             Project: Flink
>          Issue Type: Bug
>          Components: Runtime / State Backends
>    Affects Versions: 2.3.0, 2.3.1
>            Reporter: Yordan Pavlov
>            Priority: Major
>              Labels: pull-request-available
>
> h3. Problem
> {{AsyncRequestBuffer}} schedules a periodic timeout check on 
> {{AsyncRequestBuffer.DELAYER}}, a static single-thread 
> {{ScheduledThreadPoolExecutor}} shared by the whole TaskManager 
> ({{scheduleAtFixedRate}} in the constructor, whenever 
> {{execution.async-state.active-buffer-timeout}} > 0, which is the default of 
> 1000 ms). The only code that cancels it is {{AsyncRequestBuffer.close()}}, 
> reached through {{AsyncExecutionController.close()}}.
> Nothing in flink-runtime calls {{AsyncExecutionController.close()}}. 
> {{AbstractAsyncKeyOrderedStreamOperator.close()}} and 
> {{AbstractAsyncStateStreamOperatorV2.close()}} only drain in-flight records. 
> So each ended task attempt (failover, restart, rescale, cancel) leaves its 
> periodic task queued on DELAYER, and that task keeps reachable:
> {noformat}
> DELAYER thread (asyncRequestBuffer-timeout-scheduler-thread-1)
>  -> DelayedWorkQueue -> ScheduledFutureTask
>  -> AsyncRequestBuffer (timeoutHandler lambda)
>  -> AsyncExecutionController -> StateExecutionController
>  -> keyed state backend (already disposed), operator, ...
> {noformat}
> until the TaskManager exits. Heap use grows with every ended attempt of every 
> async-state (State V2) operator, regardless of state backend. The queued 
> tasks also keep running their periodic check for dead attempts.
> h3. Reproduction
> MiniCluster job, parallelism 2, a {{KeyedProcessFunction}} on 
> {{enableAsyncState()}} doing {{asyncUpdate}}, fixed-delay restarts. Subtask 0 
> throws on attempts 0-4, so the job ends up on attempt 5. The program then 
> reads the size of the {{AsyncRequestBuffer.DELAYER}} queue via reflection 
> (attached: {{DelayerLeakCheck.java}}).
> ||Build||Backend||DELAYER queued tasks (2 live subtasks)||
> |2.3.0|hashmap|12|
> |2.3.0|forst|4|
> |2.3.0 + fix below|hashmap / forst|2|
> h3. Observed in production
> Flink 2.3.0, ForSt backend with async state, a TaskManager (3 slots) up ~47h 
> through repeated job restarts. Heap dump of that TM:
> * 10 {{ForStKeyedStateBackend}} instances, all {{disposed=true}}, each 
> reachable only through the DELAYER thread along the path above. The DELAYER 
> queue held exactly 10 tasks.
> * 22 {{OneInputStreamTask}} instances, all {{isRunning=false}}.
> We did not measure the retained size per pinned attempt; it depends on the 
> backend and operator. On this TM the pinned backends also kept leaked S3 
> streams from FLINK-40644 on the heap, which is how we found it.
> h3. Proposed fix
> Close the controller in {{close()}} of 
> {{AbstractAsyncKeyOrderedStreamOperator}} and 
> {{AbstractAsyncStateStreamOperatorV2}}, the two operator bases that create 
> it. Do it in a {{finally}} so it also runs when {{super.close()}} or the 
> drain throws:
> {code:java}
> @Override
> public void close() throws Exception {
>     try {
>         super.close();
>         closeIfNeeded();
>     } finally {
>         if (asyncExecutionController != null) {
>             asyncExecutionController.close();
>         }
>     }
> }
> {code}



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to