This is an automated email from the ASF dual-hosted git repository.
jiabaosun pushed a commit to branch release-3.1
in repository https://gitbox.apache.org/repos/asf/flink-cdc.git
The following commit(s) were added to refs/heads/release-3.1 by this push:
new e18e7a252 [FLINK-35295][mysql] Improve jdbc connection pool
initialization failure message
e18e7a252 is described below
commit e18e7a2523ac1ea59471e5714eb60f544e9f4a04
Author: yux <[email protected]>
AuthorDate: Fri May 31 11:12:25 2024 +0800
[FLINK-35295][mysql] Improve jdbc connection pool initialization failure
message
---
.../mysql/source/connection/PooledDataSourceFactory.java | 12 +++++++++++-
1 file changed, 11 insertions(+), 1 deletion(-)
diff --git
a/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/connection/PooledDataSourceFactory.java
b/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/connection/PooledDataSourceFactory.java
index 424a17b0f..14ec7cd8b 100644
---
a/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/connection/PooledDataSourceFactory.java
+++
b/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/connection/PooledDataSourceFactory.java
@@ -18,9 +18,11 @@
package org.apache.flink.cdc.connectors.mysql.source.connection;
import org.apache.flink.cdc.connectors.mysql.source.config.MySqlSourceConfig;
+import org.apache.flink.util.FlinkRuntimeException;
import com.zaxxer.hikari.HikariConfig;
import com.zaxxer.hikari.HikariDataSource;
+import com.zaxxer.hikari.pool.HikariPool;
import io.debezium.connector.mysql.MySqlConnectorConfig;
import java.util.Properties;
@@ -60,7 +62,15 @@ public class PooledDataSourceFactory {
config.addDataSourceProperty("prepStmtCacheSize", "250");
config.addDataSourceProperty("prepStmtCacheSqlLimit", "2048");
- return new HikariDataSource(config);
+ try {
+ return new HikariDataSource(config);
+ } catch (HikariPool.PoolInitializationException e) {
+ throw new FlinkRuntimeException(
+ "Initialize jdbc connection pool failed, this may caused
by"
+ + " wrong jdbc configurations or unstable network."
+ + " Please check your jdbc configurations and
network.",
+ e);
+ }
}
private static String formatJdbcUrl(String hostName, int port, Properties
jdbcProperties) {