Yordan Pavlov created FLINK-40952:
-------------------------------------
Summary: ForSt: rescale restore never closes the temporary DB's
DBOptions copy, leaking the slot-shared block cache
Key: FLINK-40952
URL: https://issues.apache.org/jira/browse/FLINK-40952
Project: Flink
Issue Type: Bug
Components: Runtime / State Backends
Affects Versions: 2.3.0, 2.2.0, 2.1.0, 2.0.0
Reporter: Yordan Pavlov
Attachments: DbOptionsCopyLeak.java,
FLINK-forst-temp-dboptions-master.diff, ForStIncrementalRestoreOperation.java,
ScaleInLeakCheck.java, scale-in-fixed.log, scale-in-stock.log
h3. Problem
{{ForStIncrementalRestoreOperation#restoreTempDBInstance}} opens every
temporary restore DB with a native copy of the backend's DB options, and
nothing ever closes it:
{code:java}
// do not relocate log file for temp db
DBOptions dbOptions = new DBOptions(this.forstHandle.getDbOptions());
dbOptions.setDbLogDir("");
RocksDB restoreDb = ForStOperationUtils.openDB(..., dbOptions);
return new RestoredDBInstance(restoreDb, ...); // dbOptions is not passed on
and not closed
{code}
{{RestoredDBInstance#close()}} closes the DB, its column family handles/options
and its {{ReadOptions}}, but not this {{DBOptions}}.
With managed memory ({{state.backend.forst.memory.managed: true}}, the
default), the copied native {{DBOptions}} shares the slot's
{{WriteBufferManager}} (a {{shared_ptr}} in the C++ {{DBOptions}}), and the WBM
holds the slot-shared block cache, because memtable memory is charged to the
cache. When the backend is later disposed, all Java handles ({{RocksDB}},
{{LRUCache}}, {{WriteBufferManager}}, {{DBOptions}}) are closed, but the leaked
copy keeps the native WBM, and through it the whole block cache with its cached
blocks, alive until the TaskManager JVM exits.
Temporary DBs are used for multi-handle restores (every scale-in and some
scale-outs), so each such restore pins the block cache of the new backend. That
memory leaks at the next restart on the same TaskManager: rescale, failover, or
an in-place restart by the adaptive scheduler or autoscaler. Each leak is up to
one slot's block cache capacity. A few rescales or restart attempts on a
long-lived TM can OOM-kill the container.
* Not affected: unmanaged ForSt memory (no WBM, the copy pins only a few KB),
same-parallelism failovers (single-handle restore, no temp DB), and the RocksDB
backend ({{RocksDBIncrementalRestoreOperation}} does not copy the options).
h3. Evidence
*Production* (ForSt + async state V2, Flink 2.3.0, one TM with 3 slots and ~0.5
GB block cache per slot, job restarted in place via {{PUT
/jobs/<id>/resource-requirements}}). "Unexplained native" = native anon RSS
outside the heap, minus the live ForSt {{block-cache-usage}} and memtables:
||Restart||Backends disposed||Their cache when disposed||Unexplained
native||Step||
|baseline p=1|-|-|3,130 MiB| |
|p=1 → 3|1 (restored by a scale-in)|full, ~420 MiB|3,576 MiB|*+446*|
|p=3 → 1|3|young, ~100–200 MiB|3,604 MiB|+28|
|p=1 → 3|1 (restored by a scale-in)|full, ~413 MiB|4,040 MiB|*+436*|
|p=3 → 1|3|young, ~100–200 MiB|4,087 MiB|+47|
* jemalloc stats on the live TM: the retained memory is live allocations, not
pages kept by the allocator (allocated 2.76 GiB, resident 3.16 GiB, dirty 34
MiB).
* A heap dump shows that every Java-side object of the disposed backends is
closed ({{owningHandle_ = false}}), so the memory is held from the native side.
* With {{state.backend.forst.memory.managed: false}} (private caches, no WBM),
the same restart freed all of the disposed backends' ForSt memory (+1 MiB).
*JNI-only reproduction* (forstjni 0.1.8, jemalloc with
{{dirty_decay_ms:0,muzzy_decay_ms:0}}): a shared {{LRUCache}} +
{{WriteBufferManager(cache)}} + {{DBOptions}} with the WBM, cache filled to
~308 MiB, then everything closed. Retained RSS: 94 MiB with no copy, *422 MiB
with an unclosed {{new DBOptions(options)}} copy*, 90 MiB when the copy is
closed too.
*MiniCluster reproduction* (attached {{ScaleInLeakCheck.java}}): a ForSt
async-state job with 400k keys, a 2-slot TM with 1 GB managed memory, restored
alternately at p=2 and p=1 from native savepoints in one JVM. RSS after each
run's MiniCluster has shut down, stock flink-dist 2.3.0 vs. the same with the
fix:
||Run||Stock RSS||Step||Fixed RSS||Step||
|base|2,150 MiB||2,150 MiB||
|0: p=2, fresh|2,307 MiB|+157|2,308 MiB|+158|
|1: p=1 (scale-in)|2,713 MiB|*+406*|2,354 MiB|+46|
|2: p=2|2,759 MiB|+46|2,403 MiB|+49|
|3: p=1 (scale-in)|3,182 MiB|*+423*|2,441 MiB|+38|
|4: p=2|3,214 MiB|+32|2,468 MiB|+27|
|5: p=1 (scale-in)|3,632 MiB|*+418*|2,504 MiB|+36|
|total|| *+1,482 MiB* || +354 MiB|
Stock leaks about one slot's block cache (~410 MiB) on every scale-in. With the
fix, every run adds 27–49 MiB regardless of direction, with no scale-in pattern
(JVM and allocator warm-up). Run with {{-Xms2g -Xmx2g -XX:+AlwaysPreTouch}} and
{{LD_PRELOAD=libjemalloc.so.2 MALLOC_CONF=dirty_decay_ms:0,muzzy_decay_ms:0}},
so that RSS reflects live native allocations.
h3. Proposed fix
Pass the copy to {{RestoredDBInstance}}, close it in {{close()}} right after
the DB, and close it in {{restoreTempDBInstance}} if {{openDB}} fails. The copy
is created per temporary DB and used only to open it, so closing it drops only
this copy's references and doesn't affect the main DB's options, WBM or cache.
PR to follow.
We are deploying this patch to production on top of 2.3.0.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)