This is an automated email from the ASF dual-hosted git repository.
jerryshao pushed a commit to branch branch-1.3
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/branch-1.3 by this push:
new 4e514cb40f [Cherry-pick to branch-1.3] [#12929] fix(flink-connector):
Use catalog database instead of schema for PostgreSQL JDBC table scans (#12930)
(#12967)
4e514cb40f is described below
commit 4e514cb40f62f36f0d17f1e031a6a290041488db
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Tue Sep 8 19:53:11 2026 +0800
[Cherry-pick to branch-1.3] [#12929] fix(flink-connector): Use catalog
database instead of schema for PostgreSQL JDBC table scans (#12930) (#12967)
**Cherry-pick Information:**
- Original commit: 3e7b550d6a7452a96ffd195069f8443f9b97263a
- Target branch: `branch-1.3`
- Status: ✅ Clean cherry-pick (no conflicts)
---------
Co-authored-by: Yuhui <[email protected]>
Co-authored-by: Claude Sonnet 5 <[email protected]>
Co-authored-by: diqiu50 <[email protected]>
---
docs/flink-connector/flink-catalog-jdbc.md | 16 ++--
.../connector/jdbc/JdbcPropertiesConstants.java | 1 +
.../connector/jdbc/JdbcPropertiesConverter.java | 39 +++++++-
.../postgresql/PostgresqlPropertiesConverter.java | 25 ++++++
.../jdbc/TestMysqlPropertiesConverter.java | 23 +++++
.../jdbc/TestPostgresqlPropertiesConverter.java | 100 +++++++++++++++++++++
6 files changed, 195 insertions(+), 9 deletions(-)
diff --git a/docs/flink-connector/flink-catalog-jdbc.md
b/docs/flink-connector/flink-catalog-jdbc.md
index 6aa5b49272..8d6abbfe59 100644
--- a/docs/flink-connector/flink-catalog-jdbc.md
+++ b/docs/flink-connector/flink-catalog-jdbc.md
@@ -14,6 +14,7 @@ This document provides a comprehensive guide on configuring
and using Apache Gra
### JDBC Types
* MYSQL
+* POSTGRESQL
## Getting Started
@@ -38,6 +39,8 @@ Next, when you create the JDBC catalog in Gravitino, add the
`flink.bypass.defau
flink.bypass.default-database=db
```
+For a PostgreSQL catalog, if `flink.bypass.default-database` is not set, it
falls back to the catalog's `jdbc-database` property. This distinction matters
for PostgreSQL because a Flink "database" corresponds to a PostgreSQL schema,
not the PostgreSQL database itself; the JDBC connection always targets
`jdbc-database` (or the `flink.bypass.default-database` override), while `SHOW
TABLES FROM <schema>` and table scans address tables within that database using
the schema name.
+
### SQL Example
```sql
@@ -127,9 +130,10 @@ SELECT * FROM jdbc_table_a;
Gravitino Flink connector will transform below property names which are
defined in catalog properties to Flink JDBC connector configuration.
-| Gravitino catalog property name | Flink JDBC connector configuration |
Description |
-|:--------------------------------|------------------------------------|--------------------------------|
-| `jdbc-url` | `base-url` | JDBC
URL for MYSQL |
-| `username` | `username` |
Username of MySQL account |
-| `password` | `password` |
Password of the account |
-| `flink.bypass.default-database` | `default-database` |
Default database to connect to |
+| Gravitino catalog property name | Flink JDBC connector configuration |
Description
|
+|:--------------------------------|:-----------------------------------|:--------------------------------------------------------------------------------------------|
+| `jdbc-url` | `base-url` | JDBC
URL for the catalog
|
+| `username` | `username` |
Username of the account
|
+| `password` | `password` |
Password of the account
|
+| `flink.bypass.default-database` | `default-database` |
Default database to connect to. For PostgreSQL, falls back to `jdbc-database`
when not set. |
+| `jdbc-database` | (see above) |
Required for PostgreSQL catalogs; the PostgreSQL database the catalog connects
to. |
diff --git
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/JdbcPropertiesConstants.java
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/JdbcPropertiesConstants.java
index 4bea348e20..3a6de0ebd9 100644
---
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/JdbcPropertiesConstants.java
+++
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/JdbcPropertiesConstants.java
@@ -30,6 +30,7 @@ public class JdbcPropertiesConstants {
public static final String GRAVITINO_JDBC_PASSWORD = "jdbc-password";
public static final String GRAVITINO_JDBC_URL = "jdbc-url";
public static final String GRAVITINO_JDBC_DRIVER = "jdbc-driver";
+ public static final String GRAVITINO_JDBC_DATABASE = "jdbc-database";
public static final String FLINK_JDBC_URL = "base-url";
public static final String FLINK_JDBC_USER = "username";
diff --git
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/JdbcPropertiesConverter.java
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/JdbcPropertiesConverter.java
index da98a804c7..611952e02d 100644
---
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/JdbcPropertiesConverter.java
+++
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/JdbcPropertiesConverter.java
@@ -23,6 +23,7 @@ import java.util.HashMap;
import java.util.Map;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
+import org.apache.commons.lang3.StringUtils;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.table.catalog.ObjectPath;
import org.apache.flink.util.Preconditions;
@@ -59,6 +60,19 @@ public abstract class JdbcPropertiesConverter
// The URL in FlinkJdbcCatalog does not support database and other
parameters.
flinkCatalogProperties.put(
JdbcPropertiesConstants.FLINK_JDBC_URL,
getBaseUrlFromJdbcUrl(gravitinoJdbcUrl));
+
+ // Default the Flink catalog's 'default-database' option from the
Gravitino jdbc-database
+ // property when it was not already supplied (e.g. via
flink.bypass.default-database). This
+ // makes the value available when building per-table scan properties,
since only the Flink
+ // catalog properties (not the Gravitino table properties) reach
toFlinkTableProperties, and
+ // 'default-database' -- unlike 'jdbc-database' -- is an option Flink's
JdbcCatalogFactory
+ // already recognizes, so it survives that factory's strict option
validation.
+ String gravitinoJdbcDatabase =
+
gravitinoProperties.get(JdbcPropertiesConstants.GRAVITINO_JDBC_DATABASE);
+ if (StringUtils.isNotBlank(gravitinoJdbcDatabase)) {
+ flinkCatalogProperties.putIfAbsent(
+ JdbcPropertiesConstants.FLINK_JDBC_DEFAULT_DATABASE,
gravitinoJdbcDatabase);
+ }
return flinkCatalogProperties;
}
@@ -75,7 +89,7 @@ public abstract class JdbcPropertiesConverter
@Override
public Map<String, String> toFlinkTableProperties(
Map<String, String> flinkCatalogProperties,
- Map<String, String> gravitinoProperties,
+ Map<String, String> gravitinoTableProperties,
ObjectPath tablePath) {
String jdbcUser =
flinkCatalogProperties.get(JdbcPropertiesConstants.FLINK_JDBC_USER);
String jdbcPassword =
flinkCatalogProperties.get(JdbcPropertiesConstants.FLINK_JDBC_PASSWORD);
@@ -89,13 +103,32 @@ public abstract class JdbcPropertiesConverter
Map<String, String> tableOptions = new HashMap<>();
tableOptions.put(
JdbcPropertiesConstants.FLINK_JDBC_TABLE_DATABASE_URL,
- jdbcBaseUrl + tablePath.getDatabaseName());
- tableOptions.put(JdbcPropertiesConstants.FLINK_JDBC_TABLE_NAME,
tablePath.getObjectName());
+ jdbcBaseUrl + getConnectionDatabase(flinkCatalogProperties,
tablePath));
+ tableOptions.put(JdbcPropertiesConstants.FLINK_JDBC_TABLE_NAME,
getTableName(tablePath));
tableOptions.put(JdbcPropertiesConstants.FLINK_JDBC_USER, jdbcUser);
tableOptions.put(JdbcPropertiesConstants.FLINK_JDBC_PASSWORD,
jdbcPassword);
return tableOptions;
}
+ /**
+ * Returns the database name used to build the JDBC connection URL for a
table scan. Defaults to
+ * the Flink "database", which corresponds to the Gravitino schema name.
+ *
+ * @param flinkCatalogProperties the Flink catalog properties produced by
{@link
+ * #toFlinkCatalogProperties(Map)}.
+ * @param tablePath the Flink table path being scanned.
+ * @return the database name to use when building the JDBC connection URL.
+ */
+ protected String getConnectionDatabase(
+ Map<String, String> flinkCatalogProperties, ObjectPath tablePath) {
+ return tablePath.getDatabaseName();
+ }
+
+ /** Returns the table name used in the JDBC connector's {@code table-name}
option. */
+ protected String getTableName(ObjectPath tablePath) {
+ return tablePath.getObjectName();
+ }
+
protected abstract String defaultDriverName();
private static String getBaseUrlFromJdbcUrl(String jdbcUrl) {
diff --git
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/postgresql/PostgresqlPropertiesConverter.java
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/postgresql/PostgresqlPropertiesConverter.java
index ffa3fd8237..ea6b351572 100644
---
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/postgresql/PostgresqlPropertiesConverter.java
+++
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/postgresql/PostgresqlPropertiesConverter.java
@@ -19,7 +19,12 @@
package org.apache.gravitino.flink.connector.jdbc.postgresql;
+import com.google.common.base.Strings;
+import java.util.Map;
+import org.apache.flink.table.catalog.ObjectPath;
+import org.apache.flink.util.Preconditions;
import
org.apache.gravitino.flink.connector.jdbc.GravitinoJdbcCatalogFactoryOptions;
+import org.apache.gravitino.flink.connector.jdbc.JdbcPropertiesConstants;
import org.apache.gravitino.flink.connector.jdbc.JdbcPropertiesConverter;
public class PostgresqlPropertiesConverter extends JdbcPropertiesConverter {
@@ -37,4 +42,24 @@ public class PostgresqlPropertiesConverter extends
JdbcPropertiesConverter {
public String getFlinkCatalogType() {
return GravitinoJdbcCatalogFactoryOptions.POSTGRESQL_IDENTIFIER;
}
+
+ @Override
+ protected String getConnectionDatabase(
+ Map<String, String> flinkCatalogProperties, ObjectPath tablePath) {
+ // For PostgreSQL, the Flink "database" is a Gravitino schema, not the
PostgreSQL database
+ // that the JDBC connection must target. The connection must use the
catalog's configured
+ // jdbc-database instead, which
JdbcPropertiesConverter#toFlinkCatalogProperties defaults into
+ // the 'default-database' option.
+ String database =
+
flinkCatalogProperties.get(JdbcPropertiesConstants.FLINK_JDBC_DEFAULT_DATABASE);
+ Preconditions.checkArgument(
+ !Strings.isNullOrEmpty(database),
+ JdbcPropertiesConstants.FLINK_JDBC_DEFAULT_DATABASE + " should not be
null or empty.");
+ return database;
+ }
+
+ @Override
+ protected String getTableName(ObjectPath tablePath) {
+ return tablePath.getDatabaseName() + "." + tablePath.getObjectName();
+ }
}
diff --git
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/jdbc/TestMysqlPropertiesConverter.java
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/jdbc/TestMysqlPropertiesConverter.java
index 263575ddd9..2ad10aefa4 100644
---
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/jdbc/TestMysqlPropertiesConverter.java
+++
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/jdbc/TestMysqlPropertiesConverter.java
@@ -19,8 +19,12 @@
package org.apache.gravitino.flink.connector.jdbc;
+import com.google.common.collect.ImmutableMap;
import java.util.Map;
+import org.apache.flink.table.catalog.ObjectPath;
import
org.apache.gravitino.flink.connector.jdbc.mysql.MysqlPropertiesConverter;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
public class TestMysqlPropertiesConverter extends
AbstractJdbcPropertiesConverterTestSuite {
@@ -28,4 +32,23 @@ public class TestMysqlPropertiesConverter extends
AbstractJdbcPropertiesConverte
protected JdbcPropertiesConverter getConverter(Map<String, String>
catalogOptions) {
return MysqlPropertiesConverter.INSTANCE;
}
+
+ @Test
+ public void testToFlinkTableProperties() {
+ String schema = "myDatabase";
+ String tableName = "myTable";
+ Map<String, String> flinkCatalogProperties =
+
getConverter(catalogProperties).toFlinkCatalogProperties(catalogProperties);
+ Map<String, String> tableProperties =
+ getConverter(catalogProperties)
+ .toFlinkTableProperties(
+ flinkCatalogProperties, ImmutableMap.of(), new
ObjectPath(schema, tableName));
+
+ // For MySQL, the Flink "database" (schema) is itself the connection
database.
+ Assertions.assertEquals(
+ flinkUrl + schema,
+
tableProperties.get(JdbcPropertiesConstants.FLINK_JDBC_TABLE_DATABASE_URL));
+ Assertions.assertEquals(
+ tableName,
tableProperties.get(JdbcPropertiesConstants.FLINK_JDBC_TABLE_NAME));
+ }
}
diff --git
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/jdbc/TestPostgresqlPropertiesConverter.java
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/jdbc/TestPostgresqlPropertiesConverter.java
index 4487d14319..7c0a8023ea 100644
---
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/jdbc/TestPostgresqlPropertiesConverter.java
+++
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/jdbc/TestPostgresqlPropertiesConverter.java
@@ -19,13 +19,113 @@
package org.apache.gravitino.flink.connector.jdbc;
+import com.google.common.collect.ImmutableMap;
+import java.util.HashMap;
import java.util.Map;
+import org.apache.flink.table.catalog.ObjectPath;
import
org.apache.gravitino.flink.connector.jdbc.postgresql.PostgresqlPropertiesConverter;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
public class TestPostgresqlPropertiesConverter extends
AbstractJdbcPropertiesConverterTestSuite {
+ private static final String FLINK_BYPASS_DEFAULT_DATABASE =
"flink.bypass.default-database";
+
@Override
protected JdbcPropertiesConverter getConverter(Map<String, String>
catalogOptions) {
return PostgresqlPropertiesConverter.INSTANCE;
}
+
+ @Test
+ public void testToFlinkTableProperties() {
+ String jdbcDatabase = "gravitino";
+ String schema = "public";
+ String tableName = "table_meta";
+ // No 'flink.bypass.default-database' is set, so the connection database
must fall back to
+ // the catalog's jdbc-database.
+ Map<String, String> catalogPropertiesWithDatabase = new
HashMap<>(catalogProperties);
+ catalogPropertiesWithDatabase.remove(FLINK_BYPASS_DEFAULT_DATABASE);
+ catalogPropertiesWithDatabase.put(
+ JdbcPropertiesConstants.GRAVITINO_JDBC_URL, gravitinoUrlWithDomain);
+ catalogPropertiesWithDatabase.put(
+ JdbcPropertiesConstants.GRAVITINO_JDBC_DATABASE, jdbcDatabase);
+
+ // Mirrors the production call in BaseCatalog#toFlinkTable: the first
argument is the Flink
+ // catalog properties (as produced by toFlinkCatalogProperties), the
second is the Gravitino
+ // table's own properties, which do not carry catalog-level properties
like jdbc-database.
+ Map<String, String> flinkCatalogProperties =
+ getConverter(catalogPropertiesWithDatabase)
+ .toFlinkCatalogProperties(catalogPropertiesWithDatabase);
+ Map<String, String> tableProperties =
+ getConverter(catalogPropertiesWithDatabase)
+ .toFlinkTableProperties(
+ flinkCatalogProperties, ImmutableMap.of(), new
ObjectPath(schema, tableName));
+
+ // The connection URL must target the PostgreSQL database (jdbc-database),
not the schema.
+ Assertions.assertEquals(
+ flinkUrlWithDomain + jdbcDatabase,
+
tableProperties.get(JdbcPropertiesConstants.FLINK_JDBC_TABLE_DATABASE_URL));
+ // The schema must be carried via the schema-qualified table name instead.
+ Assertions.assertEquals(
+ schema + "." + tableName,
+ tableProperties.get(JdbcPropertiesConstants.FLINK_JDBC_TABLE_NAME));
+ }
+
+ @Test
+ public void testToFlinkTablePropertiesPrefersExplicitDefaultDatabase() {
+ // When 'flink.bypass.default-database' is explicitly set, it takes
precedence over
+ // jdbc-database.
+ String explicitDefaultDatabase = "explicit_db";
+ Map<String, String> catalogPropertiesWithBoth = new
HashMap<>(catalogProperties);
+ catalogPropertiesWithBoth.put(FLINK_BYPASS_DEFAULT_DATABASE,
explicitDefaultDatabase);
+ catalogPropertiesWithBoth.put(
+ JdbcPropertiesConstants.GRAVITINO_JDBC_URL, gravitinoUrlWithDomain);
+
catalogPropertiesWithBoth.put(JdbcPropertiesConstants.GRAVITINO_JDBC_DATABASE,
"gravitino");
+
+ Map<String, String> flinkCatalogProperties =
+
getConverter(catalogPropertiesWithBoth).toFlinkCatalogProperties(catalogPropertiesWithBoth);
+ Map<String, String> tableProperties =
+ getConverter(catalogPropertiesWithBoth)
+ .toFlinkTableProperties(
+ flinkCatalogProperties, ImmutableMap.of(), new
ObjectPath("public", "t"));
+
+ Assertions.assertEquals(
+ flinkUrlWithDomain + explicitDefaultDatabase,
+
tableProperties.get(JdbcPropertiesConstants.FLINK_JDBC_TABLE_DATABASE_URL));
+ }
+
+ @Test
+ public void testToFlinkTablePropertiesWithoutJdbcDatabase() {
+ Map<String, String> catalogPropertiesWithoutDatabase = new
HashMap<>(catalogProperties);
+ catalogPropertiesWithoutDatabase.remove(FLINK_BYPASS_DEFAULT_DATABASE);
+
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () -> {
+ Map<String, String> flinkCatalogProperties =
+ getConverter(catalogPropertiesWithoutDatabase)
+ .toFlinkCatalogProperties(catalogPropertiesWithoutDatabase);
+ getConverter(catalogPropertiesWithoutDatabase)
+ .toFlinkTableProperties(
+ flinkCatalogProperties, ImmutableMap.of(), new
ObjectPath("public", "t"));
+ });
+ }
+
+ @Test
+ public void testToFlinkTablePropertiesWithEmptyJdbcDatabase() {
+ Map<String, String> catalogPropertiesWithEmptyDatabase = new
HashMap<>(catalogProperties);
+ catalogPropertiesWithEmptyDatabase.remove(FLINK_BYPASS_DEFAULT_DATABASE);
+
catalogPropertiesWithEmptyDatabase.put(JdbcPropertiesConstants.GRAVITINO_JDBC_DATABASE,
"");
+
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () -> {
+ Map<String, String> flinkCatalogProperties =
+ getConverter(catalogPropertiesWithEmptyDatabase)
+
.toFlinkCatalogProperties(catalogPropertiesWithEmptyDatabase);
+ getConverter(catalogPropertiesWithEmptyDatabase)
+ .toFlinkTableProperties(
+ flinkCatalogProperties, ImmutableMap.of(), new
ObjectPath("public", "t"));
+ });
+ }
}