源码解析-- 海豚调度-master与worker的交互过程
yuyutoo 2024-10-21 12:00 9 浏览 0 评论
海豚调度dolphinscheduler目前是 Apache 顶级项目,作为国内优秀的开源项目,它的架构设计理念会有很多值得我们学习和借鉴。
海豚调度dolphinscheduler是分布式易扩展的可视化DAG工作流任务调度系统
本文会包含如下内容:
- 海豚调度任务执行过程中master与worker的交互过程
- 如何处理过程中的异常
本篇文章适合人群:架构师、技术专家以及对任务调度非常感兴趣的高级工程师
本文以海豚1.3.5的源代码进行分析的。
1. master与worker的消息处理器
DolphinScheduler的master与worker是不同的JVM进程,正常情况下部署在不同的服务器中,master与worker是基于netty实现RPC交互的,共用到7个消息处理器
处理器在WorkerServer和MasterServer启动时,注册到NettyRemotingServer的NettyServerHandler中processors集合中
在netty channel的接收到channelRead数据,并转换为Command时,由NettyServerHandler中的processReceived方法,根据commandType交给对应的处理器处理。
所属进程名称 | 处理器名称 | 功能描述 |
MasterServer | TaskAckProcessor | 处理TaskExecuteAckCommand消息,将消息添加到TaskResponseService的任务响应队列中 |
MasterServer | TaskResponseProcessor | 处理TaskExecuteResponseCommand消息,将消息添加到TaskResponseService的任务响应队列中 |
MasterServer | TaskKillResponseProcessor | 处理TaskKillResponseCommand消息,并在日志中打印消息内容 |
WorkerServer | TaskExecuteProcessor | 处理TaskExecuteRequestCommand消息,并发送TaskExecuteAckCommand到master,提交任务执行 |
WorkerServer | DBTaskAckProcessor | 处理DBTaskAckCommand消息,针对执行成功的,从ResponceCache中删除 |
WorkerServer | DBTaskResponseProcessor | 处理DBTaskResponseCommand消息,针对执行成功的,从ResponceCache中删除 |
WorkerServer | TaskKillProcessor | 处理TaskKillRequestCommand消息,调用kill -9 pid杀死任务对应的进程,并向master发送TaskKillResponseCommand消息 |
2. master与worker的交互
DolphinScheduler的master与worker是不同的JVM进程,正常情况下部署在不同的服务器中,master与worker是基于netty实现RPC交互的,整个过程都是异步的。
正常的交互流程如下图:
- MasterServer根据流程定义产生的DAG进行任务切分,将每个任务放在MasterBaseTaskExecThread的子类中进行执行。
只有MasterTaskExecThread子类在将Callable提交到线程池,调用call方法时,执行submitWaitComplete方法,
提交就是将任务加到任务优先级队列中taskPriorityQueue - TaskPriorityQueueConsumer线程从taskPriorityQueue队列,根据fetchTaskNum【通过master.dispatch.task.num指定】参数,获取fetchTaskNum个任务信息,然后针对这批任务进行派发。在任务派发前构建ExecutionContext
- ExecutorDispatcher派发任务
- 使用配置文件中指定的负载均衡算法,选择执行任务的worker
- 将任务下发(send)到指定的worker执行。
- TaskExecuteProcessor接收到TaskExecuteRequestCommand后,发送任务确认消息【TaskExecuteAckCommand】到master,并启动一个TaskExecuteThread执行这个任务,任务执行完成后发送执行结果到master
- TaskAckProcessor接收到TaskExecuteAckCommand后,构建一个ACK类型的TaskResponseEvent,放到TaskResponseService的任务响应队列中
- TaskResponseProcessor接收到TaskExecuteResponseCommand后,构建一个ACK类型的TaskResponseEvent,放到TaskResponseService的任务响应队列中
- TaskResponseService线程将接收到的TaskResponseEvent放到任务响应队列中,并从列表中take TaskResponseEvent, 根据Event任务,进行任务状态的更新
如果是ACK EVENT,则更新任务实例的执行开始时间、执行worker、ExecutePath及logPath等信息,更新完成后发送DBTaskAckCommand到worker
如果是RESULT EVENT,则更新任务实例的执行结束时间、ProcessId、AppIds等信息,更新完成后发送DBTaskResponseCommand到worker
10. DBTaskAckProcessor接收到DBTaskAckCommand消息后,如果ExecutionStatus为SUCCESS,则将任务从ackCache中删除
11. DBTaskResponseProcessor接收到DBTaskResponseCommand消息后,如果ExecutionStatus为SUCCESS,则将任务从responeCache中删除
12. 在MasterTaskExecThread的submitWaitComplete方法中,会循环检查任务实例是否完成,在执行完步骤11后,则此任务实例完成,将循环也将退出
13. 针对派发失败的任务,添加到failedDispatchTasks,在这批任务派发完毕后,重新将失败的任务添加到taskPriorityQueue,如果taskPriorityQueue中的任务数小于失败的任务数,则程序休眠1秒
3. 交互异常情况处理
3.1 任务下发到worker失败与重试
- 将任务send到具体的worker时,如果失败,会重试3次,如果三次都失败,则将此worker节点从任务对应的任务组worker列表中删除,并从剩下的worker中选择第一
- 如果失败,则重复步骤1
- 如果任务组worker列表中所有worker都在重试3次后失败,则任务下发失败
- 将失败的任务添加到failedDispatchTasks列表中,因为master是按批处理任务,在这批任务派发完毕后,重新将失败的任务添加到taskPriorityQueue,如果taskPriorityQueue中的任务数小于失败的任务数,则程序休眠1秒
//TaskPriorityQueueConsumer 117行
if (!failedDispatchTasks.isEmpty()) {
for (String dispatchFailedTask : failedDispatchTasks) {
taskPriorityQueue.put(dispatchFailedTask);
}
// If there are tasks in a cycle that cannot find the worker group,
// sleep for 1 second
if (taskPriorityQueue.size() <= failedDispatchTasks.size()) {
TimeUnit.MILLISECONDS.sleep(Constants.SLEEP_TIME_MILLIS);
}
}
//NettyExecutorManager 109行
Host host = context.getHost();
boolean success = false;
while (!success) {
try {
doExecute(host,command);
success = true;
context.setHost(host);
} catch (ExecuteException ex) {
logger.error(String.format("execute command : %s error", command), ex);
try {
failNodeSet.add(host.getAddress());
Set<String> tmpAllIps = new HashSet<>(allNodes);
Collection<String> remained = CollectionUtils.subtract(tmpAllIps, failNodeSet);
if (remained != null && remained.size() > 0) {
host = Host.of(remained.iterator().next());
logger.error("retry execute command : {} host : {}", command, host);
} else {
throw new ExecuteException("fail after try all nodes");
}
} catch (Throwable t) {
throw new ExecuteException("fail after try all nodes");
}
}
}
3.2 任务ack及result上报重试
worker在接收到TaskExecuteRequestCommand命令后,会向master发送任务确认消息;worker在任务执行完成后,也会向master发送任务执行完成消息;在发送消息前,会调用ResponceCache的cache方法将方法缓存。
RetryReportTaskStatusThread线程每隔5分钟,判断responceCache中ackCache和responseCache中是否为空,如果不为空,则将命令TaskExecuteAckCommand或TaskExecuteResponseCommand重新发送到master节点
相关推荐
- 当 Linux 根分区 (/) 已满时如何释放空间?
-
根分区(/)是Linux文件系统的核心,包含操作系统核心文件、配置文件、日志文件、缓存和用户数据等。当根分区满载时,系统可能出现无法写入新文件、应用程序崩溃甚至无法启动的情况。常见原因包括:...
- 玩转 Linux 之:磁盘分区、挂载知多少?
-
今天来聊聊linux下磁盘分区、挂载的问题,篇幅所限,不会聊的太底层,纯当科普!!1、Linux分区简介1.1主分区vs扩展分区硬盘分区表中最多能存储四个分区,但我们实际使用时一般只分为两...
- Linux 文件搜索神器 find 实战详解,建议收藏
-
在Linux系统使用中,作为一个管理员,我希望能查找系统中所有的大小超过200M文件,查看近7天系统中哪些文件被修改过,找出所有子目录中的可执行文件,这些任务需求...
- Linux 操作系统磁盘操作(linux 磁盘命令)
-
一、文档介绍本文档描述Linux操作系统下多种场景下的磁盘操作情况。二、名词解释...
- Win10新版19603推送:一键清理磁盘空间、首次集成Linux文件管理器
-
继上周四的Build19592后,微软今晨面向快速通道的Insider会员推送Windows10新预览版,操作系统版本号Build19603。除了一些常规修复,本次更新还带了不少新功能,一起来了...
- Android 16允许Linux终端使用手机全部存储空间
-
IT之家4月20日消息,谷歌Pixel手机正朝着成为强大便携式计算设备的目标迈进。2025年3月的更新中,Linux终端应用的推出为这一转变奠定了重要基础。该应用允许兼容的安卓设备...
- Linux 系统管理大容量磁盘(2TB+)操作指南
-
对于容量超过2TB的磁盘,传统MBR分区表的32位寻址机制存在限制(最大支持2.2TB)。需采用GPT(GUIDPartitionTable)分区方案,其支持64位寻址,理论上限为9.4ZB(9....
- Linux 服务器上查看磁盘类型的方法
-
方法1:使用lsblk命令lsblk输出说明:TYPE列显示设备类型,如disk(物理磁盘)、part(分区)、rom(只读存储)等。...
- ESXI7虚机上的Ubuntu Linux 22.04 LVM空间扩容操作记录
-
本人在实际的使用中经常遇到Vmware上安装的Linux虚机的LVM扩容情况,最终实现lv的扩容,大多数情况因为虚机都是有备用或者可停机的情况,一般情况下通过添加一块物理盘再加入vg,然后扩容lv来实...
- 5.4K Star很容易!Windows读取Linux磁盘格式工具
-
[开源日记],分享10k+Star的优质开源项目...
- Linux 文件系统监控:用脚本自动化磁盘空间管理
-
在Linux系统中,文件系统监控是一项非常重要的任务,它可以帮助我们及时发现磁盘空间不足的问题,避免因磁盘满而导致的系统服务不可用。通过编写脚本自动化磁盘空间管理,我们可以更加高效地处理这一问题。下面...
- Linux磁盘管理LVM实战(linux实验磁盘管理)
-
LVM(逻辑卷管理器,LogicalVolumeManager)是一种在Linux系统中用于灵活管理磁盘空间的技术,通过将物理磁盘抽象为逻辑卷,实现动态调整存储容量、跨磁盘扩展等功能。本章节...
- Linux查看文件大小:`ls`和`du`为何结果不同?一文讲透原理!
-
Linux查看文件大小:ls和du为何结果不同?一文讲透原理!在Linux运维中,查看文件大小是日常高频操作。但你是否遇到过以下困惑?...
- 使用 df 命令检查服务器磁盘满了,但用 du 命令发现实际小于磁盘容量
-
在Linux系统中,管理员或开发者经常会遇到一个令人困惑的问题:使用...
- Linux磁盘爆满紧急救援指南:5步清理释放50GB+小白也能轻松搞定
-
“服务器卡死?网站崩溃?当Linux系统弹出‘Nospaceleft’的红色警报,别慌!本文手把手教你从‘删库到跑路’进阶为‘磁盘清理大师’,5个关键步骤+30条救命命令,快速释放磁盘空间,拯救你...
你 发表评论:
欢迎- 一周热门
- 最近发表
- 标签列表
-
- mybatis plus (70)
- scheduledtask (71)
- css滚动条 (60)
- java学生成绩管理系统 (59)
- 结构体数组 (69)
- databasemetadata (64)
- javastatic (68)
- jsp实用教程 (53)
- fontawesome (57)
- widget开发 (57)
- vb net教程 (62)
- hibernate 教程 (63)
- case语句 (57)
- svn连接 (74)
- directoryindex (69)
- session timeout (58)
- textbox换行 (67)
- extension_dir (64)
- linearlayout (58)
- vba高级教程 (75)
- iframe用法 (58)
- sqlparameter (59)
- trim函数 (59)
- flex布局 (63)
- contextloaderlistener (56)