Apache DolphinScheduler 架构设计深度解析:从调度核心名词到分布式容错与日志原理
2026/9/24 14:51:34 网站建设 项目流程
  • 任务调度
  • 大数据
  • 后端
  • 前端

【免费下载链接】dolphinscheduler

Apache 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 官方贡献指南中的架构设计文档(architecture-design.md)为主体,系统梳理这套大数据分布式工作流调度系统的核心概念、Master/Worker 架构职责、去中心化设计思想、分布式锁、容错机制、任务优先级以及基于 Logback 与 gRPC 的日志访问原理,并结合当前仓库源码给出可验证的实现证据。读完本文,你将掌握 DolphinScheduler 从"名词体系"到"调度内核"再到"故障恢复"的完整技术脉络,能够从原理层面理解工作流为何能在大规模集群中稳定调度与自愈。

1. 调度系统的核心名词

在深入架构之前,先统一调度系统中最常出现的一组名词。它们贯穿于 DolphinScheduler 的数据库表设计、调度命令与前端页面,理解它们是读懂后续所有设计的前提。

1.1 DAG(有向无环图)

工作流中的任务以有向无环图(Directed Acyclic Graph,DAG)的形式组装:从入度为 0 的节点开始做拓扑遍历,直到不存在后继节点为止。DolphinScheduler 的编排页面正是基于 DAG 的思想,将任务节点拖拽连线形成执行依赖关系:

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

  • 流程定义(Process definition):通过拖拽任务节点并建立节点关联,将DAG可视化出来得到的产物;
  • 流程实例(Process instance):流程定义的实例化。可通过手动启动或定时调度触发,流程定义每运行一次就产生一个新的流程实例;
  • 任务实例(Task instance):流程实例运行时,其中某个具体任务节点的实例化,用于记录该任务的具体执行状态。

三者是"定义 → 实例 → 任务执行"的层级关系,也是数据库中t_ds_process_definitiont_ds_process_instancet_ds_task_instance等表设计的心智模型。

1.3 任务类型

文档最初描述的系统已支持 SHELL、SQL、SUB_PROCESS(子流程)、PROCEDURE、MR、SPARK、PYTHON、DEPENDENT(依赖)等任务类型,并计划支持动态插件化扩展。需要特别说明的是:SUB_PROCESS 本身也是一个可独立启动的流程定义。从当前仓库源码结构看,这一设计已经落地为庞大的插件体系,dolphinscheduler-task-plugin 目录下按任务类型拆分模块(如dolphinscheduler-task-shelldolphinscheduler-task-sparkdolphinscheduler-task-flinkdolphinscheduler-task-httpdolphinscheduler-task-datax等数十种),任务类型由 dolphinscheduler-task-api 提供统一抽象,真正做到"新增一种任务类型只需新增一个插件模块"的动态扩展。

1.4 调度方式与命令类型

系统支持基于 cron 表达式的定时调度手动调度。命令(Command)类型覆盖了流程运行的各类触发与干预动作。文档中提到:其中"恢复容错工作流"与"恢复等待线程"两种命令类型由调度系统内部控制,外部不可调用。

从当前源码 CommandType.java 可以看到完整的命令枚举(共 14 种,代码编号 0~13):

编号命令类型含义
0START_PROCESS启动新流程
1START_CURRENT_TASK_PROCESS从当前节点启动流程
2RECOVER_TOLERANCE_FAULT_PROCESS恢复容错流程(内部使用)
3RECOVER_SUSPENDED_PROCESS恢复暂停的流程
4START_FAILURE_TASK_PROCESS从失败节点启动流程
5COMPLEMENT_DATA补数
6SCHEDULER由定时调度启动新流程
7REPEAT_RUNNING重复运行流程
8PAUSE暂停流程
9STOP停止流程
10RECOVER_WAITING_THREAD恢复等待线程(内部使用)
11RECOVER_SERIAL_WAIT恢复串行等待
12EXECUTE_TASK在流程实例中启动某个任务节点
13DYNAMIC_GENERATION动态生成

Master 的调度线程正是周期性扫描数据库中的command 表,再根据不同的命令类型分派不同的业务处理逻辑(详见第 2 节 MasterServer 部分)。

1.5 定时调度

系统使用Quartz 分布式调度器,并支持 cron 表达式的可视化生成。当前仓库将该能力抽为调度插件:dolphinscheduler-scheduler-plugin 下的dolphinscheduler-scheduler-quartz模块实现了 QuartzScheduler.java,通过SchedulerApi接口对外统一提供调度能力,便于未来替换其他调度内核。

1.6 任务依赖

系统不仅支持 DAG 中前后继节点的简单依赖,还提供任务依赖(DEPENDENT)节点,支持跨流程的自定义任务依赖——例如流程 A 的某任务依赖流程 B 的执行结果,这种依赖可以跨越流程边界进行编排。

1.7 优先级

支持流程实例优先级任务实例优先级两级设置;若均未设置,默认按先进先出(FIFO)执行。当前源码中的枚举定义见 Priority.java,分为五级:

编号优先级说明
0HIGHEST最高
1HIGH
2MEDIUM
3LOW
4LOWEST最低

1.8 邮件告警

支持三类告警场景:SQL Task 查询结果的邮件发送流程实例运行结果的邮件告警以及容错告警通知

1.9 失败策略

对并行运行的任务,若其中出现失败任务,提供两种失败策略:

  • Continue(继续):并行任务继续运行,直到流程最终失败;
  • End(结束):一旦发现失败任务,立即 Kill 掉正在运行的并行任务,流程直接结束。

1.10 补数

用于补齐历史数据,支持区间并行串行两种补数方式。

2. 系统整体架构

2.1 系统架构图

DolphinScheduler 的经典架构由 UI、API、MasterServer、WorkerServer、ZooKeeper、任务队列、Alert 等部分组成:

2.2 MasterServer

MasterServer 采用分布式非中心化设计理念,主要负责DAG 任务拆分、任务提交监控,以及监控其他 MasterServer 与 WorkerServer 的健康状态。MasterServer 服务启动时会在 ZooKeeper 上注册临时节点,并监听 ZooKeeper 临时节点的状态变化以进行容错处理。

Master 服务内部主要包含以下核心组件(文档描述的历史组件命名,在当前代码中已演进,见下文对照):

  • 分布式 Quartz 调度组件:主要负责定时任务的启停操作;Quartz 拉起任务后,Master 内部由线程池负责任务的后续操作;
  • MasterSchedulerThread:周期性扫描数据库command 表,根据不同command 类型执行不同的业务操作。当前代码中对应 MasterSchedulerBootstrap.java,其主循环逻辑是:通过commandFetcher.fetchCommands()拉取 command → 校验 Master 负载保护(serverLoadProtection.isOverload)→ 并行调用workflowExecuteRunnableFactory.createWorkflowExecuteRunnable(command)构造工作流执行体 → 放入ProcessInstanceExecCacheManager缓存并投递START_WORKFLOW事件;无 command 时 sleep 1 秒,避免空转打爆数据库;
  • MasterExecThread:负责 DAG 任务切分、任务提交监控、各类命令类型的逻辑处理。当前对应 WorkflowExecuteRunnable.java 及 WorkflowGraph.java 等运行期组件,负责将 DAG 图结构转化为可执行的任务实例;
  • MasterTaskExecThread:负责任务持久化。当前对应 TaskExecuteThreadPool.java 与 TaskExecuteRunnable.java 构成的任务事件处理链(派发事件、运行中事件、结果事件、重试事件等均有独立 Handler)。

从源码结构看,当前版本的 Master 侧已经由早期"调度扫描线程 + 执行线程"的朴素模型,演进为"命令拉取 → 事件队列(WorkflowEventQueue)→ 状态事件处理(StateEventHandlerManager)→ 任务事件处理"的事件驱动模型,但"扫描 command 表驱动流程执行"的核心思想一脉相承。

2.3 WorkerServer

WorkerServer 同样采用分布式、非中心化设计,主要负责任务执行日志服务。WorkerServer 服务启动时在 ZooKeeper 注册临时节点并维持心跳。

Worker 服务包含:

  • FetchTaskThread:持续从任务队列接收任务,并根据任务类型调用对应的执行器(TaskScheduleThread)。当前对应 dolphinscheduler-worker 模块中的任务执行与日志相关组件;
  • 日志服务:为任务实例提供按需切分的日志文件与远程读取能力(详见第 8 节)。

2.4 ZooKeeper(注册中心)

MasterServer 与 WorkerServer 节点均使用 ZooKeeper 进行集群管理与容错。此外,系统还基于 ZooKeeper 做事件监控分布式锁。文档同时透露一个设计取舍:曾基于 Redis 实现队列,但为了让 DolphinScheduler 依赖尽可能少的组件,最终移除了 Redis 实现。从当前仓库的 dolphinscheduler-registry 模块结构看,注册中心也已插件化,除 ZooKeeper 外还提供 JDBC、etcd 等注册中心实现,通过 Registry.java 抽象接口统一管理节点注册、事件监听与分布式锁。

2.5 任务队列

提供任务队列操作,当前同样基于 ZooKeeper 实现。由于队列中存储的信息量较少,无需担心队列数据过大——文档指出已进行过百万级数据量队列的压测,对系统稳定性与性能无影响。

2.6 Alert(告警)

提供告警相关接口,主要包括告警数据的存储、查询与通知三类功能。通知功能包含邮件通知SNMP(尚未实现)两种。当前仓库已将该模块独立为 dolphinscheduler-alert 目录,并插件化支持钉钉、飞书、企业微信、Slack、Telegram、Webex Teams、HTTP、脚本等多种渠道。

2.7 API 与 UI

  • API:接口层,负责处理来自前端 UI 层的请求,对外提供RESTful API,覆盖流程的创建、定义、查询、修改、上线、下线、手动启动、停止、暂停、恢复、从当前节点开始执行等操作。对应模块为 dolphinscheduler-api;
  • UI:系统前端页面,提供各类可视化操作界面,见仓库 docs/docs/en/guide 下的使用指南,对应前端工程为 dolphinscheduler-ui。

3. 架构设计思想:去中心化 vs 中心化

3.1 中心化设计及其问题

中心化设计思路相对简单:集群节点按角色分为 Master 与 Slave 两类。Master 负责任务分发并监督 Slave 健康状态,可动态均衡地将任务分配至各 Slave,避免节点"忙闲不均";Worker(Slave)负责任务执行,并与 Master 保持心跳以便其分配任务。

但中心化设计存在两个突出问题:

  1. 单点故障:一旦 Master 出问题,集群失去领导者,整个集群崩溃。多数 Master/Slave 架构采用主备 Master 方案缓解(热备或冷备、自动或手动切换),越来越多的系统具备自动选举切换 Master 的能力以提升可用性;
  2. 调度器位置的两难:若 Scheduler 放在 Master 上,虽能支持同一 DAG 中不同任务运行在不同机器,但会加重 Master 负载;若放在 Slave 上,一个 DAG 中的所有任务只能提交到同一台机器,并行任务多时该 Slave 压力过大。

3.2 去中心化设计

去中心化设计通常没有 Master/Slave 概念,所有角色地位平等——互联网本身就是典型的去中心化分布式系统:任意节点宕机,只会影响小范围功能。其核心在于整个分布式系统中不存在"管理者"节点,因此没有单点故障问题;但代价是每个节点都需要与其他节点通信获取必要信息,分布式通信链路的不可靠性大幅增加了实现难度。

实际上,真正完全去中心化的系统很少见,取而代之的是动态中心化系统:集群中的管理者动态选举产生而非预设;集群故障时,节点自发"开会"选出新"管理者"主持工作,最典型的案例即 ZooKeeper 与 Go 实现的 Etcd。

DolphinScheduler 的去中心化实践:Master/Worker 注册到 ZooKeeper,Master 集群与 Worker 集群均无中心,并通过 ZooKeeper 分布式锁选举出某个 Master 或 Worker 作为"管理者"执行任务。当前代码中,参与选举的实现位于 AbstractHAServer.java:其participateElection()直接调用registry.acquireLock(serverPath, 3_000)抢占分布式锁,抢锁成功即切换为ACTIVE状态,从而保证同一时刻只有一个节点以主身份执行调度职责。

4. 分布式锁实践

DolphinScheduler 使用 ZooKeeper 分布式锁实现:同一时刻只有一个 Master 执行 Scheduler,或只有一个 Worker 执行任务提交

获取分布式锁的核心流程算法如下:

Master 中 Scheduler 线程的分布式锁实现流程图如下:

结合源码可见,锁能力经由 RegistryClient.java 的acquireLock(key, timeout)统一暴露,底层由各注册中心插件实现,从而支持 ZooKeeper、etcd、JDBC 等多种实现下的锁语义。

5. 线程不足的循环等待问题

这是调度系统一个非常经典且隐蔽的坑:

  • 若一个 DAG 中没有子流程,当 command 表数据量大于线程池阈值时,直接等待或失败即可;
  • 若一个大 DAG 中嵌套了大量子流程,则可能出现"死锁"状态:

如上图:MainFlowThread 等待 SubFlowThread1 结束,SubFlowThread1 等待 SubFlowThread2,SubFlowThread2 等待 SubFlowThread3,而 SubFlowThread3 等待线程池分配新线程——整个 DAG 永远无法结束,线程也无法释放,形成子父流程循环等待。此时除非启动新的 Master 增加线程来打破"卡死",否则调度集群将不可用。

启动新 Master 显然不是优雅方案,文档给出了三种候选解法:

  1. 预计算线程数:先求所有 Master 线程总数,再计算每个 DAG 所需线程数,在 DAG 执行前预判。但多 Master 线程池的总线程数难以实时获取;
  2. 单 Master 线程池满则直接失败:若线程池已满,让线程直接失败;
  3. 新增"资源不足"命令类型:线程池不足时将主流程挂起,待线程池有新的线程后,再唤醒资源不足的流程。

注意:Master Scheduler 线程获取 Command 时是 FIFO(先进先出)的。

最终 DolphinScheduler 选择了第三种方案解决线程不足问题。结合当前源码,这一思路在命令枚举中沉淀为RECOVER_WAITING_THREAD(恢复等待线程)这一内部命令类型(见第 1.4 节 CommandType 表),印证了"挂起-唤醒"闭环的设计落地。

6. 容错设计

容错分为服务容错任务重试,其中服务容错又分为Master 容错Worker 容错

6.1 宕机容错

服务容错设计依赖 ZooKeeper 的Watcher 机制,实现原理如下:

Master 监听其他 Master 与 Worker 的目录节点,一旦检测到remove 事件,就根据具体业务逻辑进行流程实例容错或任务实例容错。

Master 容错流程:ZooKeeper Master 容错后,由 DolphinScheduler 的 Scheduler 线程重新调度。它遍历 DAG,找出处于"Running"(运行中)"Submit Successful"(提交成功)状态的任务,并监控其任务实例状态:对于 "Running" 任务,需要判断任务队列中是否已存在该任务——若已存在则持续监控任务实例状态,若不存在则重新提交任务实例

Worker 容错流程:一旦 Master Scheduler 线程发现任务实例"需要容错",便接管该任务并重新提交。

补充一个关键工程细节:由于"网络抖动"可能导致节点短时间内丢失 ZooKeeper 心跳、从而触发 remove 事件,DolphinScheduler 采用最直接的处理方式——节点一旦与 ZooKeeper 连接超时,直接停止 Master 或 Worker 服务,宁可自停也不在集群中留下状态不一致的"僵尸"节点。

从当前源码看,容错逻辑沉淀为 MasterFailoverService.java、WorkerFailoverService.java 与 FailoverService.java,其中checkMasterFailover()通过ds.master.scheduler.failover.check.count等指标暴露容错检查频率,可被监控系统采集。

6.2 任务失败重试

先厘清三个容易混淆的概念:

  • 任务失败重试(Task failure Retry)任务级,由调度系统自动执行。例如某 Shell 任务设置重试 3 次,则失败后最多自动运行 3 次;
  • 流程失败恢复(Process failure recovery)流程级,由人工操作,只能从失败节点从当前节点恢复;
  • 流程失败重跑(Process failure rerun)流程级,由人工操作,从起始节点重新运行。

据此,将工作流中的任务节点分为两类:

  • 业务节点(service node):对应实际脚本或处理语句,如 Shell 节点、MR 节点、Spark 节点、依赖节点等;
  • 逻辑节点(logic node):不做实际脚本或语句处理,而是对整体流程做逻辑处理,如子流程节点。

每个业务节点可配置失败重试次数:任务节点失败后自动重试,直到成功或超过配置次数。逻辑节点不支持失败重试,但逻辑节点内的任务支持重试。若工作流中的任务失败达到最大重试次数,工作流将失败停止,可人工重跑或恢复流程。

7. 任务优先级设计

早期调度设计中没有优先级与公平调度设计时,先提交的任务可能与后提交的任务同时完成,且无法设置流程或任务的优先级。DolphinScheduler 重新设计了优先级机制,处理顺序为:

不同的流程实例优先级优先于相同流程实例优先级下的任务优先级优先于同一流程中的提交顺序,即:先按流程实例优先级从高到低,再按任务优先级从高到低,最后按提交顺序处理任务。

具体实现:根据任务实例的 JSON 解析出优先级,然后在 ZooKeeper 任务队列中保存流程实例优先级_流程实例id_任务优先级_任务id格式的键;从任务队列获取时,通过字符串比较即可直接得到应优先执行的任务。

优先级分为五级:HIGHEST、HIGH、MEDIUM、LOW、LOWEST。其中流程定义优先级用于某些流程需要先于其他流程处理,可在流程启动时或定时启动时配置:

任务优先级同样分为五级(HIGHEST、HIGH、MEDIUM、LOW、LOWEST):

与第 1.7 节源码枚举 Priority.java(HIGHEST=0、HIGH=1、MEDIUM=2、LOW=3、LOWEST=4)完全对应。

8. Logback 与 gRPC 实现远程日志访问

由于 Web(UI)与 Worker 不一定在同一台机器上,查看日志不能像查询本地文件那样直接进行,存在两种可选方案:

  1. 将日志放入 ES 搜索引擎;
  2. 通过gRPC 通信获取远程日志信息。

考虑到尽可能保持 DolphinScheduler 的轻量化,最终选择 gRPC 实现远程日志访问

8.1 日志按任务实例切分

文档最初的设计是自定义 Logback 的 FileAppender 与 Filter:通过线程名解析出processDefineId_processInstanceId_taskInstanceId,生成形如/流程定义id/流程实例id/任务实例id.log的日志文件。文档中给出了早期的TaskLogAppender(自定义 Appender,从线程名解析 logId)与TaskLogFilter(匹配TaskLogInfo-前缀线程名)示例代码。

从当前仓库源码看,该能力已演进为更成熟的MDC + SiftingAppender方案:

  • TaskLogFilter.java:判断当前日志事件的 MDC 中是否存在任务实例日志全路径键(LogUtils.TASK_INSTANCE_LOG_FULL_PATH_MDC_KEY,值为"taskInstanceLogFullPath"),存在则 ACCEPT,否则 DENY;
  • TaskLogDiscriminator.java:继承 LogbackAbstractDiscriminator,从 MDC 读取taskInstanceLogFullPath作为区分值,用于将不同任务实例的日志写入不同文件;
  • LogUtils.java 提供setTaskInstanceLogFullPathMDC(...)/ 读取 / 清理 MDC 键的方法,在任务执行时把日志路径写入 MDC,任务结束后清理。

Worker 侧的 logback-spring.xml 完整展示了这套机制:TASKLOGFILEAppender 使用SiftingAppender,挂载TaskLogFilter过滤,通过TaskLogDiscriminatortaskInstanceLogFullPath动态 sift 出独立的FileAppender,将日志写入${taskInstanceLogFullPath}指定文件;同时使用SensitiveDataConverter做敏感信息脱敏(%message转换规则),主日志WORKERLOGFILE按大小与时间滚动(单文件最大 200MB、保留 168 小时、总容量上限 50GB)。

8.2 gRPC 远程读取

日志文件生成在 Worker 本地后,通过 gRPC 服务对外提供读取能力。当前 Master 侧的 MasterLogServiceImpl.java 实现了日志相关 gRPC 接口,包括:

  • pageQueryTaskInstanceLog:按页查询任务实例日志;
  • getTaskInstanceWholeLogFileBytes:获取任务实例完整日志文件字节流;
  • getAppId:从日志中解析应用 ID(如 Yarn ApplicationId);
  • removeTaskInstanceLog:删除指定路径的任务实例日志。

这样,UI 层无论与 Worker 相隔多远,都能通过 API → Master → gRPC 的链路按需、分页地拉取某个具体任务实例的日志内容。

总结

从调度入口看,DolphinScheduler 以"流程定义 DAG → command 表驱动 → Master 事件驱动拆解 → Worker 执行 → 结果回流"为主干,以 ZooKeeper 注册中心为基石,承载了去中心化集群管理、分布式锁选举、宕机容错、任务优先级排序与任务队列等核心机制;在运维体验层面,通过 Logback 按任务实例切分日志与 gRPC 远程读取,实现了分布式环境下日志的"本地化生成、跨节点访问"。本文对应完整源码可继续查阅 dolphinscheduler-master、dolphinscheduler-worker、dolphinscheduler-registry 与 dolphinscheduler-task-plugin 等模块,结合仓库中的单元测试(如 MasterTaskExecThreadTest.java)可进一步验证各执行链路的行为细节。

  • 任务调度
  • 大数据
  • 后端
  • 前端

【免费下载链接】dolphinscheduler

Apache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code

项目地址:https://gitcode.com/gh_mirrors/do/dolphinscheduler
点击查看免费下载

相关推荐

上一篇:Cerebro屏幕亮度调节终极指南:快速调整显示器亮度
下一篇:Cerebro护眼模式插件:5分钟搞定蓝光过滤保护视力

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询