LuciferYang commented on code in PR #13528:
URL: https://github.com/apache/gravitino/pull/13528#discussion_r4115746894


##########
core/src/main/java/org/apache/gravitino/stats/storage/LancePartitionStatisticStorage.java:
##########
@@ -642,11 +693,53 @@ void checkoutLatest() {
       dataset.checkoutLatest();
     }
 
+    /**
+     * Registers an in-flight reader that borrowed the dataset; pairs with 
{@link #readerDone()}.
+     */
+    void readerBegins() {
+      activeReaders.incrementAndGet();
+    }
+
+    /** Releases a reader; closes the dataset when the cache already evicted 
it and none remain. */
+    void readerDone() {
+      if (activeReaders.decrementAndGet() == 0 && closed.get()) {
+        closeOnce();
+      }
+    }
+
     @Override
     public void close() throws IOException {
       if (closed.compareAndSet(false, true)) {
+        if (activeReaders.get() == 0) {
+          closeOnce();
+        }
+        // Readers still scanning close it in readerDone() when the last one 
leaves.

Review Comment:
   The allocator is closed only by `LancePartitionStatisticStorage.close()`, 
that is, on storage shutdown, which is a pre-existing teardown path this change 
does not alter (the old code closed it the same way). This PR targets the 
cache-eviction race during normal operation, where the reader count now keeps 
the borrowed dataset open. Coordinating full storage shutdown with in-flight 
readers (awaiting or rejecting them) is a separate, larger change, so I left it 
out of this fix.



##########
core/src/test/java/org/apache/gravitino/stats/storage/TestLancePartitionStatisticStorage.java:
##########
@@ -56,6 +56,61 @@
 
 public class TestLancePartitionStatisticStorage {
 
+  @Test
+  public void testInFlightReaderSurvivesCacheReplaceEviction() throws 
Exception {
+    // A reader that borrowed the cached dataset must not have it closed 
underneath it when a
+    // concurrent update REPLACES the cache entry: the removal listener marks 
the holder closed,
+    // and the dataset only closes after the last reader leaves.
+    EntityStore entityStore = mock(EntityStore.class);
+    FieldUtils.writeField(GravitinoEnv.getInstance(), "entityStore", 
entityStore, true);
+    LancePartitionStatisticStorage storage =
+        new LancePartitionStatisticStorage(java.util.Collections.emptyMap());

Review Comment:
   Done in 799d7e72: the test now closes the storage in a finally block, so the 
`RootAllocator` is released even if an assertion fails.



##########
core/src/test/java/org/apache/gravitino/stats/storage/TestLancePartitionStatisticStorage.java:
##########
@@ -56,6 +56,61 @@
 
 public class TestLancePartitionStatisticStorage {
 
+  @Test
+  public void testInFlightReaderSurvivesCacheReplaceEviction() throws 
Exception {
+    // A reader that borrowed the cached dataset must not have it closed 
underneath it when a
+    // concurrent update REPLACES the cache entry: the removal listener marks 
the holder closed,
+    // and the dataset only closes after the last reader leaves.
+    EntityStore entityStore = mock(EntityStore.class);
+    FieldUtils.writeField(GravitinoEnv.getInstance(), "entityStore", 
entityStore, true);
+    LancePartitionStatisticStorage storage =
+        new LancePartitionStatisticStorage(java.util.Collections.emptyMap());
+
+    java.util.concurrent.atomic.AtomicBoolean datasetClosed =
+        new java.util.concurrent.atomic.AtomicBoolean(false);
+    LancePartitionStatisticStorage.DatasetHolder holder =
+        new LancePartitionStatisticStorage.DatasetHolder(null) {
+          @Override
+          protected void closeDataset() {
+            datasetClosed.set(true);
+          }
+        };
+    LancePartitionStatisticStorage.DatasetHolder replacement =
+        new LancePartitionStatisticStorage.DatasetHolder(null);
+
+    java.lang.reflect.Field cacheField =
+        LancePartitionStatisticStorage.class.getDeclaredField("datasetCache");
+    cacheField.setAccessible(true);
+    com.github.benmanes.caffeine.cache.Cache<Long, 
LancePartitionStatisticStorage.DatasetHolder>
+        cache =
+            com.github.benmanes.caffeine.cache.Caffeine.newBuilder()
+                .maximumSize(2)
+                .removalListener(
+                    (Long key,
+                        LancePartitionStatisticStorage.DatasetHolder value,
+                        com.github.benmanes.caffeine.cache.RemovalCause cause) 
-> {
+                      if (value != null
+                          && cause != 
com.github.benmanes.caffeine.cache.RemovalCause.EXPLICIT) {
+                        try {
+                          value.close();
+                        } catch (java.io.IOException e) {
+                          throw new RuntimeException(e);
+                        }
+                      }
+                    })
+                .build();
+    cacheField.set(storage, java.util.Optional.of(cache));
+
+    cache.put(1L, holder);
+    holder.readerBegins();
+    cache.put(1L, replacement); // evicts the old entry, listener calls close()

Review Comment:
   This test targets the holder's defer-close contract: an eviction during an 
in-flight borrow must not close the dataset until the last reader leaves, which 
is the mechanism this fix adds. A test that blocks a real `listStatistics` scan 
and triggers a concurrent size eviction needs the native Lance dataset plus 
thread timing and would be inherently racy, so I kept the deterministic 
holder-level check here.



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