This is an automated email from the ASF dual-hosted git repository.
yuqi1129 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 bb93bf9234 [#12416] fix(core): Cascade cache invalidation to
hierarchical schema descendants (#12417)
bb93bf9234 is described below
commit bb93bf923433096735abdca9f248a4b3adfdfbcf
Author: Qi Yu <[email protected]>
AuthorDate: Wed Aug 12 15:10:23 2026 +0800
[#12416] fix(core): Cascade cache invalidation to hierarchical schema
descendants (#12417)
### What changes were proposed in this pull request?
Make `CaffeineEntityCache.invalidateHierarchy` scan the prefix index
once per child boundary, instead of only for `"."`.
- Scan a second time for `HierarchicalSchemaUtil.schemaSeparator()`, the
boundary between nested `HierarchicalSchema` levels.
- Correct the `invalidateHierarchy` javadoc, which asserted the
incorrect invariant that every child identifier starts with `parent
identifier + "."`.
- Add cascade tests at both the cache level and the store level, for the
default and a non-default separator.
### Why are the changes needed?
`invalidateHierarchy` found cached descendants with a single prefix
scan:
```java
String childPrefix = key.identifier().toString() + ".";
```
This assumes every child identifier starts with the parent identifier
followed by `.`. A `HierarchicalSchema` breaks that assumption: its
nested levels are not extra `NameIdentifier` levels, they are joined
**inside a single name level** by the schema separator
(`Configs.SCHEMA_SEPARATOR`, default `:`, which may not be `.`). So
`raw:events:2024` is one schema name whose identifier continues past
`raw:events` with `:`, not with `.`.
These are the cache keys a real H2-backed `RelationalEntityStore`
produces for a nested schema, its table and a sibling:
```
metalake.catalog.raw:events:SCHEMA
metalake.catalog.raw:events:2024:SCHEMA
metalake.catalog.raw:events:2024.t_child:TABLE
metalake.catalog.raw:events2:SCHEMA
```
Both boundaries are real and both are needed: `.` separates
`NameIdentifier` levels (the table under the schema), the schema
separator separates nested schema levels.
The name is logical at this layer, not physical. `SchemaMetaService`
converts between the logical name and the physical one (ASCII-1) in
`HierarchicalConversionPOStorageOps`, i.e. at the PO boundary, while the
cache lives in `RelationalEntityStore` above it. The physical separator
therefore only ever appears in backend rows, never in a cache key;
`TestRelationalEntityStoreHierarchicalCache` asserts exactly that.
Dropping or renaming `raw:events` therefore left `raw:events:2024`, any
deeper nesting, and every table/view/fileset/topic below them in the
cache until TTL:
```
--- after invalidate(raw:events schema) ---
parent cached = false (expected)
tblInParent cached = false (expected)
child cached = true <-- stale
tblInChild cached = true <-- stale
```
This is not multi-node specific. `RelationalEntityStore#delete` and the
rename paths invalidate through the same method, so the node performing
the mutation keeps stale descendants too; it reproduces on a single
node.
The trailing separator itself was correct and is kept — it is the guard
that stops `catalog1` from matching `catalog10`, and `raw:events` from
matching `raw:events2`. The defect was that there is more than one valid
child boundary and only one was handled. Because the radix index matches
on whole key strings, the added pass reaches descendants at any depth,
so no recursion is needed.
The second pass runs for schema keys only. `EntityCacheKey.toString()`
joins the identifier and the entity type with `":"`, which is also the
default schema separator, so an unguarded pass over a table key would
also match the topic or fileset of the same name. Only a schema can
carry nested levels, and a catalog still reaches its nested schemas
through the `"."` pass, so restricting it costs nothing and keeps the
scan exact.
Fix: #12416
### Does this PR introduce _any_ user-facing change?
No. No public API or configuration property changes. Cascading
invalidation of hierarchical schemas now behaves as already documented.
### How was this patch tested?
`TestCaffeineEntityCacheInvalidation` — cache-level cases:
| Test | Before |
| --- | --- |
| `testInvalidateHierarchicalSchemaCascadesToNestedSchemas` | fails |
| `testInvalidateHierarchicalSchemaCascadesToAnyDepth` (four levels) |
fails |
| `testInvalidateHierarchicalSchemaCascadesWithNonDefaultSeparator`
(separator `\|`) | fails |
| `testInvalidateLeafDoesNotEvictSameNameEntityOfAnotherType` (table vs
topic) | fails without the schema-key guard |
| `testInvalidateHierarchicalSchemaDoesNotTouchSiblings` (`raw:events`
vs `raw:events2`) | passes |
| `testInvalidateCatalogCascadesToHierarchicalSchemas` | passes |
`TestRelationalEntityStoreHierarchicalCache` — store-level case, a real
H2-backed `RelationalEntityStore` with the cache enabled, run for the
default separator `:` and for `|`. It writes `raw:events`, the nested
`raw:events:2024`, a table inside the nested schema and the sibling
`raw:events2`, reads them back through the store, drops `raw:events`
with cascade, and asserts the nested schema and its table are gone from
the cache while the sibling survives. It also asserts no cache key
contains the physical separator. Both parameter sets fail before this
change.
Commands:
- `./gradlew :core:test --tests
"org.apache.gravitino.cache.TestCaffeineEntityCacheInvalidation" --tests
"org.apache.gravitino.storage.relational.TestRelationalEntityStoreHierarchicalCache"
-PskipITs -PskipDockerTests`
- `./gradlew :core:test -PskipITs -PskipDockerTests`
- `./gradlew :core:spotlessCheck :core:javadoc -PskipITs
-PskipDockerTests`
The store-level test runs against H2 by default; MySQL/PostgreSQL
coverage comes from CI.
One case from the issue is intentionally not included here: a cross-node
assertion in `TestEntityCacheCrossNodeInvalidation`, which does not
exist on `main` yet (it is introduced by #12374). Cross-node replay
funnels into the same `cache.invalidate` -> `invalidateHierarchy` entry
point that these tests cover, so the mechanism is exercised; the
cross-node case is worth adding once #12374 lands.
---
.../gravitino/cache/CaffeineEntityCache.java | 45 +++-
.../cache/TestCaffeineEntityCacheInvalidation.java | 177 ++++++++++++++++
...TestRelationalEntityStoreHierarchicalCache.java | 231 +++++++++++++++++++++
3 files changed, 448 insertions(+), 5 deletions(-)
diff --git
a/core/src/main/java/org/apache/gravitino/cache/CaffeineEntityCache.java
b/core/src/main/java/org/apache/gravitino/cache/CaffeineEntityCache.java
index bf433fcf4d..71e2592f4c 100644
--- a/core/src/main/java/org/apache/gravitino/cache/CaffeineEntityCache.java
+++ b/core/src/main/java/org/apache/gravitino/cache/CaffeineEntityCache.java
@@ -44,6 +44,7 @@ import org.apache.gravitino.Entity;
import org.apache.gravitino.HasIdentifier;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.meta.ModelVersionEntity;
+import org.apache.gravitino.utils.HierarchicalSchemaUtil;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -83,6 +84,12 @@ public class CaffeineEntityCache extends BaseEntityCache {
private static final Logger LOG =
LoggerFactory.getLogger(CaffeineEntityCache.class.getName());
+ /**
+ * Separates {@link NameIdentifier} levels in a cache key. See {@link
+ * #invalidateHierarchy(EntityCacheKey)} for why it is not the only child
boundary.
+ */
+ private static final String NAME_LEVEL_BOUNDARY = ".";
+
/** Segmented locking for better concurrency */
private final SegmentedLock segmentedLock;
@@ -267,9 +274,25 @@ public class CaffeineEntityCache extends BaseEntityCache {
/**
* Removes the entry for the given key and all cached descendant entries.
Descendants are found
- * through the prefix index: every child identifier starts with {@code
parent identifier + "."},
- * so the scan is exact for children and never matches siblings sharing a
name prefix (e.g. {@code
- * catalog1} vs {@code catalog10}).
+ * through the prefix index, scanning once per child boundary:
+ *
+ * <ul>
+ * <li>{@code "."} separates {@link NameIdentifier} levels, so it matches
ordinary children such
+ * as the tables of a schema.
+ * <li>The {@link HierarchicalSchemaUtil#schemaSeparator() schema
separator} joins nested {@code
+ * HierarchicalSchema} levels <em>inside</em> a single name level, so
it matches nested
+ * schemas such as {@code raw:events:2024} under {@code raw:events}.
Without this pass those
+ * descendants would survive until their TTL expires. The cache sits
above the storage
+ * layer, where schema names are still logical, so the boundary is the
configured external
+ * separator and not the physical one the entity store writes to the
backend. Only a schema
+ * can carry nested levels, so this pass is skipped for every other
entity type; a catalog
+ * still reaches its nested schemas through the {@code "."} pass above.
+ * </ul>
+ *
+ * <p>Matching on a boundary rather than the bare identifier is what keeps
the scan exact: it
+ * never matches siblings sharing a name prefix, neither {@code catalog1} vs
{@code catalog10} nor
+ * {@code raw:events} vs {@code raw:events2}. Because the index matches on
the whole key string,
+ * the separator pass already collects descendants at any depth, so no
recursion is needed.
*
* @param key The key of the entity whose subtree should be invalidated
*/
@@ -277,9 +300,21 @@ public class CaffeineEntityCache extends BaseEntityCache {
cacheData.invalidate(key);
cacheIndex.remove(key.toString());
- String childPrefix = key.identifier().toString() + ".";
+ String identifier = key.identifier().toString();
+ invalidateDescendants(identifier + NAME_LEVEL_BOUNDARY);
+ if (key.entityType() == Entity.EntityType.SCHEMA) {
+ invalidateDescendants(identifier +
HierarchicalSchemaUtil.schemaSeparator());
+ }
+ }
+
+ /**
+ * Removes every cached entry whose key starts with the given prefix.
+ *
+ * @param keyPrefix The prefix that identifies the descendants to remove
+ */
+ private void invalidateDescendants(String keyPrefix) {
List<EntityCacheKey> childKeys =
-
Lists.newArrayList(cacheIndex.getValuesForKeysStartingWith(childPrefix));
+ Lists.newArrayList(cacheIndex.getValuesForKeysStartingWith(keyPrefix));
for (EntityCacheKey childKey : childKeys) {
cacheData.invalidate(childKey);
cacheIndex.remove(childKey.toString());
diff --git
a/core/src/test/java/org/apache/gravitino/cache/TestCaffeineEntityCacheInvalidation.java
b/core/src/test/java/org/apache/gravitino/cache/TestCaffeineEntityCacheInvalidation.java
index 520eb42c11..d30228e5d9 100644
---
a/core/src/test/java/org/apache/gravitino/cache/TestCaffeineEntityCacheInvalidation.java
+++
b/core/src/test/java/org/apache/gravitino/cache/TestCaffeineEntityCacheInvalidation.java
@@ -21,8 +21,11 @@ package org.apache.gravitino.cache;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
+import org.apache.commons.lang3.reflect.FieldUtils;
import org.apache.gravitino.Config;
+import org.apache.gravitino.Configs;
import org.apache.gravitino.Entity;
+import org.apache.gravitino.GravitinoEnv;
import org.apache.gravitino.Namespace;
import org.apache.gravitino.meta.CatalogEntity;
import org.apache.gravitino.meta.GroupEntity;
@@ -31,7 +34,9 @@ import org.apache.gravitino.meta.ModelVersionEntity;
import org.apache.gravitino.meta.RoleEntity;
import org.apache.gravitino.meta.SchemaEntity;
import org.apache.gravitino.meta.TableEntity;
+import org.apache.gravitino.meta.TopicEntity;
import org.apache.gravitino.meta.UserEntity;
+import org.apache.gravitino.utils.HierarchicalSchemaUtil;
import org.apache.gravitino.utils.TestUtil;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
@@ -40,6 +45,8 @@ import org.junit.jupiter.api.Test;
/** Tests hierarchical invalidation and cacheability rules of {@link
CaffeineEntityCache}. */
public class TestCaffeineEntityCacheInvalidation {
+ private static final Namespace CATALOG_NS = Namespace.of("metalake",
"catalog1");
+
private CaffeineEntityCache cache;
@BeforeEach
@@ -47,6 +54,18 @@ public class TestCaffeineEntityCacheInvalidation {
cache = new CaffeineEntityCache(new Config() {});
}
+ /**
+ * Joins nested schema levels the way they reach the cache. The cache sits
above the storage
+ * layer, so nested schema names still carry the configured external
separator.
+ */
+ private static String hierarchicalName(String... levels) {
+ return String.join(HierarchicalSchemaUtil.schemaSeparator(), levels);
+ }
+
+ private static Namespace schemaNamespace(String schemaName) {
+ return Namespace.of(CATALOG_NS.level(0), CATALOG_NS.level(1), schemaName);
+ }
+
@Test
void testInvalidateCatalogCascadesToChildren() {
CatalogEntity catalog =
@@ -155,6 +174,164 @@ public class TestCaffeineEntityCacheInvalidation {
cache.getIfPresent(catalog.nameIdentifier(),
Entity.EntityType.SCHEMA).isEmpty());
}
+ @Test
+ void testInvalidateHierarchicalSchemaCascadesToNestedSchemas() {
+ // A HierarchicalSchema nests inside a single name level, joined by the
schema separator, so
+ // "raw:events:2024" is a child of "raw:events" without adding a
NameIdentifier level.
+ String parentName = hierarchicalName("raw", "events");
+ String childName = hierarchicalName("raw", "events", "2024");
+
+ SchemaEntity parent = TestUtil.getTestSchemaEntity(2L, parentName,
CATALOG_NS, "cmt");
+ SchemaEntity child = TestUtil.getTestSchemaEntity(3L, childName,
CATALOG_NS, "cmt");
+ TableEntity tableInParent =
+ TestUtil.getTestTableEntity(4L, "t_parent",
schemaNamespace(parentName));
+ TableEntity tableInChild =
+ TestUtil.getTestTableEntity(5L, "t_child", schemaNamespace(childName));
+
+ cache.put(parent);
+ cache.put(child);
+ cache.put(tableInParent);
+ cache.put(tableInChild);
+ Assertions.assertEquals(4, cache.size());
+
+ cache.invalidate(parent.nameIdentifier(), Entity.EntityType.SCHEMA);
+
+ Assertions.assertFalse(cache.contains(parent.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ Assertions.assertFalse(cache.contains(tableInParent.nameIdentifier(),
Entity.EntityType.TABLE));
+ Assertions.assertFalse(cache.contains(child.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ Assertions.assertFalse(cache.contains(tableInChild.nameIdentifier(),
Entity.EntityType.TABLE));
+ Assertions.assertEquals(0, cache.size());
+ }
+
+ @Test
+ void testInvalidateHierarchicalSchemaCascadesToAnyDepth() {
+ String level1 = hierarchicalName("raw");
+ String level2 = hierarchicalName("raw", "events");
+ String level3 = hierarchicalName("raw", "events", "2024");
+ String level4 = hierarchicalName("raw", "events", "2024", "q1");
+
+ cache.put(TestUtil.getTestSchemaEntity(2L, level1, CATALOG_NS, "cmt"));
+ cache.put(TestUtil.getTestSchemaEntity(3L, level2, CATALOG_NS, "cmt"));
+ cache.put(TestUtil.getTestSchemaEntity(4L, level3, CATALOG_NS, "cmt"));
+ SchemaEntity deepest = TestUtil.getTestSchemaEntity(5L, level4,
CATALOG_NS, "cmt");
+ cache.put(deepest);
+ TableEntity deepestTable = TestUtil.getTestTableEntity(6L, "t_deep",
schemaNamespace(level4));
+ cache.put(deepestTable);
+ Assertions.assertEquals(5, cache.size());
+
+ cache.invalidate(
+ TestUtil.getTestSchemaEntity(2L, level1, CATALOG_NS,
"cmt").nameIdentifier(),
+ Entity.EntityType.SCHEMA);
+
+ Assertions.assertFalse(cache.contains(deepest.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ Assertions.assertFalse(cache.contains(deepestTable.nameIdentifier(),
Entity.EntityType.TABLE));
+ Assertions.assertEquals(0, cache.size());
+ }
+
+ @Test
+ void testInvalidateHierarchicalSchemaDoesNotTouchSiblings() {
+ // Guards against over-matching: "raw:events2" is a sibling of
"raw:events", not a descendant,
+ // exactly like the catalog1 / catalog10 case the "." boundary already
protects against.
+ String target = hierarchicalName("raw", "events");
+ String sibling = hierarchicalName("raw", "events2");
+ String siblingOfParent = hierarchicalName("raw2", "events");
+
+ SchemaEntity targetSchema = TestUtil.getTestSchemaEntity(2L, target,
CATALOG_NS, "cmt");
+ SchemaEntity siblingSchema = TestUtil.getTestSchemaEntity(3L, sibling,
CATALOG_NS, "cmt");
+ SchemaEntity otherBranch = TestUtil.getTestSchemaEntity(4L,
siblingOfParent, CATALOG_NS, "cmt");
+ TableEntity siblingTable =
+ TestUtil.getTestTableEntity(5L, "t_sibling", schemaNamespace(sibling));
+
+ cache.put(targetSchema);
+ cache.put(siblingSchema);
+ cache.put(otherBranch);
+ cache.put(siblingTable);
+
+ cache.invalidate(targetSchema.nameIdentifier(), Entity.EntityType.SCHEMA);
+
+ Assertions.assertFalse(cache.contains(targetSchema.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ Assertions.assertTrue(cache.contains(siblingSchema.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ Assertions.assertTrue(cache.contains(otherBranch.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ Assertions.assertTrue(cache.contains(siblingTable.nameIdentifier(),
Entity.EntityType.TABLE));
+ Assertions.assertEquals(3, cache.size());
+ }
+
+ @Test
+ void testInvalidateCatalogCascadesToHierarchicalSchemas() {
+ CatalogEntity catalog =
+ TestUtil.getTestCatalogEntity(1L, "catalog1",
Namespace.of("metalake"), "hive", "cmt");
+ String nested = hierarchicalName("raw", "events", "2024");
+ SchemaEntity schema = TestUtil.getTestSchemaEntity(2L, nested, CATALOG_NS,
"cmt");
+ TableEntity table = TestUtil.getTestTableEntity(3L, "t1",
schemaNamespace(nested));
+
+ cache.put(catalog);
+ cache.put(schema);
+ cache.put(table);
+ Assertions.assertEquals(3, cache.size());
+
+ cache.invalidate(catalog.nameIdentifier(), Entity.EntityType.CATALOG);
+
+ Assertions.assertFalse(cache.contains(schema.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ Assertions.assertFalse(cache.contains(table.nameIdentifier(),
Entity.EntityType.TABLE));
+ Assertions.assertEquals(0, cache.size());
+ }
+
+ @Test
+ void testInvalidateHierarchicalSchemaCascadesWithNonDefaultSeparator()
throws Exception {
+ Config separatorConfig = new Config(false) {};
+ separatorConfig.set(Configs.SCHEMA_SEPARATOR, "|");
+ Object previousConfig = FieldUtils.readField(GravitinoEnv.getInstance(),
"config", true);
+ FieldUtils.writeField(GravitinoEnv.getInstance(), "config",
separatorConfig, true);
+
+ try {
+ Assertions.assertEquals("|", HierarchicalSchemaUtil.schemaSeparator());
+
+ String parentName = hierarchicalName("raw", "events");
+ String childName = hierarchicalName("raw", "events", "2024");
+ String siblingName = hierarchicalName("raw", "events2");
+
+ SchemaEntity parent = TestUtil.getTestSchemaEntity(2L, parentName,
CATALOG_NS, "cmt");
+ SchemaEntity child = TestUtil.getTestSchemaEntity(3L, childName,
CATALOG_NS, "cmt");
+ SchemaEntity sibling = TestUtil.getTestSchemaEntity(4L, siblingName,
CATALOG_NS, "cmt");
+ TableEntity tableInChild =
+ TestUtil.getTestTableEntity(5L, "t_child",
schemaNamespace(childName));
+
+ cache.put(parent);
+ cache.put(child);
+ cache.put(sibling);
+ cache.put(tableInChild);
+
+ cache.invalidate(parent.nameIdentifier(), Entity.EntityType.SCHEMA);
+
+ Assertions.assertFalse(cache.contains(parent.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ Assertions.assertFalse(cache.contains(child.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ Assertions.assertFalse(
+ cache.contains(tableInChild.nameIdentifier(),
Entity.EntityType.TABLE));
+ Assertions.assertTrue(cache.contains(sibling.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ Assertions.assertEquals(1, cache.size());
+ } finally {
+ FieldUtils.writeField(GravitinoEnv.getInstance(), "config",
previousConfig, true);
+ }
+ }
+
+ @Test
+ void testInvalidateLeafDoesNotEvictSameNameEntityOfAnotherType() {
+ // A cache key is "<identifier>:<type>", and ":" is also the default
schema separator. Only a
+ // schema can nest, so the schema-separator scan must not run for other
types, otherwise
+ // invalidating a table would also drop the topic of the same name.
+ Namespace schemaNs = schemaNamespace("schema1");
+ TableEntity table = TestUtil.getTestTableEntity(2L, "shared_name",
schemaNs);
+ TopicEntity topic = TestUtil.getTestTopicEntity(3L, "shared_name",
schemaNs, "cmt");
+
+ cache.put(table);
+ cache.put(topic);
+
+ cache.invalidate(table.nameIdentifier(), Entity.EntityType.TABLE);
+
+ Assertions.assertFalse(cache.contains(table.nameIdentifier(),
Entity.EntityType.TABLE));
+ Assertions.assertTrue(cache.contains(topic.nameIdentifier(),
Entity.EntityType.TOPIC));
+ }
+
@Test
void testClearResetsSizeAndIndex() {
CatalogEntity catalog =
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/TestRelationalEntityStoreHierarchicalCache.java
b/core/src/test/java/org/apache/gravitino/storage/relational/TestRelationalEntityStoreHierarchicalCache.java
new file mode 100644
index 0000000000..bb1bfbd3e9
--- /dev/null
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/TestRelationalEntityStoreHierarchicalCache.java
@@ -0,0 +1,231 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.gravitino.storage.relational;
+
+import java.io.File;
+import java.io.IOException;
+import java.time.Instant;
+import java.util.Set;
+import java.util.UUID;
+import java.util.stream.Collectors;
+import org.apache.commons.io.FileUtils;
+import org.apache.commons.lang3.reflect.FieldUtils;
+import org.apache.gravitino.Catalog;
+import org.apache.gravitino.Config;
+import org.apache.gravitino.Configs;
+import org.apache.gravitino.Entity;
+import org.apache.gravitino.GravitinoEnv;
+import org.apache.gravitino.Namespace;
+import org.apache.gravitino.cache.CaffeineEntityCache;
+import org.apache.gravitino.meta.AuditInfo;
+import org.apache.gravitino.meta.BaseMetalake;
+import org.apache.gravitino.meta.CatalogEntity;
+import org.apache.gravitino.meta.SchemaEntity;
+import org.apache.gravitino.meta.SchemaVersion;
+import org.apache.gravitino.meta.TableEntity;
+import org.apache.gravitino.storage.RandomIdGenerator;
+import org.apache.gravitino.utils.HierarchicalSchemaUtil;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+import org.mockito.Mockito;
+
+/**
+ * Verifies that dropping a hierarchical schema through a real {@link
RelationalEntityStore} also
+ * drops its nested descendants from the entity cache, for the default and a
non-default schema
+ * separator.
+ */
+public class TestRelationalEntityStoreHierarchicalCache {
+
+ private static final String METALAKE = "metalake_hs";
+ private static final String CATALOG = "catalog_hs";
+ private static final AuditInfo AUDIT_INFO =
+
AuditInfo.builder().withCreator("creator").withCreateTime(Instant.now()).build();
+
+ private RelationalEntityStore store;
+ private String dbPath;
+ private Object previousConfig;
+
+ @AfterEach
+ void tearDown() throws Exception {
+ if (store != null) {
+ store.close();
+ store = null;
+ }
+ if (dbPath != null) {
+ FileUtils.deleteQuietly(new File(dbPath));
+ dbPath = null;
+ }
+ FieldUtils.writeField(GravitinoEnv.getInstance(), "config",
previousConfig, true);
+ }
+
+ @ParameterizedTest
+ @ValueSource(strings = {":", "|"})
+ void testDropHierarchicalSchemaEvictsNestedDescendantsFromCache(String
separator)
+ throws Exception {
+ initStore(separator);
+
+ String parentName = String.join(separator, "raw", "events");
+ String childName = String.join(separator, "raw", "events", "2024");
+ String siblingName = String.join(separator, "raw", "events2");
+
+ store.put(metalake(), false);
+ store.put(catalog(), false);
+ SchemaEntity parent = schema(parentName);
+ SchemaEntity child = schema(childName);
+ SchemaEntity sibling = schema(siblingName);
+ store.put(parent, false);
+ store.put(child, false);
+ store.put(sibling, false);
+ TableEntity tableInChild = table("t_child", childName);
+ store.put(tableInChild, false);
+
+ // Read everything back so the cache is populated with the names the store
actually returns.
+ store.get(parent.nameIdentifier(), Entity.EntityType.SCHEMA,
SchemaEntity.class);
+ store.get(child.nameIdentifier(), Entity.EntityType.SCHEMA,
SchemaEntity.class);
+ store.get(sibling.nameIdentifier(), Entity.EntityType.SCHEMA,
SchemaEntity.class);
+ store.get(tableInChild.nameIdentifier(), Entity.EntityType.TABLE,
TableEntity.class);
+ Assertions.assertTrue(
+ store.getCache().contains(child.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ Assertions.assertTrue(
+ store.getCache().contains(tableInChild.nameIdentifier(),
Entity.EntityType.TABLE));
+
+ // The cache is keyed by the identifier that reaches the store, which
still carries the
+ // configured external separator; the physical separator only exists in
the backend rows.
+ Set<String> cacheKeys =
+ ((CaffeineEntityCache) store.getCache())
+ .getCacheData().asMap().keySet().stream()
+ .map(Object::toString)
+ .collect(Collectors.toSet());
+ Assertions.assertTrue(
+ cacheKeys.stream().anyMatch(key -> key.contains(childName)),
+ "nested schema must be cached under its logical name, keys: " +
cacheKeys);
+ Assertions.assertTrue(
+ cacheKeys.stream()
+ .noneMatch(key ->
key.contains(HierarchicalSchemaUtil.physicalSeparator())),
+ "no cache key may carry the physical separator, keys: " + cacheKeys);
+
+ store.delete(parent.nameIdentifier(), Entity.EntityType.SCHEMA, true);
+
+ Assertions.assertFalse(
+ store.getCache().contains(parent.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ Assertions.assertFalse(
+ store.getCache().contains(child.nameIdentifier(),
Entity.EntityType.SCHEMA),
+ "nested schema must not survive the drop of its parent");
+ Assertions.assertFalse(
+ store.getCache().contains(tableInChild.nameIdentifier(),
Entity.EntityType.TABLE),
+ "table of a nested schema must not survive the drop of the parent
schema");
+ Assertions.assertTrue(
+ store.getCache().contains(sibling.nameIdentifier(),
Entity.EntityType.SCHEMA),
+ "a sibling sharing a name prefix must not be invalidated");
+ }
+
+ private void initStore(String separator) throws Exception {
+ dbPath = "/tmp/gravitino_hs_cache_test_" +
UUID.randomUUID().toString().replace("-", "");
+ File dir = new File(dbPath);
+ if (!dir.exists() && !dir.mkdirs()) {
+ throw new IOException("Failed to create test directory " + dbPath);
+ }
+
+ Config config = Mockito.mock(Config.class);
+
Mockito.when(config.get(Configs.ENTITY_STORE)).thenReturn(Configs.RELATIONAL_ENTITY_STORE);
+ Mockito.when(config.get(Configs.ENTITY_RELATIONAL_STORE))
+ .thenReturn(Configs.DEFAULT_ENTITY_RELATIONAL_STORE);
+ Mockito.when(config.get(Configs.ENTITY_RELATIONAL_JDBC_BACKEND_URL))
+
.thenReturn(String.format("jdbc:h2:file:%s;DB_CLOSE_DELAY=-1;MODE=MYSQL",
dbPath));
+
Mockito.when(config.get(Configs.ENTITY_RELATIONAL_JDBC_BACKEND_USER)).thenReturn("root");
+
Mockito.when(config.get(Configs.ENTITY_RELATIONAL_JDBC_BACKEND_PASSWORD)).thenReturn("123456");
+ Mockito.when(config.get(Configs.ENTITY_RELATIONAL_JDBC_BACKEND_DRIVER))
+ .thenReturn("org.h2.Driver");
+
Mockito.when(config.get(Configs.ENTITY_RELATIONAL_JDBC_BACKEND_MAX_CONNECTIONS))
+ .thenReturn(Configs.DEFAULT_RELATIONAL_JDBC_BACKEND_MAX_CONNECTIONS);
+
Mockito.when(config.get(Configs.ENTITY_RELATIONAL_JDBC_BACKEND_WAIT_MILLISECONDS))
+
.thenReturn(Configs.DEFAULT_RELATIONAL_JDBC_BACKEND_MAX_WAIT_MILLISECONDS);
+ Mockito.when(config.get(Configs.STORE_DELETE_AFTER_TIME)).thenReturn(20 *
60 * 1000L);
+ Mockito.when(config.get(Configs.VERSION_RETENTION_COUNT)).thenReturn(1L);
+
Mockito.when(config.get(Configs.ENTITY_CHANGE_LOG_POLL_INTERVAL_SECS)).thenReturn(3L);
+
Mockito.when(config.get(Configs.ENTITY_CHANGE_LOG_LISTENER_MAX_RETRIES)).thenReturn(10);
+
Mockito.when(config.get(Configs.ENTITY_CHANGE_LOG_LISTENER_FAILURE_ACTION)).thenReturn("SKIP");
+
Mockito.when(config.get(Configs.ENTITY_CHANGE_LOG_RETENTION_SECS)).thenReturn(24
* 60 * 60L);
+
Mockito.when(config.get(Configs.ENTITY_CHANGE_LOG_CLEANUP_INTERVAL_SECS)).thenReturn(60
* 60L);
+ Mockito.when(config.get(Configs.CACHE_ENABLED)).thenReturn(true);
+
Mockito.when(config.get(Configs.CACHE_IMPLEMENTATION)).thenReturn("caffeine");
+ Mockito.when(config.get(Configs.CACHE_MAX_ENTRIES)).thenReturn(10_000);
+
Mockito.when(config.get(Configs.CACHE_EXPIRATION_TIME)).thenReturn(3_600_000L);
+ Mockito.when(config.get(Configs.CACHE_WEIGHER_ENABLED)).thenReturn(true);
+ Mockito.when(config.get(Configs.CACHE_STATS_ENABLED)).thenReturn(false);
+ Mockito.when(config.get(Configs.CACHE_LOCK_SEGMENTS)).thenReturn(16);
+ Mockito.when(config.get(Configs.SCHEMA_SEPARATOR)).thenReturn(separator);
+
+ previousConfig = FieldUtils.readField(GravitinoEnv.getInstance(),
"config", true);
+ FieldUtils.writeField(GravitinoEnv.getInstance(), "config", config, true);
+ FieldUtils.writeField(
+ GravitinoEnv.getInstance(), "idGenerator", RandomIdGenerator.INSTANCE,
true);
+ Assertions.assertEquals(separator,
HierarchicalSchemaUtil.schemaSeparator());
+
+ store = new RelationalEntityStore();
+ store.initialize(config);
+ }
+
+ private static BaseMetalake metalake() {
+ return BaseMetalake.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName(METALAKE)
+ .withAuditInfo(AUDIT_INFO)
+ .withComment("")
+ .withProperties(null)
+ .withVersion(SchemaVersion.V_0_1)
+ .build();
+ }
+
+ private static CatalogEntity catalog() {
+ return CatalogEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName(CATALOG)
+ .withNamespace(Namespace.of(METALAKE))
+ .withType(Catalog.Type.RELATIONAL)
+ .withProvider("test")
+ .withComment("")
+ .withProperties(null)
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ }
+
+ private static SchemaEntity schema(String name) {
+ return SchemaEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName(name)
+ .withNamespace(Namespace.of(METALAKE, CATALOG))
+ .withComment("")
+ .withProperties(null)
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ }
+
+ private static TableEntity table(String name, String schemaName) {
+ return TableEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName(name)
+ .withNamespace(Namespace.of(METALAKE, CATALOG, schemaName))
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ }
+}