Zakelly commented on code in PR #29169:
URL: https://github.com/apache/flink/pull/29169#discussion_r4194102587


##########
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:
   Could we move this line immediately after `closed = true;` and encapsulate 
the removal in a new `FileCacheEntry.unregisterStream()` method to avoid direct 
field access?



-- 
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