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]