首页
/ Apache DolphinScheduler 核心名词与架构模块全解析:从 DAG、调度命令到补数机制

Apache DolphinScheduler 核心名词与架构模块全解析:从 DAG、调度命令到补数机制

2026-09-14 09:24:15作者:邬祺芯Juliet

本文以 Apache DolphinScheduler 官方《名词解释》文档为主体,结合当前仓库源码,系统梳理调度系统中的核心概念——DAG、流程定义/实例、任务实例、任务类型、调度方式、定时调度、依赖、优先级、失败策略与补数,并逐一介绍 master、worker、alert、api 等核心模块的职责。读完本文,你将准确掌握 DolphinScheduler 的领域术语体系与底层实现机制,为后续流程编排、参数配置和故障排查打下坚实基础。

Apache DolphinScheduler 中一个典型 DAG 工作流示意图

一、DAG:调度系统的基石

DAG 全称 Directed Acyclic Graph(有向无环图)。在 DolphinScheduler 中,工作流里的 Task 任务以有向无环图的形式组装起来:从入度为零的节点开始进行拓扑遍历,逐级执行后继任务,直到不存在后继节点为止。

  • 有向(Directed):任务之间存在明确的方向性依赖,前驱节点执行完成后才会触发后继节点;
  • 无环(Acyclic):图中不允许存在环路,从而保证拓扑排序必然存在、执行必然终止。

上图展示的就是一个典型的 DAG 示例:task-shell 作为起始节点并行分支到 task-spark01task-sql,随后经过 task-spark02task-subprocesstask-proceduretask-python 等多个节点,最终汇聚到 task-mr。所有节点仅存在单向依赖、无循环引用,符合 DAG 定义。

二、流程定义、流程实例与任务实例

流程定义(Process Definition)

通过拖拽任务节点并建立任务节点之间的关联所形成的可视化 DAG,即为一条流程定义。流程定义是静态的、可复用的"模板",描述了任务之间的拓扑关系与每个任务节点的配置。

流程实例(Process Instance)

流程实例是流程定义的实例化,可以通过手动启动定时调度生成。每运行一次流程定义,就会产生一个流程实例。流程实例具有独立的状态生命周期,其状态在仓库中由 WorkflowExecutionStatus 枚举定义,主要包括:

状态 code 含义
SUBMITTED_SUCCESS 0 已提交成功
RUNNING_EXECUTION 1 运行中
READY_PAUSE / PAUSE 2 / 3 准备暂停 / 已暂停
READY_STOP / STOP 4 / 5 准备停止 / 已停止
FAILURE 6 失败(终态)
SUCCESS 7 成功(终态)
SERIAL_WAIT 14 串行等待中
FAILOVER 18 容错中

从源码结构看,每个状态还携带 canStopcanPauseneedFailoverisFinalState 等能力位,用于判断该状态下是否允许停止、暂停、需要容错或已处于终态。

任务实例(Task Instance)

任务实例是流程定义中任务节点的实例化,标识着某个具体的任务。一次流程实例的运行会生成其 DAG 内各个任务节点的任务实例,每个任务实例拥有独立的执行状态与日志。

三、任务类型(Task Type)

DolphinScheduler 采用插件化的任务体系,官方文档最早列举了 SHELL、SQL、SUB_WORKFLOW(子工作流)、PROCEDURE、MR、SPARK、PYTHON、DEPENDENT(依赖) 等任务类型,并计划支持动态插件扩展。

从当前仓库的 dolphinscheduler-task-plugin 目录看,任务类型插件已扩展为 30+ 个模块,除上述基础类型外,还包括:HTTP、FLINK、FLINK_STREAM、DATAX、SQOOP、SEATUNNEL、CHUNJUN、DINKY、JAVA、JUPYTER、K8S、KUBEFLOW、LINKIS、MLFLOW、OPENMLDB、REMOTESHELL、SAGEMAKER、ZEPPELIN、HIVE_CLI、DVC、EMR、EMR_SERVERLESS、DATASYNC、DATAFACTORY、DMS、GRPC 以及云厂商相关(aliyunserverlessspark)等,印证了"插件化、可动态扩展"的设计目标。

需要注意:其中 SUB_WORKFLOW(子工作流) 类型的任务需要关联另外一个流程定义,且被关联的流程定义是可以单独启动执行的——这是子工作流与普通任务在生命周期上的关键区别。

四、调度方式与命令类型(Command Type)

DolphinScheduler 支持基于 cron 表达式定时调度手动调度两种方式。所有调度动作最终都会落为一条命令(Command),由 master 消费并驱动流程实例运转。

在源码 CommandType 枚举中,命令类型完整定义如下:

命令类型 code 含义
START_PROCESS 0 启动工作流(生成新的流程实例,从开始节点运行)
START_CURRENT_TASK_PROCESS 1 从当前节点开始执行
RECOVER_TOLERANCE_FAULT_PROCESS 2 恢复被容错的工作流(master 宕机后从最后一个运行节点恢复)
RECOVER_SUSPENDED_PROCESS 3 恢复暂停的流程
START_FAILURE_TASK_PROCESS 4 从失败节点开始执行
COMPLEMENT_DATA 5 补数(按补数日期列表生成流程实例)
SCHEDULER 6 定时调度触发启动
REPEAT_RUNNING 7 重跑工作流
PAUSE 8 暂停工作流
STOP 9 停止工作流(Kill 运行中的任务)
RECOVER_SERIAL_WAIT 11 恢复串行等待
EXECUTE_TASK 12 触发流程实例中的指定任务节点
DYNAMIC_GENERATION 13 动态逻辑任务实例生成(内部使用)

其中,恢复被容错的工作流(RECOVER_TOLERANCE_FAULT_PROCESS)恢复等待线程(RECOVER_SERIAL_WAIT) 两种命令类型由调度内部控制使用,外部无法直接调用——前者用于 master 异常宕机后的容错恢复,后者用于串行等待状态的流程恢复。

五、定时调度:基于 Quartz 的分布式调度

DolphinScheduler 采用 Quartz 分布式调度器,并支持 cron 表达式的可视化生成(在 UI 上可通过表单点选生成,无需手动记忆 cron 语法)。

从源码结构看,Quartz 相关实现位于 dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz 模块,核心类包括:

  • QuartzScheduler.java:调度器主入口;
  • ProcessScheduleTask.java:调度触发的任务执行体;
  • QuartzSchedulerAutoConfigurationQuartzSchedulerDataSourceAutoConfiguration:负责调度器的自动装配与调度数据源的持久化配置。

Quartz 调度产生的触发会转换为 CommandType.SCHEDULER(6) 命令,与手动 START_PROCESS 逻辑相同、仅触发来源不同(参见源码注释 "This command is same with START_PROCESS but with different trigger source")。

六、依赖:DAG 依赖之外的流程间依赖

DolphinScheduler 不仅支持 DAG 中简单的前驱节点与后继节点之间的依赖,还额外提供任务依赖(DEPENDENT) 节点,支持流程间的自定义任务依赖——即一个任务可以跨流程、跨周期地依赖其他流程定义中指定任务实例的运行结果,从而实现更复杂的上下游联动编排。

七、优先级(Priority)

系统支持流程实例任务实例两个维度的优先级控制。若流程实例和任务实例的优先级均未设置,则默认按先进先出(FIFO) 顺序执行。

优先级在源码 Priority 枚举中划分为五级:

优先级 code 含义
HIGHEST 0 最高
HIGH 1
MEDIUM 2
LOW 3
LOWEST 4 最低

八、邮件告警(Email Alert)

DolphinScheduler 的告警体系支持以下三类邮件场景:

  1. SQL 任务查询结果邮件发送:SQL 任务可将查询结果直接以邮件形式发送给指定收件人;
  2. 流程实例运行结果邮件告警:流程运行成功、失败等结果状态变化时触发邮件通知;
  3. 容错告警通知:发生容错事件时向相关人员发送告警。

邮件告警能力由 dolphinscheduler-alert-plugins/dolphinscheduler-alert-email 插件实现,属于告警插件体系(Alert Plugin)中的一员。除邮件外,仓库中的 dolphinscheduler-alert-plugins 还包含钉钉、飞书、企业微信、Slack、Telegram、WebexTeams、HTTP、Script、PagerDuty、Prometheus、阿里云语音等告警通道,告警类型与级别分别在 AlertTypeAlertWarnLevel 等枚举中定义。

九、失败策略(Failure Strategy)

对于并行运行的任务,当其中某个任务失败时,可配置两种失败策略:

  • 继续(CONTINUE):不管并行运行任务的状态如何,流程继续执行,直到整体流程失败结束;
  • 结束(END):一旦发现失败任务,立即 Kill 掉正在运行的并行任务,流程失败结束。

对应源码 FailureStrategy 枚举:

失败策略 code 含义
END 0 某任务失败即结束流程
CONTINUE 1 某任务失败后继续运行

该枚举通过 @EnumValue 注解与数据库存储值(0/1)绑定,可在流程定义配置中直接选择。

十、补数(Complement Data)

补数用于补跑历史数据,是数据平台运维中的高频场景,支持:

  • 两种补数方式区间并行串行
  • 两种日期选择方式日期范围(如指定起止日期区间)与日期枚举(逐个勾选具体日期)。

补数在命令层面由 CommandType.COMPLEMENT_DATA(5) 承载,源码注释表明其会使用 complementScheduleDateList(补数调度日期列表)批量生成流程实例。从代码结构看,补数相关的依赖模式还通过 ComplementDependentMode 枚举定义,用于约束补数场景下依赖任务的处理方式。

十一、核心模块介绍

DolphinScheduler 采用主从(Master-Worker)加插件化的多模块架构,官方文档列出的核心模块如下:

模块 职责
dolphinscheduler-master master 模块,提供工作流管理和编排服务
dolphinscheduler-worker worker 模块,提供任务执行管理服务
dolphinscheduler-alert 告警模块,提供 AlertServer 服务
dolphinscheduler-api web 应用模块,提供 ApiServer 服务
dolphinscheduler-common 通用的常量枚举、工具类、数据结构或基类
dolphinscheduler-dao 提供数据库访问等操作
dolphinscheduler-extract extract 模块,包含 master/worker/alert 的 SDK
dolphinscheduler-service service 模块,包含 Quartz、Zookeeper、日志客户端访问服务,便于 server 模块和 api 模块调用
dolphinscheduler-ui 前端模块

从当前仓库结构看,模块体系在此基础上进一步细化为更多独立工程,例如:dolphinscheduler-registry(注册中心,含 Zookeeper、JDBC、etcd 等插件)、dolphinscheduler-scheduler-plugin(调度器插件,含 Quartz 实现)、dolphinscheduler-datasource-plugin(30+ 种数据源插件)、dolphinscheduler-storage-plugin(存储插件,含 HDFS、S3、OSS、COS、OBS、GCS、ABS)、dolphinscheduler-task-plugin(30+ 种任务插件)、dolphinscheduler-task-executor(任务执行器)、dolphinscheduler-eventbus(事件总线)与 dolphinscheduler-meter(指标埋点)等,整体呈现"核心引擎 + 可插拔扩展"的现代数据编排平台形态。

十二、总结

本文围绕官方《名词解释》梳理了 DolphinScheduler 的完整领域概念:以 DAG 为组织形态的流程定义,经实例化产生流程实例与任务实例;通过丰富的命令类型驱动手动调度、定时调度、容错恢复、重跑、补数等各类运行场景;以 Quartz 支撑 cron 定时调度;以依赖节点扩展跨流程编排能力;以优先级、失败策略控制运行秩序与故障行为;以告警插件闭环通知结果;最终由 master、worker、alert、api 等模块协作完成端到端的调度执行。掌握这些名词及其源码实现,是深入使用与二次开发 DolphinScheduler 的第一步。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
34
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.21 K
2.81 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
945
1.86 K
docsdocs
暂无描述
Markdown
906
5.84 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
537
607
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
864
1.36 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
4.28 K
1.03 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.39 K
1.48 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
550
401
flutter_flutterflutter_flutter
本仓库是 Flutter SDK 与 Flutter Engine 的 OpenHarmony 适配版本,由 CPF-Flutter 团队维护。开发者可使用熟悉的 Flutter 技术栈开发 OpenHarmony 应用,3.35.7 及以后的适配版本可基于本仓库源码构建支持 OpenHarmony 的 Flutter Engine。
Dart
1.19 K
347