ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache DolphinScheduler 架构设计深度解析:从调度核心名词到分布式容错与日志原理

Apache DolphinScheduler 架构设计深度解析:从调度核心名词到分布式容错与日志原理 任务调度大数据后端前端【免费下载链接】dolphinschedulerApache 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 GraphDAG的形式组装从入度为 0 的节点开始做拓扑遍历直到不存在后继节点为止。DolphinScheduler 的编排页面正是基于 DAG 的思想将任务节点拖拽连线形成执行依赖关系1.2 流程定义、流程实例与任务实例流程定义Process definition通过拖拽任务节点并建立节点关联将DAG可视化出来得到的产物流程实例Process instance流程定义的实例化。可通过手动启动或定时调度触发流程定义每运行一次就产生一个新的流程实例任务实例Task instance流程实例运行时其中某个具体任务节点的实例化用于记录该任务的具体执行状态。三者是定义 → 实例 → 任务执行的层级关系也是数据库中t_ds_process_definition、t_ds_process_instance、t_ds_task_instance等表设计的心智模型。1.3 任务类型文档最初描述的系统已支持 SHELL、SQL、SUB_PROCESS子流程、PROCEDURE、MR、SPARK、PYTHON、DEPENDENT依赖等任务类型并计划支持动态插件化扩展。需要特别说明的是SUB_PROCESS 本身也是一个可独立启动的流程定义。从当前仓库源码结构看这一设计已经落地为庞大的插件体系dolphinscheduler-task-plugin 目录下按任务类型拆分模块如dolphinscheduler-task-shell、dolphinscheduler-task-spark、dolphinscheduler-task-flink、dolphinscheduler-task-http、dolphinscheduler-task-datax等数十种任务类型由 dolphinscheduler-task-api 提供统一抽象真正做到新增一种任务类型只需新增一个插件模块的动态扩展。1.4 调度方式与命令类型系统支持基于 cron 表达式的定时调度与手动调度。命令Command类型覆盖了流程运行的各类触发与干预动作。文档中提到其中恢复容错工作流与恢复等待线程两种命令类型由调度系统内部控制外部不可调用。从当前源码 CommandType.java 可以看到完整的命令枚举共 14 种代码编号 013编号命令类型含义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 MasterServerMasterServer 采用分布式非中心化设计理念主要负责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 WorkerServerWorkerServer 同样采用分布式、非中心化设计主要负责任务执行与日志服务。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 与 UIAPI接口层负责处理来自前端 UI 层的请求对外提供RESTful API覆盖流程的创建、定义、查询、修改、上线、下线、手动启动、停止、暂停、恢复、从当前节点开始执行等操作。对应模块为 dolphinscheduler-apiUI系统前端页面提供各类可视化操作界面见仓库 docs/docs/en/guide 下的使用指南对应前端工程为 dolphinscheduler-ui。3. 架构设计思想去中心化 vs 中心化3.1 中心化设计及其问题中心化设计思路相对简单集群节点按角色分为 Master 与 Slave 两类。Master 负责任务分发并监督 Slave 健康状态可动态均衡地将任务分配至各 Slave避免节点忙闲不均WorkerSlave负责任务执行并与 Master 保持心跳以便其分配任务。但中心化设计存在两个突出问题单点故障一旦 Master 出问题集群失去领导者整个集群崩溃。多数 Master/Slave 架构采用主备 Master 方案缓解热备或冷备、自动或手动切换越来越多的系统具备自动选举切换 Master 的能力以提升可用性调度器位置的两难若 Scheduler 放在 Master 上虽能支持同一 DAG 中不同任务运行在不同机器但会加重 Master 负载若放在 Slave 上一个 DAG 中的所有任务只能提交到同一台机器并行任务多时该 Slave 压力过大。3.2 去中心化设计去中心化设计通常没有 Master/Slave 概念所有角色地位平等——互联网本身就是典型的去中心化分布式系统任意节点宕机只会影响小范围功能。其核心在于整个分布式系统中不存在管理者节点因此没有单点故障问题但代价是每个节点都需要与其他节点通信获取必要信息分布式通信链路的不可靠性大幅增加了实现难度。实际上真正完全去中心化的系统很少见取而代之的是动态中心化系统集群中的管理者动态选举产生而非预设集群故障时节点自发开会选出新管理者主持工作最典型的案例即 ZooKeeper 与 Go 实现的 Etcd。DolphinScheduler 的去中心化实践Master/Worker 注册到 ZooKeeperMaster 集群与 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 等待 SubFlowThread2SubFlowThread2 等待 SubFlowThread3而 SubFlowThread3 等待线程池分配新线程——整个 DAG 永远无法结束线程也无法释放形成子父流程循环等待。此时除非启动新的 Master 增加线程来打破卡死否则调度集群将不可用。启动新 Master 显然不是优雅方案文档给出了三种候选解法预计算线程数先求所有 Master 线程总数再计算每个 DAG 所需线程数在 DAG 执行前预判。但多 Master 线程池的总线程数难以实时获取单 Master 线程池满则直接失败若线程池已满让线程直接失败新增资源不足命令类型线程池不足时将主流程挂起待线程池有新的线程后再唤醒资源不足的流程。注意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.javaHIGHEST0、HIGH1、MEDIUM2、LOW3、LOWEST4完全对应。8. Logback 与 gRPC 实现远程日志访问由于 WebUI与 Worker 不一定在同一台机器上查看日志不能像查询本地文件那样直接进行存在两种可选方案将日志放入 ES 搜索引擎通过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否则 DENYTaskLogDiscriminator.java继承 LogbackAbstractDiscriminator从 MDC 读取taskInstanceLogFullPath作为区分值用于将不同任务实例的日志写入不同文件LogUtils.java 提供setTaskInstanceLogFullPathMDC(...)/ 读取 / 清理 MDC 键的方法在任务执行时把日志路径写入 MDC任务结束后清理。Worker 侧的 logback-spring.xml 完整展示了这套机制TASKLOGFILEAppender 使用SiftingAppender挂载TaskLogFilter过滤通过TaskLogDiscriminator按taskInstanceLogFullPath动态 sift 出独立的FileAppender将日志写入${taskInstanceLogFullPath}指定文件同时使用SensitiveDataConverter做敏感信息脱敏%message转换规则主日志WORKERLOGFILE按大小与时间滚动单文件最大 200MB、保留 168 小时、总容量上限 50GB。8.2 gRPC 远程读取日志文件生成在 Worker 本地后通过 gRPC 服务对外提供读取能力。当前 Master 侧的 MasterLogServiceImpl.java 实现了日志相关 gRPC 接口包括pageQueryTaskInstanceLog按页查询任务实例日志getTaskInstanceWholeLogFileBytes获取任务实例完整日志文件字节流getAppId从日志中解析应用 ID如 Yarn ApplicationIdremoveTaskInstanceLog删除指定路径的任务实例日志。这样UI 层无论与 Worker 相隔多远都能通过 API → Master → gRPC 的链路按需、分页地拉取某个具体任务实例的日志内容。总结从调度入口看DolphinScheduler 以流程定义 DAG → command 表驱动 → Master 事件驱动拆解 → Worker 执行 → 结果回流为主干以 ZooKeeper 注册中心为基石承载了去中心化集群管理、分布式锁选举、宕机容错、任务优先级排序与任务队列等核心机制在运维体验层面通过 Logback 按任务实例切分日志与 gRPC 远程读取实现了分布式环境下日志的本地化生成、跨节点访问。本文对应完整源码可继续查阅 dolphinscheduler-master、dolphinscheduler-worker、dolphinscheduler-registry 与 dolphinscheduler-task-plugin 等模块结合仓库中的单元测试如 MasterTaskExecThreadTest.java可进一步验证各执行链路的行为细节。赞分享任务调度大数据后端前端【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址https://gitcode.com/gh_mirrors/do/dolphinscheduler点击查看免费下载相关推荐DolphinScheduler分布式调度架构深度解析从核心设计到企业级实践DolphinScheduler分布式调度架构深度解析从核心设计到企业级实践 Apache DolphinScheduler作为一款开源的分布式工作流任务调度任务调度数据编排工作流自动化后端大数据Apache DolphinScheduler 系统架构设计深度解析去中心化调度、分布式锁、容错与任务优先级Apache DolphinScheduler 系统架构设计深度解析去中心化调度、分布式锁、容错与任务优先级 本篇文章以 Apache DolphinSche任务调度数据编排工作流自动化后端大数据AI Scientist-v2零基础上手:从一段主题描述自动写出工作坊论文AI Scientist v2零基础上手:从一段主题描述自动写出工作坊论文 AI Scientist v2 是一套自动化科研系统:给它一份主题描述,它就能自主产人工智能AI Agent自主智能体科研Agent 工作流上一篇Cerebro屏幕亮度调节终极指南快速调整显示器亮度下一篇Cerebro护眼模式插件5分钟搞定蓝光过滤保护视力创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表