[
https://issues.apache.org/jira/browse/FLINK-40872?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated FLINK-40872:
-----------------------------------
Labels: pull-request-available (was: )
> 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)