• 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


本文基于Spark 1.3.1

先上一些stage相关的知识点:

stage的产生和提交的大致步骤如下:

RDD.action
    SparkContext.runJob
        DAGScheduler.runJob
        DAGScheduler.submitJob
        DAGScheduler.eventProcessLoop.post( JobSubmitted(...) )
        DAGScheduler.handleJobSubmitted
        DAGScheduler.submitStage

可以看到,在DAGScheduler内部通过post一个JobSubmitted事件来触发Job的提交

DAGScheduler.eventProcessLoop.post( JobSubmitted(...) )
DAGScheduler.handleJobSubmitted

既然这两个方法都是DAGScheduler内部的实现,为什么不直接调用函数而要这样“多此一举”呢?我猜想这是为了保证整个系统事件模型的完整性。


DAGScheduler.handleJobSubmitted部分源码及如下

private[scheduler] def handleJobSubmitted(jobId: Int,
      finalRDD: RDD[_],
      func: (TaskContext, Iterator[_]) => _,
      partitions: Array[Int],
      allowLocal: Boolean,
      callSite: CallSite,
      listener: JobListener,
      properties: Properties) {
    var finalStage: ResultStage = null
    try {
      //< 创建finalStage可能会抛出一个异常, 比如, jobs是基于一个HadoopRDD的但这个HadoopRDD已被删除
      finalStage = newResultStage(finalRDD, partitions.size, jobId, callSite)
    } catch {
      case e: Exception => 
        return
    }

    //< 此处省略n行代码
  }

handleJobSubmitted首先通过调用newResultStage来创建了一个ResultStage,即finalStage。这里先给出一个结论:在finalStage = newResultStage(finalRDD, partitions.size, jobId, callSite)调用完之后,stage的划分其实已经完成,并已将所有的ShuffleMapStage保存在成员shuffleToMapStage: HashMap[Int, ShuffleMapStage]中。

shuffleToMapStage的key的含义是什么?key和value又是怎么确定的呢?往下看~


跟进到DAGScheduler.newResultStage

  private def newResultStage(
      rdd: RDD[_],
      numTasks: Int,
      jobId: Int,
      callSite: CallSite): ResultStage = {
    val (parentStages: List[Stage], id: Int) = getParentStagesAndId(rdd, jobId)
    val stage: ResultStage = new ResultStage(id, rdd, numTasks, parentStages, jobId, callSite)

    stageIdToStage(id) = stage
    updateJobIdStageIdMaps(jobId, stage)
    stage
  }

DAGScheduler.newResultStage首先调用val (parentStages: List[Stage], id: Int) = getParentStagesAndId(rdd, jobId),这个调用看起来像是要先确定好该ResultStage依赖的所有父stages。

接下来,调用val stage: ResultStage = new ResultStage(id, rdd, numTasks, parentStages, jobId, callSite)来创建DAG图中唯一的一个ResultStage,new ResultStage值得注意的有两点
1. ResultStage构造函数参数numTasks的值等于newResultStage调用时传入的finalRDD的partitions.size,那么我们就知道了“对于ResultStage,task数目就是它finalRDD的partition个数”
2. new ResultStage时会使用val (parentStages: List[Stage], id: Int) = getParentStagesAndId(rdd, jobId)获得的该stage依赖的父stages来构造ResultStage,说明了ResultStage对象内部保存了它依赖的stages,可以在需要的时候使用

创建ResultStage的调用比较简单,得到ResultStage对象后存入成员stageIdToStage中。接下来,看看getParentStagesAndId都干了点啥

  private def getParentStagesAndId(rdd: RDD[_], jobId: Int): (List[Stage], Int) = {
    val parentStages = getParentStages(rdd, jobId)
    val id = nextStageId.getAndIncrement()
    (parentStages, id)
  }

跟到getParentStages里

  //< getParentStages函数以参数rdd为起点,通过一层一层遍历依赖来划分stage,划分stage的依据是ShuffleDependency,即宽依赖。
  //< 这个函数的实现方式非常巧妙
  private def getParentStages(rdd: RDD[_], jobId: Int): List[Stage] = {
    //< 通过vist一级一级vist得到的父stage
    val parents = new HashSet[Stage]
    //< 已经visted的rdd
    val visited = new HashSet[RDD[_]]
    val waitingForVisit = new Stack[RDD[_]]
    def visit(r: RDD[_]) {
      if (!visited(r)) {
        visited += r

        for (dep <- r.dependencies) {
          dep match {
            //< 若为宽依赖,调用getShuffleMapStage
            case shufDep: ShuffleDependency[_, _, _] => parents += getShuffleMapStage(shufDep, jobId)
            case _ =>
              //< 若为窄依赖,将该依赖中的rdd加入到待vist列表,以保证能一级一级往上vist,直至遍历整个DAG图
              waitingForVisit.push(dep.rdd)
          }
        }
      }
    }
    waitingForVisit.push(rdd)
    while (waitingForVisit.nonEmpty) {
      visit(waitingForVisit.pop())
    }
    parents.toList
  }

case shufDep: ShuffleDependency[_, _, _] => parents += getShuffleMapStage(shufDep, jobId)若为宽依赖,则调用getShuffleMapStage并将返回的stage加入到ResultStage依赖的父stage中,那么getShuffleMapStage又是干嘛的?

  private def getShuffleMapStage(
      shuffleDep: ShuffleDependency[_, _, _],
      jobId: Int): ShuffleMapStage = {
    shuffleToMapStage.get(shuffleDep.shuffleId) match {
      case Some(stage) => stage
      case None =>
        // We are going to register ancestor shuffle dependencies
        registerShuffleDependencies(shuffleDep, jobId)

        //< 然后创建新的ShuffleMapStage
        val stage = newOrUsedShuffleStage(shuffleDep, jobId)
        shuffleToMapStage(shuffleDep.shuffleId) = stage

        stage
    }
  }

该函数首先判断成员shuffleToMapStage中是否包含了参数shuffleDep.shuffleId为key的stage,在划分stage的过程中,在以ResultStage的finalRDD为起点向前遍历RDD依赖的过程中,每次碰到的宽依赖都是不同的,它们的id自然也是不同,所以调用到这个函数的时候,都会进入到 case None的分支中

该分支中首先调用了registerShuffleDependencies(shuffleDep, jobId),该函数比较关键,为了方便阅读,我将该函数的实现及其调用的函数的源码及注释都贴在一起,如下:

  private def registerShuffleDependencies(shuffleDep: ShuffleDependency[_, _, _], jobId: Int) {
    


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

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

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

相关文章

  • 超级权限容器,权限容器
  • Hive 合并输入输出文件,
  • HIVE动态分区实战,hive动态实战
  • Spark性能优化:shuffle调优
  • myeclipse2013 建的web项目没有web.xml,myeclipse中web.xml
  • RabbitMQ(python实现)学习之一:简单两点传输“Hello World”的实现,rabbitmqpython
  • Apache Flink简介,apacheflink简介
  • spark core源码分析9 从简单例子看action操作,sparkcore
  • 【HBase】how many zookeepers should i run?,hbasezookeepers
  • HBase常用操作之namespace,hbasenamespace

文章分类

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

最近更新的内容

    • MapReduce排序及实例,MapReduce排序实例
    • 玩转OpenStack网络Neutron(1)--热身,openstackneutron
    • flume 组件概述与列表
    • ElasticSearch的分布式安装
    • Hbase详解—–管理 Splitting,hbase详解splitting
    • 【Spark1.3官方翻译】Spark集群模式概览,spark1.3spark
    • 数学定理证明机械化的中国学派(I),机械化学派
    • Spark Streaming和Flume集成指南V1.4.1,flumev1.4.1
    • Storm计算结果是如何存放的,Storm计算结果存放
    • NOSQL(一)为什么选用NoSQL?,选用nosql

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

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