百度360必应搜狗淘宝本站头条
当前位置:网站首页 > 编程网 > 正文

源码解析-- 海豚调度-master与worker的交互过程

yuyutoo 2024-10-21 12:00 10 浏览 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交互的,整个过程都是异步的。

正常的交互流程如下图:

  1. MasterServer根据流程定义产生的DAG进行任务切分,将每个任务放在MasterBaseTaskExecThread的子类中进行执行。
    只有MasterTaskExecThread子类在将Callable提交到线程池,调用call方法时,执行submitWaitComplete方法,
    提交就是将任务加到任务优先级队列中taskPriorityQueue
  2. TaskPriorityQueueConsumer线程从taskPriorityQueue队列,根据fetchTaskNum【通过master.dispatch.task.num指定】参数,获取fetchTaskNum个任务信息,然后针对这批任务进行派发。在任务派发前构建ExecutionContext
  3. ExecutorDispatcher派发任务
  4. 使用配置文件中指定的负载均衡算法,选择执行任务的worker
  5. 将任务下发(send)到指定的worker执行。
  6. TaskExecuteProcessor接收到TaskExecuteRequestCommand后,发送任务确认消息【TaskExecuteAckCommand】到master,并启动一个TaskExecuteThread执行这个任务,任务执行完成后发送执行结果到master
  7. TaskAckProcessor接收到TaskExecuteAckCommand后,构建一个ACK类型的TaskResponseEvent,放到TaskResponseService的任务响应队列中
  8. TaskResponseProcessor接收到TaskExecuteResponseCommand后,构建一个ACK类型的TaskResponseEvent,放到TaskResponseService的任务响应队列中
  9. 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失败与重试

  1. 将任务send到具体的worker时,如果失败,会重试3次,如果三次都失败,则将此worker节点从任务对应的任务组worker列表中删除,并从剩下的worker中选择第一
  2. 如果失败,则重复步骤1
  3. 如果任务组worker列表中所有worker都在重试3次后失败,则任务下发失败
  4. 将失败的任务添加到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节点

相关推荐

墨尔本一华裔男子与亚裔男子分别失踪数日 警方寻人

中新网5月15日电据澳洲新快网报道,据澳大利亚维州警察局网站消息,22岁的华裔男子邓跃(Yue‘Peter’Deng,音译)失踪已6天,维州警方于当地时间13日发布寻人通告,寻求公众协助寻找邓跃。华...

网络交友须谨慎!美国犹他州一男子因涉嫌杀害女网友被捕

伊森·洪克斯克(图源网络,侵删)据美国广播公司(ABC)25日报道,美国犹他州一名男子于24日因涉嫌谋杀被捕。警方表示,这名男子主动告知警局,称其杀害了一名在网络交友软件上认识的25岁女子。雷顿警...

一课译词:来龙去脉(来龙去脉 的意思解释)

Mountainranges[Photo/SIPA]“来龙去脉”,汉语成语,本指山脉的走势和去向,现比喻一件事的前因后果(causeandeffectofanevent),可以翻译为“i...

高考重要考点:range(range高考用法)

range可以用作动词,也可以用作名词,含义特别多,在阅读理解中出现的频率很高,还经常作为完形填空的选项,而且在作文中使用是非常好的高级词汇。...

C++20 Ranges:现代范围操作(现代c++白皮书)

1.引言:C++20Ranges库简介C++20引入的Ranges库是C++标准库的重要更新,旨在提供更现代化、表达力更强的方式来处理数据序列(范围,range)。Ranges库基于...

学习VBA,报表做到飞 第二章 数组 2.4 Filter函数

第二章数组2.4Filter函数Filter函数功能与autofilter函数类似,它对一个一维数组进行筛选,返回一个从0开始的数组。...

VBA学习笔记:数组:数组相关函数—Split,Join

Split拆分字符串函数,语法Split(expression,字符,Limit,compare),第1参数为必写,后面3个参数都是可选项。Expression为需要拆分的数据,“字符”就是以哪个字...

VBA如何自定义序列,学会这些方法,让你工作更轻松

No.1在Excel中,自定义序列是一种快速填表机制,如何有效地利用这个方法,可以大大增加工作效率。通常在操作工作表的时候,可能会输入一些很有序的序列,如果一一录入就显得十分笨拙。Excel给出了一种...

Excel VBA入门教程1.3 数组基础(vba数组详解)

1.3数组使用数组和对象时,也要声明,这里说下数组的声明:'确定范围的数组,可以存储b-a+1个数,a、b为整数Dim数组名称(aTob)As数据类型Dimarr...

远程网络调试工具百宝箱-MobaXterm

MobaXterm是一个功能强大的远程网络工具百宝箱,它将所有重要的远程网络工具(SSH、Telnet、X11、RDP、VNC、FTP、MOSH、Serial等)和Unix命令(bash、ls、cat...

AREX:携程新一代自动化回归测试工具的设计与实现

一、背景随着携程机票BU业务规模的不断提高,业务系统日趋复杂,各种问题和挑战也随之而来。对于研发测试团队,面临着各种效能困境,包括业务复杂度高、数据构造工作量大、回归测试全量回归、沟通成本高、测试用例...

Windows、Android、IOS、Web自动化工具选择策略

Windows平台中应用UI自动化测试解决方案AutoIT是开源工具,该工具识别windows的标准控件效果不错,但是当它遇到应用中非标准控件定义的UI元素时往往就无能为力了,这个时候选择silkte...

python自动化工具:pywinauto(python快速上手 自动化)

简介Pywinauto是完全由Python构建的一个模块,可以用于自动化Windows上的GUI应用程序。同时,它支持鼠标、键盘操作,在元素控件树较复杂的界面,可以辅助我们完成自动化操作。我在...

时下最火的 Airtest 如何测试手机 APP?

引言Airtest是网易出品的一款基于图像识别的自动化测试工具,主要应用在手机APP和游戏的测试。一旦使用了这个工具进行APP的自动化,你就会发现自动化测试原来是如此简单!!连接手机要进行...

【推荐】7个最强Appium替代工具,移动App自动化测试必备!

在移动应用开发日益火爆的今天,自动化测试成为了确保应用质量和用户体验的关键环节。Appium作为一款广泛应用的移动应用自动化测试工具,为测试人员所熟知。然而,在不同的测试场景和需求下,还有许多其他优...

取消回复欢迎 发表评论: