十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

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、方向directIN/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并收集这些前置节点产出的varPoolListProperty随后合并为一个 varPool。文档明确描述了合并过程中同名变量的处理逻辑这一逻辑在代码中有两处印证文档描述的规则同名变量冲突时若所有值均为 null则合并后的值为 null若有且仅有一个值为非 null则合并后的值为该非 null 值若所有值均非 null则取这些 varPool 所属 taskInstance 中endtime 最早的那一个合并过程中所有合并进来的 Property 的方向都会被更新为IN合并结果保存在taskInstance.varPool中。对应的核心实现位于 WorkflowExecuteRunnable.initializeTaskInstanceVarPool()// 获取当前任务的前置节点 preTasks String preTasks workflowExecuteContext.getWorkflowGraph() .getTaskNodeByCode(taskInstance.getTaskCode()).getPreTasks(); SetLong preTaskList new HashSet(JSONUtils.toList(preTasks, Long.class)); ProcessInstance workflowInstance workflowExecuteContext.getWorkflowInstance(); if (CollectionUtils.isEmpty(preTaskList)) { // 无前置节点时直接继承流程实例的 varPool taskInstance.setVarPool(workflowInstance.getVarPool()); return; } // 收集所有前置节点实例的 varPool并按 endTime 排序 ListString 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 合并工具VarPoolUtilsvarPool 的合并、反序列化与减法操作统一封装在 VarPoolUtils.javadeserializeVarPool(String varPoolJson)将 JSON 字符串反序列化为ListPropertymergeVarPoolJsonString(ListString varPoolJsons)批量合并多个 varPool 的 JSON 串返回合并后的 JSON 串空集合返回 nullmergeVarPool(ListListProperty varPools)核心合并逻辑仅处理方向为 OUT 的 Property以property.getProp()变量名为 key 存入 Map后写入的覆盖先写入的subtractVarPool / subtractVarPoolJson从 varPool 中剔除指定变量用于失败重跑等场景下清理过期数据。public ListProperty mergeVarPool(ListListProperty varPools) { if (CollectionUtils.isEmpty(varPools)) return null; if (varPools.size() 1) return varPools.get(0); MapString, Property result new HashMap(); for (ListProperty 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解析为MapString, Property格式其中map 的 key 为property.prop即变量名。在 processor 处理参数时会将三个变量池合并处理变量池说明合并优先级globalParam流程级全局参数高保留varPool前置节点传递来的输出参数中localParam节点自身定义的参数低被替换合并时若存在同名参数高优先级保留低优先级被替换即同名时优先取globalParam其次varPool最后才是localParam。参数会在节点内容执行之前通过正则表达式匹配${变量名}并替换为对应的值。因此一个典型的使用场景是上游 SQL 节点产出 OUT 参数后下游 Shell 节点可以直接在脚本中书写${变量名}引用上游产出值。三、参数的设置SQL 与 SHELL 节点的输出回写文档明确指出目前仅支持 SQL 和 SHELL 节点的参数获取从localParam中获取方向为 OUT 的参数再根据节点类型做不同处理。3.1 SQL 节点ListMapString, String 结构SQL 节点参数返回的结构为ListMapString, StringList 的元素为每行数据Map 的 key为列名value为该列对应的值。匹配规则如下若 SQL 语句返回一行数据根据用户在定义 task 时定义的 OUT 参数名匹配列名匹配成功则取值未匹配到则放弃若 SQL 语句返回多行根据用户定义的类型为LIST的 OUT 参数名匹配列名将对应列的所有行数据转换为ListString作为该参数的值未匹配到则放弃。这一逻辑与 SqlParameters.dealOutParam() 的实现完全对应public void dealOutParam(String result) { if (CollectionUtils.isEmpty(localParams)) return; ListProperty outProperty getOutProperty(localParams); if (CollectionUtils.isEmpty(outProperty)) return; if (StringUtils.isEmpty(result)) { varPool VarPoolUtils.mergeVarPool(Lists.newArrayList(varPool, outProperty)); return; } ListMapString, String sqlResult getListMapByString(result); // 多行按列聚合为 ListString if (sqlResult.size() 1) { // ... 按列名聚合所有行数据 ... for (Property info : outProperty) { if (info.getType() DataType.LIST) { info.setValue(JSONUtils.toJsonString(sqlResultFormat.get(info.getProp()))); } } } else { // 单行直接按列名取值 MapString, 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(keyvalue)} 语法SHELL 节点 processor 执行后的结果返回为MapString, String。用户在编写 shell 脚本时需要在输出中定义${setValue(keyvalue)}格式的标记例如echo ${setValue(custom_keycustom_value)}参数处理时框架会去掉${setValue(...)}外壳按照进行拆分第 0 个为 key第 1 个为 value随后用用户定义 task 时配置的 OUT 参数名与该 key 匹配将 value 作为该参数的值写入。该语法由 TaskOutputParameterParser.java 负责解析值得注意的实现要点同时支持${setValue(...)}与#{setValue(...)}两种写法解析要求表达式必须以setValue(开头、以)}结尾且内部必须包含否则视为非法并告警跳过为了防止构造的超长参数导致 OOMparser 设置了参数行数上限默认单参数最多1024 行与长度上限超出即跳过并打印告警日志解析结果以 key-value 形式写入MapString, 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 为ListMapString, StringSHELL 为MapString, String将匹配好值的参数赋值给 varPoolListProperty其中包含原有 IN 的参数。随后 varPool 被格式化为 JSON 传递给 MasterMaster 接收到 varPool 后会将其中方向为 OUT 的参数回写到 localParam中从而完成“输出参数沉淀回节点定义”的闭环。这也解释了为什么下游节点能够稳定地通过 varPool 读取到上游的 OUT 参数——数据始终以 JSON 形式在节点间显式传递而非依赖共享内存或外部存储。四、全链路时序总结综合文档与源码一次全局参数OUT 参数的完整生命周期如下1. 用户定义 task 时配置 OUT 参数 → 存入 localParam 2. 任务执行SQL/SHELL产出输出 ├─ SQL结果格式化为 ListMapString,String按 OUT 参数名列名匹配取值 └─ SHELL解析 ${setValue(keyvalue)}按 OUT 参数名与 key 匹配取值 3. Worker 将匹配结果与原有 varPool 合并ListProperty含原有 IN 参数→ 序列化为 JSON 4. Master 接收 varPool将 OUT 参数回写到 localParam 5. 下游任务实例创建时Master 收集所有直接前置节点的 varPool → 同名冲突时按endtime 最早者优先合并 → 合并后方向更新为 IN → 存入 taskInstance.varPool 6. Worker 将 varPool 解析为 MapString,Propertykey 为变量名 → 与 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.javaSQL 输出处理见 SqlParameters.javaSHELL 输出解析见 TaskOutputParameterParser.javaMaster 侧合并见 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),仅供参考
返回列表