[ 
https://issues.apache.org/jira/browse/FLINK-40952?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18124914#comment-18124914
 ] 

Yordan Pavlov commented on FLINK-40952:
---------------------------------------

PR is up, can I please be assigned on this

> 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.0.0, 2.1.0, 2.2.0, 2.3.0
>            Reporter: Yordan Pavlov
>            Priority: Major
>              Labels: pull-request-available
>         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