精选·Spark

Spark 的宽依赖和窄依赖怎么区分?Stage 是怎么划分的?

91学AI·2026/7/13·10 阅读

考察点

这题是 Spark 调度体系的入口,答不好后面 shuffle、AQE 都聊不下去。面试官想确认三点:你能不能不看代码只凭算子语义判断依赖类型(比如 coalesce(true/false)、join 什么情况下是窄依赖);stage 划分的规则和「为什么以宽依赖为界」;以及这个划分对流水线执行、容错重算的影响。追问常往 shuffle 边界、stage 内 task 的流水线化、一个 action 产生几个 job 这些方向走。

参考答案

宽窄依赖的本质区别

判断标准只有一条:子 RDD 的每个 partition 依赖父 RDD 的多少 partition

  • 窄依赖:子 partition 依赖父 RDD 固定的一到几个 partition。包括 OneToOneDependency(map、filter、flatMap、mapPartitions)和 RangeDependency(union),以及 N:1 的 coalesce(false) 这类 co-partition 合并。特点是不需要跨节点搬数据,父 partition 算完本地就能出子 partition。
  • 宽依赖(ShuffleDependency):子 partition 依赖父 RDD 的全部 partition。groupByKey、reduceByKey、sortByKey、repartition、coalesce(true) 都是。因为聚合/重排需要看到全量数据,必然发生 shuffle——数据要按新的分区规则跨网络搬运。

有几个容易答错的边界情况要心里有数:

  • join 不固定。两个 RDD 如果已经被同一个 partitioner 分区好(比如上游都做过同样参数的 partitionBy),join 就是窄依赖(OneToOneDependency);否则是宽依赖,而且通常是双方都 shuffle。
  • mapPartitions、filter 是窄的,虽然它们一次看一个 partition 的全部数据,但依赖关系仍然是一对一。
  • coalesce 的两个重载语义相反:默认 shuffle=false 是窄依赖(父 partition 合并到子 partition);coalesce(n, shuffle=true) 或 repartition 是宽依赖。

Stage 的划分规则

一个 action 触发 DAGScheduler 做 stage 划分,算法是一个从结果 RDD 出发的反向遍历:

  1. 以最终 RDD 为根,建一个 ResultStage;
  2. 沿血缘反向走,遇到宽依赖就切一刀,为父 RDD 新建 ShuffleMapStage;
  3. 窄依赖一路向上合并进同一个 stage,直到血缘起点或遇到已存在的 stage。

划完得到一个 stage 的有向无环图。每个 ShuffleMapStage 内部是一条纯窄依赖的流水线。直觉理解:shuffle 是天然的物化点——数据要落盘(map 端写 shuffle file)、要排序、要按分区归并,到了这里流水线必须断开,前一段的结果必须全部就绪才能开始后一段。所以「以宽依赖为 stage 边界」不是人为规定,是 shuffle 的物化语义决定的。

为什么这个划分重要

流水线执行:同一 stage 内的窄依赖算子会被串成一条计算链,task 拿到一条记录可以一口气 map → filter → map 走完,中间不落地、不物化。这是 Spark 比 MapReduce 快的根本原因之一——MR 每个算子之间都要落盘,Spark 只在 shuffle 边界落地。Stage 数 = shuffle 次数 + 1(近似),所以调优的第一原则就是减少不必要的 shuffle。

调度与并行度:DAGScheduler 按 stage 依赖拓扑提交,父 stage 全部完成才提交子 stage。每个 stage 内按 partition 数切 task,task 并行跑、互不通信——这就是 Spark 任务模型简单的原因,也是它能容忍 task 级失败的基础(失败 task 在父 stage 可用的情况下直接重跑)。

失败恢复的粒度:stage 边界是恢复的记账点。某个 task 失败,先重跑本 task;如果是 shuffle fetch 失败(map 端数据丢了),调度器会把对应的父 ShuffleMapStage 重新提交,只重算丢失的 map 输出。窄依赖场景重算范围天然小,宽依赖的恢复才需要回溯到上一个物化点。

一个具体的例子

val r1 = sc.textFile("hdfs://logs/2026/")            // RDD0
val r2 = r1.filter(_.contains("ERROR"))              // 窄
             .map(l => (l.split(",")(0), 1))          // 窄
val r3 = r2.reduceByKey(_ + _)                       // 宽,shuffle 1
val r4 = r3.mapValues(_ * 2)                         // 窄
val r5 = r4.sortByKey()                              // 宽,shuffle 2
r5.collect()                                         // action

划分结果:Stage 0(RDD0→r2,到 shuffle 1 的 map 端)、Stage 1(r3→r4,到 shuffle 2 的 map 端)、Stage 2(r5 的 ResultStage)。两次 shuffle 切出三个 stage,Stage 0、1 无依赖关系?不——有依赖,Stage 1 依赖 Stage 0 的 shuffle 输出,串行提交。

可能的追问

  • shuffle 时 map 端和 reduce 端各做什么? map 端按分区器把记录写进本地 shuffle 文件(SortShuffleManager 下按 partition 排序/分桶),生成 index + data 文件;reduce 端 task 从所有 map 端拉属于自己的那段数据,做归并聚合。中间靠 MapOutputTracker 汇报位置。
  • 一个 application 里 job、stage、task 的数量关系? 一个 action 一个 job;job 内按 shuffle 切 stage;stage 内按 partition 切 task。注意 cache 命中的 RDD 会让 DAGScheduler 跳过它上游的 stage(跳过或缩短血缘)。
  • 为什么 reduceByKey 和 sortByKey 挨着写会产生两次 shuffle,能优化吗? 两次宽依赖切两个 stage,确实两次 shuffle。调优思路是合并语义:如果排序后还要聚合,考虑直接在一个 shuffle 里完成(自定义分区器 + mapPartitions 内排序聚合),或者用 DataFrame API 让 Catalyst 帮忙合并。

评论 (0)

暂无评论,快来抢沙发吧!

91学AI

© 2026 91学AI · 按岗位学 AI 与大数据. All rights reserved.