diff --git a/docs/configs/docsdev.js b/docs/configs/docsdev.js
index 323e8b21d9..233f6d20e7 100644
--- a/docs/configs/docsdev.js
+++ b/docs/configs/docsdev.js
@@ -366,6 +366,10 @@ export default {
{
title: 'doris',
link: '/en-us/docs/dev/user_doc/guide/datasource/doris.html',
+ },
+ {
+ title: 'flink',
+ link: '/en-us/docs/dev/user_doc/guide/datasource/flink.html',
}
],
},
@@ -1086,6 +1090,10 @@ export default {
{
title: 'Doris',
link: '/zh-cn/docs/dev/user_doc/guide/datasource/doris.html',
+ },
+ {
+ title: 'Flink',
+ link: '/zh-cn/docs/dev/user_doc/guide/datasource/flink.html',
}
],
},
diff --git a/docs/docs/en/guide/datasource/flink.md b/docs/docs/en/guide/datasource/flink.md
new file mode 100644
index 0000000000..e69de29bb2
diff --git a/docs/docs/zh/guide/datasource/flink.md b/docs/docs/zh/guide/datasource/flink.md
new file mode 100644
index 0000000000..e69de29bb2
diff --git a/dolphinscheduler-bom/pom.xml b/dolphinscheduler-bom/pom.xml
index 18e75de57b..d1561f0901 100644
--- a/dolphinscheduler-bom/pom.xml
+++ b/dolphinscheduler-bom/pom.xml
@@ -124,6 +124,7 @@
0.10.1
2.1.4
0.3.2
+ 1.18.1
@@ -578,6 +579,12 @@
${commons-io.version}
+
+ org.apache.flink
+ flink-sql-jdbc-driver-bundle
+ ${flink-jdbc.version}
+
+
com.github.oshi
oshi-core
diff --git a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/constants/DataSourceConstants.java b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/constants/DataSourceConstants.java
index 11347942bd..97dc2235e8 100644
--- a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/constants/DataSourceConstants.java
+++ b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/constants/DataSourceConstants.java
@@ -39,6 +39,7 @@ public class DataSourceConstants {
public static final String COM_TRINO_JDBC_DRIVER = "io.trino.jdbc.TrinoDriver";
public static final String COM_DAMENG_JDBC_DRIVER = "dm.jdbc.driver.DmDriver";
public static final String ORG_APACHE_KYUUBI_JDBC_DRIVER = "org.apache.kyuubi.jdbc.KyuubiHiveDriver";
+ public static final String ORG_APACHE_FLINK_JDBC_DRIVER = "org.apache.flink.table.jdbc.FlinkDriver";
public static final String COM_OCEANBASE_JDBC_DRIVER = "com.oceanbase.jdbc.Driver";
public static final String NET_SNOWFLAKE_JDBC_DRIVER = "net.snowflake.client.jdbc.SnowflakeDriver";
public static final String COM_VERTICA_JDBC_DRIVER = "com.vertica.jdbc.Driver";
@@ -63,6 +64,7 @@ public class DataSourceConstants {
public static final String SNOWFLAKE_VALIDATION_QUERY = "select 1";
public static final String KYUUBI_VALIDATION_QUERY = "select 1";
+ public static final String FLINK_VALIDATION_QUERY = "select 1";
public static final String VERTICA_VALIDATION_QUERY = "select 1";
public static final String HANA_VALIDATION_QUERY = "select 1 from DUMMY";
@@ -91,6 +93,7 @@ public class DataSourceConstants {
public static final String JDBC_SNOWFLAKE = "jdbc:snowflake://";
public static final String JDBC_VERTICA = "jdbc:vertica://";
public static final String JDBC_HANA = "jdbc:sap://";
+ public static final String JDBC_FLINK = "jdbc:flink://";
/**
* database type
diff --git a/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-all/pom.xml b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-all/pom.xml
index effe3c9abb..b60c4ff31e 100644
--- a/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-all/pom.xml
+++ b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-all/pom.xml
@@ -158,5 +158,10 @@
dolphinscheduler-datasource-hana
${project.version}
+
+ org.apache.dolphinscheduler
+ dolphinscheduler-datasource-flink
+ ${project.version}
+
diff --git a/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/pom.xml b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/pom.xml
new file mode 100644
index 0000000000..b743cae8d4
--- /dev/null
+++ b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/pom.xml
@@ -0,0 +1,61 @@
+
+
+
+ 4.0.0
+
+ org.apache.dolphinscheduler
+ dolphinscheduler-datasource-plugin
+ dev-SNAPSHOT
+
+
+ dolphinscheduler-datasource-flink
+ jar
+ ${project.artifactId}
+
+
+
+ org.apache.dolphinscheduler
+ dolphinscheduler-spi
+ provided
+
+
+ org.apache.dolphinscheduler
+ dolphinscheduler-task-api
+ provided
+
+
+
+ org.apache.dolphinscheduler
+ dolphinscheduler-datasource-hive
+ ${project.version}
+ provided
+
+
+
+ org.apache.dolphinscheduler
+ dolphinscheduler-datasource-api
+ ${project.version}
+
+
+
+ org.apache.flink
+ flink-sql-jdbc-driver-bundle
+
+
+
diff --git a/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/FlinkAdHocDataSourceClient.java b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/FlinkAdHocDataSourceClient.java
new file mode 100644
index 0000000000..66160066e3
--- /dev/null
+++ b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/FlinkAdHocDataSourceClient.java
@@ -0,0 +1,12 @@
+package org.apache.dolphinscheduler.plugin.datasource.flink;
+
+import org.apache.dolphinscheduler.plugin.datasource.api.client.BaseAdHocDataSourceClient;
+import org.apache.dolphinscheduler.spi.datasource.BaseConnectionParam;
+import org.apache.dolphinscheduler.spi.enums.DbType;
+
+public class FlinkAdHocDataSourceClient extends BaseAdHocDataSourceClient {
+
+ public FlinkAdHocDataSourceClient(BaseConnectionParam baseConnectionParam, DbType dbType) {
+ super(baseConnectionParam, dbType);
+ }
+}
diff --git a/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/FlinkDataSourceChannel.java b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/FlinkDataSourceChannel.java
new file mode 100644
index 0000000000..657b0113c3
--- /dev/null
+++ b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/FlinkDataSourceChannel.java
@@ -0,0 +1,20 @@
+package org.apache.dolphinscheduler.plugin.datasource.flink;
+
+import org.apache.dolphinscheduler.spi.datasource.AdHocDataSourceClient;
+import org.apache.dolphinscheduler.spi.datasource.BaseConnectionParam;
+import org.apache.dolphinscheduler.spi.datasource.DataSourceChannel;
+import org.apache.dolphinscheduler.spi.datasource.PooledDataSourceClient;
+import org.apache.dolphinscheduler.spi.enums.DbType;
+
+public class FlinkDataSourceChannel implements DataSourceChannel {
+
+ @Override
+ public AdHocDataSourceClient createAdHocDataSourceClient(BaseConnectionParam baseConnectionParam, DbType dbType) {
+ return new FlinkAdHocDataSourceClient(baseConnectionParam, dbType);
+ }
+
+ @Override
+ public PooledDataSourceClient createPooledDataSourceClient(BaseConnectionParam baseConnectionParam, DbType dbType) {
+ return new FlinkPooledDataSourceClient(baseConnectionParam, dbType);
+ }
+}
\ No newline at end of file
diff --git a/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/FlinkDataSourceChannelFactory.java b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/FlinkDataSourceChannelFactory.java
new file mode 100644
index 0000000000..4372bb2039
--- /dev/null
+++ b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/FlinkDataSourceChannelFactory.java
@@ -0,0 +1,19 @@
+package org.apache.dolphinscheduler.plugin.datasource.flink;
+
+import com.google.auto.service.AutoService;
+import org.apache.dolphinscheduler.spi.datasource.DataSourceChannel;
+import org.apache.dolphinscheduler.spi.datasource.DataSourceChannelFactory;
+
+@AutoService(DataSourceChannelFactory.class)
+public class FlinkDataSourceChannelFactory implements DataSourceChannelFactory {
+
+ @Override
+ public String getName() {
+ return "flink";
+ }
+
+ @Override
+ public DataSourceChannel create() {
+ return new FlinkDataSourceChannel();
+ }
+}
\ No newline at end of file
diff --git a/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/FlinkPooledDataSourceClient.java b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/FlinkPooledDataSourceClient.java
new file mode 100644
index 0000000000..1820d3e608
--- /dev/null
+++ b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/FlinkPooledDataSourceClient.java
@@ -0,0 +1,28 @@
+package org.apache.dolphinscheduler.plugin.datasource.flink;
+
+import lombok.extern.slf4j.Slf4j;
+import org.apache.dolphinscheduler.plugin.datasource.api.client.BasePooledDataSourceClient;
+import org.apache.dolphinscheduler.spi.datasource.BaseConnectionParam;
+import org.apache.dolphinscheduler.spi.enums.DbType;
+
+import java.sql.Connection;
+import java.sql.SQLException;
+
+@Slf4j
+public class FlinkPooledDataSourceClient extends BasePooledDataSourceClient {
+
+ public FlinkPooledDataSourceClient(BaseConnectionParam baseConnectionParam, DbType dbType) {
+ super(baseConnectionParam, dbType);
+ }
+
+ @Override
+ public Connection getConnection() throws SQLException {
+ return dataSource.getConnection();
+ }
+
+ @Override
+ public void close() {
+ super.close();
+ log.info("Closed Flink datasource client.");
+ }
+}
\ No newline at end of file
diff --git a/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/param/FlinkConnectionParam.java b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/param/FlinkConnectionParam.java
new file mode 100644
index 0000000000..11ff12a621
--- /dev/null
+++ b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/param/FlinkConnectionParam.java
@@ -0,0 +1,21 @@
+package org.apache.dolphinscheduler.plugin.datasource.flink.param;
+
+import org.apache.dolphinscheduler.spi.datasource.BaseConnectionParam;
+
+public class FlinkConnectionParam extends BaseConnectionParam {
+
+ @Override
+ public String toString() {
+ return "FlinkConnectionParam{"
+ + "user='" + user + '\''
+ + ", password='" + password + '\''
+ + ", address='" + address + '\''
+ + ", database='" + database + '\''
+ + ", jdbcUrl='" + jdbcUrl + '\''
+ + ", driverLocation='" + driverLocation + '\''
+ + ", driverClassName='" + driverClassName + '\''
+ + ", validationQuery='" + validationQuery + '\''
+ + ", other='" + other + '\''
+ + '}';
+ }
+}
diff --git a/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/param/FlinkDataSourceParamDTO.java b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/param/FlinkDataSourceParamDTO.java
new file mode 100644
index 0000000000..2382148510
--- /dev/null
+++ b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/param/FlinkDataSourceParamDTO.java
@@ -0,0 +1,24 @@
+package org.apache.dolphinscheduler.plugin.datasource.flink.param;
+
+import org.apache.dolphinscheduler.plugin.datasource.api.datasource.BaseDataSourceParamDTO;
+import org.apache.dolphinscheduler.spi.enums.DbType;
+
+public class FlinkDataSourceParamDTO extends BaseDataSourceParamDTO {
+
+ @Override
+ public String toString() {
+ return "FlinkDataSourceParamDTO{"
+ + "host='" + host + '\''
+ + ", port=" + port
+ + ", database='" + database + '\''
+ + ", userName='" + userName + '\''
+ + ", password='" + password + '\''
+ + ", other='" + other + '\''
+ + '}';
+ }
+
+ @Override
+ public DbType getType() {
+ return DbType.FLINK;
+ }
+}
diff --git a/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/param/FlinkDataSourceProcessor.java b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/param/FlinkDataSourceProcessor.java
new file mode 100644
index 0000000000..88e77483c0
--- /dev/null
+++ b/dolphinscheduler-datasource-plugin/dolphinscheduler-datasource-flink/src/main/java/org/apache/dolphinscheduler/plugin/datasource/flink/param/FlinkDataSourceProcessor.java
@@ -0,0 +1,152 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.dolphinscheduler.plugin.datasource.flink.param;
+
+import org.apache.dolphinscheduler.common.constants.Constants;
+import org.apache.dolphinscheduler.common.constants.DataSourceConstants;
+import org.apache.dolphinscheduler.common.utils.JSONUtils;
+import org.apache.dolphinscheduler.plugin.datasource.api.datasource.AbstractDataSourceProcessor;
+import org.apache.dolphinscheduler.plugin.datasource.api.datasource.BaseDataSourceParamDTO;
+import org.apache.dolphinscheduler.plugin.datasource.api.datasource.DataSourceProcessor;
+import org.apache.dolphinscheduler.plugin.datasource.api.utils.PasswordUtils;
+import org.apache.dolphinscheduler.spi.datasource.BaseConnectionParam;
+import org.apache.dolphinscheduler.spi.datasource.ConnectionParam;
+import org.apache.dolphinscheduler.spi.enums.DbType;
+
+import org.apache.commons.collections4.MapUtils;
+
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.SQLException;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+
+import com.google.auto.service.AutoService;
+
+@AutoService(DataSourceProcessor.class)
+public class FlinkDataSourceProcessor extends AbstractDataSourceProcessor {
+
+ @Override
+ public void checkDatasourceParam(BaseDataSourceParamDTO baseDataSourceParamDTO) {
+ }
+
+ @Override
+ public BaseDataSourceParamDTO castDatasourceParamDTO(String paramJson) {
+ return JSONUtils.parseObject(paramJson, FlinkDataSourceParamDTO.class);
+ }
+
+ @Override
+ public BaseDataSourceParamDTO createDatasourceParamDTO(String connectionJson) {
+ FlinkDataSourceParamDTO flinkDataSourceParamDTO = new FlinkDataSourceParamDTO();
+ FlinkConnectionParam kyuubiConnectionParam = (FlinkConnectionParam) createConnectionParams(connectionJson);
+ flinkDataSourceParamDTO.setDatabase(kyuubiConnectionParam.getDatabase());
+ flinkDataSourceParamDTO.setUserName(kyuubiConnectionParam.getUser());
+ flinkDataSourceParamDTO.setOther(kyuubiConnectionParam.getOther());
+
+ String[] tmpArray = kyuubiConnectionParam.getAddress().split(Constants.DOUBLE_SLASH);
+ StringBuilder hosts = new StringBuilder();
+ String[] hostPortArray = tmpArray[tmpArray.length - 1].split(Constants.COMMA);
+ for (String hostPort : hostPortArray) {
+ hosts.append(hostPort.split(Constants.COLON)[0]).append(Constants.COMMA);
+ }
+ hosts.deleteCharAt(hosts.length() - 1);
+ flinkDataSourceParamDTO.setHost(hosts.toString());
+ flinkDataSourceParamDTO.setPort(Integer.parseInt(hostPortArray[0].split(Constants.COLON)[1]));
+
+ return flinkDataSourceParamDTO;
+ }
+
+ @Override
+ public BaseConnectionParam createConnectionParams(BaseDataSourceParamDTO datasourceParam) {
+ FlinkDataSourceParamDTO flinkParam = (FlinkDataSourceParamDTO) datasourceParam;
+ StringBuilder address = new StringBuilder();
+ address.append(DataSourceConstants.JDBC_FLINK);
+ for (String zkHost : flinkParam.getHost().split(",")) {
+ address.append(String.format("%s:%s,", zkHost, flinkParam.getPort()));
+ }
+ address.deleteCharAt(address.length() - 1);
+ String jdbcUrl = address + "/" + flinkParam.getDatabase();
+ FlinkConnectionParam kyuubiConnectionParam = new FlinkConnectionParam();
+ kyuubiConnectionParam.setDatabase(flinkParam.getDatabase());
+ kyuubiConnectionParam.setAddress(address.toString());
+ kyuubiConnectionParam.setJdbcUrl(jdbcUrl);
+ kyuubiConnectionParam.setUser(flinkParam.getUserName());
+ kyuubiConnectionParam.setPassword(PasswordUtils.encodePassword(flinkParam.getPassword()));
+ kyuubiConnectionParam.setDriverClassName(getDatasourceDriver());
+ kyuubiConnectionParam.setValidationQuery(getValidationQuery());
+ kyuubiConnectionParam.setOther(flinkParam.getOther());
+ return kyuubiConnectionParam;
+ }
+
+ @Override
+ public ConnectionParam createConnectionParams(String connectionJson) {
+ return JSONUtils.parseObject(connectionJson, FlinkConnectionParam.class);
+
+ }
+
+ @Override
+ public String getDatasourceDriver() {
+ return DataSourceConstants.ORG_APACHE_FLINK_JDBC_DRIVER;
+ }
+
+ @Override
+ public String getValidationQuery() {
+ return DataSourceConstants.FLINK_VALIDATION_QUERY;
+ }
+
+ @Override
+ public String getJdbcUrl(ConnectionParam connectionParam) {
+ FlinkConnectionParam flinkConnectionParam = (FlinkConnectionParam) connectionParam;
+ String jdbcUrl = flinkConnectionParam.getJdbcUrl();
+
+ if (MapUtils.isNotEmpty(flinkConnectionParam.getOther())) {
+ return jdbcUrl + "?" + transformOther(flinkConnectionParam.getOther());
+ }
+ return jdbcUrl;
+ }
+
+ @Override
+ public Connection getConnection(ConnectionParam connectionParam) throws ClassNotFoundException, SQLException {
+ FlinkConnectionParam flinkConnectionParam = (FlinkConnectionParam) connectionParam;
+ Class.forName(getDatasourceDriver());
+ // todo:
+ return DriverManager.getConnection(getJdbcUrl(connectionParam),
+ flinkConnectionParam.getUser(), PasswordUtils.decodePassword(flinkConnectionParam.getPassword()));
+ }
+
+ @Override
+ public DbType getDbType() {
+ return DbType.FLINK;
+ }
+
+ @Override
+ public DataSourceProcessor create() {
+ return new FlinkDataSourceProcessor();
+ }
+
+ private String transformOther(Map otherMap) {
+ if (MapUtils.isEmpty(otherMap)) {
+ return null;
+ }
+ List otherList = new ArrayList<>();
+ otherMap.forEach((key, value) -> otherList.add(String.format("%s=%s", key, value)));
+ return String.join(";", otherList);
+ }
+
+}
\ No newline at end of file
diff --git a/dolphinscheduler-datasource-plugin/pom.xml b/dolphinscheduler-datasource-plugin/pom.xml
index c30a6b4258..6450ae3c14 100644
--- a/dolphinscheduler-datasource-plugin/pom.xml
+++ b/dolphinscheduler-datasource-plugin/pom.xml
@@ -56,6 +56,7 @@
dolphinscheduler-datasource-sagemaker
dolphinscheduler-datasource-k8s
dolphinscheduler-datasource-hana
+ dolphinscheduler-datasource-flink
diff --git a/dolphinscheduler-spi/src/main/java/org/apache/dolphinscheduler/spi/enums/DbType.java b/dolphinscheduler-spi/src/main/java/org/apache/dolphinscheduler/spi/enums/DbType.java
index e7ebbeee0a..61f458c287 100644
--- a/dolphinscheduler-spi/src/main/java/org/apache/dolphinscheduler/spi/enums/DbType.java
+++ b/dolphinscheduler-spi/src/main/java/org/apache/dolphinscheduler/spi/enums/DbType.java
@@ -55,7 +55,8 @@ public enum DbType {
ZEPPELIN(24, "zeppelin"),
SAGEMAKER(25, "sagemaker"),
- K8S(26, "k8s");
+ K8S(26, "k8s"),
+ FLINK(27, "flink");
private static final Map DB_TYPE_MAP =
Arrays.stream(DbType.values()).collect(toMap(DbType::getCode, Functions.identity()));
@EnumValue
diff --git a/dolphinscheduler-task-plugin/dolphinscheduler-task-sql/src/main/java/org/apache/dolphinscheduler/plugin/task/sql/SqlTask.java b/dolphinscheduler-task-plugin/dolphinscheduler-task-sql/src/main/java/org/apache/dolphinscheduler/plugin/task/sql/SqlTask.java
index 1c8fd4519f..e70a4311ce 100644
--- a/dolphinscheduler-task-plugin/dolphinscheduler-task-sql/src/main/java/org/apache/dolphinscheduler/plugin/task/sql/SqlTask.java
+++ b/dolphinscheduler-task-plugin/dolphinscheduler-task-sql/src/main/java/org/apache/dolphinscheduler/plugin/task/sql/SqlTask.java
@@ -188,6 +188,22 @@ public class SqlTask extends AbstractTask {
Connection connection =
DataSourceClientProvider.getAdHocConnection(DbType.valueOf(sqlParameters.getType()),
baseConnectionParam)) {
+ // main execute
+ String result = null;
+ if (sqlParameters.getType().equals("FLINK")) {
+ // pre execute
+ executeFlinkSQLQuery(connection, preStatementsBinds, "pre");
+
+ // query statements need to be convert to JsonArray and inserted into Alert to send
+ result = executeFlinkSQLQuery(connection, mainStatementsBinds.get(0), "main");
+
+ sqlParameters.dealOutParam(result);
+
+ // post execute
+ executeFlinkSQLQuery(connection, postStatementsBinds, "post");
+
+ return;
+ }
// create temp function
if (CollectionUtils.isNotEmpty(createFuncs)) {
@@ -197,8 +213,6 @@ public class SqlTask extends AbstractTask {
// pre execute
executeUpdate(connection, preStatementsBinds, "pre");
- // main execute
- String result = null;
// decide whether to executeQuery or executeUpdate based on sqlType
if (sqlParameters.getSqlType() == SqlType.QUERY.ordinal()) {
// query statements need to be convert to JsonArray and inserted into Alert to send
@@ -208,6 +222,7 @@ public class SqlTask extends AbstractTask {
String updateResult = executeUpdate(connection, mainStatementsBinds, "main");
result = setNonQuerySqlReturn(updateResult, sqlParameters.getLocalParams());
}
+
// deal out params
sqlParameters.dealOutParam(result);
@@ -319,6 +334,25 @@ public class SqlTask extends AbstractTask {
}
}
+ // see https://nightlies.apache.org/flink/flink-docs-master/docs/dev/table/jdbcdriver/#java
+ private String executeFlinkSQLQuery(Connection connection, SqlBinds sqlBinds, String handlerType) throws Exception {
+ try (Statement statement = connection.createStatement()) {
+ log.info("{} statement execute flink sql query, for sql: {}", handlerType, sqlBinds.getSql());
+ ResultSet resultSet = statement.executeQuery(sqlBinds.getSql());
+ return resultProcess(resultSet);
+ }
+ }
+
+ private String executeFlinkSQLQuery(Connection connection, List statementsBinds, String handlerType) throws Exception {
+ String result = "";
+ for (SqlBinds sqlBind : statementsBinds) {
+ result = executeFlinkSQLQuery(connection, sqlBind, handlerType);
+ log.info("{} statement execute update result: {}, for sql: {}", handlerType, result,
+ sqlBind.getSql());
+ }
+ return String.valueOf(result);
+ }
+
private String executeUpdate(Connection connection, List statementsBinds,
String handlerType) throws Exception {
int result = 0;
diff --git a/dolphinscheduler-ui/src/service/modules/data-source/types.ts b/dolphinscheduler-ui/src/service/modules/data-source/types.ts
index 444f5293dd..cb3b4a59c9 100644
--- a/dolphinscheduler-ui/src/service/modules/data-source/types.ts
+++ b/dolphinscheduler-ui/src/service/modules/data-source/types.ts
@@ -42,6 +42,7 @@ type IDataBase =
| 'ZEPPELIN'
| 'SAGEMAKER'
| 'K8S'
+ | 'FLINK'
type IDataBaseLabel =
| 'MYSQL'
@@ -65,6 +66,7 @@ type IDataBaseLabel =
| 'ZEPPELIN'
| 'SAGEMAKER'
| 'K8S'
+ | 'FLINK'
interface IDataSource {
id?: number
diff --git a/dolphinscheduler-ui/src/views/datasource/list/detail.tsx b/dolphinscheduler-ui/src/views/datasource/list/detail.tsx
index 4842651290..0aed3487ba 100644
--- a/dolphinscheduler-ui/src/views/datasource/list/detail.tsx
+++ b/dolphinscheduler-ui/src/views/datasource/list/detail.tsx
@@ -546,7 +546,8 @@ const DetailModal = defineComponent({