|
|
@ -17,6 +17,7 @@ |
|
|
|
|
|
|
|
|
|
|
|
package org.apache.dolphinscheduler.plugin.datasource.api.client; |
|
|
|
package org.apache.dolphinscheduler.plugin.datasource.api.client; |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
import org.apache.dolphinscheduler.plugin.datasource.api.provider.JDBCDataSourceProvider; |
|
|
|
import org.apache.dolphinscheduler.spi.datasource.BaseConnectionParam; |
|
|
|
import org.apache.dolphinscheduler.spi.datasource.BaseConnectionParam; |
|
|
|
import org.apache.dolphinscheduler.spi.datasource.DataSourceClient; |
|
|
|
import org.apache.dolphinscheduler.spi.datasource.DataSourceClient; |
|
|
|
import org.apache.dolphinscheduler.spi.enums.DbType; |
|
|
|
import org.apache.dolphinscheduler.spi.enums.DbType; |
|
|
@ -24,13 +25,15 @@ import org.apache.dolphinscheduler.spi.enums.DbType; |
|
|
|
import org.apache.commons.lang3.StringUtils; |
|
|
|
import org.apache.commons.lang3.StringUtils; |
|
|
|
|
|
|
|
|
|
|
|
import java.sql.Connection; |
|
|
|
import java.sql.Connection; |
|
|
|
import java.sql.DriverManager; |
|
|
|
|
|
|
|
import java.sql.SQLException; |
|
|
|
import java.sql.SQLException; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
|
|
|
|
|
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
import org.springframework.jdbc.core.JdbcTemplate; |
|
|
|
|
|
|
|
|
|
|
|
import com.google.common.base.Stopwatch; |
|
|
|
import com.google.common.base.Stopwatch; |
|
|
|
|
|
|
|
import com.zaxxer.hikari.HikariDataSource; |
|
|
|
|
|
|
|
|
|
|
|
@Slf4j |
|
|
|
@Slf4j |
|
|
|
public class CommonDataSourceClient implements DataSourceClient { |
|
|
|
public class CommonDataSourceClient implements DataSourceClient { |
|
|
@ -39,7 +42,8 @@ public class CommonDataSourceClient implements DataSourceClient { |
|
|
|
public static final String COMMON_VALIDATION_QUERY = "select 1"; |
|
|
|
public static final String COMMON_VALIDATION_QUERY = "select 1"; |
|
|
|
|
|
|
|
|
|
|
|
protected final BaseConnectionParam baseConnectionParam; |
|
|
|
protected final BaseConnectionParam baseConnectionParam; |
|
|
|
protected Connection connection; |
|
|
|
protected HikariDataSource dataSource; |
|
|
|
|
|
|
|
protected JdbcTemplate jdbcTemplate; |
|
|
|
|
|
|
|
|
|
|
|
public CommonDataSourceClient(BaseConnectionParam baseConnectionParam, DbType dbType) { |
|
|
|
public CommonDataSourceClient(BaseConnectionParam baseConnectionParam, DbType dbType) { |
|
|
|
this.baseConnectionParam = baseConnectionParam; |
|
|
|
this.baseConnectionParam = baseConnectionParam; |
|
|
@ -59,7 +63,8 @@ public class CommonDataSourceClient implements DataSourceClient { |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
protected void initClient(BaseConnectionParam baseConnectionParam, DbType dbType) { |
|
|
|
protected void initClient(BaseConnectionParam baseConnectionParam, DbType dbType) { |
|
|
|
this.connection = buildConn(baseConnectionParam); |
|
|
|
this.dataSource = JDBCDataSourceProvider.createJdbcDataSource(baseConnectionParam, dbType); |
|
|
|
|
|
|
|
this.jdbcTemplate = new JdbcTemplate(dataSource); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
protected void checkUser(BaseConnectionParam baseConnectionParam) { |
|
|
|
protected void checkUser(BaseConnectionParam baseConnectionParam) { |
|
|
@ -68,20 +73,6 @@ public class CommonDataSourceClient implements DataSourceClient { |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private Connection buildConn(BaseConnectionParam baseConnectionParam) { |
|
|
|
|
|
|
|
Connection conn = null; |
|
|
|
|
|
|
|
try { |
|
|
|
|
|
|
|
Class.forName(baseConnectionParam.getDriverClassName()); |
|
|
|
|
|
|
|
conn = DriverManager.getConnection(baseConnectionParam.getJdbcUrl(), baseConnectionParam.getUser(), |
|
|
|
|
|
|
|
baseConnectionParam.getPassword()); |
|
|
|
|
|
|
|
} catch (ClassNotFoundException e) { |
|
|
|
|
|
|
|
throw new RuntimeException("Driver load fail", e); |
|
|
|
|
|
|
|
} catch (SQLException e) { |
|
|
|
|
|
|
|
throw new RuntimeException("JDBC connect failed", e); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
return conn; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
protected void setDefaultUsername(BaseConnectionParam baseConnectionParam) { |
|
|
|
protected void setDefaultUsername(BaseConnectionParam baseConnectionParam) { |
|
|
|
baseConnectionParam.setUser(COMMON_USER); |
|
|
|
baseConnectionParam.setUser(COMMON_USER); |
|
|
|
} |
|
|
|
} |
|
|
@ -101,7 +92,7 @@ public class CommonDataSourceClient implements DataSourceClient { |
|
|
|
// Checking data source client
|
|
|
|
// Checking data source client
|
|
|
|
Stopwatch stopwatch = Stopwatch.createStarted(); |
|
|
|
Stopwatch stopwatch = Stopwatch.createStarted(); |
|
|
|
try { |
|
|
|
try { |
|
|
|
this.connection.prepareStatement(this.baseConnectionParam.getValidationQuery()).executeQuery(); |
|
|
|
this.jdbcTemplate.execute(this.baseConnectionParam.getValidationQuery()); |
|
|
|
} catch (Exception e) { |
|
|
|
} catch (Exception e) { |
|
|
|
throw new RuntimeException("JDBC connect failed", e); |
|
|
|
throw new RuntimeException("JDBC connect failed", e); |
|
|
|
} finally { |
|
|
|
} finally { |
|
|
@ -113,21 +104,20 @@ public class CommonDataSourceClient implements DataSourceClient { |
|
|
|
@Override |
|
|
|
@Override |
|
|
|
public Connection getConnection() { |
|
|
|
public Connection getConnection() { |
|
|
|
try { |
|
|
|
try { |
|
|
|
return connection.isClosed() ? buildConn(baseConnectionParam) : connection; |
|
|
|
return this.dataSource.getConnection(); |
|
|
|
} catch (SQLException e) { |
|
|
|
} catch (SQLException e) { |
|
|
|
throw new RuntimeException("get conn is fail", e); |
|
|
|
log.error("get druidDataSource Connection fail SQLException: {}", e.getMessage(), e); |
|
|
|
|
|
|
|
return null; |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
@Override |
|
|
|
public void close() { |
|
|
|
public void close() { |
|
|
|
log.info("do close connection {}.", baseConnectionParam.getDatabase()); |
|
|
|
log.info("do close dataSource {}.", baseConnectionParam.getDatabase()); |
|
|
|
try { |
|
|
|
try (HikariDataSource closedDatasource = dataSource) { |
|
|
|
connection.close(); |
|
|
|
// only close the resource
|
|
|
|
} catch (SQLException e) { |
|
|
|
|
|
|
|
log.info("colse connection fail"); |
|
|
|
|
|
|
|
throw new RuntimeException(e); |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
this.jdbcTemplate = null; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|