• Spark任务调度概述_大数据培训


    Spark 任务调度机制

    在工厂环境下,Spark集群的部署方式一般为YARN-Cluster模式,之后的内核分析内容中我们默认集群的部署方式为YARN-Cluster模式。在上一章中我们讲解了Spark YARN-Cluster模式下的任务提交流程,但是我们并没有具体说明Driver的工作流程, Driver线程主要是初始化SparkContext对象,准备运行所需的上下文,然后一方面保持与ApplicationMaster的RPC连接,通过ApplicationMaster申请资源,另一方面根据用户业务逻辑开始调度任务,将任务下发到已有的空闲Executor上。

    当ResourceManager向ApplicationMaster返回Container资源时,ApplicationMaster就尝试在对应的Container上启动Executor进程,Executor进程起来后,会向Driver反向注册,注册成功后保持与Driver的心跳,同时等待Driver分发任务,当分发的任务执行完毕后,将任务状态上报给Driver。

    Spark任务调度概述

    当Driver起来后,Driver则会根据用户程序逻辑准备任务,并根据Executor资源情况逐步分发任务。在详细阐述任务调度前,首先说明下Spark里的几个概念。一个Spark应用程序包括Job、Stage以及Task三个概念:

    (1)Job是以Action方法为界,遇到一个Action方法则触发一个Job;

    (2)Stage是Job的子集,以RDD宽依赖(即Shuffle)为界,遇到Shuffle做一次划分;

    (3)Task是Stage的子集,以并行度(分区数)来衡量,分区数是多少,则有多少个task。

    Spark的任务调度总体来说分两路进行,一路是Stage级的调度,一路是Task级的调度,总体调度流程如下图所示:

    图4-1 Spark任务调度概览

    Spark RDD通过其Transactions操作,形成了RDD血缘关系图,即DAG,最后通过Action的调用,触发Job并调度执行。DAGScheduler负责Stage级的调度,主要是将job切分成若干Stages,并将每个Stage打包成TaskSet交给TaskScheduler调度。TaskScheduler负责Task级的调度,将DAGScheduler给过来的TaskSet按照指定的调度策略分发到Executor上执行,调度过程中SchedulerBackend负责提供可用资源,其中SchedulerBackend有多种实现,分别对接不同的资源管理系统。

    图4-2 Job提交和Task拆分

    Driver初始化SparkContext过程中,会分别初始化DAGScheduler、TaskScheduler、SchedulerBackend以及HeartbeatReceiver,并启动SchedulerBackend以及HeartbeatReceiver。SchedulerBackend通过ApplicationMaster申请资源,并不断从TaskScheduler中拿到合适的Task分发到Executor执行。HeartbeatReceiver负责接收Executor的心跳信息,监控Executor的存活状况,并通知到TaskScheduler。

  • 相关阅读:
    Fortify-解决中文乱码
    三维模型3DTile格式轻量化的跨平台兼容性问题分析
    初始MySQL(七)(MySQL表类型和存储引擎,MySQL视图,MySQL用户管理)
    Linux系统管理指南:用户权限、进程管理和网络配置精解
    vue2 vue3 各传值汇总
    [极客大挑战 2019]Http1(BUCTF在线评测)
    SQLMAP自动注入-优化参数
    Python最新学习路线
    【云原生之kubernetes实战】在k8s集群环境下部署Tomcat应用
    【SA8295P 源码分析】103 - QNX DDR RAM 内存布局分析
  • 原文地址:https://blog.csdn.net/zjjcchina/article/details/126165156