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)