diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/catalog/FlinkCatalog.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/catalog/FlinkCatalog.java index 40e9781cd6..181a55b268 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/catalog/FlinkCatalog.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/catalog/FlinkCatalog.java @@ -802,11 +802,17 @@ public void alterPartition( @Override public List listFunctions(String dbName) throws DatabaseNotExistException, CatalogException { + if (!databaseExists(dbName)) { + throw new DatabaseNotExistException(getName(), dbName); + } return new ArrayList<>(BUILTIN_BITMAP_FUNCTIONS.keySet()); } @Override public boolean functionExists(ObjectPath objectPath) throws CatalogException { + if (!databaseExists(objectPath.getDatabaseName())) { + return false; + } return BUILTIN_BITMAP_FUNCTIONS.containsKey( objectPath.getObjectName().toLowerCase(Locale.ROOT)); } @@ -814,6 +820,9 @@ public boolean functionExists(ObjectPath objectPath) throws CatalogException { @Override public CatalogFunction getFunction(ObjectPath functionPath) throws FunctionNotExistException, CatalogException { + if (!databaseExists(functionPath.getDatabaseName())) { + throw new FunctionNotExistException(getName(), functionPath); + } String className = BUILTIN_BITMAP_FUNCTIONS.get(functionPath.getObjectName().toLowerCase(Locale.ROOT)); if (className == null) { diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/catalog/FlinkCatalogITCase.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/catalog/FlinkCatalogITCase.java index 5762ec7658..9924cfbc4c 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/catalog/FlinkCatalogITCase.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/catalog/FlinkCatalogITCase.java @@ -1004,6 +1004,17 @@ void testCreateCatalogWithUnexistedDatabase() { "The configured default-database 'non-exist' does not exist in the Fluss cluster."); } + @Test + void testBitmapFunctionFailsForFullyQualifiedNonexistentDatabase() { + assertThatThrownBy( + () -> + tEnv.executeSql( + "SELECT " + + CATALOG_NAME + + ".nonexistent_db.rb_build(ARRAY[1,2])")) + .hasMessageContaining("No match found for function signature"); + } + @Test void testCreateCatalogWithLakeProperties() throws Exception { Map properties = new HashMap<>(); diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/catalog/FlinkCatalogTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/catalog/FlinkCatalogTest.java index a72595dce2..088574e766 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/catalog/FlinkCatalogTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/catalog/FlinkCatalogTest.java @@ -1221,4 +1221,37 @@ void testBitmapFunctionsRegistered() throws Exception { // verify unknown still returns false assertThat(catalog.functionExists(new ObjectPath(DEFAULT_DB, "unknown_fn"))).isFalse(); } + + @Test + void testBuiltinFunctionsRequireExistingDatabase() throws Exception { + String nonexistentDb = "nonexistent_db_for_functions"; + assertThat(catalog.databaseExists(nonexistentDb)).isFalse(); + + // listFunctions on a nonexistent database throws + assertThatThrownBy(() -> catalog.listFunctions(nonexistentDb)) + .isInstanceOf(DatabaseNotExistException.class) + .hasMessage( + "Database %s does not exist in Catalog %s.", nonexistentDb, CATALOG_NAME); + + // functionExists on a nonexistent database returns false, not throw + ObjectPath qualifiedInNonexistentDb = new ObjectPath(nonexistentDb, "rb_build"); + assertThat(catalog.functionExists(qualifiedInNonexistentDb)).isFalse(); + + // getFunction on a nonexistent database throws FunctionNotExistException + assertThatThrownBy(() -> catalog.getFunction(qualifiedInNonexistentDb)) + .isInstanceOf(FunctionNotExistException.class); + + // built-in functions still resolve from every existing database, not just DEFAULT_DB + String secondDb = "second_db_for_functions"; + catalog.createDatabase( + secondDb, new CatalogDatabaseImpl(Collections.emptyMap(), null), true); + try { + assertThat(catalog.functionExists(new ObjectPath(secondDb, "rb_build"))).isTrue(); + assertThat(catalog.getFunction(new ObjectPath(secondDb, "rb_build"))).isNotNull(); + assertThat(catalog.listFunctions(secondDb)) + .contains("rb_build_agg", "rb_or_agg", "rb_and_agg", "rb_xor_agg"); + } finally { + catalog.dropDatabase(secondDb, true, true); + } + } }