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写数据



