This is an automated email from the ASF dual-hosted git repository.
FANNG1 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new 97bb4a28c8 [#13523] fix(core): replace the Lance dataset cache with a
shared Session (#13528)
97bb4a28c8 is described below
commit 97bb4a28c8076610dd61372f4e6bfebd3e81d8dc
Author: YangJie <[email protected]>
AuthorDate: Tue Sep 29 04:17:47 2026 -0400
[#13523] fix(core): replace the Lance dataset cache with a shared Session
(#13528)
### What changes were proposed in this pull request?
Replace the per-storage Lance dataset cache with a shared `Session`.
Each list/update/drop now opens its own `Dataset` and closes it in a
`finally`, and the storage holds one `Session` that carries Lance's
index and metadata caches across those opens (`ReadOptions.setSession`).
This removes the `DatasetHolder` refcounting, the Caffeine cache, and
the cleanup thread. The `datasetCacheSize` option and its docs row go
with it.
### Why are the changes needed?
With the cache enabled, a cross-table SIZE eviction could close a
`Dataset` while another thread still held it for an in-flight scan, a
use-after-free on the native handle. Opening and closing per operation
removes the shared handle an eviction could pull out from under a
reader, and the `Session` keeps the cache sharing that the dataset cache
was originally added for.
Fix: #13523
### Does this PR introduce any user-facing change?
Removes the `gravitino.stats.partition.storageOption.datasetCacheSize`
option. It defaulted to `0` (cache disabled), so default behavior is
unchanged.
### How was this patch tested?
Removed the cache-specific tests; the end-to-end functional test now
exercises the Session-backed read, update and drop path.
---
.../storage/LancePartitionStatisticStorage.java | 182 +++-----------
.../TestLancePartitionStatisticStorage.java | 265 ---------------------
docs/manage-statistics-in-gravitino.md | 1 -
3 files changed, 27 insertions(+), 421 deletions(-)
diff --git
a/core/src/main/java/org/apache/gravitino/stats/storage/LancePartitionStatisticStorage.java
b/core/src/main/java/org/apache/gravitino/stats/storage/LancePartitionStatisticStorage.java
index 98a4d91559..2dd4788d26 100644
---
a/core/src/main/java/org/apache/gravitino/stats/storage/LancePartitionStatisticStorage.java
+++
b/core/src/main/java/org/apache/gravitino/stats/storage/LancePartitionStatisticStorage.java
@@ -19,17 +19,9 @@
package org.apache.gravitino.stats.storage;
import com.fasterxml.jackson.core.JsonProcessingException;
-import com.github.benmanes.caffeine.cache.Cache;
-import com.github.benmanes.caffeine.cache.Caffeine;
-import com.github.benmanes.caffeine.cache.RemovalCause;
-import com.github.benmanes.caffeine.cache.RemovalListener;
-import com.github.benmanes.caffeine.cache.Scheduler;
-import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Preconditions;
import com.google.common.collect.Lists;
import com.google.common.collect.Maps;
-import com.google.common.util.concurrent.ThreadFactoryBuilder;
-import java.io.Closeable;
import java.io.File;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
@@ -38,10 +30,6 @@ import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.Map;
-import java.util.Optional;
-import java.util.concurrent.ScheduledThreadPoolExecutor;
-import java.util.concurrent.ThreadFactory;
-import java.util.concurrent.atomic.AtomicBoolean;
import java.util.stream.Collectors;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;
@@ -73,13 +61,12 @@ import org.lance.Dataset;
import org.lance.Fragment;
import org.lance.FragmentMetadata;
import org.lance.ReadOptions;
+import org.lance.Session;
import org.lance.SourcedTransaction;
import org.lance.WriteParams;
import org.lance.ipc.LanceScanner;
import org.lance.ipc.ScanOptions;
import org.lance.operation.Append;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
/** LancePartitionStatisticStorage is based on Lance format files. */
public class LancePartitionStatisticStorage implements
PartitionStatisticStorage {
@@ -95,8 +82,6 @@ public class LancePartitionStatisticStorage implements
PartitionStatisticStorage
private static final int DEFAULT_MAX_ROWS_PER_GROUP = 1000000; // 1M
private static final String READ_BATCH_SIZE = "readBatchSize";
private static final int DEFAULT_READ_BATCH_SIZE = 10000; // 10K
- private static final String DATASET_CACHE_SIZE = "datasetCacheSize";
- private static final int DEFAULT_DATASET_CACHE_SIZE = 0;
private static final String METADATA_FILE_CACHE_SIZE =
"metadataFileCacheSizeBytes";
private static final long DEFAULT_METADATA_FILE_CACHE_SIZE = 100L * 1024; //
100KB
private static final String INDEX_CACHE_SIZE = "indexCacheSizeBytes";
@@ -110,7 +95,7 @@ public class LancePartitionStatisticStorage implements
PartitionStatisticStorage
private static final String STATISTIC_VALUE_COLUMN = "statistic_value";
private static final String AUDIT_INFO_COLUMN = "audit_info";
- private final Optional<Cache<Long, DatasetHolder>> datasetCache;
+ private final Session session;
private static final Schema SCHEMA =
new Schema(
@@ -131,12 +116,9 @@ public class LancePartitionStatisticStorage implements
PartitionStatisticStorage
private final long metadataFileCacheSize;
private final long indexCacheSize;
private final int maxStatisticsPerUpdate;
- private final ScheduledThreadPoolExecutor scheduler;
private final EntityStore entityStore =
GravitinoEnv.getInstance().entityStore();
- private static final Logger LOG =
LoggerFactory.getLogger(LancePartitionStatisticStorage.class);
-
public LancePartitionStatisticStorage(Map<String, String> properties) {
this.allocator = new RootAllocator();
this.location = properties.getOrDefault(LOCATION, DEFAULT_LOCATION);
@@ -165,13 +147,6 @@ public class LancePartitionStatisticStorage implements
PartitionStatisticStorage
properties.getOrDefault(READ_BATCH_SIZE,
String.valueOf(DEFAULT_READ_BATCH_SIZE)));
Preconditions.checkArgument(
readBatchSize > 0, "Lance partition statistics storage readBatchSize
must be positive");
- int datasetCacheSize =
- Integer.parseInt(
- properties.getOrDefault(
- DATASET_CACHE_SIZE,
String.valueOf(DEFAULT_DATASET_CACHE_SIZE)));
- Preconditions.checkArgument(
- datasetCacheSize >= 0,
- "Lance partition statistics storage datasetCacheSize must be greater
than or equal to 0");
this.metadataFileCacheSize =
Long.parseLong(
properties.getOrDefault(
@@ -195,32 +170,15 @@ public class LancePartitionStatisticStorage implements
PartitionStatisticStorage
"Lance partition statistics storage maxStatisticsPerUpdate must be
positive");
this.properties = properties;
- if (datasetCacheSize != 0) {
- this.scheduler =
- new ScheduledThreadPoolExecutor(
- 1,
newDaemonThreadFactory("lance-partition-statistic-storage-cache-cleaner"));
-
- this.datasetCache =
- Optional.of(
- Caffeine.newBuilder()
- .maximumSize(datasetCacheSize)
-
.scheduler(Scheduler.forScheduledExecutorService(this.scheduler))
- .removalListener(
- (RemovalListener<Long, DatasetHolder>)
- (key, value, cause) -> {
- LOG.debug(
- "Removed Lance dataset cache entry,
tableId={}, cause={}",
- key,
- cause);
- if (value != null && cause !=
RemovalCause.EXPLICIT) {
- closeDatasetHolder(value);
- }
- })
- .build());
- } else {
- this.datasetCache = Optional.empty();
- this.scheduler = null;
- }
+ // Share Lance's index and metadata caches across per-operation datasets
through a single
+ // storage-scoped session. Each operation opens and closes its own
short-lived dataset, so
+ // there is no cached dataset that a concurrent eviction could close
underneath an in-flight
+ // scan.
+ this.session =
+ Session.builder()
+ .indexCacheSizeBytes(indexCacheSize)
+ .metadataCacheSizeBytes(metadataFileCacheSize)
+ .build();
}
@Override
@@ -300,7 +258,7 @@ public class LancePartitionStatisticStorage implements
PartitionStatisticStorage
Dataset datasetRead = null;
Dataset newDataset = null;
try {
- datasetRead = getDataset(tableId);
+ datasetRead = open(getFilePath(tableId));
List<FragmentMetadata> fragmentMetas = createFragmentMetadata(tableId,
updates);
SourcedTransaction appendTxn =
@@ -310,23 +268,18 @@ public class LancePartitionStatisticStorage implements
PartitionStatisticStorage
.transactionProperties(Collections.emptyMap())
.build();
newDataset = appendTxn.commit();
-
- Dataset finalNewDataset = newDataset;
- datasetCache.ifPresent(cache -> cache.put(tableId, new
DatasetHolder(finalNewDataset)));
} finally {
- if (!datasetCache.isPresent()) {
- if (datasetRead != null) {
- datasetRead.close();
- }
- if (newDataset != null) {
- newDataset.close();
- }
+ if (datasetRead != null) {
+ datasetRead.close();
+ }
+ if (newDataset != null) {
+ newDataset.close();
}
}
}
private void dropStatisticsImpl(Long tableId, List<PartitionStatisticsDrop>
drops) {
- Dataset dataset = getDataset(tableId);
+ Dataset dataset = open(getFilePath(tableId));
try {
List<String> partitionSQLs = Lists.newArrayList();
for (PartitionStatisticsDrop drop : drops) {
@@ -352,7 +305,7 @@ public class LancePartitionStatisticStorage implements
PartitionStatisticStorage
dataset.delete(filterSQL);
}
} finally {
- if (!datasetCache.isPresent() && dataset != null) {
+ if (dataset != null) {
dataset.close();
}
}
@@ -360,25 +313,12 @@ public class LancePartitionStatisticStorage implements
PartitionStatisticStorage
@Override
public void close() throws IOException {
- if (datasetCache.isPresent()) {
- Cache<Long, DatasetHolder> cache = datasetCache.get();
-
cache.asMap().values().forEach(LancePartitionStatisticStorage::closeDatasetHolder);
- cache.invalidateAll();
- cache.cleanUp();
+ if (session != null) {
+ session.close();
}
-
if (allocator != null) {
allocator.close();
}
-
- if (scheduler != null) {
- scheduler.shutdown();
- }
- }
-
- @VisibleForTesting
- Cache<Long, DatasetHolder> getDatasetCache() {
- return datasetCache.orElse(null);
}
private String getFilePath(Long tableId) {
@@ -496,9 +436,11 @@ public class LancePartitionStatisticStorage implements
PartitionStatisticStorage
private List<PersistedPartitionStatistics> listStatisticsImpl(
Long tableId, String partitionFilter) {
+ return listStatisticsWith(open(getFilePath(tableId)), tableId,
partitionFilter);
+ }
- Dataset dataset = getDataset(tableId);
-
+ private List<PersistedPartitionStatistics> listStatisticsWith(
+ Dataset dataset, Long tableId, String partitionFilter) {
String filter = "table_id = " + tableId + partitionFilter;
try (LanceScanner scanner =
@@ -553,45 +495,18 @@ public class LancePartitionStatisticStorage implements
PartitionStatisticStorage
} catch (Exception e) {
throw new RuntimeException(e);
} finally {
- if (!datasetCache.isPresent() && dataset != null) {
+ if (dataset != null) {
dataset.close();
}
}
}
- private Dataset getDataset(Long tableId) {
- AtomicBoolean newlyCreated = new AtomicBoolean(false);
- return datasetCache
- .map(
- cache -> {
- DatasetHolder holder =
- cache.get(
- tableId,
- id -> {
- newlyCreated.set(true);
- return new DatasetHolder(open(getFilePath(id)));
- });
-
- // Ensure dataset uses the latest version
- if (!newlyCreated.get()) {
- holder.checkoutLatest();
- }
-
- return holder.getDataset();
- })
- .orElse(open(getFilePath(tableId)));
- }
-
private Dataset open(String fileName) {
try {
return Dataset.open()
.allocator(allocator)
.uri(fileName)
- .readOptions(
- new ReadOptions.Builder()
- .setMetadataCacheSizeBytes(metadataFileCacheSize)
- .setIndexCacheSizeBytes(indexCacheSize)
- .build())
+ .readOptions(new ReadOptions.Builder().setSession(session).build())
.build();
} catch (IllegalArgumentException illegalArgumentException) {
if (illegalArgumentException.getMessage().contains("was not found")) {
@@ -606,47 +521,4 @@ public class LancePartitionStatisticStorage implements
PartitionStatisticStorage
}
}
}
-
- private ThreadFactory newDaemonThreadFactory(String name) {
- return new ThreadFactoryBuilder().setDaemon(true).setNameFormat(name +
"-%d").build();
- }
-
- private static void closeDatasetHolder(DatasetHolder holder) {
- try {
- holder.close();
- } catch (IOException | RuntimeException e) {
- LOG.warn("Failed to close cached Lance dataset", e);
- }
- }
-
- /**
- * Package-private wrapper around a {@link Dataset} stored in the dataset
cache. Exists solely to
- * allow test code to mock this holder (and thus verify close-ordering)
without requiring Mockito
- * to instrument the JNI-heavy {@link Dataset} class itself.
- */
- static class DatasetHolder implements Closeable {
-
- private final Dataset dataset;
-
- private final AtomicBoolean closed = new AtomicBoolean(false);
-
- DatasetHolder(Dataset dataset) {
- this.dataset = dataset;
- }
-
- Dataset getDataset() {
- return dataset;
- }
-
- void checkoutLatest() {
- dataset.checkoutLatest();
- }
-
- @Override
- public void close() throws IOException {
- if (closed.compareAndSet(false, true)) {
- dataset.close();
- }
- }
- }
}
diff --git
a/core/src/test/java/org/apache/gravitino/stats/storage/TestLancePartitionStatisticStorage.java
b/core/src/test/java/org/apache/gravitino/stats/storage/TestLancePartitionStatisticStorage.java
index 79d3525700..0203c77799 100644
---
a/core/src/test/java/org/apache/gravitino/stats/storage/TestLancePartitionStatisticStorage.java
+++
b/core/src/test/java/org/apache/gravitino/stats/storage/TestLancePartitionStatisticStorage.java
@@ -19,24 +19,15 @@
package org.apache.gravitino.stats.storage;
import static org.mockito.ArgumentMatchers.any;
-import static org.mockito.Mockito.doAnswer;
-import static org.mockito.Mockito.inOrder;
import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.spy;
-import static org.mockito.Mockito.timeout;
-import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
-import com.github.benmanes.caffeine.cache.Cache;
import com.google.common.collect.Lists;
import com.google.common.collect.Maps;
import java.io.File;
import java.nio.file.Files;
import java.util.List;
import java.util.Map;
-import org.apache.arrow.memory.BufferAllocator;
-import org.apache.arrow.memory.RootAllocator;
-import org.apache.arrow.vector.VarCharVector;
import org.apache.commons.io.FileUtils;
import org.apache.commons.lang3.reflect.FieldUtils;
import org.apache.gravitino.EntityStore;
@@ -49,10 +40,8 @@ import
org.apache.gravitino.stats.PartitionStatisticsModification;
import org.apache.gravitino.stats.PartitionStatisticsUpdate;
import org.apache.gravitino.stats.StatisticValue;
import org.apache.gravitino.stats.StatisticValues;
-import
org.apache.gravitino.stats.storage.LancePartitionStatisticStorage.DatasetHolder;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
-import org.mockito.InOrder;
public class TestLancePartitionStatisticStorage {
@@ -219,176 +208,6 @@ public class TestLancePartitionStatisticStorage {
storage.close();
}
- @Test
- public void testLancePartitionStatisticStorageWithCache() throws Exception {
- PartitionStatisticStorageFactory factory = new
LancePartitionStatisticStorageFactory();
-
- // Prepare table entity
- String metalakeName = "metalake";
- String catalogName = "catalog";
- String schemaName = "schema";
- String tableName = "table";
-
- MetadataObject metadataObject =
- MetadataObjects.of(
- Lists.newArrayList(catalogName, schemaName, tableName),
MetadataObject.Type.TABLE);
-
- EntityStore entityStore = mock(EntityStore.class);
- TableEntity tableEntity = mock(TableEntity.class);
- when(entityStore.get(any(), any(), any())).thenReturn(tableEntity);
- when(tableEntity.id()).thenReturn(1L);
- FieldUtils.writeField(GravitinoEnv.getInstance(), "entityStore",
entityStore, true);
-
- String location = Files.createTempDirectory("lance_stats_test").toString();
- Map<String, String> properties = Maps.newHashMap();
- properties.put("location", location);
- properties.put("datasetCacheSize", "1000");
-
- LancePartitionStatisticStorage storage =
- (LancePartitionStatisticStorage) factory.create(properties);
-
- int count = 100;
- int partitions = 10;
- Map<MetadataObject, Map<String, Map<String, StatisticValue<?>>>>
originData =
- generateData(metadataObject, count, partitions);
- Map<MetadataObject, List<PartitionStatisticsUpdate>> statisticsToUpdate =
- convertData(originData);
-
- List<MetadataObjectStatisticsUpdate> objectUpdates = Lists.newArrayList();
- for (Map.Entry<MetadataObject, List<PartitionStatisticsUpdate>> entry :
- statisticsToUpdate.entrySet()) {
- MetadataObject metadata = entry.getKey();
- List<PartitionStatisticsUpdate> updates = entry.getValue();
- objectUpdates.add(MetadataObjectStatisticsUpdate.of(metadata, updates));
- }
- storage.updateStatistics(metalakeName, objectUpdates);
- Assertions.assertEquals(1, storage.getDatasetCache().estimatedSize());
-
- String fromPartitionName =
- "partition" + String.format("%0" + String.valueOf(partitions).length()
+ "d", 0);
- String toPartitionName =
- "partition" + String.format("%0" + String.valueOf(partitions).length()
+ "d", 1);
-
- List<PersistedPartitionStatistics> listedStats =
- storage.listStatistics(
- metalakeName,
- metadataObject,
- PartitionRange.between(
- fromPartitionName,
- PartitionRange.BoundType.CLOSED,
- toPartitionName,
- PartitionRange.BoundType.OPEN));
- Assertions.assertEquals(1, listedStats.size());
- Assertions.assertEquals(1, storage.getDatasetCache().estimatedSize());
-
- String targetPartitionName = "partition00";
- for (PersistedPartitionStatistics persistStat : listedStats) {
- String partitionName = persistStat.partitionName();
- List<PersistedStatistic> stats = persistStat.statistics();
- Assertions.assertEquals(targetPartitionName, partitionName);
- Assertions.assertEquals(10, stats.size());
-
- for (PersistedStatistic statistic : stats) {
- String statisticName = statistic.name();
- StatisticValue<?> statisticValue = statistic.value();
-
- Assertions.assertTrue(
-
originData.get(metadataObject).get(targetPartitionName).containsKey(statisticName));
- Assertions.assertEquals(
-
originData.get(metadataObject).get(targetPartitionName).get(statisticName).value(),
- statisticValue.value());
- Assertions.assertNotNull(statistic.auditInfo());
- }
- }
-
- // Drop one statistic from partition00
- List<MetadataObjectStatisticsDrop> tableStatisticsToDrop =
- Lists.newArrayList(
- MetadataObjectStatisticsDrop.of(
- metadataObject,
- Lists.newArrayList(
- PartitionStatisticsModification.drop(
- targetPartitionName,
Lists.newArrayList("statistic0")))));
-
- storage.dropStatistics(metalakeName, tableStatisticsToDrop);
- Assertions.assertEquals(1, storage.getDatasetCache().estimatedSize());
-
- listedStats =
- storage.listStatistics(
- metalakeName,
- metadataObject,
- PartitionRange.between(
- fromPartitionName,
- PartitionRange.BoundType.CLOSED,
- toPartitionName,
- PartitionRange.BoundType.OPEN));
- Assertions.assertEquals(1, listedStats.size());
- Assertions.assertEquals(1, storage.getDatasetCache().estimatedSize());
-
- for (PersistedPartitionStatistics partitionStat : listedStats) {
- String partitionName = partitionStat.partitionName();
- List<PersistedStatistic> stats = partitionStat.statistics();
- Assertions.assertEquals(targetPartitionName, partitionName);
- Assertions.assertEquals(9, stats.size());
-
- for (PersistedStatistic statistic : stats) {
- String statisticName = statistic.name();
- StatisticValue<?> statisticValue = statistic.value();
-
- Assertions.assertTrue(
-
originData.get(metadataObject).get(targetPartitionName).containsKey(statisticName));
- Assertions.assertEquals(
-
originData.get(metadataObject).get(targetPartitionName).get(statisticName).value(),
- statisticValue.value());
- Assertions.assertNotNull(statistic.auditInfo());
- }
-
- // Drop one statistics from partition01 and partition02
- tableStatisticsToDrop =
- Lists.newArrayList(
- MetadataObjectStatisticsDrop.of(
- metadataObject,
- Lists.newArrayList(
- PartitionStatisticsModification.drop(
- "partition01", Lists.newArrayList("statistic1")),
- PartitionStatisticsModification.drop(
- "partition02", Lists.newArrayList("statistic2")))));
- storage.dropStatistics(metalakeName, tableStatisticsToDrop);
- Assertions.assertEquals(1, storage.getDatasetCache().estimatedSize());
-
- listedStats =
- storage.listStatistics(
- metalakeName,
- metadataObject,
- PartitionRange.between(
- fromPartitionName,
- PartitionRange.BoundType.CLOSED,
- "partition03",
- PartitionRange.BoundType.OPEN));
- Assertions.assertEquals(3, listedStats.size());
- Assertions.assertEquals(1, storage.getDatasetCache().estimatedSize());
- for (PersistedPartitionStatistics persistPartStat : listedStats) {
- stats = persistPartStat.statistics();
- Assertions.assertEquals(9, stats.size());
- for (PersistedStatistic statistic : stats) {
- partitionName = persistPartStat.partitionName();
- String statisticName = statistic.name();
- StatisticValue<?> statisticValue = statistic.value();
-
- Assertions.assertTrue(
-
originData.get(metadataObject).get(partitionName).containsKey(statisticName));
- Assertions.assertEquals(
-
originData.get(metadataObject).get(partitionName).get(statisticName).value(),
- statisticValue.value());
- Assertions.assertNotNull(statistic.auditInfo());
- }
- }
- }
-
- FileUtils.deleteDirectory(new File(location + "/" + tableEntity.id() +
".lance"));
- storage.close();
- }
-
@Test
public void testExceedMaxStatisticsPerUpdateLimit() throws Exception {
PartitionStatisticStorageFactory factory = new
LancePartitionStatisticStorageFactory();
@@ -605,88 +424,4 @@ public class TestLancePartitionStatisticStorage {
storage.close();
}
}
-
- @Test
- public void testCloseReleasesCachedDatasetBeforeAllocator() throws Exception
{
- String location =
Files.createTempDirectory("lance_stats_close_cache").toString();
- Map<String, String> properties = Maps.newHashMap();
- properties.put("location", location);
- properties.put("datasetCacheSize", "10");
-
- EntityStore entityStore = org.mockito.Mockito.mock(EntityStore.class);
- TableEntity tableEntity = org.mockito.Mockito.mock(TableEntity.class);
- when(entityStore.get(any(), any(), any())).thenReturn(tableEntity);
- when(tableEntity.id()).thenReturn(1L);
- FieldUtils.writeField(GravitinoEnv.getInstance(), "entityStore",
entityStore, true);
-
- LancePartitionStatisticStorage storage = new
LancePartitionStatisticStorage(properties);
-
- try {
- BufferAllocator allocator = spy(new RootAllocator(Long.MAX_VALUE));
- FieldUtils.writeField(storage, "allocator", allocator, true);
-
- Cache<Long, DatasetHolder> datasetCache = storage.getDatasetCache();
- Assertions.assertNotNull(datasetCache);
-
- DatasetHolder holder = mock(DatasetHolder.class);
- VarCharVector buffer = new VarCharVector("test", allocator);
- buffer.allocateNew(1024);
-
- doAnswer(
- invocation -> {
- buffer.close();
- return null;
- })
- .when(holder)
- .close();
-
- datasetCache.put(1L, holder);
-
- storage.close();
-
- Assertions.assertEquals(0, allocator.getAllocatedMemory());
-
- InOrder inOrder = inOrder(holder, allocator);
- inOrder.verify(holder).close();
- inOrder.verify(allocator).close();
-
- } finally {
- FileUtils.deleteDirectory(new File(location));
- }
- }
-
- @Test
- public void testDatasetCacheClosesPreviousHolderOnReplacement() throws
Exception {
- String location =
Files.createTempDirectory("lance_stats_replace_cache").toString();
- Map<String, String> properties = Maps.newHashMap();
- properties.put("location", location);
- properties.put("datasetCacheSize", "10");
-
- EntityStore entityStore = mock(EntityStore.class);
- TableEntity tableEntity = mock(TableEntity.class);
- when(entityStore.get(any(), any(), any())).thenReturn(tableEntity);
- when(tableEntity.id()).thenReturn(1L);
- FieldUtils.writeField(GravitinoEnv.getInstance(), "entityStore",
entityStore, true);
-
- LancePartitionStatisticStorage storage = new
LancePartitionStatisticStorage(properties);
-
- try {
- Cache<Long, DatasetHolder> datasetCache = storage.getDatasetCache();
- Assertions.assertNotNull(datasetCache);
-
- DatasetHolder previousHolder = mock(DatasetHolder.class);
- DatasetHolder newHolder = mock(DatasetHolder.class);
-
- datasetCache.put(1L, previousHolder);
- datasetCache.put(1L, newHolder);
-
- verify(previousHolder, timeout(5000)).close();
-
- storage.close();
-
- verify(newHolder).close();
- } finally {
- FileUtils.deleteDirectory(new File(location));
- }
- }
}
diff --git a/docs/manage-statistics-in-gravitino.md
b/docs/manage-statistics-in-gravitino.md
index 2037c43b60..05676e8fba 100644
--- a/docs/manage-statistics-in-gravitino.md
+++ b/docs/manage-statistics-in-gravitino.md
@@ -274,7 +274,6 @@ For Lance remote storage, you can refer to the document
[here](https://lancedb.g
| `gravitino.stats.partition.storageOption.maxBytesPerFile` | The
maximum bytes per file | `104857600`
| No |
| `gravitino.stats.partition.storageOption.maxRowsPerGroup` | The
maximum rows per group | `1000000`
| No |
| `gravitino.stats.partition.storageOption.readBatchSize` | The
batch record number when reading | `10000`
| No |
-| `gravitino.stats.partition.storageOption.datasetCacheSize` | size
of dataset cache for Lance | `0`, It means we don't
use the cache | No |
| `gravitino.stats.partition.storageOption.metadataFileCacheSizeBytes` | The
Lance's metadata file cache size | `102400`
| No |
| `gravitino.stats.partition.storageOption.indexCacheSizeBytes` | The
Lance's index cache size | `102400`
| No |
| `gravitino.stats.partition.storageOption.maxStatisticsPerUpdate` |
Maximum number of statistics allowed per update operation | `100`
| No |