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);