This is an automated email from the ASF dual-hosted git repository.

fanrui pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-jdbc.git

commit 5a8060b39d50490dc317ac941782df36d9a1d42d
Author: Roc Marshal <[email protected]>
AuthorDate: Fri Apr 19 23:37:48 2024 +0800

    [FLINK-35176][Connector/JDBC] Support property authentication connection 
for JDBC catalog
---
 .../connector/jdbc/JdbcConnectionOptions.java      | 12 ++++++
 .../jdbc/catalog/AbstractJdbcCatalog.java          | 46 ++++++++++++++++------
 .../flink/connector/jdbc/catalog/JdbcCatalog.java  | 45 +++++++++++++++++----
 .../connector/jdbc/catalog/JdbcCatalogUtils.java   | 27 +++++++++++--
 .../jdbc/catalog/factory/JdbcCatalogFactory.java   |  6 +--
 .../catalog/factory/JdbcCatalogFactoryOptions.java |  7 ++--
 .../databases/cratedb/catalog/CrateDBCatalog.java  | 24 +++++++++--
 .../jdbc/databases/mysql/catalog/MySqlCatalog.java | 25 ++++++++++--
 .../postgres/catalog/PostgresCatalog.java          | 45 ++++++++++++++++++---
 .../flink/connector/jdbc/utils/JdbcUtils.java      | 28 +++++++++++++
 10 files changed, 223 insertions(+), 42 deletions(-)

diff --git 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/JdbcConnectionOptions.java
 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/JdbcConnectionOptions.java
index 85b12951..f9f7eb53 100644
--- 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/JdbcConnectionOptions.java
+++ 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/JdbcConnectionOptions.java
@@ -81,6 +81,18 @@ public class JdbcConnectionOptions implements Serializable {
         return properties;
     }
 
+    @Nonnull
+    public static Properties getBriefAuthProperties(String user, String 
password) {
+        final Properties result = new Properties();
+        if (Objects.nonNull(user)) {
+            result.put(USER_KEY, user);
+        }
+        if (Objects.nonNull(password)) {
+            result.put(PASSWORD_KEY, password);
+        }
+        return result;
+    }
+
     /** Builder for {@link JdbcConnectionOptions}. */
     public static class JdbcConnectionOptionsBuilder {
         private String url;
diff --git 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/AbstractJdbcCatalog.java
 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/AbstractJdbcCatalog.java
index 7ba0c06d..a0521484 100644
--- 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/AbstractJdbcCatalog.java
+++ 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/AbstractJdbcCatalog.java
@@ -71,8 +71,12 @@ import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Optional;
+import java.util.Properties;
 import java.util.function.Predicate;
 
+import static 
org.apache.flink.connector.jdbc.JdbcConnectionOptions.PASSWORD_KEY;
+import static org.apache.flink.connector.jdbc.JdbcConnectionOptions.USER_KEY;
+import static 
org.apache.flink.connector.jdbc.JdbcConnectionOptions.getBriefAuthProperties;
 import static 
org.apache.flink.connector.jdbc.table.JdbcConnectorOptions.PASSWORD;
 import static 
org.apache.flink.connector.jdbc.table.JdbcConnectorOptions.TABLE_NAME;
 import static org.apache.flink.connector.jdbc.table.JdbcConnectorOptions.URL;
@@ -88,11 +92,11 @@ public abstract class AbstractJdbcCatalog extends 
AbstractCatalog {
     private static final Logger LOG = 
LoggerFactory.getLogger(AbstractJdbcCatalog.class);
 
     protected final ClassLoader userClassLoader;
-    protected final String username;
-    protected final String pwd;
     protected final String baseUrl;
     protected final String defaultUrl;
+    protected final Properties connectionProperties;
 
+    @Deprecated
     public AbstractJdbcCatalog(
             ClassLoader userClassLoader,
             String catalogName,
@@ -100,20 +104,36 @@ public abstract class AbstractJdbcCatalog extends 
AbstractCatalog {
             String username,
             String pwd,
             String baseUrl) {
+        this(
+                userClassLoader,
+                catalogName,
+                defaultDatabase,
+                baseUrl,
+                getBriefAuthProperties(username, pwd));
+    }
+
+    public AbstractJdbcCatalog(
+            ClassLoader userClassLoader,
+            String catalogName,
+            String defaultDatabase,
+            String baseUrl,
+            Properties connectionProperties) {
         super(catalogName, defaultDatabase);
 
         checkNotNull(userClassLoader);
-        checkArgument(!StringUtils.isNullOrWhitespaceOnly(username));
-        checkArgument(!StringUtils.isNullOrWhitespaceOnly(pwd));
         checkArgument(!StringUtils.isNullOrWhitespaceOnly(baseUrl));
 
         JdbcCatalogUtils.validateJdbcUrl(baseUrl);
 
         this.userClassLoader = userClassLoader;
-        this.username = username;
-        this.pwd = pwd;
         this.baseUrl = baseUrl.endsWith("/") ? baseUrl : baseUrl + "/";
         this.defaultUrl = this.baseUrl + defaultDatabase;
+        this.connectionProperties = 
Preconditions.checkNotNull(connectionProperties);
+        checkArgument(
+                
!StringUtils.isNullOrWhitespaceOnly(connectionProperties.getProperty(USER_KEY)));
+        checkArgument(
+                !StringUtils.isNullOrWhitespaceOnly(
+                        connectionProperties.getProperty(PASSWORD_KEY)));
     }
 
     @Override
@@ -122,7 +142,7 @@ public abstract class AbstractJdbcCatalog extends 
AbstractCatalog {
         try (TemporaryClassLoaderContext ignored =
                 TemporaryClassLoaderContext.of(userClassLoader)) {
             // test connection, fail early if we cannot connect to database
-            try (Connection conn = DriverManager.getConnection(defaultUrl, 
username, pwd)) {
+            try (Connection conn = DriverManager.getConnection(defaultUrl, 
connectionProperties)) {
             } catch (SQLException e) {
                 throw new ValidationException(
                         String.format("Failed connecting to %s via JDBC.", 
defaultUrl), e);
@@ -139,11 +159,11 @@ public abstract class AbstractJdbcCatalog extends 
AbstractCatalog {
     // ----- getters ------
 
     public String getUsername() {
-        return username;
+        return connectionProperties.getProperty(USER_KEY);
     }
 
     public String getPassword() {
-        return pwd;
+        return connectionProperties.getProperty(PASSWORD_KEY);
     }
 
     public String getBaseUrl() {
@@ -248,7 +268,7 @@ public abstract class AbstractJdbcCatalog extends 
AbstractCatalog {
         String databaseName = tablePath.getDatabaseName();
         String dbUrl = baseUrl + databaseName;
 
-        try (Connection conn = DriverManager.getConnection(dbUrl, username, 
pwd)) {
+        try (Connection conn = DriverManager.getConnection(dbUrl, 
connectionProperties)) {
             DatabaseMetaData metaData = conn.getMetaData();
             Optional<UniqueConstraint> primaryKey =
                     getPrimaryKey(
@@ -282,8 +302,8 @@ public abstract class AbstractJdbcCatalog extends 
AbstractCatalog {
             Map<String, String> props = new HashMap<>();
             props.put(CONNECTOR.key(), IDENTIFIER);
             props.put(URL.key(), dbUrl);
-            props.put(USERNAME.key(), username);
-            props.put(PASSWORD.key(), pwd);
+            props.put(USERNAME.key(), 
connectionProperties.getProperty(USER_KEY));
+            props.put(PASSWORD.key(), 
connectionProperties.getProperty(PASSWORD_KEY));
             props.put(TABLE_NAME.key(), getSchemaTableName(tablePath));
             return CatalogTable.of(tableSchema, null, Lists.newArrayList(), 
props);
         } catch (Exception e) {
@@ -497,7 +517,7 @@ public abstract class AbstractJdbcCatalog extends 
AbstractCatalog {
             Predicate<String> filterFunc,
             Object... params) {
 
-        try (Connection conn = DriverManager.getConnection(connUrl, username, 
pwd);
+        try (Connection conn = DriverManager.getConnection(connUrl, 
connectionProperties);
                 PreparedStatement ps = conn.prepareStatement(sql)) {
             return extractColumnValuesByStatement(ps, columnIndex, filterFunc, 
params);
 
diff --git 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/JdbcCatalog.java
 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/JdbcCatalog.java
index 3f6e28fa..44c0fa64 100644
--- 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/JdbcCatalog.java
+++ 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/JdbcCatalog.java
@@ -29,6 +29,9 @@ import 
org.apache.flink.table.catalog.exceptions.DatabaseNotExistException;
 import org.apache.flink.table.catalog.exceptions.TableNotExistException;
 
 import java.util.List;
+import java.util.Properties;
+
+import static 
org.apache.flink.connector.jdbc.JdbcConnectionOptions.getBriefAuthProperties;
 
 /** Catalogs for relational databases via JDBC. */
 @PublicEvolving
@@ -36,11 +39,12 @@ public class JdbcCatalog extends AbstractJdbcCatalog {
 
     private final AbstractJdbcCatalog internal;
 
+    @Deprecated
     /**
      * Creates a JdbcCatalog.
      *
      * @deprecated please use {@link JdbcCatalog#JdbcCatalog(ClassLoader, 
String, String, String,
-     *     String, String, String)} instead.
+     *     String, Properties)} instead.
      */
     public JdbcCatalog(
             String catalogName,
@@ -52,12 +56,12 @@ public class JdbcCatalog extends AbstractJdbcCatalog {
                 Thread.currentThread().getContextClassLoader(),
                 catalogName,
                 defaultDatabase,
-                username,
-                pwd,
                 baseUrl,
-                null);
+                null,
+                getBriefAuthProperties(username, pwd));
     }
 
+    @VisibleForTesting
     /**
      * Creates a JdbcCatalog.
      *
@@ -77,17 +81,42 @@ public class JdbcCatalog extends AbstractJdbcCatalog {
             String pwd,
             String baseUrl,
             String compatibleMode) {
-        super(userClassLoader, catalogName, defaultDatabase, username, pwd, 
baseUrl);
+        this(
+                userClassLoader,
+                catalogName,
+                defaultDatabase,
+                baseUrl,
+                compatibleMode,
+                getBriefAuthProperties(username, pwd));
+    }
+
+    /**
+     * Creates a JdbcCatalog.
+     *
+     * @param userClassLoader the classloader used to load JDBC driver
+     * @param catalogName the registered catalog name
+     * @param defaultDatabase the default database name
+     * @param connectProperties the properties used to connect the database
+     * @param baseUrl the base URL of the database, e.g. 
jdbc:mysql://localhost:3306
+     * @param compatibleMode the compatible mode of the database
+     */
+    public JdbcCatalog(
+            ClassLoader userClassLoader,
+            String catalogName,
+            String defaultDatabase,
+            String baseUrl,
+            String compatibleMode,
+            Properties connectProperties) {
+        super(userClassLoader, catalogName, defaultDatabase, baseUrl, 
connectProperties);
 
         internal =
                 JdbcCatalogUtils.createCatalog(
                         userClassLoader,
                         catalogName,
                         defaultDatabase,
-                        username,
-                        pwd,
                         baseUrl,
-                        compatibleMode);
+                        compatibleMode,
+                        connectProperties);
     }
 
     // ------ databases -----
diff --git 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/JdbcCatalogUtils.java
 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/JdbcCatalogUtils.java
index 09d4d924..a1dbf5af 100644
--- 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/JdbcCatalogUtils.java
+++ 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/JdbcCatalogUtils.java
@@ -27,6 +27,9 @@ import 
org.apache.flink.connector.jdbc.databases.postgres.dialect.PostgresDialec
 import org.apache.flink.connector.jdbc.dialect.JdbcDialect;
 import org.apache.flink.connector.jdbc.dialect.JdbcDialectLoader;
 
+import java.util.Properties;
+
+import static 
org.apache.flink.connector.jdbc.JdbcConnectionOptions.getBriefAuthProperties;
 import static org.apache.flink.util.Preconditions.checkArgument;
 
 /** Utils for {@link JdbcCatalog}. */
@@ -41,6 +44,7 @@ public class JdbcCatalogUtils {
         checkArgument(parts.length == 2);
     }
 
+    @Deprecated
     /** Create catalog instance from given information. */
     public static AbstractJdbcCatalog createCatalog(
             ClassLoader userClassLoader,
@@ -50,17 +54,34 @@ public class JdbcCatalogUtils {
             String pwd,
             String baseUrl,
             String compatibleMode) {
+        return createCatalog(
+                userClassLoader,
+                catalogName,
+                defaultDatabase,
+                baseUrl,
+                compatibleMode,
+                getBriefAuthProperties(username, pwd));
+    }
+
+    /** Create catalog instance from given information. */
+    public static AbstractJdbcCatalog createCatalog(
+            ClassLoader userClassLoader,
+            String catalogName,
+            String defaultDatabase,
+            String baseUrl,
+            String compatibleMode,
+            Properties connectionProperties) {
         JdbcDialect dialect = JdbcDialectLoader.load(baseUrl, compatibleMode, 
userClassLoader);
 
         if (dialect instanceof PostgresDialect) {
             return new PostgresCatalog(
-                    userClassLoader, catalogName, defaultDatabase, username, 
pwd, baseUrl);
+                    userClassLoader, catalogName, defaultDatabase, baseUrl, 
connectionProperties);
         } else if (dialect instanceof CrateDBDialect) {
             return new CrateDBCatalog(
-                    userClassLoader, catalogName, defaultDatabase, username, 
pwd, baseUrl);
+                    userClassLoader, catalogName, defaultDatabase, baseUrl, 
connectionProperties);
         } else if (dialect instanceof MySqlDialect) {
             return new MySqlCatalog(
-                    userClassLoader, catalogName, defaultDatabase, username, 
pwd, baseUrl);
+                    userClassLoader, catalogName, defaultDatabase, baseUrl, 
connectionProperties);
         } else {
             throw new UnsupportedOperationException(
                     String.format("Catalog for '%s' is not supported yet.", 
dialect));
diff --git 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/factory/JdbcCatalogFactory.java
 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/factory/JdbcCatalogFactory.java
index 15225c4f..a19e0bae 100644
--- 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/factory/JdbcCatalogFactory.java
+++ 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/factory/JdbcCatalogFactory.java
@@ -35,6 +35,7 @@ import static 
org.apache.flink.connector.jdbc.catalog.factory.JdbcCatalogFactory
 import static 
org.apache.flink.connector.jdbc.catalog.factory.JdbcCatalogFactoryOptions.DEFAULT_DATABASE;
 import static 
org.apache.flink.connector.jdbc.catalog.factory.JdbcCatalogFactoryOptions.PASSWORD;
 import static 
org.apache.flink.connector.jdbc.catalog.factory.JdbcCatalogFactoryOptions.USERNAME;
+import static 
org.apache.flink.connector.jdbc.utils.JdbcUtils.getConnectionProperties;
 import static org.apache.flink.table.factories.FactoryUtil.PROPERTY_VERSION;
 
 /** Factory for {@link JdbcCatalog}. */
@@ -75,9 +76,8 @@ public class JdbcCatalogFactory implements CatalogFactory {
                 context.getClassLoader(),
                 context.getName(),
                 helper.getOptions().get(DEFAULT_DATABASE),
-                helper.getOptions().get(USERNAME),
-                helper.getOptions().get(PASSWORD),
                 helper.getOptions().get(BASE_URL),
-                helper.getOptions().get(COMPATIBLE_MODE));
+                helper.getOptions().get(COMPATIBLE_MODE),
+                getConnectionProperties(helper.getOptions()));
     }
 }
diff --git 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/factory/JdbcCatalogFactoryOptions.java
 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/factory/JdbcCatalogFactoryOptions.java
index ab1ce130..197f6b9f 100644
--- 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/factory/JdbcCatalogFactoryOptions.java
+++ 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/catalog/factory/JdbcCatalogFactoryOptions.java
@@ -22,6 +22,7 @@ import org.apache.flink.annotation.Internal;
 import org.apache.flink.configuration.ConfigOption;
 import org.apache.flink.configuration.ConfigOptions;
 import org.apache.flink.connector.jdbc.catalog.JdbcCatalog;
+import org.apache.flink.connector.jdbc.table.JdbcConnectorOptions;
 import org.apache.flink.table.catalog.CommonCatalogOptions;
 
 /** {@link ConfigOption}s for {@link JdbcCatalog}. */
@@ -35,11 +36,9 @@ public class JdbcCatalogFactoryOptions {
                     .stringType()
                     .noDefaultValue();
 
-    public static final ConfigOption<String> USERNAME =
-            ConfigOptions.key("username").stringType().noDefaultValue();
+    public static final ConfigOption<String> USERNAME = 
JdbcConnectorOptions.USERNAME;
 
-    public static final ConfigOption<String> PASSWORD =
-            ConfigOptions.key("password").stringType().noDefaultValue();
+    public static final ConfigOption<String> PASSWORD = 
JdbcConnectorOptions.PASSWORD;
 
     public static final ConfigOption<String> BASE_URL =
             ConfigOptions.key("base-url").stringType().noDefaultValue();
diff --git 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/databases/cratedb/catalog/CrateDBCatalog.java
 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/databases/cratedb/catalog/CrateDBCatalog.java
index 04ed2a4f..3b00495f 100644
--- 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/databases/cratedb/catalog/CrateDBCatalog.java
+++ 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/databases/cratedb/catalog/CrateDBCatalog.java
@@ -19,6 +19,7 @@
 package org.apache.flink.connector.jdbc.databases.cratedb.catalog;
 
 import org.apache.flink.annotation.Internal;
+import org.apache.flink.annotation.VisibleForTesting;
 import 
org.apache.flink.connector.jdbc.databases.postgres.catalog.PostgresCatalog;
 import org.apache.flink.table.catalog.ObjectPath;
 import org.apache.flink.table.catalog.exceptions.CatalogException;
@@ -34,8 +35,11 @@ import java.sql.SQLException;
 import java.util.Collections;
 import java.util.HashSet;
 import java.util.List;
+import java.util.Properties;
 import java.util.Set;
 
+import static 
org.apache.flink.connector.jdbc.JdbcConnectionOptions.getBriefAuthProperties;
+
 /** Catalog for CrateDB. */
 @Internal
 public class CrateDBCatalog extends PostgresCatalog {
@@ -53,6 +57,7 @@ public class CrateDBCatalog extends PostgresCatalog {
                 }
             };
 
+    @VisibleForTesting
     public CrateDBCatalog(
             ClassLoader userClassLoader,
             String catalogName,
@@ -60,14 +65,27 @@ public class CrateDBCatalog extends PostgresCatalog {
             String username,
             String pwd,
             String baseUrl) {
+        this(
+                userClassLoader,
+                catalogName,
+                defaultDatabase,
+                baseUrl,
+                getBriefAuthProperties(username, pwd));
+    }
+
+    public CrateDBCatalog(
+            ClassLoader userClassLoader,
+            String catalogName,
+            String defaultDatabase,
+            String baseUrl,
+            Properties connecProperties) {
         super(
                 userClassLoader,
                 catalogName,
                 defaultDatabase,
-                username,
-                pwd,
                 baseUrl,
-                new CrateDBTypeMapper());
+                new CrateDBTypeMapper(),
+                connecProperties);
     }
 
     // ------ databases ------
diff --git 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/databases/mysql/catalog/MySqlCatalog.java
 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/databases/mysql/catalog/MySqlCatalog.java
index b54aa23d..a511e77c 100644
--- 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/databases/mysql/catalog/MySqlCatalog.java
+++ 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/databases/mysql/catalog/MySqlCatalog.java
@@ -19,6 +19,7 @@
 package org.apache.flink.connector.jdbc.databases.mysql.catalog;
 
 import org.apache.flink.annotation.Internal;
+import org.apache.flink.annotation.VisibleForTesting;
 import org.apache.flink.connector.jdbc.catalog.AbstractJdbcCatalog;
 import org.apache.flink.connector.jdbc.dialect.JdbcDialectTypeMapper;
 import org.apache.flink.table.catalog.ObjectPath;
@@ -38,10 +39,13 @@ import java.sql.ResultSetMetaData;
 import java.sql.SQLException;
 import java.util.HashSet;
 import java.util.List;
+import java.util.Properties;
 import java.util.Set;
 import java.util.regex.Matcher;
 import java.util.regex.Pattern;
 
+import static 
org.apache.flink.connector.jdbc.JdbcConnectionOptions.getBriefAuthProperties;
+
 /** Catalog for MySQL. */
 @Internal
 public class MySqlCatalog extends AbstractJdbcCatalog {
@@ -60,6 +64,7 @@ public class MySqlCatalog extends AbstractJdbcCatalog {
                 }
             };
 
+    @VisibleForTesting
     public MySqlCatalog(
             ClassLoader userClassLoader,
             String catalogName,
@@ -67,7 +72,21 @@ public class MySqlCatalog extends AbstractJdbcCatalog {
             String username,
             String pwd,
             String baseUrl) {
-        super(userClassLoader, catalogName, defaultDatabase, username, pwd, 
baseUrl);
+        this(
+                userClassLoader,
+                catalogName,
+                defaultDatabase,
+                baseUrl,
+                getBriefAuthProperties(username, pwd));
+    }
+
+    public MySqlCatalog(
+            ClassLoader userClassLoader,
+            String catalogName,
+            String defaultDatabase,
+            String baseUrl,
+            Properties connectionProperties) {
+        super(userClassLoader, catalogName, defaultDatabase, baseUrl, 
connectionProperties);
 
         String driverVersion =
                 Preconditions.checkNotNull(getDriverVersion(), "Driver version 
must not be null.");
@@ -122,7 +141,7 @@ public class MySqlCatalog extends AbstractJdbcCatalog {
     private String getDatabaseVersion() {
         try (TemporaryClassLoaderContext ignored =
                 TemporaryClassLoaderContext.of(userClassLoader)) {
-            try (Connection conn = DriverManager.getConnection(defaultUrl, 
username, pwd)) {
+            try (Connection conn = DriverManager.getConnection(defaultUrl, 
connectionProperties)) {
                 return conn.getMetaData().getDatabaseProductVersion();
             } catch (Exception e) {
                 throw new CatalogException(
@@ -134,7 +153,7 @@ public class MySqlCatalog extends AbstractJdbcCatalog {
     private String getDriverVersion() {
         try (TemporaryClassLoaderContext ignored =
                 TemporaryClassLoaderContext.of(userClassLoader)) {
-            try (Connection conn = DriverManager.getConnection(defaultUrl, 
username, pwd)) {
+            try (Connection conn = DriverManager.getConnection(defaultUrl, 
connectionProperties)) {
                 String driverVersion = conn.getMetaData().getDriverVersion();
                 Pattern regexp = Pattern.compile("\\d+?\\.\\d+?\\.\\d+");
                 Matcher matcher = regexp.matcher(driverVersion);
diff --git 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/databases/postgres/catalog/PostgresCatalog.java
 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/databases/postgres/catalog/PostgresCatalog.java
index 7ece9b64..27aa1f4d 100644
--- 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/databases/postgres/catalog/PostgresCatalog.java
+++ 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/databases/postgres/catalog/PostgresCatalog.java
@@ -19,6 +19,7 @@
 package org.apache.flink.connector.jdbc.databases.postgres.catalog;
 
 import org.apache.flink.annotation.Internal;
+import org.apache.flink.annotation.VisibleForTesting;
 import org.apache.flink.connector.jdbc.catalog.AbstractJdbcCatalog;
 import org.apache.flink.connector.jdbc.dialect.JdbcDialectTypeMapper;
 import org.apache.flink.table.catalog.ObjectPath;
@@ -39,8 +40,11 @@ import java.sql.ResultSetMetaData;
 import java.sql.SQLException;
 import java.util.HashSet;
 import java.util.List;
+import java.util.Properties;
 import java.util.Set;
 
+import static 
org.apache.flink.connector.jdbc.JdbcConnectionOptions.getBriefAuthProperties;
+
 /** Catalog for PostgreSQL. */
 @Internal
 public class PostgresCatalog extends AbstractJdbcCatalog {
@@ -72,6 +76,7 @@ public class PostgresCatalog extends AbstractJdbcCatalog {
 
     protected final JdbcDialectTypeMapper dialectTypeMapper;
 
+    @VisibleForTesting
     public PostgresCatalog(
             ClassLoader userClassLoader,
             String catalogName,
@@ -83,12 +88,26 @@ public class PostgresCatalog extends AbstractJdbcCatalog {
                 userClassLoader,
                 catalogName,
                 defaultDatabase,
-                username,
-                pwd,
                 baseUrl,
-                new PostgresTypeMapper());
+                getBriefAuthProperties(username, pwd));
+    }
+
+    public PostgresCatalog(
+            ClassLoader userClassLoader,
+            String catalogName,
+            String defaultDatabase,
+            String baseUrl,
+            Properties connectProperties) {
+        this(
+                userClassLoader,
+                catalogName,
+                defaultDatabase,
+                baseUrl,
+                new PostgresTypeMapper(),
+                connectProperties);
     }
 
+    @Deprecated
     protected PostgresCatalog(
             ClassLoader userClassLoader,
             String catalogName,
@@ -97,7 +116,23 @@ public class PostgresCatalog extends AbstractJdbcCatalog {
             String pwd,
             String baseUrl,
             JdbcDialectTypeMapper dialectTypeMapper) {
-        super(userClassLoader, catalogName, defaultDatabase, username, pwd, 
baseUrl);
+        this(
+                userClassLoader,
+                catalogName,
+                defaultDatabase,
+                baseUrl,
+                dialectTypeMapper,
+                getBriefAuthProperties(username, pwd));
+    }
+
+    protected PostgresCatalog(
+            ClassLoader userClassLoader,
+            String catalogName,
+            String defaultDatabase,
+            String baseUrl,
+            JdbcDialectTypeMapper dialectTypeMapper,
+            Properties connectProperties) {
+        super(userClassLoader, catalogName, defaultDatabase, baseUrl, 
connectProperties);
         this.dialectTypeMapper = dialectTypeMapper;
     }
 
@@ -153,7 +188,7 @@ public class PostgresCatalog extends AbstractJdbcCatalog {
         }
 
         final String url = baseUrl + databaseName;
-        try (Connection conn = DriverManager.getConnection(url, username, 
pwd)) {
+        try (Connection conn = DriverManager.getConnection(url, 
connectionProperties)) {
             // get all schemas
             List<String> schemas;
             try (PreparedStatement ps =
diff --git 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/utils/JdbcUtils.java
 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/utils/JdbcUtils.java
index 501a7a63..d86c7319 100644
--- 
a/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/utils/JdbcUtils.java
+++ 
b/flink-connector-jdbc/src/main/java/org/apache/flink/connector/jdbc/utils/JdbcUtils.java
@@ -18,19 +18,47 @@
 
 package org.apache.flink.connector.jdbc.utils;
 
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.ReadableConfig;
 import org.apache.flink.types.Row;
+import org.apache.flink.util.Preconditions;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.sql.PreparedStatement;
 import java.sql.SQLException;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Properties;
+
+import static 
org.apache.flink.connector.jdbc.JdbcConnectionOptions.getBriefAuthProperties;
+import static 
org.apache.flink.connector.jdbc.catalog.factory.JdbcCatalogFactoryOptions.PASSWORD;
+import static 
org.apache.flink.connector.jdbc.catalog.factory.JdbcCatalogFactoryOptions.USERNAME;
 
 /** Utils for jdbc connectors. */
 public class JdbcUtils {
 
     private static final Logger LOG = LoggerFactory.getLogger(JdbcUtils.class);
 
+    public static final String PROPERTIES_PREFIX = "connection-properties.";
+
+    public static Properties getConnectionProperties(ReadableConfig config) {
+        final Properties result =
+                getBriefAuthProperties(config.get(USERNAME), 
config.get(PASSWORD));
+        Preconditions.checkArgument(config instanceof Configuration);
+        Map<String, String> configMap = ((Configuration) config).toMap();
+        configMap.forEach(
+                (k, v) -> {
+                    if (Objects.nonNull(k)
+                            && Objects.nonNull(v)
+                            && k.startsWith(PROPERTIES_PREFIX)) {
+                        
result.setProperty(k.substring(PROPERTIES_PREFIX.length()), v);
+                    }
+                });
+        return result;
+    }
+
     /**
      * Adds a record to the prepared statement.
      *

Reply via email to