Apache DolphinScheduler SeaTunnel 任务类型实战指南:配置、执行与源码原理
【免费下载链接】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 官方文档《Apache SeaTunnel》任务章节为主体,结合仓库内
dolphinscheduler-task-seatunnel插件源码与测试用例,完整讲解如何在 DolphinScheduler 中创建、配置与运行 SeaTunnel 数据同步任务,覆盖 Flink / Spark / SeaTunnel Engine 三种引擎的参数体系、配置文件的两种组织方式、参数传递机制以及任务底层的命令构建逻辑。读完本文,你将能够在 DolphinScheduler 的 DAG 中直接编排 SeaTunnel 任务,并理解其「包装 SeaTunnel CLI」的本质实现。
一、任务概述:DolphinScheduler 如何执行 SeaTunnel 任务
SeaTunnel是 Apache DolphinScheduler 内置的一种任务类型(Task Type),用于创建并执行 SeaTunnel 数据集成任务。当 Worker 执行该任务时,DolphinScheduler 并不会自行解析 SeaTunnel 的配置语法,而是将任务包装为一条 Shell 命令,通过${SEATUNNEL_HOME}/bin/目录下的启动脚本(如seatunnel.sh/start-seatunnel-*-connector-v2.sh)解析并运行用户提供的配置文件。
这一设计在源码中有清晰的体现:任务的核心执行类 SeatunnelTask.java 中,buildCommand()方法将${SEATUNNEL_HOME}/bin/与用户选择的启动脚本名拼接成命令的第一段,再追加buildOptions()生成的参数列表,最终交给ShellCommandExecutor以 Shell 方式执行:
private static final String SEATUNNEL_BIN_DIR = "${SEATUNNEL_HOME}/bin/"; private String buildCommand() throws Exception { List<String> args = new ArrayList<>(); args.add(SEATUNNEL_BIN_DIR + seatunnelParameters.getStartupScript()); args.addAll(buildOptions()); String command = String.join(" ", args); log.info("SeaTunnel task command: {}", command); return command; }因此可以明确:SeaTunnel 任务类型本质上是 SeaTunnel CLI 的一层封装,Worker 节点必须预先安装 SeaTunnel 发行版,并正确设置SEATUNNEL_HOME环境变量,任务才能正常运行。
二、创建 SeaTunnel 任务
在 DolphinScheduler 前端界面中创建 SeaTunnel 任务的步骤如下:
- 进入项目管理 -> 项目名称 -> 工作流定义,点击「创建工作流」按钮进入 DAG 编辑页面;
- 从左侧工具栏将 SeaTunnel 任务图标(
)拖拽到画布上;
- 双击任务节点,在弹出的配置面板中填写任务参数;
- 配置完成后,与工作流中的其他任务连线,保存并上线工作流。
SeaTunnel 任务节点默认参数(如任务名称、运行环境、失败重试次数、超时时间、优先级等)与其他任务类型一致,请参考 DolphinScheduler 任务参数附录 中的「默认任务参数」一节。
三、任务参数详解:三大引擎分支
SeaTunnel 任务的核心参数是「启动脚本」与「引擎相关参数」。从源码结构看,该插件将参数按引擎拆分为三个子类,分别对应三种运行方式:
| 源码类 | 对应引擎 | 文件 |
|---|---|---|
SeatunnelFlinkParameters | Flink | flink/SeatunnelFlinkParameters.java |
SeatunnelSparkParameters | Spark | spark/SeatunnelSparkParameters.java |
SeatunnelEngineParameters | SeaTunnel Engine | self/SeatunnelEngineParameters.java |
3.1 启动脚本(Startup Script)
选择用于启动任务的脚本名称。脚本存放于${SEATUNNEL_HOME}/bin/目录下,具体名称可能因 SeaTunnel 发行版而异,请以实际安装目录为准。常见脚本包括:
seatunnel.shstart-seatunnel-flink-13-connector-v2.shstart-seatunnel-flink-15-connector-v2.shstart-seatunnel-flink-connector-v2.shstart-seatunnel-flink.shstart-seatunnel-spark-2-connector-v2.shstart-seatunnel-spark-3-connector-v2.shstart-seatunnel-spark-connector-v2.shstart-seatunnel-spark.sh
需要注意的是,启动脚本名并非可以随意填写。在 SeatunnelParameters.java 中,参数校验逻辑使用正则^[A-Za-z0-9][A-Za-z0-9._-]*\.sh$对脚本名进行合法性校验——必须以字母或数字开头、以.sh结尾,且不能包含路径分隔符。同时checkParameters()还要求配置必须满足以下条件之一,否则任务初始化时会抛出TaskException:
- 选择「自定义配置」(
useCustom = true)且脚本内容(rawScript)非空;或 - 选择「从资源中心选择配置文件」(
useCustom = false)且资源列表恰好包含 1 个配置文件(resourceList.size() == 1)。
3.2 FLINK 引擎参数
- Run model(运行模式):支持
run与run-application两种模式。在 SeatunnelFlinkParameters.java 中,两种模式被映射为 SeaTunnel 命令行参数--deploy-mode run与--deploy-mode run-application,未选择时对应none(不追加该参数)。 - Option parameters(可选参数):用于追加 Flink 引擎的自定义参数,例如
-m yarn-cluster -ynm seatunnel。该参数在 SeatunnelFlinkTask.buildOptions() 中作为一段原始字符串直接拼接到命令末尾,因此多个参数之间请使用空格分隔。
3.3 SPARK 引擎参数
- Deployment mode(部署模式):指定部署模式,可选
cluster、client(源码枚举定义见 DeployModeEnum.java,该枚举同时服务于 Spark 与 SeaTunnel Engine 两个分支)。 - Master:指定 Master 模式,可选
yarn、local、spark、mesos。其中spark与mesos需要额外指定 Master 服务地址,例如127.0.0.1:7077。
在 SeatunnelSparkParameters.checkParameters() 中有更细致的校验规则:部署模式必填;当部署模式不是local时,Master 必填;而当 Master 为spark或mesos时,Master 地址(masterUrl)必填。命令构造逻辑在 SeatunnelSparkTask.buildOptions() 中:部署模式为local时 Master 会被强制取local值,spark/mesos模式下则拼接为spark://127.0.0.1:7077、mesos://127.0.0.1:7077形式的地址。底层命令选项--deploy-mode、--master定义于 Constants.java。
3.4 SEATUNNEL_ENGINE 引擎参数
- Deployment mode(部署模式):指定部署模式,可选
cluster、local。该参数通过--deploy-mode选项传入,逻辑见 SeatunnelEngineTask.buildOptions(),未选择时不会追加该参数。
关于 Apache SeaTunnel 命令行使用的更多信息,可参考 SeaTunnel 官方文档中
2.3.3/command/usage一节。
四、自定义配置(Custom Configuration)
SeaTunnel 任务支持两种方式提供运行所需的配置文件:
- 自定义配置:直接在任务节点上编写配置内容;
- 从资源中心选择:从 DolphinScheduler 资源中心(Resource Center)上传一个 SeaTunnel 配置文件,然后在任务节点上引用。
两种方式在源码层面的处理路径不同。在 SeatunnelTask.buildOptions() 中可以看到:
if (BooleanUtils.isTrue(seatunnelParameters.getUseCustom())) { scriptContent = buildCustomConfigContent(); } else { String resourceFileName = seatunnelParameters.getResourceList().get(0).getResourceName(); ResourceContext resourceContext = taskRequest.getResourceContext(); scriptContent = FileUtils.readFileToString( new File(resourceContext.getResourceItem(resourceFileName).getResourceAbsolutePathInLocal()), StandardCharsets.UTF_8); }- 自定义配置模式下,任务会把节点上的脚本内容直接作为配置文件写入;
- 资源中心模式下,任务会通过
ResourceContext将所选资源文件内容读取出来,再写入本地执行目录。
值得关注的是配置文件的格式与命名。SeatunnelTask会根据脚本内容是否为合法 JSON 自动决定生成.json还是.conf文件(formatDetector()):
private String formatDetector() { return JSONUtils.checkJsonValid(seatunnelParameters.getRawScript(), false) ? Constants.JSON_SUFFIX : Constants.CONF_SUFFIX; }最终配置文件写入{执行目录}/seatunnel_{taskAppId}.{json|conf},随后以--config选项传入启动脚本。这一行为在测试类 SeatunnelTaskTest.java 中得到了验证——同一份脚本分别以 HOCON 与 JSON 语法书写时,生成的命令分别指向.conf与.json文件。
4.1 Script 配置结构
自定义配置的脚本内容通常由四个部分组成:env、source、transform、sink。
env:环境配置,如并行度等引擎级参数;source:数据源插件配置,声明输入来源(数据库、消息队列、文件等);transform:转换插件配置,对数据进行加工处理(可选);sink:目标端插件配置,声明数据输出位置。
关于 Apache SeaTunnel 配置文件的完整语法,可参考 SeaTunnel 官方文档中2.3.3/concept/config一节。
五、自定义参数与全局参数传递
当在 SeaTunnel 任务节点上定义了自定义参数(Custom Parameters),或工作流定义了全局参数(Global Parameters)时,这些参数会被传递给 SeaTunnel 任务,并可在配置文件中通过${}方式引用,在任务执行时完成动态替换。
从源码看,参数传递在 SeatunnelTask.generateTaskParameters() 中实现,具体逻辑是:
- 将工作流全局参数反序列化后,从
paramsMap中取值,以参数名 -> 参数值的形式放入变量表; - 遍历任务本地参数,仅取方向为
IN(输入)的参数加入变量表; - 将变量表转换为 SeaTunnel CLI 的
-i参数列表,格式为-i 'key=value',并追加到命令末尾。
其中-i参数的取值会经过 bash 安全转义(quoteForBash()),避免参数值中的单引号等特殊字符破坏命令。此外,脚本中出现的${}占位符也会在parseScript()中通过ParameterUtils.convertParameterPlaceholders()做进一步替换处理。
关于 Apache SeaTunnel 变量替换(Variable Substitution)的更多细节,可参考 SeaTunnel 官方文档中
2.3.3/concept/config/#config-variable-substitution一节。
六、任务示例:Flink 引擎读取 FakeSource 输出到控制台
下面通过一个完整示例,演示如何在 DolphinScheduler 中使用 Flink 引擎运行 SeaTunnel 任务:从一个 Fake 数据源读取数据,打印到控制台。
6.1 配置 DolphinScheduler 中的 SeaTunnel 环境
在生产环境中使用 SeaTunnel 任务类型之前,需要先在 Worker 所在机器上配置好运行环境。DolphinScheduler 的环境配置文件位于/dolphinscheduler/conf/env/dolphinscheduler_env.sh,需要在该文件中设置 SeaTunnel 相关环境变量(如SEATUNNEL_HOME),并确保其bin目录下的启动脚本具备执行权限。
6.2 配置 SeaTunnel 任务节点
进入任务节点的编辑面板,根据上文「任务参数详解」一节的内容进行配置:选择启动脚本(如 Flink 对应的start-seatunnel-flink-connector-v2.sh),选择运行模式,并填写自定义配置内容。
6.3 配置示例
在「自定义配置」中粘贴以下 SeaTunnel 配置(HOCON 格式):
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中设置并行度为 1;source使用FakeSource插件模拟产生包含name、age两个字段的数据,并注册为名为fake的结果表;transform通过 SQL 从fake表查询两个字段;sink使用ConsoleSink将结果打印到控制台。整个链路演示了 SeaTunnel 任务「读取 -> 转换 -> 输出」的最小可运行闭环,可作为验证环境连通性的基准示例。
七、支持的 SeaTunnel 版本
- 本文示例基于 SeaTunnel
2.3.x的 CLI 选项与启动脚本编写; - 已验证版本:v2.3.1、v2.3.2、v2.3.3;
- 其他版本:由于该任务类型本质上是 SeaTunnel CLI 的包装器,只要启动脚本与 CLI 选项保持兼容,通常新版本也能正常工作(升级后建议运行回归测试确认)。
八、延伸阅读:任务执行的完整链路
结合上文源码分析,可以梳理出 SeaTunnel 任务在 Worker 上的完整执行链路:
- 参数解析与校验:
SeatunnelTask.init()解析任务参数,调用checkParameters()校验启动脚本名与配置来源是否合法; - 命令构建:
buildCommand()将${SEATUNNEL_HOME}/bin/{startupScript}与buildOptions()结果拼接为完整命令; - 配置落盘:根据脚本内容格式生成
seatunnel_{taskAppId}.conf或.json文件到任务执行目录; - 参数注入:全局参数与本地
IN参数通过-i 'key=value'追加到命令; - Shell 执行:
ShellCommandExecutor通过 Shell 拦截器执行命令,记录进程 ID 与退出码; - 结果回收:
dealOutParam()收集输出参数,供下游任务引用。
SeaTunnel 任务插件本身不提供 Application 级的状态跟踪(submitApplication/trackApplicationStatus均为空实现,SeatunnelTask.java),任务是否成功完全取决于 Shell 命令的退出状态码,这一点在将 SeaTunnel 任务接入复杂 DAG 时值得注意。
如果你希望进一步扩展 SeaTunnel 任务的能力(例如新增启动脚本选项或引擎参数),可以阅读插件模块的完整实现与测试代码:dolphinscheduler-task-seatunnel,其中 SeatunnelTaskTest.java 覆盖了配置格式检测、资源中心读取与参数传递三条核心路径,可作为理解插件行为的参考。
【免费下载链接】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),仅供参考