尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Apache DolphinScheduler Python 任务节点(Python Task)使用详解:脚本编排、自定义参数与执行原理

Apache DolphinScheduler Python 任务节点(Python Task)使用详解:脚本编排、自定义参数与执行原理 Apache DolphinScheduler Python 任务节点Python Task使用详解脚本编排、自定义参数与执行原理【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/GitHub_Trending/dol/dolphinscheduler导读Python 任务是 Apache DolphinScheduler 中最常用的任务类型之一用于在数据编排工作流中直接执行一段或多段 Python 脚本。本文以 Python 节点官方文档 为主体结合仓库中dolphinscheduler-task-plugin/dolphinscheduler-task-python模块的源码与测试用例系统讲解 Python 任务的创建流程、核心参数、自定义参数的传递机制以及 Worker 端生成临时脚本并交由租户同名 Linux 用户执行的底层实现原理。阅读完本文你将能够独立在工作流 DAG 中配置 Python 任务、通过${param}占位符实现脚本复用并理解任务运行失败时可排查的工程细节。综述Python 任务类型是什么Python 任务类型用于创建 Python 类型的任务并执行一系列的 Python 脚本。其核心工作机制是当 Worker 执行该任务时会生成一个临时的 Python 脚本文件并使用与租户同名的 Linux 用户来执行这个脚本。这一机制由源码直接印证PythonTask.java 中的handle()方法依次完成三步动作调用buildPythonScriptContent()基于任务配置的原始脚本生成脚本内容调用buildPythonCommandFilePath()生成脚本文件的落盘路径调用createPythonCommandFileIfNotExists()将脚本内容写入临时文件随后通过 Shell 命令执行。其中脚本文件的命名与落盘位置固定为// PythonTask.java protected String buildPythonCommandFilePath() { return String.format(%s/py_%s.py, taskRequest.getExecutePath(), taskRequest.getTaskAppId()); }即脚本会被写入 Worker 的执行目录executePath下文件名为py_任务实例ID.py。写入时还会自动在脚本头部添加#-*- encodingutf8 -*-编码声明避免因编码问题导致脚本解析异常sb.append(#-*- encodingutf8 -*-).append(System.lineSeparator());Python 任务通过插件化机制注册PythonTaskChannelFactory使用AutoService(TaskChannelFactory.class)自动注册任务类型名为PYTHON见 PythonTaskChannelFactory.java这也是前端任务类型下拉框与后端分发识别的统一标识。创建任务在 DolphinScheduler 的 Web UI 中创建 Python 任务的步骤如下点击项目管理 - 项目名称 - 工作流定义点击创建工作流按钮进入 DAG 编辑页面。在左侧工具栏中拖动 Python 图标docs/img/tasks/icons/python.png到画板上即可完成节点的创建。上图展示了 Python 任务节点的核心配置界面左侧为任务通用参数区域右侧脚本输入框用于填写待执行的 Python 代码。创建完成后即可在节点配置面板中填写任务参数。需要说明的是Python 节点还支持在创建时关联资源文件见下文参数说明用于加载外部 Python 模块或数据文件。任务参数Python 任务节点的参数分为两部分默认任务参数所有任务类型共有的公共参数与Python 任务特有参数。默认任务参数公共参数所有 DolphinScheduler 任务类型共享一组公共参数Python 任务同样适用。完整定义见 DolphinScheduler任务参数附录常用项摘录如下任务参数描述任务名称任务的名称同一个工作流定义中的节点名称不能重复。运行标志标识这个节点是否需要调度执行如果不需要执行可以打开禁止执行开关。缓存执行标识这个节点是否需要进行缓存。若开启缓存对于相同标识相同任务版本、相同任务定义、相同参数传入的任务运行时若已存在缓存过的任务则直接复用结果不再重复执行。描述当前节点的功能描述。任务优先级Worker 线程数不足时根据优先级从高到低依次执行任务优先级相同时按先到先得原则执行。Worker 分组设置分组后任务会被分配给对应 Worker 组的机器执行若选择 Default则随机选择一个 Worker 执行。任务组名称任务资源组未配置则不生效。组内优先级一个任务组内此任务的优先级。环境名称配置任务执行的环境。失败重试次数任务失败后重新提交的次数可在下拉菜单中选择或手动填充。失败重试间隔任务失败后重新提交任务的时间间隔可在下拉菜单中选择或手动填充。CPU 配额为执行的任务分配指定的 CPU 时间配额单位为百分比默认 -1 代表不限制例如 1 个核心满载为 100%16 个核心为 1600%。该功能由task.resource.limit.state配置控制见架构配置文档。最大内存为执行的任务分配指定的内存大小超过会触发 OOM 被 Kill 且不会自动重试单位 MB默认 -1 代表不限制。同样由task.resource.limit.state控制。超时告警设置超时告警、超时失败。当任务超过超时时长后会发送告警邮件并且任务执行失败。资源任务执行时所需资源文件。前置任务设置当前任务的前置上游任务用于构建 DAG 依赖关系。延时执行时间任务延迟执行的时间以分钟为单位。Python 任务特有参数任务参数描述脚本用户开发的 PYTHON 程序。自定义参数PYTHON 局部的用户自定义参数会替换脚本中以${变量}形式出现的占位内容。其中自定义参数支持配置参数名、参数类型如IN、数据类型如VARCHAR与参数值可在界面中通过号新增多条。从源码角度看这两个参数对应 PythonParameters.java 中的字段Data public class PythonParameters extends AbstractParameters { /** origin python script */ private String rawScript; // 对应脚本参数 private ListResourceInfo resourceList; // 对应资源参数 Override public boolean checkParameters() { return rawScript ! null !rawScript.isEmpty(); } }值得注意的校验逻辑任务参数是否合法只检查rawScript是否非空。也就是说脚本内容为空的任务在初始化阶段就会抛出python task params is not valid异常见PythonTask.init()而自定义参数、资源文件均为可选项。此外任务参数在运行时通过JSONUtils.parseObject(taskRequest.getTaskParams(), PythonParameters.class)反序列化而来即界面上的所有配置最终以 JSON 形式随任务上下文传递到 Worker。任务样例简单打印一行文字该样例模拟了常见的简单任务——只需要一两行脚本就能运行。这里以打印一行日志为例该任务仅会在日志文件中输出一行文本This is a demo of python taskprint(This is a demo of python task)运行后可在任务实例的日志中看到对应输出。使用自定义参数该样例模拟了带自定义参数的任务。为了更便捷地复用已有任务、或应对动态需求通常会使用变量来保证脚本的复用性。操作步骤如下在节点配置的自定义参数区域新增一行参数参数名param_key、参数类型IN、数据类型VARCHAR、参数值param_val在脚本区域使用print函数并通过${param_key}占位符引用该参数print(${param_key})保存并运行任务后在日志中可以看到占位符被替换为参数值打印出param_val。上图展示了自定义参数与脚本占位符的组合用法自定义参数栏中配置了param_key param_val脚本中通过${param_key}引用。自定义参数的替换原理占位符替换并非简单的前端文本处理而是由 Worker 端在生成脚本内容时完成。PythonTask.buildPythonScriptContent()的实现如下protected String buildPythonScriptContent() { log.info(raw python script : {}, pythonParameters.getRawScript()); String rawPythonScript pythonParameters.getRawScript().replaceAll(\\r\\n, System.lineSeparator()); MapString, Property paramsMap mergeParamsWithContext(pythonParameters); return ParameterUtils.convertParameterPlaceholders(rawPythonScript, ParameterUtils.convert(paramsMap)); }从这段代码可以看出两个关键细节跨平台行尾统一脚本中的\r\nWindows 换行会被统一替换为系统默认换行符避免在 Linux 执行环境下因行尾不一致导致脚本异常占位符替换convertParameterPlaceholders()会把rawScript中的${param_key}替换为来自mergeParamsWithContext()合并后的参数值。mergeParamsWithContext()返回的是taskRequest.getPrepareParamsMap()即该参数是任务上下文工作流参数、全局参数、上游任务输出等与节点自定义参数合并后的结果——因此${param_key}不仅能引用节点自定义参数也能引用工作流级别的全局参数与上下游传递的变量。此外PythonTask.handle()中还调用了pythonParameters.dealOutParam(...)与taskRequest.setVarPool(...)这意味着 Python 脚本运行后产生的输出参数OUT 类型参数会被写回任务变量池供下游任务继续引用形成上游 Python 任务输出 → 下游任务输入的参数链路。Python 解释器选择PYTHON_LAUNCHERWorker 执行 Python 脚本时使用哪个解释器由环境变量PYTHON_LAUNCHER决定。源码中buildPythonExecuteCommand()的注释与实现说明了这一策略/** * If user have set the PYTHON_LAUNCHER environment, we will use the PYTHON_LAUNCHER, * if not, we will default use python. */ protected String buildPythonExecuteCommand(String pythonFile) { Preconditions.checkNotNull(pythonFile, Python file cannot be null); String pythonHome String.format(${%s}, PYTHON_LAUNCHER); return pythonHome pythonFile; }即最终执行的命令形如python /path/to/py_taskAppId.py若在 Worker 的租户环境或任务环境名称中配置了PYTHON_LAUNCHER环境变量则优先使用该变量指向的解释器例如指定 conda 虚拟环境中的 Python否则回退使用系统默认的python命令。该行为也由单元测试直接验证PythonTaskTest.javaTest public void buildPythonExecuteCommand() throws Exception { PythonTask pythonTask createPythonTask(); Assertions.assertEquals(${PYTHON_LAUNCHER} test.py, pythonTask.buildPythonExecuteCommand(test.py)); }因此当生产环境中需要指定 Python 版本或使用特定虚拟环境如内置 pandas、numpy 的环境时应通过配置PYTHON_LAUNCHER环境变量来实现而不是依赖 Worker 机器默认的 Python。执行流程与故障排查视角综合以上源码分析一个 Python 任务从提交到完成的核心调用链可归纳如下任务分发Master 依据 Worker 分组将任务实例下发到 Worker参数解析PythonTaskChannel.parseParameters()将任务参数 JSON 反序列化为PythonParameters初始化校验PythonTask.init()校验rawScript非空否则抛出python task params is not valid脚本生成buildPythonScriptContent()完成换行统一与${变量}占位符替换脚本落盘createPythonCommandFileIfNotExists()在 Worker 执行目录生成py_taskAppId.py并写入 UTF-8 编码声明与脚本内容命令执行buildPythonExecuteCommand()组装${PYTHON_LAUNCHER} 脚本路径命令通过ShellCommandExecutor以租户同名 Linux 用户执行结果回写根据退出码设置任务状态输出参数写回变量池供下游任务使用执行过程中可随时调用cancel()终止进程。据此任务失败时可从以下角度排查租户/权限问题脚本由与租户同名的 Linux 用户执行若该用户对执行目录无写权限、或缺少对应 Python 解释器的执行权限任务会失败解释器缺失未配置PYTHON_LAUNCHER且系统默认python不存在或配置的解释器路径无效脚本语法/依赖脚本本身语法错误或依赖第三方库未安装此时应通过PYTHON_LAUNCHER指向已装依赖的解释器参数替换异常${变量}占位符在运行环境中未定义时convertParameterPlaceholders的处理结果可能与预期不一致建议先在自定义参数或工作流参数中显式声明默认值。注意事项原文档注意事项一节内容为 None即 Python 任务目前没有额外限制。但从源码实现可以补充以下实操提醒脚本内容为空时任务无法通过参数校验会直接初始化失败任务脚本会被自动加上 UTF-8 编码声明#-*- encodingutf8 -*-脚本中不建议再重复声明或使用与 UTF-8 冲突的编码脚本文件位于 Worker 的执行目录下命名格式为py_taskAppId.py排查问题时可直接到对应 Worker 机器查看生成的临时脚本内容使用自定义参数时注意${变量}替换发生在脚本生成阶段即 Worker 执行前并非 Python 运行时语法因此不要在 Python 字符串内部刻意使用未定义的${...}文本以免被意外替换。参考与延伸阅读Python 任务官方文档docs/docs/zh/guide/task/python.mdPython 任务源码模块dolphinscheduler-task-plugin/dolphinscheduler-task-python核心实现PythonTask.java参数模型PythonParameters.java插件注册PythonTaskChannelFactory.java单元测试PythonTaskTest.java任务公共参数附录DolphinScheduler任务参数附录相关配置项说明架构配置文档Python 任务图标docs/img/tasks/icons/python.png【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/GitHub_Trending/dol/dolphinscheduler创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表