spark dag切分

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快的原因

最后编辑于 :
©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

友情链接更多精彩内容