Copilot commented on code in PR #13528:
URL: https://github.com/apache/gravitino/pull/13528#discussion_r4115263760
##########
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:
The new test constructs a `LancePartitionStatisticStorage`, which owns a
`RootAllocator`, but never closes the storage. This leaks the allocator on
every run (and also leaks it if an assertion fails); wrap the test body in
try-with-resources or a finally block that closes `storage`.
##########
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:
Deferring `Dataset.close()` is not sufficient during storage shutdown:
`LancePartitionStatisticStorage.close()` calls `closeDatasetHolder` (which now
intentionally leaves active holders open) and then closes the shared Arrow
allocator immediately. An in-flight scanner still uses that allocator, so
shutdown can invalidate or fail the scan before `readerDone()` gets a chance to
close the dataset. Coordinate allocator shutdown with the reader count, or
explicitly disallow/await shutdown while readers are active.
##########
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:
The regression is in `withDataset`: its `asMap().compute` must pin the
holder before a real `listStatistics` scan can begin. This test bypasses that
path by calling `readerBegins()` directly and then invoking `cache.put`, so it
would still pass if `withDataset` looked up the holder non-atomically or failed
to call `readerDone`; please add coverage that blocks an actual scan and
triggers a replacement/size eviction while it is in flight.
--
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]