seungjoo-choi-bucketplace commented on code in PR #29169:
URL: https://github.com/apache/flink/pull/29169#discussion_r4201288630
##########
flink-state-backends/flink-statebackend-forst/src/main/java/org/apache/flink/state/forst/fs/cache/CachedDataInputStream.java:
##########
@@ -305,6 +305,11 @@ public void close() throws IOException {
if (closed) {
return;
}
+ // Unregister from the cache entry right away.
FileCacheEntry#doRemoveFile only drops
+ // closed streams when the file itself is evicted or deleted, so while
a file stays cached
+ // every open/close cycle of a stream on it would otherwise leave a
closed stream (and the
+ // remote stream it wraps) on the heap. Done first so that a failing
close cannot skip it.
+ cacheEntry.openedStreams.remove(this);
Review Comment:
Done in 552fda25 — `closed = true` now comes first and the removal goes
through a new `FileCacheEntry#unregisterStream()`, so no more direct field
access from the stream.
One decision worth flagging: I left the new method unsynchronized (noted in
a comment on it). `LinkedBlockingQueue#remove` is safe against the concurrent
iteration in `doRemoveFile()`, whereas taking the entry lock there would
serialize every stream close against cache eviction.
Kept it as a separate commit so the delta is easy to review — happy to
squash it into the first commit before merge if you prefer.
--
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]