2-1 Job Stage Task的划分与执行

  • Worker Node:物理节点,上面执行executor进程
  • Executor:Worker Node为某应用启动的一个进程,执行多个tasks
  • Jobs:action 的触发会生成一个job, Job会提交给DAGScheduler,分解成Stage,
  • Stage:DAGScheduler 根据shuffle将job划分为不同的stage,同一个stage中包含多个task,这些tasks有相同的 shuffle dependencies。

有两类shuffle map stage和result stage:

shuffle map stage:case its tasks' results are input for other stage(s)

result stage:case its tasks directly compute a Spark action (e.g. count(), save(), etc) by running a function on an RDD,输入与结果间划分stage

  • Task:被送到executor上的工作单元,task简单的说就是在一个数据partition上的单个数据处理流程。

小结:

action触发一个job (task对应在一个partition上的数据处理流程)

------stage1(多个tasks 有相同的shuffle依赖)------【map--shuffle】------- stage2---- 【result--shuffle】-----

样例说明

已下面的一段样例分析下具体的执行过程,为了分析,将map和reduce没有写成链式
data.txt数据如下,分别表示event uuid pv
1 10010
1 10011
1 10020
1 10031
2 10021
2 10031
2 10030
3 10010
3 10010
计算每个事件的uv和pv,即实现
event维度下,count(distinct if(pv > 0,uuid,null)) ,sum(pv)

package com.dt.spark.main.DAG


import org.apache.spark.{SparkConf, SparkContext}

/**
  * Created by hjw on 17/9/11.
  *
  **/
/*
event uuid  pv data.txt的数据
1        1001    0
1        1001    1
1        1002    0
1        1003    1
2        1002    1
2        1003    1
2        1003    0
3        1001    0
3        1001    0
计算 event,count(distinct if(pv > 0,uuid,null)) ,sum(pv)
2     UV=2     PV=2
3     UV=0     PV=0
1     UV=2     PV=2
为了分析,将map和reduce没有写成链式
 */
object DAG {
  def main(args: Array[String]) : Unit = {
    val conf = new SparkConf()

    conf.setAppName("test")
    conf.setMaster("local")

    val sc = new SparkContext(conf)

    val txtFile = sc.textFile("./src/main/java/com/dt/spark/main/DAG/srcFile/data.txt")
    val inputRDD = txtFile.map(x => (x.split("\t")(0), x.split("\t")(1), x.split("\t")(2).toInt))
    //val partitionsSzie = inputRDD.partitions.size

    //这里为了分析task数先重分区,分区前partitions.size = 1,下面每个stage的task数为1
    val inputPartionRDD = inputRDD.repartition(2)

    //------map_shuffle stage 有shuffle Read
    //结果:(事件-用户,pv)
    val eventUser2PV = inputPartionRDD.map(x => (x._1 + "-" + x._2, x._3))

    //结果: (事件,(用户,pv))
    val PvRDDTemp1 = eventUser2PV.reduceByKey(_ + _).map(x =>
      (x._1.split("-")(0), (x._1.split("-")(1), x._2))
    )

    //-------map_shuffle stage   有shuffle Read 和 有shuffle Write
    //结果: (事件, Tuple2(Tuple2(用户,是否出现),该用户的pv) )
    val PvUvRDDTemp2 = PvRDDTemp1.map(
      x => x match {
        case x if x._2._2 > 0 => (x._1, (1, x._2._2))
        case x if x._2._2 == 0 => (x._1, (0, x._2._2))
      }
    )

    //结果:(事件,Tuple2(uv,pv))
    val PVUVRDD = PvUvRDDTemp2.reduceByKey(
      (a, b) => (a._1 + b._1, a._2 + b._2)
    )

    //------result_shuffle stage 有shuffle Read
    //--------触发一个job
    val res = PVUVRDD.collect();


    //------result_shuffle stage 有shuffle Read
    //--------触发一个job
    PVUVRDD.foreach(a => println(a._1 + "\t UV=" + a._2._1 + "\t PV=" + a._2._2))
    //    2     UV=2     PV=2
    //    3     UV=0     PV=0
    //    1     UV=2     PV=2

    while (true) {
    ;
    }
    sc.stop()
  }
}

ine 69 collect() 和 line 74 行 foreach()触发两个job

  //------result_shuffle stage 有shuffle Read
    //--------触发一个job
    val res = PVUVRDD.collect();


    //------result_shuffle stage 有shuffle Read
    //--------触发一个job
    PVUVRDD.foreach(a => println(a._1 + "\t UV=" + a._2._1 + "\t PV=" + a._2._2))

基于下图知

job0有4个stage,共7个task

  • stage0:1个分区有1个task执行
  • stage1--3:各有2个分区,各有2个task执行,共6个task

job1有1个stage(复用job0的RDD),有2个task具体job0-stage0

1 1001 0
1 1001 1
1 1002 0
1 1003 1
2 1002 1
2 1003 1
2 1003 0
3 1001 0
3 1001 0

输入9条记录,只要一个分区,有一个task执行,想repartion进行shuffle写数据

results matching ""

    No results matching ""