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. *
