Apache DolphinScheduler:分布式可视化 DAG 工作流调度平台的定位、架构与核心能力解析
【免费下载链接】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(下称 DolphinScheduler)是一个面向企业级场景的分布式、易扩展的可视化工作流任务调度开源平台。本文基于当前仓库中的官方介绍文档(docs/docs/en/about/introduction.md)展开,结合根目录 CLAUDE.md、README.md 与dolphinscheduler-common中的核心枚举源码,系统讲解 DolphinScheduler 要解决的 DataOps 编排问题、其 DAG 流式任务模型的运行机理,以及重试、失败恢复、暂停、恢复、终止等关键操作在源码层面的真实映射,帮助读者在动手部署前先建立对该平台完整而准确的技术认知。
一、定位:解决复杂大数据任务依赖与 DataOps 编排问题
官方介绍文档对 DolphinScheduler 的定位可以归纳为三句话:
- 分布式且易扩展的可视化工作流调度平台——任务、工作流、整个数据处理过程都可以可视化操作;
- 解决复杂的大数据任务依赖与触发关系——专为 DataOps 编排设计,应对数据研发 ETL 中"依赖错综复杂"的痛点;
- 可实时监测任务健康状态——针对传统手工调度"无法监控任务健康状态"的问题,提供统一的任务状态监控视图。
在大数据应用中,一条典型 ETL 流水线往往横跨 Shell 脚本、SQL、Spark/Flink 作业、数据同步(DataX/Sqoop)等多种异构任务类型,并跨多台机器执行。如果依赖关系靠文档维护、任务状态靠人工 SSH 查看,那么任务失败后"从哪里恢复、如何恢复"就成了运维黑洞。DolphinScheduler 的解决方案是:把任务以 DAG(Directed Acyclic Graph,有向无环图)的流式方式组装起来,由平台统一负责依赖计算、状态流转与健康监控。这一点可以从仓库的模块结构直接印证——dolphinscheduler-task-plugin/下并列存放了 shell、sql、spark、flink、datasync、datax、sqoop、k8s、emr 等 30 余个任务插件模块,dolphinscheduler-datasource-plugin/下则并列存放了 mysql、hive、doris、snowflake、trino 等 20 余个数据源插件模块,正是"异构任务 + 异构数据源统一编排"这一定位的落地形态。
二、DAG 流式组装:任务如何被依赖图驱动执行
官方文档强调 DolphinScheduler 以DAG 流式方式组装任务,这一表述包含两层含义:
- 依赖即拓扑:工作流定义保存为 JSON(核心字段
process_definition_json),节点之间的preTasks数组显式描述前驱任务,执行引擎据此做拓扑排序与依赖计算,保证"前驱成功才提交后继"; - 状态即流转:任务提交到执行队列后,其运行状态(运行中、成功、失败、暂停、容错等)被持续跟踪,任何一个节点状态变化都会驱动后续节点的触发或阻塞,因此可以"及时监控任务的执行状态"。
从源码结构看,DAG 的"流式"推进由 Master 端的多条职责线程协同完成:Master 内置的调度线程定期扫描数据库中的指令表(t_ds_command),按指令类型执行不同业务操作;随后由 DAG 状态机线程负责任务切分与提交监控,事件轮询线程与状态轮询线程分别处理实例事件队列和任务超时/重试/依赖轮询。这些线程的职责划分在架构设计文档 docs/docs/en/architecture/design.md 中有完整描述(MasterSchedulerService、WorkflowExecuteRunnable、EventExecuteService、StateWheelExecuteThread等),本文不再展开,但值得记住:DAG 的每一步推进本质上都是"状态事件驱动",这正是"流式"一词的含义。
三、灵活的状态控制:重试、恢复、暂停、终止的源码级映射
介绍文档中最具实操价值的一句话是:DolphinScheduler"支持重试、从指定节点恢复失败、暂停、恢复、终止任务等操作"。这些操作不是文档修辞,而是平台指令模型中明确定义的一等公民。当前仓库中,所有针对工作流实例的操作指令统一收敛在枚举 CommandType(dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/enums/CommandType.java)中,其定义与介绍文档中的操作一一对应:
| 指令 | 代码值 | 语义(对应介绍文档中的能力) |
|---|---|---|
START_PROCESS | 0 | 启动一个新的工作流实例,从起始节点开始执行 |
START_CURRENT_TASK_PROCESS | 1 | 从指定当前节点启动(标记为待移除的旧指令) |
RECOVER_TOLERANCE_FAULT_PROCESS | 2 | 从容错状态恢复——Master 故障后从最后运行节点接管实例 |
RECOVER_SUSPENDED_PROCESS | 3 | 从暂停状态恢复,从被暂停且未触发的任务实例继续 |
START_FAILURE_TASK_PROCESS | 4 | 从失败节点恢复(即"recovery from specified nodes") |
COMPLEMENT_DATA | 5 | 补数(backfill),按补数日期列表批量生成实例 |
SCHEDULER | 6 | 定时调度触发(与手动启动逻辑相同,仅触发源不同) |
REPEAT_RUNNING | 7 | 重复运行,将历史任务实例标记为失效后从首节点重跑 |
PAUSE | 8 | 暂停工作流实例(pause a workflow) |
STOP | 9 | 终止工作流实例,会 kill 正在运行的任务(kill tasks) |
RECOVER_SERIAL_WAIT | 11 | 从串行等待状态恢复 |
EXECUTE_TASK | 12 | 在实例内从指定任务节点触发执行 |
DYNAMIC_GENERATION | 13 | 动态逻辑任务实例使用 |
源码注释同时给出了各指令的精确行为边界,例如:RECOVER_SUSPENDED_PROCESS是"从被暂停且未触发的任务实例继续";PAUSE是"暂停运行中的任务,但并非所有任务都会被暂停";STOP是"kill 运行中的任务"。这些注释让"暂停/恢复/终止"从界面按钮变成了可审计的语义契约。
在 Master 侧,指令的落地路径是:用户操作(UI / Open API / Python SDK)先写入指令记录,Master 的扫描线程按指令类型分发处理。此外,文档提到的"任务重试"发生在任务级别,由任务参数maxRetryTimes/retryInterval控制——在各任务类型的节点 JSON 结构中(见 task-structure.md)这两个字段是每个任务节点的标配字段,且源码明确区分了业务任务(可配置失败重试次数)与逻辑任务(如子流程、依赖任务,不支持失败重试)。
与此配套,工作流与任务都支持五级优先级:Priority(dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/enums/Priority.java)定义了HIGHEST(0) → LOWEST(4)五个档位。任务提交到执行队列时,队列键按"工作流实例优先级 + 实例 ID + 任务优先级 + 任务 ID"的字符串格式拼接(详见 design.md 的 Task Priority Design 一节),取队列时通过字符串比较即可得到最高优先级任务,从而保证不同工作流之间、同一工作流内不同任务之间的调度顺序符合优先级预期。
四、分布式架构:四个可独立扩展的服务
"分布式且易扩展"是介绍文档的第一关键词。从源码结构看,生产部署由四个相互独立的服务构成(外加外部注册中心与元数据库),每个服务都是独立的 Spring Boot 应用,主类与默认端口在根目录 CLAUDE.md 的 "Runnable services" 一节中给出:
| 服务 | 所在模块 | 主类 | 默认端口 |
|---|---|---|---|
| API | dolphinscheduler-api | org.apache.dolphinscheduler.api.ApiApplicationServer | 12345(HTTP:UI + REST) |
| Master | dolphinscheduler-master | org.apache.dolphinscheduler.server.master.MasterServer | 5679(RPC) |
| Worker | dolphinscheduler-worker | org.apache.dolphinscheduler.server.worker.WorkerServer | 1235(RPC) |
| Alert | dolphinscheduler-alert(→-alert-server) | org.apache.dolphinscheduler.alert.AlertServer | 50053(HTTP)、50052(RPC) |
此外还有一个第五入口StandaloneServer(dolphinscheduler-standalone-server模块),把上述四个服务内嵌到同一个 JVM 中,面向开发与冒烟测试场景。
各服务的职责边界可以从调用链中直接读出:UI → API 服务 → 元数据库(读写定义数据)+ Master(RPC 下发运行时指令:启动/暂停/终止)→ Master 消费指令行、驱动 DAG 状态机、把任务分发给 Worker → Worker 加载对应任务插件执行,并把生命周期事件流式回传 Master → 失败或 SLA 违规事件流向 Alert 服务,经告警插件(邮件、飞书、钉钉、Slack、Telegram 等,见 dolphinscheduler-alert-plugins 下的各插件模块)扇出。Master / Worker / Alert 均可水平扩展,服务发现、领导者选举与分布式锁统一由注册中心承担——注册中心是可插拔的:默认 Zookeeper,同时提供 Etcd 与 JDBC 实现(dolphinscheduler-registry/dolphinscheduler-registry-plugins/下三个插件模块),其中 JDBC 实现利用关系数据库本身完成事件监听与分布式锁,适合不愿额外引入中间件的环境。
这套去中心化设计的直接收益是介绍文档中"易扩展"的工程含义:Master 与 Worker 集群内部没有"主从"之分,节点宕机不会导致整个集群瘫痪,而是通过注册中心的节点移除事件触发故障转移逻辑——Worker 掉线由其上的运行任务被 Master 接管重新提交;Master 掉线则触发工作流实例级别的容错恢复(对应上表的RECOVER_TOLERANCE_FAULT_PROCESS指令)。完整的容错时序、故障范围划分(Master 容错 vs Worker 容错)与网络抖动处理策略,建议进一步阅读 docs/docs/en/architecture/design.md。
五、面向场景的能力清单:从介绍文档到仓库证据
把介绍文档与 README.md 中的特性描述交叉对照,可以整理出一份"文档宣称 ↔ 仓库证据"的能力清单,便于读者按需检索:
- 低代码创建工作流:拖拽式 Web UI 构建 DAG,同时提供 Open API 与 Python SDK 编程式管理——对应
dolphinscheduler-ui/(Vue 3 前端,页面源码位于dolphinscheduler-ui/src/views/<feature>/,API 调用封装在dolphinscheduler-ui/src/service/modules/)与dolphinscheduler-api/的 REST 控制器(dolphinscheduler-api/src/main/java/.../api/controller/); - 开箱即用的丰富任务类型:
dolphinscheduler-task-plugin/下 30 余个任务插件模块,覆盖 Shell、SQL、Spark、Flink、DataX、Sqoop、HTTP、K8s、EMR、MLflow、SeaTunnel 等; - 统一数据源接入:
dolphinscheduler-datasource-plugin/下 20 余个数据源插件,覆盖 MySQL、PostgreSQL、Hive、Trino、Doris、Snowflake、ClickHouse、StarRocks 等,为 SQL 类任务与数据同步任务提供统一数据访问能力; - 多租户支持:工作流定义中的
tenant_id字段(见 task-structure.md 中t_ds_process_definition表结构)为任务在 Worker 上以租户身份执行提供隔离基础; - 补数(backfill)原生支持:对应
COMPLEMENT_DATA(5)指令,按补数日期列表批量生成工作流实例; - 工作流版本控制:定义与实例两级版本管理,对应表结构中的
version字段与版本化的发布状态(release_state); - 权限控制:项目、数据源等资源维度的权限体系(
dolphinscheduler-api/下的权限相关控制器与服务)。
部署形态方面,README.md 列出了四种模式:Standalone(单机)、Cluster(集群)、Docker、Kubernetes,仓库中对应deploy/docker/(含 docker-compose.yml)、deploy/kubernetes/(Helm Chart)与deploy/terraform/aws/(Terraform 模板);本地快速体验则推荐构建后进入dolphinscheduler-standalone-server/target执行./bin/start.sh启动一体式服务(构建命令./mvnw clean install -Prelease,详见根目录 CLAUDE.md 的 "Build & run" 一节)。
六、总结
回顾官方介绍文档的核心命题——"分布式、易扩展的可视化 DAG 工作流调度平台,解决复杂任务依赖与 DataOps 编排问题,可监控、可恢复、可暂停、可终止"——在当前仓库中都能找到明确的工程落点:
- "DAG 流式组装"落在
process_definition_json的节点/依赖结构(task-structure.md)与 Master 端的 DAG 状态机线程上; - "重试、恢复、暂停、恢复、终止"落在 CommandType 指令模型与任务级
maxRetryTimes配置上; - "分布式易扩展"落在 API / Master / Worker / Alert 四个可独立水平扩展的服务与可插拔注册中心(Zookeeper / Etcd / JDBC)上;
- "企业级场景"落在多租户、版本控制、项目级权限、补数、28+ 数据源与 30+ 任务插件的生态上。
对于准备采用 DolphinScheduler 的团队,建议的阅读与操作路径是:先通过本文与 架构设计文档 建立架构认知,再查阅 功能介绍 熟悉各功能界面,最后按 README.md 的 QuickStart 选择 Standalone / Docker / Kubernetes 任一形态完成部署验证。
【免费下载链接】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),仅供参考