seungjoo-choi-bucketplace commented on PR #29169:
URL: https://github.com/apache/flink/pull/29169#issuecomment-5987091170

   Hi @Zakelly, @fredia and @Myasuka, could you take a look at this when you 
have a chance?
   
   It fixes a heap leak in the ForSt file cache: `FileCacheEntry#open` 
registers every `CachedDataInputStream` in `FileCacheEntry#openedStreams`, but 
the only place that removes them is `doRemoveFile()`, i.e. when the file is 
evicted from or deleted in the cache. While a file stays cached, every 
open/close cycle leaves a closed stream — together with the remote stream and 
its statistics objects it wraps — in the queue. In production (Flink 2.1.3, 
ForSt on S3A, `files.open=64`) this accumulated 83,403 `S3AInputStream` 
instances (~650 MB) in 24 minutes and the TaskManagers died of heap exhaustion. 
The fix unregisters the stream as the first step of `close()`, so a failing 
close cannot skip the removal.
   
   It's small and self-contained (+5 lines in `CachedDataInputStream`), has 
unit test coverage (`FileCacheEntryTest`), CI is green on `4eb8efb`, and it 
already carries approvals from @Jackeyzhe and @YordanPavlov. Given your work on 
ForSt's file cache and stream lifecycle, I'd really appreciate a committer 
review so this can land. Thanks!
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to