《楚门的世界》
如果DolphinScheduler官方的案例没有演示如何动态传,我们开发者应该如何去处理这种需求?
SqlTask,一个是
HttpTask。我的观点是:总不能做完SQL查询,或者做完HTTP请求后就不管结果吧?
分析HttpTask源码
分析HttpTask源码,直接找到HttpTask的handle方法,DolphinScheduler中,任何Task的具体执行逻辑都在这个handle方法中。
handle方法分析
1@Override
2public void handle(TaskCallBack taskCallBack) throws TaskException {
3 long startTime = System.currentTimeMillis();
4 String formatTimeStamp = DateUtils.formatTimeStamp(startTime);
5 String statusCode = null;
6 String body = null;
7
8 try (
9 CloseableHttpClient client = createHttpClient();
10 CloseableHttpResponse response = sendRequest(client)) {
11 statusCode = String.valueOf(getStatusCode(response));
12 body = getResponseBody(response);
13 exitStatusCode = validResponse(body, statusCode);
14 // 看名字应该就能猜到是处理请求结果的
15 addDefaultOutput(body);
16 long costTime = System.currentTimeMillis() - startTime;
17 log.info(
18 "startTime: {}, httpUrl: {}, httpMethod: {}, costTime : {} milliseconds, statusCode : {}, body : {}, log : {}",
19 formatTimeStamp, httpParameters.getUrl(),
20 httpParameters.getHttpMethod(), costTime, statusCode, body, output);
21 } catch (Exception e) {
22 appendMessage(e.toString());
23 exitStatusCode = -1;
24 log.error("httpUrl[" + httpParameters.getUrl() + "] connection failed:" + output, e);
25 throw new TaskException("Execute http task failed", e);
26 }
27
28}
addDefaultOutput方法
1public void addDefaultOutput(String response) {
2 // put response in output
3 // 创建Property对象
4 Property outputProperty = new Property();
5 // 设置Prop,也就是设置Key
6 outputProperty.setProp(String.format("%s.%s", taskExecutionContext.getTaskName(), "response"));
7 // 设置是入参还是出参,这里是出参,因为是将结果给下游任务
8 outputProperty.setDirect(Direct.OUT);
9 // 设置参数类型,VARCHAR表示就是字符串
10 outputProperty.setType(DataType.VARCHAR);
11 // 设置Value,就是http请求结果
12 outputProperty.setValue(response);
13 // 重点:将Property添加到varPool中
14 httpParameters.addPropertyToValPool(outputProperty);
15}
分析SqlTask源码
handler方法分析
1@Override
2public void handle(TaskCallBack taskCallBack) throws TaskException {
3 log.info("Full sql parameters: {}", sqlParameters);
4 log.info(
5 "sql type : {}, datasource : {}, sql : {} , localParams : {},udfs : {},showType : {},connParams : {},varPool : {} ,query max result limit {}",
6 sqlParameters.getType(),
7 sqlParameters.getDatasource(),
8 sqlParameters.getSql(),
9 sqlParameters.getLocalParams(),
10 sqlParameters.getUdfs(),
11 sqlParameters.getShowType(),
12 sqlParameters.getConnParams(),
13 sqlParameters.getVarPool(),
14 sqlParameters.getLimit());
15 try {
16
17 // get datasource
18 baseConnectionParam = (BaseConnectionParam) DataSourceUtils.buildConnectionParams(dbType,
19 sqlTaskExecutionContext.getConnectionParams());
20 List<String> subSqls = DataSourceProcessorProvider.getDataSourceProcessor(dbType)
21 .splitAndRemoveComment(sqlParameters.getSql());
22
23 // ready to execute SQL and parameter entity Map
24 List<SqlBinds> mainStatementSqlBinds = subSqls
25 .stream()
26 .map(this::getSqlAndSqlParamsMap)
27 .collect(Collectors.toList());
28
29 List<SqlBinds> preStatementSqlBinds = Optional.ofNullable(sqlParameters.getPreStatements())
30 .orElse(new ArrayList<>())
31 .stream()
32 .map(this::getSqlAndSqlParamsMap)
33 .collect(Collectors.toList());
34 List<SqlBinds> postStatementSqlBinds = Optional.ofNullable(sqlParameters.getPostStatements())
35 .orElse(new ArrayList<>())
36 .stream()
37 .map(this::getSqlAndSqlParamsMap)
38 .collect(Collectors.toList());
39
40 List<String> createFuncs = createFuncs(sqlTaskExecutionContext.getUdfFuncParametersList());
41
42 // execute sql task
43 // 这个方法就是处理sql结果的
44 executeFuncAndSql(mainStatementSqlBinds, preStatementSqlBinds, postStatementSqlBinds, createFuncs);
45
46 setExitStatusCode(TaskConstants.EXIT_CODE_SUCCESS);
47
48 } catch (Exception e) {
49 setExitStatusCode(TaskConstants.EXIT_CODE_FAILURE);
50 log.error("sql task error", e);
51 throw new TaskException("Execute sql task failed", e);
52 }
53}
所以我们在看下executeFuncAndSql
方法内部实现
1public void executeFuncAndSql(List<SqlBinds> mainStatementsBinds,
2 List<SqlBinds> preStatementsBinds,
3 List<SqlBinds> postStatementsBinds,
4 List<String> createFuncs) throws Exception {
5 try (
6 Connection connection =
7 DataSourceClientProvider.getAdHocConnection(DbType.valueOf(sqlParameters.getType()),
8 baseConnectionParam)) {
9
10 // create temp function
11 if (CollectionUtils.isNotEmpty(createFuncs)) {
12 createTempFunction(connection, createFuncs);
13 }
14
15 // pre execute
16 executeUpdate(connection, preStatementsBinds, "pre");
17
18 // main execute
19 String result = null;
20 // decide whether to executeQuery or executeUpdate based on sqlType
21 if (sqlParameters.getSqlType() == SqlType.QUERY.ordinal()) {
22 // query statements need to be convert to JsonArray and inserted into Alert to send
23 result = executeQuery(connection, mainStatementsBinds.get(0), "main");
24 } else if (sqlParameters.getSqlType() == SqlType.NON_QUERY.ordinal()) {
25 // non query statement
26 String updateResult = executeUpdate(connection, mainStatementsBinds, "main");
27 result = setNonQuerySqlReturn(updateResult, sqlParameters.getLocalParams());
28 }
29 // deal out params
30 // 这个方法就是来处理结果的
31 sqlParameters.dealOutParam(result);
32
33 // post execute
34 executeUpdate(connection, postStatementsBinds, "post");
35 } catch (Exception e) {
36 log.error("execute sql error: {}", e.getMessage());
37 throw e;
38 }
39}
通过dealOutParam
看具体处理细节
1public void dealOutParam(String result) {
2 if (CollectionUtils.isEmpty(localParams)) {
3 return;
4 }
5 List<Property> outProperty = getOutProperty(localParams);
6 if (CollectionUtils.isEmpty(outProperty)) {
7 return;
8 }
9 if (StringUtils.isEmpty(result)) {
10 varPool = VarPoolUtils.mergeVarPool(Lists.newArrayList(varPool, outProperty));
11 return;
12 }
13 List<Map<String, String>> sqlResult = getListMapByString(result);
14 if (CollectionUtils.isEmpty(sqlResult)) {
15 return;
16 }
17 // if sql return more than one line
18 if (sqlResult.size() > 1) {
19 Map<String, List<String>> sqlResultFormat = new HashMap<>();
20 // init sqlResultFormat
21 Set<String> keySet = sqlResult.get(0).keySet();
22 for (String key : keySet) {
23 sqlResultFormat.put(key, new ArrayList<>());
24 }
25 for (Map<String, String> info : sqlResult) {
26 for (String key : info.keySet()) {
27 sqlResultFormat.get(key).add(String.valueOf(info.get(key)));
28 }
29 }
30 for (Property info : outProperty) {
31 if (info.getType() == DataType.LIST) {
32 info.setValue(JSONUtils.toJsonString(sqlResultFormat.get(info.getProp())));
33 }
34 }
35 } else {
36 // result only one line
37 Map<String, String> firstRow = sqlResult.get(0);
38 for (Property info : outProperty) {
39 info.setValue(String.valueOf(firstRow.get(info.getProp())));
40 }
41 }
42
43 // 本质还是将sql结果处理后保存在varPool中,varPool才是关键所在
44 varPool = VarPoolUtils.mergeVarPool(Lists.newArrayList(varPool, outProperty));
45
46}
所以,源代码分析到这,我们就知道了:如果想实现动态传参,那么我们需要将传递的数据封装成org.apache.dolphinscheduler.plugin.task.api.model.Property
,然后添加到内置集合变量org.apache.dolphinscheduler.plugin.task.api.parameters.AbstractParameters#varPool
中
1Property outputProperty = new Property();
2// 添加我们要传递的数据Key
3outputProperty.setProp("xxxxKey"));
4// OUT
5outputProperty.setDirect(Direct.OUT);
6// 这里传递的数据是什么类型就写什么类型,建议通过json字符串处理数据
7outputProperty.setType(DataType.VARCHAR);
8// 添加我们要传递的数据Key
9outputProperty.setValue("xxxxValue");
10// 这里的xxxxParameters是我们自己自定义的,一般情况下,一个Task对应一个Parameters
11xxxxParameters.addPropertyToValPool(outputProperty);
DolphinScheduler内部有将List<Property> varPool
转换成Map<String, Property> varParams
的逻辑,然后会将varParams
与其他的参数合并,最后通过taskExecutionContext.setPrepareParamsMap(propertyMap)
将数据设置给Map<String, Property> prepareParamsMap
。
参与Apache DolphinScheduler 社区有非常多的参与贡献的方式,包括:





