news 2026/9/23 18:04:51

Apache DolphinScheduler 全局参数(OUT 参数)机制详解:varPool 合并、跨节点传递与 ${setValue} 输出解析

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache DolphinScheduler 全局参数(OUT 参数)机制详解:varPool 合并、跨节点传递与 ${setValue} 输出解析

Apache DolphinScheduler 全局参数(OUT 参数)机制详解:varPool 合并、跨节点传递与 ${setValue} 输出解析

【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/gh_mirrors/do/dolphinscheduler

本文是 Apache DolphinScheduler 全局参数开发机制的深度技术指南,以仓库内 全局参数开发文档 为主体骨架,结合 Master、Worker 与任务插件模块的源码实现展开。读者读完将掌握 OUT 参数从定义、跨节点合并传递、Worker 侧三池合并到 SQL/SHELL 节点输出回写的完整生命周期,可直接据此理解全局参数行为或在此基础上进行二次开发。

一、全局参数的定位:从 localParam 到 varPool

在 Apache DolphinScheduler 中,用户在定义任务(task)时配置的参数分为两类:

  • IN 方向参数:作为当前任务节点的输入,在节点执行前完成替换;
  • OUT 方向参数:作为当前任务节点的输出,执行完成后回传给流程,供后续节点作为输入使用。

按照 全局参数开发文档 的说明,用户在定义方向为 OUT 的参数后,该参数会保存在 task 的localParam中。这里的localParam与流程级的globalParam(全局参数)不同:localParam是节点私有的参数列表,其中方向为 OUT 的条目承载了“本节点向外输出变量”的语义。

从源码结构看,localParam对应的模型是 Property.java,其核心字段包括变量名prop、方向direct(IN/OUT)、类型type(如 VARCHAR、LIST 等)与值valuevarPool则是一组Property的序列化形态(JSON 字符串),在 Master 与 Worker 之间、节点与节点之间流转。整条数据链路可以概括为:

task 定义 OUT 参数 (localParam) → 任务执行产出 → varPool (JSON) → Master 合并传递 → Worker 解析合并 → 节点执行替换 → 输出参数回写 localParam → 继续传递给下游

二、参数的使用:Master 侧 varPool 的获取与合并

2.1 从前置节点收集 varPool

当一个任务实例(taskInstance)需要创建时,Master 会从 DAG 中获取该节点的直接前置节点 preTasks,并收集这些前置节点产出的varPoolList<Property>),随后合并为一个 varPool。

文档明确描述了合并过程中同名变量的处理逻辑,这一逻辑在代码中有两处印证:

  • 文档描述的规则(同名变量冲突时):
    • 若所有值均为 null,则合并后的值为 null;
    • 若有且仅有一个值为非 null,则合并后的值为该非 null 值;
    • 若所有值均非 null,则取这些 varPool 所属 taskInstance 中endtime 最早的那一个;
  • 合并过程中,所有合并进来的 Property 的方向都会被更新为IN
  • 合并结果保存在taskInstance.varPool中。

对应的核心实现位于 WorkflowExecuteRunnable.initializeTaskInstanceVarPool():

// 获取当前任务的前置节点 preTasks String preTasks = workflowExecuteContext.getWorkflowGraph() .getTaskNodeByCode(taskInstance.getTaskCode()).getPreTasks(); Set<Long> preTaskList = new HashSet<>(JSONUtils.toList(preTasks, Long.class)); ProcessInstance workflowInstance = workflowExecuteContext.getWorkflowInstance(); if (CollectionUtils.isEmpty(preTaskList)) { // 无前置节点时,直接继承流程实例的 varPool taskInstance.setVarPool(workflowInstance.getVarPool()); return; } // 收集所有前置节点实例的 varPool,并按 endTime 排序 List<String> preTaskInstanceVarPools = preTaskList .stream() .map(taskCode -> getTaskInstance(taskCode).orElse(null)) .filter(Objects::nonNull) .sorted(Comparator.comparing(TaskInstance::getEndTime)) .map(TaskInstance::getVarPool) .collect(Collectors.toList()); taskInstance.setVarPool(VarPoolUtils.mergeVarPoolJsonString(preTaskInstanceVarPools));

从这段实现可以推断文档中“取 endtime 最早的一个”这一规则的落地方式:先按endTime升序排序,再交由VarPoolUtils.mergeVarPool合并,同名变量在合并时以先到(endtime 更早)的值为准。而在流程实例层面,WorkflowExecuteRunnable中还有对整条流程 varPool 的维护(例如 mergeVarPoolJsonString 调用 与失败重跑时对 varPool 的清理逻辑),保证流程级数据的一致性。

2.2 合并工具:VarPoolUtils

varPool 的合并、反序列化与减法操作统一封装在 VarPoolUtils.java:

  • deserializeVarPool(String varPoolJson):将 JSON 字符串反序列化为List<Property>
  • mergeVarPoolJsonString(List<String> varPoolJsons):批量合并多个 varPool 的 JSON 串,返回合并后的 JSON 串(空集合返回 null);
  • mergeVarPool(List<List<Property>> varPools):核心合并逻辑,仅处理方向为 OUT 的 Property,以property.getProp()(变量名)为 key 存入 Map,后写入的覆盖先写入的;
  • subtractVarPool / subtractVarPoolJson:从 varPool 中剔除指定变量(用于失败重跑等场景下清理过期数据)。
public List<Property> mergeVarPool(List<List<Property>> varPools) { if (CollectionUtils.isEmpty(varPools)) return null; if (varPools.size() == 1) return varPools.get(0); Map<String, Property> result = new HashMap<>(); for (List<Property> varPool : varPools) { if (CollectionUtils.isEmpty(varPool)) continue; for (Property property : varPool) { if (!Direct.OUT.equals(property.getDirect())) { log.info("The direct should be OUT in varPool, but got {}", property.getDirect()); continue; } result.put(property.getProp(), property); } } return new ArrayList<>(result.values()); }

注意:mergeVarPool只接受方向为 OUT 的 Property,且以“后写覆盖先写”为语义;而 WorkflowExecuteRunnable 在收集前置节点 varPool 时已按 endtime 升序排序,因此“endtime 最早者优先”与“先写后覆盖”正好吻合。该工具类的合并行为有对应单元测试覆盖,可参考 VarPoolUtilsTest.java。

2.3 Worker 侧解析与三池合并优先级

Worker 收到任务后,会将varPool解析为Map<String, Property>格式,其中map 的 key 为property.prop,即变量名

在 processor 处理参数时,会将三个变量池合并处理:

变量池说明合并优先级
globalParam流程级全局参数(保留)
varPool前置节点传递来的输出参数
localParam节点自身定义的参数低(被替换)

合并时若存在同名参数,高优先级保留,低优先级被替换,即同名时优先取globalParam,其次varPool,最后才是localParam

参数会在节点内容执行之前,通过正则表达式匹配${变量名}并替换为对应的值。因此,一个典型的使用场景是:上游 SQL 节点产出 OUT 参数后,下游 Shell 节点可以直接在脚本中书写${变量名}引用上游产出值。

三、参数的设置:SQL 与 SHELL 节点的输出回写

文档明确指出,目前仅支持 SQL 和 SHELL 节点的参数获取:从localParam中获取方向为 OUT 的参数,再根据节点类型做不同处理。

3.1 SQL 节点:List<Map<String, String>> 结构

SQL 节点参数返回的结构为List<Map<String, String>>

  • List 的元素为每行数据;
  • Map 的 key为列名,value为该列对应的值。

匹配规则如下:

  • 若 SQL 语句返回一行数据:根据用户在定义 task 时定义的 OUT 参数名匹配列名,匹配成功则取值,未匹配到则放弃;
  • 若 SQL 语句返回多行:根据用户定义的类型为LIST的 OUT 参数名匹配列名,将对应列的所有行数据转换为List<String>作为该参数的值,未匹配到则放弃。

这一逻辑与 SqlParameters.dealOutParam() 的实现完全对应:

public void dealOutParam(String result) { if (CollectionUtils.isEmpty(localParams)) return; List<Property> outProperty = getOutProperty(localParams); if (CollectionUtils.isEmpty(outProperty)) return; if (StringUtils.isEmpty(result)) { varPool = VarPoolUtils.mergeVarPool(Lists.newArrayList(varPool, outProperty)); return; } List<Map<String, String>> sqlResult = getListMapByString(result); // 多行:按列聚合为 List<String> if (sqlResult.size() > 1) { // ... 按列名聚合所有行数据 ... for (Property info : outProperty) { if (info.getType() == DataType.LIST) { info.setValue(JSONUtils.toJsonString(sqlResultFormat.get(info.getProp()))); } } } else { // 单行:直接按列名取值 Map<String, String> firstRow = sqlResult.get(0); for (Property info : outProperty) { info.setValue(String.valueOf(firstRow.get(info.getProp()))); } } varPool = VarPoolUtils.mergeVarPool(Lists.newArrayList(varPool, outProperty)); }

结合实现可以补充两个细节:其一,DataType.LIST是触发多行聚合的关键条件,只有 OUT 参数类型定义为 LIST 才会把多行数据聚合为 JSON 数组字符串;其二,即使结果为空,原有 varPool 与 OUT 参数占位也会被合并保留。相关单测可参考 SqlParametersTest.java。

3.2 SHELL 节点:${setValue(key=value)} 语法

SHELL 节点 processor 执行后的结果返回为Map<String, String>。用户在编写 shell 脚本时,需要在输出中定义${setValue(key=value)}格式的标记,例如:

echo "${setValue(custom_key=custom_value)}"

参数处理时,框架会去掉${setValue(...)}外壳,按照=进行拆分:第 0 个为 key,第 1 个为 value;随后用用户定义 task 时配置的 OUT 参数名与该 key 匹配,将 value 作为该参数的值写入。

该语法由 TaskOutputParameterParser.java 负责解析,值得注意的实现要点:

  • 同时支持${setValue(...)}#{setValue(...)}两种写法;
  • 解析要求表达式必须以setValue(开头、以)}结尾,且内部必须包含=,否则视为非法并告警跳过;
  • 为了防止构造的超长参数导致 OOM,parser 设置了参数行数上限(默认单参数最多1024 行)与长度上限,超出即跳过并打印告警日志;
  • 解析结果以 key-value 形式写入Map<String, String>,供后续与 OUT 参数名匹配使用。
int indexOfVarPoolBegin = logLine.indexOf("${setValue("); if (indexOfVarPoolBegin == -1) { indexOfVarPoolBegin = logLine.indexOf("#{setValue("); } ... String[] keyValue = keyValueExpression.split("=", 2); return ImmutablePair.of(keyValue[0], keyValue[1]);

该解析器的边界行为(多行参数、异常格式、超长截断)均有测试覆盖,可参考 TaskOutputParameterParserTest.java。

3.3 返回参数处理的标准流程

无论 SQL 还是 SHELL,返回参数的处理都遵循同一套流程,文档将其归纳为六个步骤:

  1. 获取到的 processor 的结果为 String;
  2. 判断 processor 结果是否为空,为空则退出;
  3. 判断 localParam 是否为空,为空则退出;
  4. 获取 localParam 中方向为 OUT 的参数,为空则退出;
  5. 将 String 按对应格式格式化(SQL 为List<Map<String, String>>,SHELL 为Map<String, String>);
  6. 将匹配好值的参数赋值给 varPool(List<Property>,其中包含原有 IN 的参数)。

随后 varPool 被格式化为 JSON 传递给 Master;Master 接收到 varPool 后,会将其中方向为 OUT 的参数回写到 localParam中,从而完成“输出参数沉淀回节点定义”的闭环。这也解释了为什么下游节点能够稳定地通过 varPool 读取到上游的 OUT 参数——数据始终以 JSON 形式在节点间显式传递,而非依赖共享内存或外部存储。

四、全链路时序总结

综合文档与源码,一次全局参数(OUT 参数)的完整生命周期如下:

1. 用户定义 task 时配置 OUT 参数 → 存入 localParam 2. 任务执行(SQL/SHELL)产出输出 ├─ SQL:结果格式化为 List<Map<String,String>>,按 OUT 参数名列名匹配取值 └─ SHELL:解析 ${setValue(key=value)},按 OUT 参数名与 key 匹配取值 3. Worker 将匹配结果与原有 varPool 合并(List<Property>,含原有 IN 参数)→ 序列化为 JSON 4. Master 接收 varPool,将 OUT 参数回写到 localParam 5. 下游任务实例创建时,Master 收集所有直接前置节点的 varPool → 同名冲突时按"endtime 最早者优先"合并 → 合并后方向更新为 IN → 存入 taskInstance.varPool 6. Worker 将 varPool 解析为 Map<String,Property>(key 为变量名) → 与 globalParam、localParam 三池合并(优先级 globalParam > varPool > localParam) → 节点内容执行前用正则替换 ${变量名}

五、开发者注意事项

  • 仅 SQL 与 SHELL 节点支持参数输出:其余节点类型即使配置了 OUT 参数,也不会触发本文所述的回写逻辑;
  • 同名变量冲突有明确优先级:三池合并时globalParam优先于varPool优先于localParam,跨节点合并时“endtime 最早”者胜出,设计多级变量时应避免依赖容易产生歧义的重复命名;
  • SHELL 输出语法必须严格${setValue(...)}必须完整闭合且包含=,超长(超过 1024 行)或格式非法的表达式会被静默跳过并告警,排查参数未生效问题时优先检查输出日志中的告警;
  • SQL 多行输出需要 LIST 类型:只有 OUT 参数类型为 LIST 时才会聚合多行数据,单行输出则直接按列名取标量值;
  • 方向语义随合并变化:合并进下游 varPool 的参数方向被更新为 IN,表示它们已成为下游的输入来源,后续任务回写时会保留原有 IN 参数。

本文全部机制描述均可在仓库源码中逐一验证:参数模型见 Property.java,合并工具见 VarPoolUtils.java,SQL 输出处理见 SqlParameters.java,SHELL 输出解析见 TaskOutputParameterParser.java,Master 侧合并见 WorkflowExecuteRunnable.java。本系列机制的文档入口位于 机制综述。

【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/gh_mirrors/do/dolphinscheduler

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/23 18:03:41

Cesium卫星雷达可视化实战:从时钟驱动到波束建模与性能调优

简介&#xff1a;面向前端开发者的 Cesium 三维开发示例包&#xff0c;聚焦卫星雷达数据在三维场景中的可视化&#xff0c;演示如何通过 HTML 页面加载遥感底图并构建动态效果。压缩包共 11 个文件&#xff0c;含 2 个 HTML 示例与 9 张 PNG 卫星/雷达影像图&#xff0c;整体约…

作者头像 李华
网站建设 2026/9/23 18:02:14

PSO优化风-水电联合调度:快速求解多约束能源调度问题

简介&#xff1a;本资源是基于粒子群算法&#xff08;PSO&#xff09;实现风电-水电&#xff08;抽水蓄能&#xff09;联合优化调度的MATLAB仿真程序&#xff0c;面向电力系统优化、新能源并网调度及智能算法应用方向的研究生、工程师与科研人员&#xff0c;解决风电出力波动大…

作者头像 李华
网站建设 2026/9/23 17:53:36

RecRecNet深度学习畸变矫正实战:推理、训练与部署指南

简介&#xff1a;基于RecRecNet深度网络实现广角图像畸变矫正&#xff0c;所附Python源码适用于高校计算机相关专业学生与教师&#xff0c;可支撑毕业设计、课程设计及初学进阶。压缩包共26个文件&#xff0c;主要包含py源码、C辅助工具、Shell脚本、Markdown说明与示例图片&am…

作者头像 李华
网站建设 2026/9/23 17:50:14

火车轨道检测数据集实战:3900张COCO标注与93.7%准确率验证

简介&#xff1a;这份火车轨道检测数据集面向计算机视觉开发者、轨道交通智能化研究者及深度学习实践者&#xff0c;用于训练和验证轨道区域与障碍物识别模型&#xff0c;可支撑列车前方障碍预警、轨道巡检自动化等场景。资源以COCO标注格式组织&#xff0c;包含3900张原始图片…

作者头像 李华