• linkedu视频
  • 平面设计
  • 电脑入门
  • 操作系统
  • 办公应用
  • 电脑硬件
  • 动画设计
  • 3D设计
  • 网页设计
  • CAD设计
  • 影音处理
  • 数据库
  • 程序设计
  • 认证考试
  • 信息管理
  • 信息安全
菜单
linkedu.com
导航菜单
  • 网页制作
  • 数据库
  • 程序设计
  • 操作系统
  • CMS教程
  • 游戏攻略
  • 脚本语言
  • 平面设计
  • 软件教程
  • 网络安全
  • 电脑知识
  • 服务器
  • 视频教程
  • windows
  • 服务器硬件
  • 服务器运维
  • 云计算
  • 虚拟化
  • IIS教程
  • Linux
  • Apache
  • Ftp
  • DNS
  • Nginx
您的位置:首页 > 服务器 >云计算 > [Spark源码剖析] DAGScheduler提交stage,sparkdagscheduler

[Spark源码剖析] DAGScheduler提交stage,sparkdagscheduler

作者:网友 字体:[增加 减小] 来源:互联网

本文主要包含spark stage,spark源码剖析,apache spark源码剖析,spark源码下载,spark源码等服务器相关知识,网友希望可以进行参考

[Spark源码剖析] DAGScheduler提交stage,sparkdagscheduler


DAGScheduler通过调用submitStage来提交stage,实现如下:

  private def submitStage(stage: Stage) {
    val jobId = activeJobForStage(stage)
    if (jobId.isDefined) {
      logDebug("submitStage(" + stage + ")")
      if (!waitingStages(stage) && !runningStages(stage) && !failedStages(stage)) {
        //< 获取该stage未提交的父stages,并按stage id从小到大排序
        val missing = getMissingParentStages(stage).sortBy(_.id)
        logDebug("missing: " + missing)
        if (missing.isEmpty) {
          logInfo("Submitting " + stage + " (" + stage.rdd + "), which has no missing parents")
          //< 若无未提交的父stage, 则提交该stage对应的tasks
          submitMissingTasks(stage, jobId.get)
        } else {
          //< 若存在未提交的父stage, 依次提交所有父stage (若父stage也存在未提交的父stage, 则提交之, 依次类推); 并把该stage添加到等待stage队列中
          for (parent <- missing) {
            submitStage(parent)
          }
          waitingStages += stage
        }
      }
    } else {
      abortStage(stage, "No active job for stage " + stage.id)
    }
  }

submitStage先调用getMissingParentStages来获取参数stageX(这里为了区分在取名为stageX)是否有未提交的父stages,若有,则依次递归(按stage id从小到大排列,也就是stage是从后往前提交的)提交父stage,并将stageX加入到waitingStages: HashSet[Stage]中。对于要依次提交的父stage,也是如此。

getMissingParentStages与DAGScheduler划分stage中介绍的getParentStages有点像,但不同的是不再需要划分stage,并对每个stage的状态做了判断,源码及注释如下:

//< 以参数stage为起点,向前遍历所有stage,判断stage是否为未提交,若使则加入missing中
  private def getMissingParentStages(stage: Stage): List[Stage] = {
    //< 未提交的stage
    val missing = new HashSet[Stage]
    //< 存储已经被访问到得RDD
    val visited = new HashSet[RDD[_]]

    val waitingForVisit = new Stack[RDD[_]]
    def visit(rdd: RDD[_]) {
      if (!visited(rdd)) {
        visited += rdd
        if (getCacheLocs(rdd).contains(Nil)) {
          for (dep <- rdd.dependencies) {
            dep match {
              //< 若为宽依赖,生成新的stage
              case shufDep: ShuffleDependency[_, _, _] =>
                //< 这里调用getShuffleMapStage不像在getParentStages时需要划分stage,而是直接根据shufDep.shuffleId获取对应的ShuffleMapStage
                val mapStage = getShuffleMapStage(shufDep, stage.jobId)
                if (!mapStage.isAvailable) {
                  // 若stage得状态为available则为未提交stage
                  missing += mapStage
                }
              //< 若为窄依赖,那就属于同一个stage。并将依赖的RDD放入waitingForVisit中,以能够在下面的while中继续向上visit,直至遍历了整个DAG图
              case narrowDep: NarrowDependency[_] =>
                waitingForVisit.push(narrowDep.rdd)
            }
          }
        }
      }
    }
    waitingForVisit.push(stage.rdd)
    while (waitingForVisit.nonEmpty) {
      visit(waitingForVisit.pop())
    }
    missing.toList
  }

上面提到,若stageX存在未提交的父stages,则先提交父stages;那么,如果stageX没有未提交的父stage呢(比如,包含从HDFS读取数据生成HadoopRDD的那个stage是没有父stage的)?

这时会调用submitMissingTasks(stage, jobId.get),参数就是stageX及其对应的jobId.get。这个函数便是我们时常在其他文章或书籍中看到的将stage与taskSet对应起来,然后DAGScheduler将taskSet提交给TaskScheduler去执行的实施者。这个函数的实现比较长,下面分段说明。

Step1: 得到RDD中需要计算的partition

对于Shuffle类型的stage,需要判断stage中是否缓存了该结果;对于Result类型的Final Stage,则判断计算Job中该partition是否已经计算完成。这么做的原因是,stage中某个task执行失败其他执行成功地时候就需要找出这个失败的task对应要计算的partition而不是要计算所有partition

  private def submitMissingTasks(stage: Stage, jobId: Int) {
    stage.pendingTasks.clear()


    //< 首先得到RDD中需要计算的partition
    //< 对于Shuffle类型的stage,需要判断stage中是否缓存了该结果;
    //< 对于Result类型的Final Stage,则判断计算Job中该partition是否已经计算完成
    //< 这么做的原因是,stage中某个task执行失败其他执行成功地时候就需要找出这个失败的task对应要计算的partition而不是要计算所有partition
    val partitionsToCompute: Seq[Int] = {
      stage match {
        case stage: ShuffleMapStage =>
          (0 until stage.numPartitions).filter(id => stage.outputLocs(id).isEmpty)
        case stage: ResultStage =>
          val job = stage.resultOfJob.get
          (0 until job.numPartitions).filter(id => !job.finished(id))
      }
    }

Step2: 序列化task的binary

Executor可以通过广播变量得到它。每个task运行的时候首先会反序列化

var taskBinary: Broadcast[Array[Byte]] = null
    try {
      // For ShuffleMapTask, serialize and broadcast (rdd, shuffleDep).
      // For ResultTask, serialize and broadcast (rdd, func).
      val taskBinaryBytes: Array[Byte] = stage match {
        case stage: ShuffleMapStage =>
          //< 对于ShuffleMapTask,将rdd及其依赖关系序列化;在Executor执行task之前会反序列化
          closureSerializer.serialize((stage.rdd, stage.shuffleDep): AnyRef).array()
          //< 对于ResultTask,对rdd及要在每个partition上执
  


 
分享到:QQ空间新浪微博腾讯微博微信百度贴吧QQ好友复制网址打印

您可能想查找下面的文章:

  • [Spark源码剖析] DAGScheduler提交stage,sparkdagscheduler
  • [Spark源码剖析] DAGScheduler划分stage,sparkdagscheduler

相关文章

  • Docker exec与Docker attach,dockerexecattach
  • 几本京东上的CISCO云计算相关书籍,cisco相关书籍
  • 《2》CentOS7.0+OpenStack+kvm云平台部署—配置Keystone,centos7.0kvm
  • HDFS小文件的合并优化
  • 开源图计算框架GraphLab介绍,开源图框架graphlab
  • Hadoop2.6 Ha 安装,hadoop2.6ha安装
  • hive创建表语句详解,hive语句详解
  • HBase扫描器与过滤器,HBase扫描器过滤器
  • HDFS读文件解析,
  • Storm简述及集群安装,Storm简述集群安装

文章分类

  • windows
  • 服务器硬件
  • 服务器运维
  • 云计算
  • 虚拟化
  • IIS教程
  • Linux
  • Apache
  • Ftp
  • DNS
  • Nginx

最近更新的内容

    • 给大数据文件的每一行产生唯一的id,数据一行id
    • 《5》CentOS7.0+OpenStack+kvm云平台部署—配置Horizon,《5》kvm
    • LXC学习,学习中国app上线
    • openstack中nova组件Hypervisors、Floating_ips的所有python API 汇总,openstackpython
    • Hadoop之——Linux基本命令回顾,hadooplinux回顾
    • 一道hadoop面试题,hadoop面试题
    • MapReduce处理表的自连接,mapreduce处理表
    • hadoop-2.6.0安装(new),hadoop-2.6.0new
    • Hadoop启动时报错:Incorrect configuration: namenode address dfs.namenode.servicerpc-address or...,hadoopnamenode
    • Centos 64下实现socket通信,centos64下

关于我们 - 联系我们 - 免责声明 - 网站地图

©2020-2025 All Rights Reserved. linkedu.com 版权所有