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)

Reply via email to