Yordan Pavlov created FLINK-40872:
-------------------------------------

             Summary: 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


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