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]
