Apache DolphinScheduler SeaTunnel 任务类型完全指南:CLI 封装原理、参数配置与实战样例
Apache DolphinScheduler SeaTunnel 任务类型完全指南CLI 封装原理、参数配置与实战样例【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/GitHub_Trending/dol/dolphinscheduler导读本文系统讲解 Apache DolphinScheduler 中的 SeaTunnel 任务类型涵盖任务创建、启动脚本选择、FLINK / SPARK / SEATUNNEL_ENGINE 三种引擎的参数配置、自定义配置与资源中心配置两种写法以及变量替换与一个完整的 Flink 引擎实战样例。读完本文你将掌握如何在 DolphinScheduler 工作流中快速接入 Apache SeaTunnel 数据同步任务并理解该任务类型在 worker 端被封装为$SEATUNNEL_HOME/bin/启动脚本调用的底层原理。文中所有结论均以当前仓库 dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel 模块的源码与官方中文文档为事实依据。SeaTunnel 任务类型综述SeaTunnel任务类型用于创建并执行 SeaTunnel 类型的任务。其核心运行机制是worker 执行该任务时会通过${SEATUNNEL_HOME}/bin/下的启动脚本如seatunnel.sh/start-seatunnel-*-connector-v2.sh解析并执行 config 文件。也就是说该任务类型本质上是对 SeaTunnel CLI 的 shell 封装——DolphinScheduler 负责生成一段 shell 命令交给 shell 执行器运行最终把 SeaTunnel 作业接入到整个工作流的调度编排中。这一封装逻辑在源码中非常清晰。核心实现位于 SeatunnelTask.javabuildCommand()首先拼接${SEATUNNEL_HOME}/bin/与用户选择的启动脚本名随后通过buildOptions()依次追加--config 配置文件路径与任务参数最终命令形如${SEATUNNEL_HOME}/bin/start-seatunnel-flink-13-connector-v2.sh --config /path/to/config.conf ...由ShellCommandExecutor在 worker 上以 shell 方式执行退出码与输出参数会回传给调度系统。创建任务在 DolphinScheduler 的 Web UI 中创建 SeaTunnel 任务只需两步进入项目管理 → 项目名称 → 工作流定义点击“创建工作流”按钮进入 DAG 编辑页面从左侧工具栏中拖拽SeaTunnel任务节点到画板中。任务节点放置到 DAG 画板后即可在右侧配置面板中填写下文介绍的各类参数并将该节点与其他任务节点连线纳入整个工作流的依赖编排。任务参数详解SeaTunnel 任务参数分为三组通用参数启动脚本、三种引擎各自的专属参数FLINK / SPARK / SEATUNNEL_ENGINE、以及配置来源自定义配置或资源中心文件与脚本内容。默认参数如任务名称、运行标志、失败重试、超时告警、优先级等请参考 DolphinScheduler 任务参数附录 中“默认任务参数”一栏。通用参数启动脚本启动脚本选择运行任务所使用的启动脚本。不同 SeaTunnel 发行包中bin/目录下的脚本可能存在差异请以实际${SEATUNNEL_HOME}/bin/下的文件为准。DolphinScheduler 支持以下启动脚本脚本名对应引擎seatunnel.shSeaTunnel 自研引擎Zeta / SEATUNNEL_ENGINEstart-seatunnel-flink-13-connector-v2.shFlink 1.13 Connector V2start-seatunnel-flink-15-connector-v2.shFlink 1.15 Connector V2start-seatunnel-flink-connector-v2.shFlink Connector V2start-seatunnel-flink.shFlink Connector V1start-seatunnel-spark-2-connector-v2.shSpark 2 Connector V2start-seatunnel-spark-3-connector-v2.shSpark 3 Connector V2start-seatunnel-spark-connector-v2.shSpark Connector V2start-seatunnel-spark.shSpark Connector V1从源码看启动脚本并非随意填写的字符串。SeatunnelParameters.java 中用正则^[A-Za-z0-9][A-Za-z0-9._-]*\.sh$校验脚本名必须以.sh结尾同时参数校验要求启动脚本必填且配置来源只能二选一自定义配置useCustom true且脚本非空或资源中心文件useCustom false且资源列表恰好 1 个文件。任一条件不满足任务初始化阶段就会抛出SeaTunnel task params is not valid异常。FLINK 引擎参数运行模式Run Mode支持run与run-application两种模式。对应源码 SeatunnelFlinkParameters.RunModeEnum当运行模式为RUN时实际追加--deploy-mode run为RUN_APPLICATION时追加--deploy-mode run-application选择NONE则不追加任何模式参数各引擎默认行为。选项参数用于添加 Flink 引擎本身的参数例如-m yarn-cluster -ynm seatunnel。该值原样拼接到命令行尾部由 SeatunnelFlinkTask.buildOptions() 在非空时追加。SPARK 引擎参数部署方式指定部署模式可选cluster、client对应 DeployModeEnum 枚举。SPARK 任务会把部署模式转换为--deploy-mode cluster|client参数。Master指定 Master 模式可选yarn、local、spark、mesos。其中spark与mesos需要额外指定 Master 服务地址例如127.0.0.1:7077。对应源码 SeatunnelSparkParameters.MasterTypeEnumyarn/local直接输出--master yarn|local而spark/mesos会拼接为--master spark://地址或--master mesos://地址。此外当部署方式为local时Master 会强制取local无需填写地址参数校验也会保证在非 local 部署下 Master 必填、且 spark/mesos 模式必须提供 masterUrl见 SeatunnelSparkParameters.checkParameters()。命令行拼装逻辑见 SeatunnelSparkTask.buildOptions()。SEATUNNEL_ENGINE 引擎参数部署方式指定部署模式可选cluster、local同样对应 DeployModeEnum。SEATUNNEL_ENGINE 任务会把部署模式转换为--deploy-mode cluster|local。选项参数与 FLINK 类似others字段用于追加 SeaTunnel 引擎自身的额外参数见 SeatunnelEngineTask.buildOptions()。配置来源与脚本内容自定义配置支持在任务节点内直接编写配置或从资源中心选择配置文件。从源码看二者最终都会在 worker 执行目录下生成一个临时配置文件再交给 SeaTunnel 执行自定义配置rawScript会经过换行归一化与参数占位符解析后写入文件资源中心配置则读取选中资源在本地落盘的绝对路径内容后写入文件见 SeatunnelTask.buildOptions()。生成文件名的后缀由配置内容自动判定内容是合法 JSON 则生成.json否则生成.conf见 formatDetector()。脚本在任务节点中自定义配置信息包括四部分env、source、transform、sink。这是 SeaTunnel 的 HOCON 风格配置文件结构详见下一节样例。自定义参数 / 全局参数当定义了自定义参数或全局参数时会将这些参数传递给 SeaTunnel 任务任务中可通过${}引用该参数在运行时动态替换参数值。其实现细节在 SeatunnelTask.generateTaskParameters()DolphinScheduler 会把本地参数Direct.IN方向与全局参数收集起来逐个转换为-i keyvalue形式的命令行参数追加到启动脚本后同时parseScript()会通过convertParameterPlaceholders对配置文件中的${}占位符做一次预处理替换保证变量在配置与命令行两个层面都能生效。值为空时会输出单引号也会做 Bash 安全的转义处理quoteForBash避免命令注入与解析错误。任务样例Flink 引擎 FakeSource → Console下面以一个完整样例演示如何用 Flink 引擎从 Fake 源读取数据并打印到控制台。在 DolphinScheduler 中配置 SeaTunnel 环境若生产环境需要使用 SeaTunnel 任务类型必须先配置好所需环境。配置文件位于/dolphinscheduler/conf/env/dolphinscheduler_env.sh需要在其中声明SEATUNNEL_HOME指向 SeaTunnel 安装目录并加入 PATH使 worker 能找到bin/下的启动脚本以 Flink 引擎为例还须保证环境变量中配置了 Flink 的HADOOP_HOME、FLINK_HOME等依赖若使用 Spark 引擎则对应配置SPARK_HOME。SeaTunnel 启动脚本会按${SEATUNNEL_HOME}/bin/定位请确保该目录下存在上表所列的启动脚本文件。配置 SeaTunnel 任务节点根据上文参数说明在任务节点配置面板中完成配置即可。下图展示了一个典型的 SeaTunnel 任务节点配置启动脚本、引擎参数、配置内容等关键配置项归纳如下配置项本例取值说明启动脚本start-seatunnel-flink-13-connector-v2.sh选择 Flink 1.13 Connector V2 启动脚本运行模式runFlink 引擎运行模式配置来源自定义配置直接在脚本框中编写下文 Config 样例自定义参数无需要变量替换时在此定义Config 样例任务节点“脚本”栏可粘贴如下 SeaTunnel 配置文件以 Flink 引擎为例从 Fake 源读取两列数据经 SQL 变换后打印到控制台env { execution.parallelism 1 } source { FakeSource { result_table_name fake field_name name,age } } transform { sql { sql select name,age from fake } } sink { ConsoleSink {} }配置说明env定义执行环境参数execution.parallelism 1表示并行度为 1source数据源插件配置。FakeSource是 SeaTunnel 内置的模拟数据源result_table_name为中间结果表名field_name声明字段name,agetransform变换插件配置示例使用sql插件对fake表执行select name,age from fakesink结果输出插件配置ConsoleSink将数据打印到控制台便于验证整条链路是否打通。在实际生产任务中可将source替换为 MySQL、Kafka 等真实数据源将sink替换为目标库并补充对应的 connector 依赖。变量替换示例假设在任务节点“自定义参数”中定义了参数var_name seatunnel_test或定义了同名全局参数那么在脚本配置中可直接用${var_name}引用env { execution.parallelism 1 } source { FakeSource { result_table_name fake field_name ${var_name},age } } transform { sql { sql select name,age from fake } } sink { ConsoleSink {} }任务运行时DolphinScheduler 会将该参数以-i var_nameseatunnel_test的形式传给 SeaTunnel CLI并对配置中的${var_name}做预替换实现参数动态化。多个参数会生成多条-i参数全局参数与本地参数同时存在时二者会一并收集见 SeatunnelTask.generateTaskParameters()。任务执行原理从配置到命令行结合上述样例SeaTunnel 任务在 worker 端的完整执行链路如下任务被调度后worker 根据任务参数实例化对应引擎的任务对象Flink / Spark / Engine并反序列化参数见 SeatunnelFlinkTask.init() 等三个引擎的init()init()中调用checkParameters()做参数合法性校验非法参数直接抛异常中止任务handle()中调用buildCommand()组装 shell 命令${SEATUNNEL_HOME}/bin/启动脚本 --config 临时配置文件 [--deploy-mode ...] [--master ...] [-i keyvalue] ...配置内容先写入 worker 执行目录下的临时文件seatunnel_taskAppId.conf|.json再交给ShellCommandExecutor以子进程方式执行执行结束后DolphinScheduler 采集退出码与输出参数写入任务实例结果供后续节点与告警使用任务被取消时调用shellCommandExecutor.cancelApplication()杀掉对应进程见 SeatunnelTask.cancelApplication()。该模块的单元测试覆盖了参数校验等关键逻辑可参考 SeatunnelParametersTest.java 与 SeatunnelTaskTest.java。支持 SeaTunnel 版本本文示例基于 SeaTunnel2.3.x版本的命令行参数与启动脚本官方已验证版本2.3.1、2.3.2、2.3.3其他版本由于该任务类型本质是对 SeaTunnel CLI 的封装只要 SeaTunnel 启动脚本与命令行参数保持兼容通常可直接使用更高版本。建议升级 SeaTunnel 版本后在测试环境先做回归验证确认bin/下脚本命名、--config、--deploy-mode等参数行为未发生变化再上线到生产工作流。【免费下载链接】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),仅供参考
上一篇/下一篇内容由系统自动关联
返回资讯列表 →