From f9bca319bcf86e49348f006057464257cc84778b Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 26 Sep 2026 16:59:54 +0800 Subject: [PATCH] [core] Invalidate partition cache when dropping a database The cascade branch of dropDatabase invalidated only the table-cache entries it happened to find, by key set. Tables whose table-cache entry had expired were missed entirely, so with the partition cache enabled a dropped database's partitions were served to a recreated same-name table until the entry itself expired. Enumerate the database's tables before the cascade drop and route each through invalidateTable, which clears the partition cache and the branch entries too. Assisted-by: GLM-5.3 --- .../apache/paimon/catalog/CachingCatalog.java | 24 ++++++++++++------- .../paimon/catalog/CachingCatalogTest.java | 20 ++++++++++++++++ 2 files changed, 35 insertions(+), 9 deletions(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/catalog/CachingCatalog.java b/paimon-core/src/main/java/org/apache/paimon/catalog/CachingCatalog.java index 037f7ea2bfba..846b14247fda 100644 --- a/paimon-core/src/main/java/org/apache/paimon/catalog/CachingCatalog.java +++ b/paimon-core/src/main/java/org/apache/paimon/catalog/CachingCatalog.java @@ -41,10 +41,11 @@ import javax.annotation.Nullable; import java.time.Duration; -import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.stream.Collectors; import static org.apache.paimon.options.CatalogOptions.CACHE_DV_MAX_NUM; import static org.apache.paimon.options.CatalogOptions.CACHE_ENABLED; @@ -197,17 +198,22 @@ public Database getDatabase(String databaseName) throws DatabaseNotExistExceptio @Override public void dropDatabase(String name, boolean ignoreIfNotExists, boolean cascade) throws DatabaseNotExistException, DatabaseNotEmptyException { - super.dropDatabase(name, ignoreIfNotExists, cascade); - databaseCache.invalidate(name); + // enumerate before the drop: cascade removes the tables from the wrapped catalog, + // and cache entries alone miss tables whose table-cache entry already expired + List tables = Collections.emptyList(); if (cascade) { - List tables = new ArrayList<>(); - for (Identifier identifier : tableCache.asMap().keySet()) { - if (identifier.getDatabaseName().equals(name)) { - tables.add(identifier); - } + try { + tables = + listTables(name).stream() + .map(tableName -> new Identifier(name, tableName)) + .collect(Collectors.toList()); + } catch (DatabaseNotExistException ignored) { + // super.dropDatabase reports the missing database on its own terms } - tables.forEach(tableCache::invalidate); } + super.dropDatabase(name, ignoreIfNotExists, cascade); + databaseCache.invalidate(name); + tables.forEach(this::invalidateTable); } @Override diff --git a/paimon-core/src/test/java/org/apache/paimon/catalog/CachingCatalogTest.java b/paimon-core/src/test/java/org/apache/paimon/catalog/CachingCatalogTest.java index 71427641b79f..2506c1719ed2 100644 --- a/paimon-core/src/test/java/org/apache/paimon/catalog/CachingCatalogTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/catalog/CachingCatalogTest.java @@ -345,6 +345,26 @@ public void testPartitionCache() throws Exception { assertThat(partitionEntryListFromCache).containsAll(partitionEntryList); } + @Test + public void testDropDatabaseCascadeInvalidatesPartitionCache() throws Exception { + Catalog wrapped = Mockito.mock(Catalog.class); + TestableCachingCatalog catalog = + new TestableCachingCatalog(wrapped, EXPIRATION_TTL, ticker); + Identifier identifier = new Identifier("db", "tbl"); + Partition dropped = new Partition(singletonMap("dt", "20260101"), 0, 0, 0, 0, -1, false); + Mockito.when(wrapped.listTables("db")).thenReturn(singletonList("tbl")); + Mockito.when(wrapped.listPartitions(identifier)) + .thenReturn(singletonList(dropped), emptyList()); + + assertThat(catalog.listPartitions(identifier)).containsExactly(dropped); + + // the table-cache key-set enumeration missed tables whose entry already expired, + // so a recreated same-name table served the dropped table's partitions + catalog.dropDatabase("db", false, true); + + assertThat(catalog.listPartitions(identifier)).isEmpty(); + } + @Test public void testCreatePartitionsInvalidatesPartitionCache() throws Exception { Catalog wrapped = Mockito.mock(Catalog.class);