暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Apache DolphinScheduler将上游Task执行结果传递给下游Task

海豚调度 2024-11-11
361

真实的世界不总是顺心的。

《楚门的世界》



01

背景



公司的数据开发平台需要用到DolphinScheduler做任务调度,其中一个场景是:上游任务执行结束后,需要将任务执行结果传递给下游任务。
Apache DolphinScheduler肯定是能实现任务之间的传参的,具体的可以看:DolphinScheduler | 文档中心 (https://dolphinscheduler.apache.org/zh-cn/docs/3.2.2/guide/parameter/context)
但是官方案例中介绍的任务之间传参是提前在管理台上配置好的,OK,那么问题来了,如何实现任务之间的动态传参呢?比如说我们自定义Task,然后在Task执行结束后将执行结果封装,传递给DAG中的下一个Task。
02

分析



如果DolphinScheduler官方的案例没有演示如何动态传,我们开发者应该如何去处理这种需求?

我是这么做的:分析DolphinScheduler内置的Task,总有一个Task是需要传递参数给下游的。我这里盲猜两个,一个是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

03

具体实现



这里我们不去讨论自定义Task的具体实现步骤,这不是本文的重点。
当我们实现自定义Task后,可以这样编码实现动态传参:

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

04

总结



关于DolphinScheduler(海豚调度器)是什么,能做什么,怎么使用等等,这里我就不再赘述,大家感兴趣的可以去看看官方文档:DolphinScheduler | 文档中心 (https://dolphinscheduler.apache.org/zh-cn/docs/3.2.2)
希望通过本篇文章能让各位读者掌握Task之间的动态传参,然后应用在实际工作中。如果本篇文章能给屏幕前的你们或多或少的一些帮助,也是我喜闻乐见的。
<🐬🐬 >

推荐阅读

用户实践案例
奇富科技  蜀海供应链 联通数科 拈花云科
蔚来汽车 长城汽车 集度 长安汽车
思科网讯 生鲜电商 联通医疗 联想
新网银行 消费金融  腾讯音乐 自如
有赞 伊利 当贝大数据
联想 传智教育 Bigo
通信行业  作业帮 太美医疗
某新能源 中电信翼康 每日互动
迁移实践
Azkaban   Ooize   
Airflow (有赞案例) Air2phin(迁移工具)
Airflow迁移实践
Apache DolphinScheduler 3.0.0 升级到 3.1.8 教程
Apache DolphinScheduler 1.3.4升级至3.1.2版本解决方案合集

新手入门
选择Apache DolphinScheduler的10个理由
Apache DolphinScheduler 3.1.8 保姆级教程【安装、介绍、项目运用、邮箱预警设置】轻松拿捏!
Apache DolphinScheduler 如何实现自动化打包+单机/集群部署?
Apache DolphinScheduler-3.1.3 版本安装部署详细教程
Apache DolphinScheduler 在大数据环境中的应用与调优

< 🐬🐬 >
参与社区


参与Apache DolphinScheduler 社区有非常多的参与贡献的方式,包括:



贡献第一个PR(文档、代码) 我们也希望是简单的,第一个PR用于熟悉提交的流程和社区协作以及感受社区的友好度。

社区汇总了以下适合新手的问题列表:https://github.com/apache/dolphinscheduler/issues/5689

非新手问题列表:https://github.com/apache/dolphinscheduler/issues?
q=is%3Aopen+is%3Aissue+label%3A%22volunteer+wanted%22

如何参与贡献链接:https://dolphinscheduler.apache.org/zh-cn/community/development/contribute.html

来吧,DolphinScheduler开源社区需要您的参与,为中国开源崛起添砖加瓦吧,哪怕只是小小的一块瓦,汇聚起来的力量也是巨大的!

文章转载自海豚调度,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论