1、RDD介绍
弹性分布式数据集。
父RDD: 生成当前RDD所以来的上游RDD,父RDD是相对的
子RDD: 经过算子计算后,有上游RDD生成的数据,子RDD是相对的
2、RDD依赖关系
依赖关系直接决定了shuffle,stage划分,是spark dag调度的基石。宽依赖要分成依赖的stage,需等待上游执行完毕
- 窄依赖:子RDD的一个分区,只依赖父RDD的1个分区数据。 简单说shuffle前阶段。子RDD值读取父RDD的一个分区数据;
- 宽依赖:子RDD的一个分区,依赖父RDD的多个分区数据。 简单说shuffle后阶段
比如a left join b on a.id=b.id
a表有4个文件块,会有4个map,每个map生成一份shuffle数据,该阶段RDD_A_shuf1,窄依赖
从不同的map读取属于自己reduce分区shuffle数据进行join运行, 依赖父RDD的多个分区RDD_A_shuf1,宽依赖,该阶段叫RDD_AB
RDD_A(读a表)
↓ 窄依赖(map转kv)
RDD_A_kv
RDD_B(读b表)
↓ 窄依赖(map转kv)
RDD_B_kv
RDD_A_kv RDD_B_kv
↘ ↙
ShuffleDependency 【唯一1个宽依赖,stage分界线】
↓
ShuffledRDD # 唯一的ShuffledRDD,内部同时shuffle read A、B两边,输出 (key,(iterA,iterB))
↓ 窄依赖(flatMapValues)
RDD_AB(join结果RDD:MapPartitionsRDD)
对应stage说明
【Stage1】
RDD_A(读a表)
↓ 窄依赖(map转kv)
RDD_A_kv
RDD_B(读b表)
↓ 窄依赖(map转kv)
RDD_B_kv
【Stage2:ReduceTask】
RDD_A_kv RDD_B_kv
↘ ↙
ShuffleDependency 【唯一1个宽依赖,stage分界线】
↓
ShuffledRDD # 唯一的ShuffledRDD,内部同时shuffle read A、B两边,输出 (key,(iterA,iterB))
↓ 窄依赖(flatMapValues)
RDD_AB(join结果RDD:MapPartitionsRDD)
RDD_A 读取a表,按id hash, 写本地文件
RDD_B 读取b表,按id hash, 写本地文件
RDD_A_kv,RDD_B_kv 是shuffle后生成的数据
ShuffledRDD 是reduce阶段拉取数据
RDD_AB是join算子合并数据
3. DAGScheduler说明
以select * from a left join b on a.id=b.id left join c on a.task_id = c.id为例
rdd过程说明
RDD_A(读a表)
↓ 窄依赖map:(id, a完整行)
RDD_A_kv
RDD_B(读b表)
↓ 窄依赖map:(id, b完整行)
RDD_B_kv
RDD_A_kv RDD_B_kv
↘ ↙
ShuffleDependency #【宽依赖,Stage切割点1】第一次shuffle,join key = id
↓
ShuffledRDD_AB # 内部:shuffle read + cogroup → (id, (iter[A], iter[B]))
↓ 窄依赖flatMapValues
RDD_AB # (id, (a_row, Option[b_row])) A left join B结果
↓ 窄依赖map重key:以a_row.task_id作为新key,value携带(a_row, Option[b_row])
RDD_AB_rekey
↘
RDD_C(读c表)
↘ ↓ 窄依赖map:(task_id, c完整行)
↘ RDD_C_kv
↘ ↙
↙
↘ ↙
ShuffleDependency #【宽依赖,Stage切割点2】第二次shuffle,join key = task_id
↓
ShuffledRDD_ABC # 内部:shuffle read + cogroup → (task_id, (iter[AB], iter[C]))
↓ 窄依赖flatMapValues
RDD_ABC # (task_id, ((a_row,Option[b_row]), Option[c_row])) 最终join结果
备注: RDD_AB是join后结果,在到RDD_AB_rekey是map阶段将数据展开,因此是窄依赖
3.1、stage切分说明
从最底层RDD 向前回溯 DAG;遇到宽依赖ShuffleDependency就切断,生成新 Stage;窄依赖不切 Stage,同属一个 Stage。
拆分成两个stage, 最上游stage阶段最小
- stage1 RDD_A->RDD_A_kv,RDD_B-> RDD_B_kv
- stage2 RDD_A_kv,RDD_B_kv -> ShuffledRDD_AB -> RDD_AB -> RDD_AB_rekey (将join后数据重新map,因此是窄依赖)
- stage3 RDD_C->RDD_C_kv
- stage4 ShuffledRDD_ABC、RDD_ABC
3.2、执行过程
DAGScheduler 规则:
从最后一个 RDD 向前回溯 DAG 生成全部 Stage;
只有某个 Stage 的所有父 Stage 全部执行完成,该 Stage 才可以提交 Task 调度;
无依赖的多个Stage,资源充足时可以并行跑;资源不足就串行排队。
这也是为啥spark 比mr快的原因