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"));
+        });
+  }
 }

Reply via email to