Apache DolphinScheduler Flink 任务节点完全指南:参数详解、部署模式与底层命令生成原理 Apache DolphinScheduler Flink 任务节点完全指南参数详解、部署模式与底层命令生成原理【免费下载链接】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 官方文档中 Flink 任务节点docs/docs/zh/guide/task/flink.md为主体结合dolphinscheduler-task-flink任务插件的真实源码与单元测试系统讲解如何在 DolphinScheduler 中创建、配置并运行 Flink含 Flink SQL任务。读完本文你将掌握 Flink 节点的全部任务参数含义、四种部署模式local / cluster / application / standalone的差异、Worker 底层生成的flink run与sql-client.sh命令细节以及从环境配置、资源上传到任务运行的一条完整实战链路。Flink 节点综述Flink 任务类型用于在 DolphinScheduler 的工作流中执行 Flink 程序。根据用户选择的程序类型Worker 会采用两种不同的提交方式程序类型为 Java、Scala 或 PythonWorker 使用 Flink 命令行工具提交任务即flink run或run-application命令程序类型为 SQLWorker 使用 Flink 自带的sql-client.sh脚本提交任务。这一逻辑在源码中有直接体现。任务插件入口 FlinkTask.java 继承自AbstractYarnTask其getScript()方法调用FlinkArgsUtils.buildRunCommandLine(...)拼装命令行而 FlinkArgsUtils.java 会按程序类型分流ProgramType.SQL→buildRunCommandLineForSql(...)使用${FLINK_HOME}/bin/sql-client.sh其余类型JAVA / SCALA / PYTHON→buildRunCommandLineForOthers(...)使用${FLINK_HOME}/bin/flink。其中${FLINK_HOME}是环境占位符具体指向哪个 Flink 安装目录由任务执行环境的配置决定见下文“环境配置”一节。创建任务创建 Flink 任务节点的操作路径与 DolphinScheduler 其他任务类型一致进入项目管理 → 项目名称 → 工作流定义点击“创建工作流”按钮进入 DAG 编辑页面从左侧工具栏中拖动 Flink 任务节点图标为flink到画板中在节点配置面板中填写下文介绍的各项任务参数保存后即可将节点接入 DAG。任务参数详解除DolphinScheduler 任务参数附录中描述的默认任务参数任务名称、运行标志、缓存执行、描述、任务优先级、Worker 分组、失败重试次数/间隔、超时告警、资源、前置任务、延时执行时间等外Flink 节点还特有如下参数任务参数描述程序类型支持 Java、Scala、Python 和 SQL 四种语言主函数的 ClassFlink 程序的入口 Main Class 的全路径如org.example.Main主程序包执行 Flink 程序的 jar 包通过资源中心上传部署方式支持 cluster、local 和 applicationFlink 1.11 及之后版本支持以及 standalone 模式初始化脚本用于初始化会话上下文的脚本文件Flink SQL 场景使用脚本用户开发的应该执行的 SQL 脚本文件Flink SQL 场景使用Flink 版本根据所需环境选择对应的版本即可任务名称选填Flink 程序的名称对应提交到 Yarn 上的应用名jobManager 内存数用于设置 jobManager 内存数可根据实际生产环境设置对应的内存数Slot 数量用于设置 Slot 的数量可根据实际生产环境设置对应的数量taskManager 内存数用于设置 taskManager 内存数可根据实际生产环境设置对应的内存数taskManager 数量用于设置 taskManager 的数量可根据实际生产环境设置对应的数量并行度用于设置执行 Flink 任务的并行度Yarn 队列用于设置 Yarn 队列默认使用 default 队列主程序参数设置 Flink 程序的输入参数支持自定义参数变量的替换选项参数设置 Flink 命令的选项参数例如-D、-C、-yt自定义参数是 Flink 局部的用户自定义参数会替换脚本中以${变量}的内容这些 UI 参数与任务插件的数据模型一一对应。查看 FlinkParameters.java可以看到后台字段定义mainJar主程序包、mainClass主函数的 Class、mainArgs主程序参数deployMode部署方式枚举FlinkDeployModeLOCAL / CLUSTER / APPLICATION / STANDALONEslotSlot 数量、parallelism并行度、taskManagertaskManager 数量、jobManagerMemory/taskManagerMemoryJobManager / TaskManager 内存appName任务名称、yarnQueueYarn 队列、others选项参数flinkVersionFlink 版本、programType程序类型、initScript初始化脚本、rawScriptSQL 脚本。参数校验规则源码级FlinkParameters.checkParameters() 定义了保存/提交任务时必填项的校验逻辑程序类型为 SQLrawScriptSQL 脚本不能为空程序类型为 Java / Scala / PythonmainJar主程序包不能为空。此外getResourceFilesList() 会自动把mainJar加入资源文件列表保证 jar 包随任务分发到 Worker 本地。源码视角Worker 如何生成提交命令理解 DolphinScheduler 如何把“参数”翻译成“命令”有助于排查线上 Flink 任务提交失败问题。核心逻辑集中在 FlinkArgsUtils.java其行为随部署方式和Flink 版本而变化命令选项常量定义在 FlinkConstants.java。Java / Scala / Python 程序的命令生成默认部署方式为CLUSTER见 FlinkArgsUtils.java。命令生成遵循以下分支CLUSTER 模式Flink 版本1.12/1.13时使用flink run -t yarn-per-job更早版本使用flink run -m yarn-clusterAPPLICATION 模式使用flink run-application -t yarn-application对应 Flink 1.11 引入的 Application ModeLOCAL / STANDALONE 模式直接使用flink run不再追加 Yarn 相关参数。随后按部署方式追加资源参数所有选项常量均定义在 FlinkConstants.java参数命令选项说明Slot 数量-ys每个 TaskManager 的 Slot 数任务名称-ynmYarn 上的应用名taskManager 数量-yn仅 Flink 1.10 之前版本支持1.10 起该参数被移除jobManager 内存-yjm如1024mtaskManager 内存-ytm如1024m并行度-p全局并行度主函数 Class-cPython 程序不追加见注意事项Python 入口-py仅 Python 程序追加附加退出关闭-sae提交命令固定追加保证 CLI 中断时集群任务同步关闭、任务状态与集群状态保持一致主程序参数会经过 ParameterUtils.convertParameterPlaceholders 做${变量}占位符替换从而实现与 DolphinScheduler 全局/局部参数的联动。Yarn 队列则通过 determinedYarnQueue 处理新版本-t目标模式使用-Dyarn.application.queuexxx旧版本使用-yqu xxx若“选项参数”中已包含队列选项则不重复追加。以上行为均有单元测试固化见 FlinkArgsUtilsTest.java。测试断言的命令示例# APPLICATION 模式 ${FLINK_HOME}/bin/flink run-application -t yarn-application -ys 4 -ynm demo-app-name \ -yjm 1024m -ytm 1024m -p 4 -sae -c org.example.Main /opt/job.jar # CLUSTER 模式Flink 1.11 / 1.10 ${FLINK_HOME}/bin/flink run -m yarn-cluster -ys 4 -ynm demo-app-name \ -yjm 1024m -ytm 1024m -p 4 -sae -c org.example.Main /opt/job.jar # CLUSTER 模式Flink 1.12 ${FLINK_HOME}/bin/flink run -t yarn-per-job -ys 4 -ynm demo-app-name \ -yjm 1024m -ytm 1024m -p 4 -sae -c org.example.Main /opt/job.jar # LOCAL 模式 ${FLINK_HOME}/bin/flink run -p 4 -sae -c org.example.Main /opt/job.jarFlink SQL 的命令生成SQL 程序的命令生成走 buildRunCommandLineForSql${FLINK_HOME}/bin/sql-client.sh -i 初始化脚本路径 -f SQL脚本路径 [选项参数]其中-i指向初始化脚本-f指向SQL 脚本两个脚本由 FileUtils.java 在 Worker 执行目录下生成命名形如${taskAppId}_init.sql与${taskAppId}_node.sql初始化脚本的内容由 buildInitOptionsForSql 自动生成根据部署方式写入一组set ...语句再拼接用户在 UI 中填写的“初始化脚本”。set语句与部署方式的关系测试用例见 FlinkArgsUtilsTest.javaLOCAL 模式set execution.targetlocalCLUSTER 模式Yarnset execution.targetyarn-per-job并视参数追加set taskmanager.numberOfTaskSlots...、set yarn.application.name...、set jobmanager.memory.process.size...、set taskmanager.memory.process.size...、set yarn.application.queue...无论何种模式只要设置了并行度都会追加set parallelism.default...从源码注释看当前 Flink SQL 在 Yarn 上仅支持 yarn-per-job 模式。任务样例一执行 WordCount 程序WordCount 是大数据生态中最常见的入门案例常用于 MapReduce、Flink、Spark 等计算框架核心是统计输入文本中相同单词的数量Flink 官方 Releases 附带了此示例作业。下面按完整链路演示在 DolphinScheduler 中配置并运行。第一步在 DolphinScheduler 中配置 Flink 环境若生产环境要使用 Flink 任务类型需要先配置好所需环境配置文件为部署目录下的bin/env/dolphinscheduler_env.sh。其中关键是设置FLINK_HOME环境变量使${FLINK_HOME}/bin/flink与${FLINK_HOME}/bin/sql-client.sh能够被 Worker 解析到真实路径。第二步上传主程序包使用 Flink 任务节点前需要利用资源中心上传执行程序的 jar 包具体操作可参考资源中心文档。配置完成资源中心后直接使用拖拽方式即可上传目标文件。第三步配置 Flink 节点根据上文参数说明配置所需内容程序类型选择 Java或 Scala部署方式选择符合集群实际情况的模式如 cluster填写主函数的 Class如 WordCount 示例的入口全路径、选择上传的主程序包并可按需设置并行度、内存、Yarn 队列等。任务样例二执行 FlinkSQL 程序Flink SQL 场景同样根据上文参数说明配置即可注意以下几点程序类型选择 SQL初始化脚本用于初始化会话上下文可选如设置 Catalog、注册 UDF 等脚本填写要执行的 SQL 语句该脚本为必填项否则无法通过参数校验见上文checkParameters源码说明部署方式与资源参数Slot、内存、并行度、应用名、Yarn 队列会被自动翻译为初始化脚本中的set语句并随sql-client.sh -i加载。注意事项Java 和 Scala 只是用来标识没有本质区别如果是 Python 开发的 Flink 程序则没有主函数的 Class其余配置与 Java/Scala 一致从 FlinkArgsUtils.java 可以看到Python 程序不会追加-c主类参数而是追加-py使用 SQL 执行 Flink SQL 任务目前只支持 Flink 1.13 及以上版本源码中通过 Flink 版本常量1.13参与命令分支判断见 FlinkArgsUtils.java部署方式选择受 Flink 版本约束application 模式需要 Flink 1.11 及之后版本Flink 1.10 之后-yntaskManager 数量参数被移除提交命令不会携带该选项SQL 任务在 Yarn 上仅支持 yarn-per-job 模式LOCAL 与 CLUSTER 之外的部署方式在生成初始化脚本时不会写入execution.target见 FlinkArgsUtils.java任务提交命令固定追加-saeattached 模式退出时尽力关闭集群保证任务状态与集群任务状态同步。深入阅读Flink 任务节点官方文档本文主体来源DolphinScheduler 任务参数附录Flink 节点通用的默认任务参数含失败重试、超时告警、Worker 分组、缓存执行等资源中心文档jar 包、脚本等资源的上传与管理方式FlinkTask.java任务插件入口负责参数解析、脚本文件生成与命令执行FlinkArgsUtils.java提交命令与 SQL 初始化选项的生成核心FlinkParameters.java任务参数模型与必填项校验FlinkConstants.javaflink run/sql-client.sh的全部命令选项常量FlinkArgsUtilsTest.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),仅供参考