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 等)与值value。varPool则是一组Property的序列化形态(JSON 字符串),在 Master 与 Worker 之间、节点与节点之间流转。整条数据链路可以概括为:
task 定义 OUT 参数 (localParam) → 任务执行产出 → varPool (JSON) → Master 合并传递 → Worker 解析合并 → 节点执行替换 → 输出参数回写 localParam → 继续传递给下游二、参数的使用:Master 侧 varPool 的获取与合并
2.1 从前置节点收集 varPool
当一个任务实例(taskInstance)需要创建时,Master 会从 DAG 中获取该节点的直接前置节点 preTasks,并收集这些前置节点产出的varPool(List<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,返回参数的处理都遵循同一套流程,文档将其归纳为六个步骤:
- 获取到的 processor 的结果为 String;
- 判断 processor 结果是否为空,为空则退出;
- 判断 localParam 是否为空,为空则退出;
- 获取 localParam 中方向为 OUT 的参数,为空则退出;
- 将 String 按对应格式格式化(SQL 为
List<Map<String, String>>,SHELL 为Map<String, String>); - 将匹配好值的参数赋值给 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),仅供参考