Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Identifier> tables = Collections.emptyList();
if (cascade) {
List<Identifier> 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Loading