This is an automated email from the ASF dual-hosted git repository.
zihaoxiang pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/dolphinscheduler.git
The following commit(s) were added to refs/heads/dev by this push:
new 7467cb2467 [Improvement-18019][task-sql] Support SQL from resource
file and parameter placeholders (#18020)
7467cb2467 is described below
commit 7467cb24678cacb76762e268aa526be7b5cec032
Author: macdoor <[email protected]>
AuthorDate: Wed Apr 8 16:26:37 2026 +0800
[Improvement-18019][task-sql] Support SQL from resource file and parameter
placeholders (#18020)
---
.../plugin/task/api/enums/SqlSourceType.java | 26 +++
.../plugin/task/api/parameters/SqlParameters.java | 45 +++++-
.../task/api/parameters/SqlParametersTest.java | 57 +++++++
.../dolphinscheduler/plugin/task/sql/SqlTask.java | 51 +++++-
.../plugin/task/sql/SqlTaskTest.java | 176 +++++++++++++++++++--
dolphinscheduler-ui/src/locales/en_US/project.ts | 4 +
dolphinscheduler-ui/src/locales/zh_CN/project.ts | 4 +
.../task/components/node/fields/use-sql.ts | 71 ++++++++-
.../projects/task/components/node/tasks/use-sql.ts | 2 +
9 files changed, 409 insertions(+), 27 deletions(-)
diff --git
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/enums/SqlSourceType.java
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/enums/SqlSourceType.java
new file mode 100644
index 0000000000..3c812049bb
--- /dev/null
+++
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/enums/SqlSourceType.java
@@ -0,0 +1,26 @@
+/*
+ * 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.task.api.enums;
+
+public enum SqlSourceType {
+ /**
+ * SCRIPT: inline sql text
+ * FILE: sql from resource center file
+ */
+ SCRIPT, FILE
+}
diff --git
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/parameters/SqlParameters.java
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/parameters/SqlParameters.java
index 8864943bac..c614bfc15e 100644
---
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/parameters/SqlParameters.java
+++
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/parameters/SqlParameters.java
@@ -21,6 +21,7 @@ import org.apache.dolphinscheduler.common.utils.JSONUtils;
import org.apache.dolphinscheduler.plugin.task.api.SQLTaskExecutionContext;
import org.apache.dolphinscheduler.plugin.task.api.enums.DataType;
import org.apache.dolphinscheduler.plugin.task.api.enums.ResourceType;
+import org.apache.dolphinscheduler.plugin.task.api.enums.SqlSourceType;
import org.apache.dolphinscheduler.plugin.task.api.model.Property;
import org.apache.dolphinscheduler.plugin.task.api.model.ResourceInfo;
import
org.apache.dolphinscheduler.plugin.task.api.parameters.resource.DataSourceParameters;
@@ -58,6 +59,16 @@ public class SqlParameters extends AbstractParameters {
*/
private String sql;
+ /**
+ * sql source
+ */
+ private SqlSourceType sqlSource;
+
+ /**
+ * sql resource file path in resource center
+ */
+ private String sqlResource;
+
/**
* sql type
* 0 query
@@ -139,6 +150,22 @@ public class SqlParameters extends AbstractParameters {
this.sql = sql;
}
+ public SqlSourceType getSqlSource() {
+ return sqlSource;
+ }
+
+ public void setSqlSource(SqlSourceType sqlSource) {
+ this.sqlSource = sqlSource;
+ }
+
+ public String getSqlResource() {
+ return sqlResource;
+ }
+
+ public void setSqlResource(String sqlResource) {
+ this.sqlResource = sqlResource;
+ }
+
public int getSqlType() {
return sqlType;
}
@@ -213,12 +240,24 @@ public class SqlParameters extends AbstractParameters {
@Override
public boolean checkParameters() {
- return datasource != 0 && StringUtils.isNotEmpty(type) &&
StringUtils.isNotEmpty(sql);
+ if (datasource == 0 || StringUtils.isEmpty(type)) {
+ return false;
+ }
+ if (StringUtils.isNotEmpty(sql)) {
+ return true;
+ }
+ return StringUtils.isNotEmpty(sqlResource);
}
@Override
public List<ResourceInfo> getResourceFilesList() {
- return new ArrayList<>();
+ List<ResourceInfo> resourceFiles = new ArrayList<>();
+ if (StringUtils.isNotEmpty(sqlResource)) {
+ ResourceInfo resourceInfo = new ResourceInfo();
+ resourceInfo.setResourceName(sqlResource);
+ resourceFiles.add(resourceInfo);
+ }
+ return resourceFiles;
}
public void dealOutParam(String result) {
@@ -272,6 +311,8 @@ public class SqlParameters extends AbstractParameters {
+ "type='" + type + '\''
+ ", datasource=" + datasource
+ ", sql='" + sql + '\''
+ + ", sqlSource='" + sqlSource + '\''
+ + ", sqlResource='" + sqlResource + '\''
+ ", sqlType=" + sqlType
+ ", sendEmail=" + sendEmail
+ ", displayRows=" + displayRows
diff --git
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/test/java/org/apache/dolphinscheduler/plugin/task/api/parameters/SqlParametersTest.java
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/test/java/org/apache/dolphinscheduler/plugin/task/api/parameters/SqlParametersTest.java
index 8f1ee76560..b9388b3ba5 100644
---
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/test/java/org/apache/dolphinscheduler/plugin/task/api/parameters/SqlParametersTest.java
+++
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/test/java/org/apache/dolphinscheduler/plugin/task/api/parameters/SqlParametersTest.java
@@ -17,9 +17,14 @@
package org.apache.dolphinscheduler.plugin.task.api.parameters;
+import org.apache.dolphinscheduler.plugin.task.api.SQLTaskExecutionContext;
import org.apache.dolphinscheduler.plugin.task.api.enums.DataType;
import org.apache.dolphinscheduler.plugin.task.api.enums.Direct;
+import org.apache.dolphinscheduler.plugin.task.api.enums.ResourceType;
import org.apache.dolphinscheduler.plugin.task.api.model.Property;
+import
org.apache.dolphinscheduler.plugin.task.api.parameters.resource.DataSourceParameters;
+import
org.apache.dolphinscheduler.plugin.task.api.parameters.resource.ResourceParametersHelper;
+import org.apache.dolphinscheduler.spi.enums.DbType;
import org.apache.commons.collections4.CollectionUtils;
@@ -87,5 +92,57 @@ public class SqlParametersTest {
sqlParameters.setLocalParams(properties);
sqlParameters.dealOutParam(sqlResult);
Assertions.assertNotNull(sqlParameters.getVarPool().get(0));
+
+ // resource files list should contain sqlResource when it is set
+ sqlParameters.setSql(null);
+ sqlParameters.setSqlResource("/sql/demo.sql");
+
Assertions.assertFalse(CollectionUtils.isEmpty(sqlParameters.getResourceFilesList()));
+ }
+
+ @Test
+ public void testCheckParameters_variants() {
+ SqlParameters p = new SqlParameters();
+
+ // datasource/type invalid
+ p.setDatasource(0);
+ p.setType("");
+ p.setSql("select 1");
+ p.setSqlResource(null);
+ Assertions.assertFalse(p.checkParameters());
+
+ // sql present -> true
+ p.setDatasource(1);
+ p.setType("MYSQL");
+ p.setSql("select 1");
+ p.setSqlResource(null);
+ Assertions.assertTrue(p.checkParameters());
+
+ // sql absent, sqlResource present -> true
+ p.setSql(null);
+ p.setSqlResource("/sql/demo.sql");
+ Assertions.assertTrue(p.checkParameters());
+
+ // both absent -> false
+ p.setSql(null);
+ p.setSqlResource(null);
+ Assertions.assertFalse(p.checkParameters());
+ }
+
+ @Test
+ public void testGenerateExtendedContext_setsConnectionParams() {
+ SqlParameters p = new SqlParameters();
+ p.setDatasource(1);
+
+ DataSourceParameters dataSourceParameters = new DataSourceParameters();
+ dataSourceParameters.setType(DbType.MYSQL);
+ dataSourceParameters.setResourceType(ResourceType.DATASOURCE.name());
+ dataSourceParameters.setConnectionParams("conn_params");
+
+ ResourceParametersHelper helper = new ResourceParametersHelper();
+ helper.put(ResourceType.DATASOURCE, 1, dataSourceParameters);
+
+ SQLTaskExecutionContext ctx = p.generateExtendedContext(helper);
+ Assertions.assertNotNull(ctx);
+ Assertions.assertEquals("conn_params", ctx.getConnectionParams());
}
}
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 33b85e1a66..113dfbcd9e 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
@@ -35,12 +35,17 @@ import
org.apache.dolphinscheduler.plugin.task.api.model.Property;
import org.apache.dolphinscheduler.plugin.task.api.model.TaskAlertInfo;
import
org.apache.dolphinscheduler.plugin.task.api.parameters.AbstractParameters;
import org.apache.dolphinscheduler.plugin.task.api.parameters.SqlParameters;
+import org.apache.dolphinscheduler.plugin.task.api.resource.ResourceContext;
import org.apache.dolphinscheduler.plugin.task.api.utils.ParameterUtils;
import org.apache.dolphinscheduler.spi.datasource.BaseConnectionParam;
import org.apache.dolphinscheduler.spi.enums.DbType;
import org.apache.commons.lang3.StringUtils;
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Paths;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
@@ -48,6 +53,7 @@ import java.sql.ResultSetMetaData;
import java.sql.SQLException;
import java.sql.Statement;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
@@ -128,6 +134,8 @@ public class SqlTask extends AbstractTask {
sqlParameters.getLimit());
try {
+ ensureSqlContent();
+
// get datasource
baseConnectionParam = (BaseConnectionParam)
DataSourceUtils.buildConnectionParams(dbType,
sqlTaskExecutionContext.getConnectionParams());
@@ -405,6 +413,32 @@ public class SqlTask extends AbstractTask {
log.info("Sql Params are {}", logPrint);
}
+ private void ensureSqlContent() {
+ if (StringUtils.isNotEmpty(sqlParameters.getSql())) {
+ return;
+ }
+ if (StringUtils.isEmpty(sqlParameters.getSqlResource())) {
+ return;
+ }
+ String resourcePathInStorage = sqlParameters.getSqlResource();
+ try {
+ ResourceContext resourceContext =
taskExecutionContext.getResourceContext();
+ ResourceContext.ResourceItem resourceItem =
+ resourceContext.getResourceItem(resourcePathInStorage);
+ String localPath = resourceItem.getResourceAbsolutePathInLocal();
+ log.info("Load sql content from resource file: {}",
resourcePathInStorage);
+ String sqlContent = new String(
+ Files.readAllBytes(Paths.get(localPath)),
+ StandardCharsets.UTF_8);
+ sqlParameters.setSql(sqlContent);
+ } catch (IOException e) {
+ log.error("Read sql content from resource file {} error",
resourcePathInStorage, e);
+ throw new TaskException(
+ String.format("Read sql content from resource file %s
error", resourcePathInStorage),
+ e);
+ }
+ }
+
/**
* ready to execute SQL and parameter entity Map
*
@@ -420,19 +454,22 @@ public class SqlTask extends AbstractTask {
Map<String, Property> paramsMap =
taskExecutionContext.getPrepareParamsMap();
- // spell SQL according to the final user-defined variable
- if (paramsMap == null) {
- sqlBuilder.append(sql);
- return new SqlBinds(sqlBuilder.toString(), sqlParamsMap);
- }
+ Map<String, String> placeholderParamsMap = paramsMap == null
+ ? Collections.emptyMap()
+ : ParameterUtils.convert(paramsMap);
if (StringUtils.isNotEmpty(sqlParameters.getTitle())) {
- String title =
ParameterUtils.convertParameterPlaceholders(sqlParameters.getTitle(),
- ParameterUtils.convert(paramsMap));
+ String title =
ParameterUtils.convertParameterPlaceholders(sqlParameters.getTitle(),
placeholderParamsMap);
log.info("SQL title : {}", title);
sqlParameters.setTitle(title);
}
+ // spell SQL according to the final user-defined variable
+ if (paramsMap == null) {
+ sqlBuilder.append(sql);
+ return new SqlBinds(sqlBuilder.toString(), sqlParamsMap);
+ }
+
// special characters need to be escaped, ${} needs to be escaped
setSqlParamsMap(sql, sqlParamsMap, paramsMap,
taskExecutionContext.getTaskInstanceId());
// Replace the original value in sql !{...} ,Does not participate in
precompilation
diff --git
a/dolphinscheduler-task-plugin/dolphinscheduler-task-sql/src/test/java/org/apache/dolphinscheduler/plugin/task/sql/SqlTaskTest.java
b/dolphinscheduler-task-plugin/dolphinscheduler-task-sql/src/test/java/org/apache/dolphinscheduler/plugin/task/sql/SqlTaskTest.java
index 3832d653a5..92f0fa0702 100644
---
a/dolphinscheduler-task-plugin/dolphinscheduler-task-sql/src/test/java/org/apache/dolphinscheduler/plugin/task/sql/SqlTaskTest.java
+++
b/dolphinscheduler-task-plugin/dolphinscheduler-task-sql/src/test/java/org/apache/dolphinscheduler/plugin/task/sql/SqlTaskTest.java
@@ -20,6 +20,7 @@ package org.apache.dolphinscheduler.plugin.task.sql;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
+import org.apache.dolphinscheduler.common.utils.DateUtils;
import org.apache.dolphinscheduler.common.utils.JSONUtils;
import org.apache.dolphinscheduler.plugin.task.api.TaskConstants;
import org.apache.dolphinscheduler.plugin.task.api.TaskException;
@@ -31,11 +32,15 @@ import
org.apache.dolphinscheduler.plugin.task.api.model.Property;
import org.apache.dolphinscheduler.plugin.task.api.parameters.SqlParameters;
import
org.apache.dolphinscheduler.plugin.task.api.parameters.resource.DataSourceParameters;
import
org.apache.dolphinscheduler.plugin.task.api.parameters.resource.ResourceParametersHelper;
+import org.apache.dolphinscheduler.plugin.task.api.resource.ResourceContext;
import org.apache.dolphinscheduler.plugin.task.api.utils.ParameterUtils;
import org.apache.dolphinscheduler.spi.enums.DbType;
import java.lang.reflect.InvocationTargetException;
import java.lang.reflect.Method;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
import java.sql.ResultSet;
import java.sql.ResultSetMetaData;
import java.util.HashMap;
@@ -44,6 +49,7 @@ import java.util.Map;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
@@ -69,6 +75,47 @@ class SqlTaskTest {
sqlTask = new SqlTask(ctx);
}
+ @Test
+ void testSqlLoadedFromResourceFileWhenSqlIsEmpty(@TempDir Path tempDir)
throws Exception {
+ Path sqlFile = tempDir.resolve("test.sql");
+ String sqlContent = "SELECT 1";
+ Files.write(sqlFile, sqlContent.getBytes(StandardCharsets.UTF_8));
+
+ SqlParameters sqlParameters = new SqlParameters();
+ sqlParameters.setType("MYSQL");
+ sqlParameters.setDatasource(1);
+ sqlParameters.setSql(null);
+ sqlParameters.setSqlResource("/sql/test.sql");
+
+ DataSourceParameters dataSourceParameters = new DataSourceParameters();
+ dataSourceParameters.setType(DbType.MYSQL);
+ dataSourceParameters.setResourceType(ResourceType.DATASOURCE.name());
+
+ ResourceParametersHelper resourceParametersHelper = new
ResourceParametersHelper();
+ resourceParametersHelper.put(ResourceType.DATASOURCE, 1,
dataSourceParameters);
+
+ TaskExecutionContext taskExecutionContext = new TaskExecutionContext();
+
taskExecutionContext.setTaskParams(JSONUtils.toJsonString(sqlParameters));
+ taskExecutionContext.setScheduleTime(System.currentTimeMillis());
+
taskExecutionContext.setResourceParametersHelper(resourceParametersHelper);
+
+ ResourceContext resourceContext = new ResourceContext();
+ resourceContext.addResourceItem(ResourceContext.ResourceItem.builder()
+ .resourceAbsolutePathInStorage(sqlParameters.getSqlResource())
+ .resourceAbsolutePathInLocal(sqlFile.toString())
+ .build());
+ taskExecutionContext.setResourceContext(resourceContext);
+
+ SqlTask task = new SqlTask(taskExecutionContext);
+
+ Method ensureSqlContent =
SqlTask.class.getDeclaredMethod("ensureSqlContent");
+ ensureSqlContent.setAccessible(true);
+ ensureSqlContent.invoke(task);
+
+ SqlParameters loadedParameters = (SqlParameters) task.getParameters();
+ Assertions.assertEquals(sqlContent, loadedParameters.getSql());
+ }
+
@Test
void testReplacingSqlWithoutParams() {
String querySql = "select 1";
@@ -210,7 +257,6 @@ class SqlTaskTest {
@Test
void
testGenerateEmptyRow_WithNonNullResultSet_ReturnsEmptyValuesForAllColumns()
throws Exception {
- // Arrange
ResultSet mockResultSet = mock(ResultSet.class);
ResultSetMetaData mockMetaData = mock(ResultSetMetaData.class);
@@ -222,10 +268,8 @@ class SqlTaskTest {
Method method = SqlTask.class.getDeclaredMethod("generateEmptyRow",
ResultSet.class);
method.setAccessible(true);
- // Act
ArrayNode result = (ArrayNode) method.invoke(sqlTask, mockResultSet);
- // Assert
Assertions.assertNotNull(result);
Assertions.assertEquals(1, result.size());
@@ -236,14 +280,11 @@ class SqlTaskTest {
@Test
void testGenerateEmptyRow_WithNullResultSet_ReturnsErrorObject() throws
Exception {
- // Arrange
Method method = SqlTask.class.getDeclaredMethod("generateEmptyRow",
ResultSet.class);
method.setAccessible(true);
- // Act
ArrayNode result = (ArrayNode) method.invoke(sqlTask, (ResultSet)
null);
- // Assert
Assertions.assertNotNull(result);
Assertions.assertEquals(1, result.size());
@@ -260,7 +301,7 @@ class SqlTaskTest {
when(mockResultSet.getMetaData()).thenReturn(mockMetaData);
when(mockMetaData.getColumnCount()).thenReturn(3);
when(mockMetaData.getColumnLabel(1)).thenReturn("id");
- when(mockMetaData.getColumnLabel(2)).thenReturn("id"); // duplicate
+ when(mockMetaData.getColumnLabel(2)).thenReturn("id");
when(mockMetaData.getColumnLabel(3)).thenReturn("name");
Method method = SqlTask.class.getDeclaredMethod("generateEmptyRow",
ResultSet.class);
@@ -281,7 +322,6 @@ class SqlTaskTest {
Method resultProcessMethod =
SqlTask.class.getDeclaredMethod("resultProcess", ResultSet.class);
resultProcessMethod.setAccessible(true);
- // Mock a null ResultSet
String result = (String) resultProcessMethod.invoke(sqlTask,
(ResultSet) null);
Assertions.assertNotNull(result);
@@ -290,7 +330,6 @@ class SqlTaskTest {
@Test
void testResultProcess_EmptyResultSet_ReturnsEmptyResult() throws
Exception {
- // Mock a non-null ResultSet that contains no data rows
ResultSet mockResultSet = mock(ResultSet.class);
ResultSetMetaData mockMetaData = mock(ResultSetMetaData.class);
@@ -298,7 +337,7 @@ class SqlTaskTest {
when(mockMetaData.getColumnCount()).thenReturn(2);
when(mockMetaData.getColumnLabel(1)).thenReturn("id");
when(mockMetaData.getColumnLabel(2)).thenReturn("name");
- when(mockResultSet.next()).thenReturn(false); // no rows available
+ when(mockResultSet.next()).thenReturn(false);
Method resultProcessMethod =
SqlTask.class.getDeclaredMethod("resultProcess", ResultSet.class);
resultProcessMethod.setAccessible(true);
@@ -306,7 +345,6 @@ class SqlTaskTest {
String result = (String) resultProcessMethod.invoke(sqlTask,
mockResultSet);
Assertions.assertNotNull(result);
- // Verify the result contains empty string values for all columns and
is a valid JSON array
Assertions.assertTrue(result.contains("\"id\":\"\""));
Assertions.assertTrue(result.contains("\"name\":\"\""));
Assertions.assertTrue(result.startsWith("[{"));
@@ -321,17 +359,15 @@ class SqlTaskTest {
when(mockRs.getMetaData()).thenReturn(mockMd);
when(mockMd.getColumnCount()).thenReturn(2);
when(mockMd.getColumnLabel(1)).thenReturn("id");
- when(mockMd.getColumnLabel(2)).thenReturn("id"); // duplicate column
name
+ when(mockMd.getColumnLabel(2)).thenReturn("id");
Method method = SqlTask.class.getDeclaredMethod("resultProcess",
ResultSet.class);
method.setAccessible(true);
- // Assert that InvocationTargetException is thrown
InvocationTargetException thrown = Assertions.assertThrows(
InvocationTargetException.class,
() -> method.invoke(sqlTask, mockRs));
- // Check the actual cause
Throwable cause = thrown.getCause();
Assertions.assertNotNull(cause);
Assertions.assertInstanceOf(TaskException.class, cause,
@@ -341,4 +377,116 @@ class SqlTaskTest {
"TaskException message should mention duplicate column name");
}
+ @Test
+ void
testGetSqlAndSqlParamsMap_nullPrepareParamsMap_replacesScheduleTimeAndTitlePlaceholder()
throws Exception {
+ long scheduleTimeMillis = 1700000000000L;
+ String expectedDate = DateUtils.format(new
java.util.Date(scheduleTimeMillis), "yyyyMMdd");
+
+ TaskExecutionContext ctx = new TaskExecutionContext();
+
ctx.setTaskParams("{\"type\":\"HIVE\",\"datasource\":1,\"sql\":\"select
1\",\"title\":\"title-$[yyyyMMdd]\"}");
+ ctx.setScheduleTime(scheduleTimeMillis);
+
ctx.setResourceParametersHelper(getResourceParametersHelperWithDatasourceType(DbType.HIVE));
+
+ // Ensure prepareParamsMap == null
+ ctx.setPrepareParamsMap(null);
+
+ SqlTask task = new SqlTask(ctx);
+
+ Method method =
SqlTask.class.getDeclaredMethod("getSqlAndSqlParamsMap", String.class);
+ method.setAccessible(true);
+
+ String inputSql = "select '" + "$[yyyyMMdd]" + "'";
+ SqlBinds binds = (SqlBinds) method.invoke(task, inputSql);
+
+ Assertions.assertEquals("select '" + expectedDate + "'",
binds.getSql());
+ SqlParameters loadedParameters = (SqlParameters) task.getParameters();
+
Assertions.assertTrue(loadedParameters.getTitle().matches("title-\\d{8}"));
+ }
+
+ @Test
+ void
testGetSqlAndSqlParamsMap_withPrepareParamsMap_coversPrintReplacedSql() throws
Exception {
+ Map<String, Property> prepareParamsMap = new HashMap<>();
+ prepareParamsMap.put("dt", new Property("dt", Direct.IN,
DataType.VARCHAR, "1970"));
+
+ TaskExecutionContext ctx = new TaskExecutionContext();
+
ctx.setTaskParams("{\"type\":\"HIVE\",\"datasource\":1,\"sql\":\"select 1\"}");
+ ctx.setScheduleTime(System.currentTimeMillis());
+ ctx.setTaskInstanceId(1);
+
ctx.setResourceParametersHelper(getResourceParametersHelperWithDatasourceType(DbType.HIVE));
+ ctx.setPrepareParamsMap(prepareParamsMap);
+
+ SqlTask task = new SqlTask(ctx);
+
+ Method method =
SqlTask.class.getDeclaredMethod("getSqlAndSqlParamsMap", String.class);
+ method.setAccessible(true);
+
+ String inputSql = "select * from student where dt=${dt}";
+ SqlBinds binds = (SqlBinds) method.invoke(task, inputSql);
+
+ Assertions.assertEquals("select * from student where dt=?",
binds.getSql());
+ Assertions.assertNotNull(binds.getParamsMap());
+ Assertions.assertEquals("1970",
binds.getParamsMap().get(1).getValue());
+ }
+
+ @Test
+ void testEnsureSqlContent_whenSqlAlreadyPresent_doesNotReadResource()
throws Exception {
+ SqlTask task = this.sqlTask;
+
+ Method ensureSqlContent =
SqlTask.class.getDeclaredMethod("ensureSqlContent");
+ ensureSqlContent.setAccessible(true);
+
+ // Should early return because sqlParameters.getSql() is not empty.
+ ensureSqlContent.invoke(task);
+ }
+
+ @Test
+ void testEnsureSqlContent_whenResourceMissing_throwsTaskException(@TempDir
Path tempDir) throws Exception {
+ SqlParameters sqlParameters = new SqlParameters();
+ sqlParameters.setType("HIVE");
+ sqlParameters.setDatasource(1);
+ sqlParameters.setSql(null);
+ sqlParameters.setSqlResource("/sql/missing.sql");
+
+ DataSourceParameters dataSourceParameters = new DataSourceParameters();
+ dataSourceParameters.setType(DbType.HIVE);
+ dataSourceParameters.setResourceType(ResourceType.DATASOURCE.name());
+
+ ResourceParametersHelper resourceParametersHelper = new
ResourceParametersHelper();
+ resourceParametersHelper.put(ResourceType.DATASOURCE, 1,
dataSourceParameters);
+
+ TaskExecutionContext taskExecutionContext = new TaskExecutionContext();
+
taskExecutionContext.setTaskParams(JSONUtils.toJsonString(sqlParameters));
+ taskExecutionContext.setScheduleTime(System.currentTimeMillis());
+
taskExecutionContext.setResourceParametersHelper(resourceParametersHelper);
+
+ ResourceContext resourceContext = new ResourceContext();
+ // Point to a file that does not exist to trigger IOException in
ensureSqlContent.
+ Path missingLocalPath = tempDir.resolve("missing.sql");
+ resourceContext.addResourceItem(ResourceContext.ResourceItem.builder()
+ .resourceAbsolutePathInStorage(sqlParameters.getSqlResource())
+ .resourceAbsolutePathInLocal(missingLocalPath.toString())
+ .build());
+ taskExecutionContext.setResourceContext(resourceContext);
+
+ SqlTask task = new SqlTask(taskExecutionContext);
+
+ Method ensureSqlContent =
SqlTask.class.getDeclaredMethod("ensureSqlContent");
+ ensureSqlContent.setAccessible(true);
+
+ InvocationTargetException thrown = Assertions.assertThrows(
+ InvocationTargetException.class,
+ () -> ensureSqlContent.invoke(task));
+ Assertions.assertInstanceOf(TaskException.class, thrown.getCause());
+ }
+
+ private ResourceParametersHelper
getResourceParametersHelperWithDatasourceType(DbType dbType) {
+ DataSourceParameters parameters = new DataSourceParameters();
+ parameters.setType(dbType);
+ parameters.setResourceType(ResourceType.DATASOURCE.name());
+
+ ResourceParametersHelper resourceParametersHelper = new
ResourceParametersHelper();
+ resourceParametersHelper.put(ResourceType.DATASOURCE, 1, parameters);
+ return resourceParametersHelper;
+ }
+
}
diff --git a/dolphinscheduler-ui/src/locales/en_US/project.ts
b/dolphinscheduler-ui/src/locales/en_US/project.ts
index 667153f215..f10e7fbdd0 100644
--- a/dolphinscheduler-ui/src/locales/en_US/project.ts
+++ b/dolphinscheduler-ui/src/locales/en_US/project.ts
@@ -552,10 +552,14 @@ export default {
sql_type_query: 'Query',
sql_type_non_query: 'Non Query',
sql_statement: 'SQL Statement',
+ sql_source: 'SQL Source',
+ sql_source_script: 'Script',
+ sql_source_file: 'Resource file',
pre_sql_statement: 'Pre SQL Statement',
post_sql_statement: 'Post SQL Statement',
sql_input_placeholder: 'Please enter non-query sql.',
sql_empty_tips: 'The sql can not be empty.',
+ sql_resource_file: 'SQL Resource File',
procedure_method: 'SQL Statement',
procedure_method_tips: 'Please enter the procedure script',
procedure_method_snippet:
diff --git a/dolphinscheduler-ui/src/locales/zh_CN/project.ts
b/dolphinscheduler-ui/src/locales/zh_CN/project.ts
index ef9634095d..a3ad17816a 100644
--- a/dolphinscheduler-ui/src/locales/zh_CN/project.ts
+++ b/dolphinscheduler-ui/src/locales/zh_CN/project.ts
@@ -534,10 +534,14 @@ export default {
sql_type_query: '查询',
sql_type_non_query: '非查询',
sql_statement: 'SQL语句',
+ sql_source: 'SQL来源',
+ sql_source_script: '脚本',
+ sql_source_file: '资源文件',
pre_sql_statement: '前置SQL语句',
post_sql_statement: '后置SQL语句',
sql_input_placeholder: '请输入非查询SQL语句',
sql_empty_tips: '语句不能为空',
+ sql_resource_file: 'SQL资源文件',
procedure_method: 'SQL语句',
procedure_method_tips: '请输入存储脚本',
procedure_method_snippet:
diff --git
a/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-sql.ts
b/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-sql.ts
index 0044fa9cf8..5374fbfa25 100644
---
a/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-sql.ts
+++
b/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-sql.ts
@@ -14,14 +14,40 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-import { computed } from 'vue'
+import { computed, ref, onMounted } from 'vue'
import { useI18n } from 'vue-i18n'
+import { queryResourceList } from '@/service/modules/resources'
+import { useTaskNodeStore } from '@/store/project/task-node'
+import utils from '@/utils'
import { useCustomParams } from '.'
-import type { IJsonItem } from '../types'
+import type { IJsonItem, IResource } from '../types'
export function useSql(model: { [field: string]: any }): IJsonItem[] {
const { t } = useI18n()
const hiveSpan = computed(() => (model.type === 'HIVE' ? 24 : 0))
+ const taskStore = useTaskNodeStore()
+ const sqlResourceOptions = ref<IResource[]>([])
+
+ const sqlEditorSpan = computed(() => (model.sqlSource === 'FILE' ? 0 : 24))
+ const sqlResourceSpan = computed(() => (model.sqlSource === 'FILE' ? 24 : 0))
+ const isScriptSource = computed(
+ () => model.sqlSource === 'SCRIPT' || !model.sqlSource
+ )
+
+ const loadSqlResourceTree = async () => {
+ if (taskStore.resources.length) {
+ sqlResourceOptions.value = taskStore.resources
+ return
+ }
+ const res = await queryResourceList({ type: 'FILE', fullName: '' })
+ utils.removeUselessChildren(res)
+ sqlResourceOptions.value = res || []
+ taskStore.updateResource(res)
+ }
+
+ onMounted(() => {
+ void loadSqlResourceTree()
+ })
return [
{
@@ -34,19 +60,56 @@ export function useSql(model: { [field: string]: any }):
IJsonItem[] {
},
span: hiveSpan
},
+ {
+ type: 'radio',
+ field: 'sqlSource',
+ name: t('project.node.sql_source'),
+ options: [
+ {
+ label: t('project.node.sql_source_script'),
+ value: 'SCRIPT'
+ },
+ {
+ label: t('project.node.sql_source_file'),
+ value: 'FILE'
+ }
+ ],
+ span: 24
+ },
{
type: 'editor',
field: 'sql',
name: t('project.node.sql_statement'),
+ span: sqlEditorSpan,
validate: {
- trigger: ['input', 'trigger'],
- required: true,
+ trigger: ['input', 'blur'],
+ required: isScriptSource,
message: t('project.node.sql_empty_tips')
},
props: {
language: 'sql'
}
},
+ {
+ type: 'tree-select',
+ field: 'sqlResource',
+ name: t('project.node.sql_resource_file'),
+ span: sqlResourceSpan,
+ options: sqlResourceOptions,
+ props: {
+ placeholder: t('project.node.resources_tips'),
+ keyField: 'fullName',
+ labelField: 'name',
+ disabledField: 'disable',
+ filterable: true,
+ showPath: true
+ },
+ validate: {
+ trigger: ['input', 'blur'],
+ required: computed(() => model.sqlSource === 'FILE'),
+ message: t('project.node.resources_tips')
+ }
+ },
...useCustomParams({
model,
field: 'localParams',
diff --git
a/dolphinscheduler-ui/src/views/projects/task/components/node/tasks/use-sql.ts
b/dolphinscheduler-ui/src/views/projects/task/components/node/tasks/use-sql.ts
index e967c2f09b..371ae737fa 100644
---
a/dolphinscheduler-ui/src/views/projects/task/components/node/tasks/use-sql.ts
+++
b/dolphinscheduler-ui/src/views/projects/task/components/node/tasks/use-sql.ts
@@ -47,6 +47,8 @@ export function useSql({
type: 'MYSQL',
displayRows: 10,
sql: '',
+ sqlSource: 'SCRIPT',
+ sqlResource: '',
sqlType: '0',
preStatements: [],
postStatements: [],