做Spark调优的人,最后多半会追到同一个问题上:并行度。任务跑得慢、Executor CPU利用率上不去、某个Stage卡半天,翻来覆去大概率都是并行度没对上。但不少朋友对并行度的理解停留在“加分区数”这个结论上,真到看源码、对着DAG图排查时,才发现宽依赖和窄依赖对并行度的作用完全是两套逻辑。这篇算Spark RDD任务并行度系列解读的Part2,从Dependency类源码、Stage切分逻辑、Shuffle分区器三个角度,把“并行度为什么长成这样”讲透,适合已经能写Spark作业、正在调优或者准备面试的朋友。Part1聊过RDD分区数与Task调度的映射,这一篇重点落在宽窄依赖上。
1. 并行度的底层来源:RDD分区数如何换算成Task数
1.1 一个分区对应一个Task:getPartitions与Task生成逻辑
要理解宽窄依赖对并行度的影响,先得把并行度到底从哪里来的搞清楚。Spark里没有任何一个全局变量叫“并行度”,真正决定Task数量的,是RDD的分区数。每个RDD内部维护一个分区数组,源码里就是def getPartitions: Array[Partition],由各个子类实现。比如ParallelCollectionRDD会把你传入的集合按切片个数切成若干分区,HadoopRDD会按文件split生成分区,ShuffledRDD则根据分区器生成对应数量的分区。
当一次action触发job时,DAGScheduler拿到最后一个RDD,向前回溯依赖关系生成Stage,每个Stage再根据它对应的RDD的分区数生成Task。这里有个关键点:一个分区最终对应一个Task,Task被调度到Executor上执行时,该分区的compute方法就是Task的业务逻辑。所以一个Stage的并行度上限,本质上等于这个Stage最后一个RDD的partitions.length。
举个实际例子。你用sc.parallelize(1 to 100, 8),这个RDD有8个分区,那么后面接map、filter,只要不触发shuffle,最终Stage的任务数就是8。即使Executor有32个核,这个Stage也只能并行跑8个Task,剩下24个核算力闲置。很多任务跑得慢,不是集群不够大,是并行度根本没给到位。
1.2 为什么Stage内可并行、Stage间必须串行
分区数决定Task数量,但Task之间能不能真正并行,还要看依赖关系。Stage内部的所有Task之间没有数据依赖,拿到哪块分区就算哪块,天然适合并行调度。Stage之间则不一样,下游Stage的Task可能需要上游Stage的shuffle输出,这种跨Stage的数据传递是硬性的先后关系。
宽窄依赖在这里起了决定性作用:窄依赖意味着子RDD的每个分区只依赖父RDD的对应分区,数据不需要跨节点重组,所以多个算子可以串联在同一个Stage里流水线执行。宽依赖则意味着子RDD的分区要依赖父RDD的多个分区,必须等到上游全部算完、数据经过shuffle落盘后才能开始,因此Stage边界恰好就切在宽依赖上。
理解这点后你再回头看DAG图就清晰了:图中每两个Stage之间必然有一次shuffle,也就是一个宽依赖。你看到的Stage小方块里的数字,就是这个Stage的Task数量——它由这个Stage内最后一个RDD的分区数决定,和上游Stage的Task数量没有直接关系。
2. 宽窄依赖源码解读:从Dependency类看懂分区血缘
2.1 窄依赖的源码形态:OneToOneDependency与RangeDependency
先看源码。依赖这个抽象概念在Spark里对应Dependency[T]基类,窄依赖是它的子类NarrowDependency:
abstract class NarrowDependency[T](rdd: RDD[T]) extends Dependency[T] { def getParents(partitionId: Int): Seq[Int] }这个抽象类只有一个方法getParents(partitionId),作用是问:子RDD的某个分区,是从父RDD的哪些分区算出来的。注意它返回的是Seq[Int],是一组父分区ID列表。
窄依赖有两个实现。OneToOneDependency最常见,map、filter、flatMap、mapPartitions这类算子产生的就是它:
class OneToOneDependency[T](rdd: RDD[T]) extends NarrowDependency[T](rdd) { override def getParents(partitionId: Int): Seq[Int] = Seq(partitionId) }getParents直接返回Seq(partitionId),意思是子RDD的第N个分区只依赖父RDD的第N个分区,一一对应。这正是窄依赖的典型血缘关系,也决定了这类算子不会改变并行度:子RDD分区数等于父RDD分区数。
另一个实现是RangeDependency,用于union操作。UnionRDD把多个RDD拼在一起,子分区会跨多个父RDD:
class RangeDependency[T](rdd: RDD[T], inStart: Int, outStart: Int, length: Int) extends NarrowDependency[T](rdd) { override def getParents(partitionId: Int): Seq[Int] = { if (partitionId >= outStart && partitionId < outStart + length) { Seq(partitionId - outStart + inStart) } else { Nil } } }虽然union之后的子RDD分区数等于所有父RDD分区数之和,但每个父RDD的每个分区最多只被子RDD的一个分区使用——前面的分区不会跑到后面的子分区里去。所以它仍然是窄依赖,不需要shuffle。从这里能看出一个判断窄依赖的关键标准:父RDD的单个分区,被子RDD至多一个分区消费。
2.2 ShuffleDependency源码:宽依赖为什么“宽”
宽依赖在源码里就是ShuffleDependency,它是Dependency的直接实现,和窄依赖完全不同的抽象路径:
class ShuffleDependency[K, V, C]( rdd: RDD[_ <: Product2[K, V]], partitioner: Partitioner, serializer: Serializer, keyOrdering: Ordering[K], aggregator: Aggregator[K, V, C], mapSideCombine: Boolean) extends Dependency[Product2[K, V]] { val shuffleId: Int = ... }构造参数里的几个字段特别值得关注。rdd是父RDD,它的元素必须是Product2[K, V],说白了就是key-value形式。partitioner是分区器,它决定了父RDD里的每个元素经过shuffle后落到下游哪个分区。aggregator是可选的聚合器,如果定义了,map端可以先做一轮combine减少传输量。
为什么它是“宽”的?因为窄依赖里父分区与子分区是一对一或一对定范围的关系,而shuffle依赖下,父RDD的同一个分区里的数据,会按key的哈希值被分发到下游多个不同的分区去。一个父分区同时被多个子分区消费,这种一对多的血缘关系就是宽依赖。每次出现ShuffleDependency,DAG就会在这里切开Stage,上游Stage的任务跑完后把结果按分区器写到磁盘,下游Stage才能拉取数据继续跑。
2.3 DAGScheduler切分Stage的源码路径
Stage切分是理解并行度的最后一环。DAGScheduler在处理job时,从final RDD开始向前遍历依赖树。遍历到ShuffleDependency就停下来,把shuffle之前的所有窄依赖归为一个Stage,继续往前切。源码核心在getShuffleDependencies和createShuffleMapStage、createResultStage这几个方法里,逻辑可以概括为一条规则:遇到宽依赖就断,窄依赖全部留在一个Stage内。
这条规则对并行度有三层实际含义。第一,同一个Stage内的多个算子,比如map之后接filter再接mapPartitions,分区数始终传递,并行度不变。第二,Stage内的计算以流水线方式执行,每个Task一口气处理完这个分区上的所有算子,中间不落盘,因此窄依赖区间的性能损耗主要就是单分区CPU计算量。第三,Stage与Stage之间的并行度是独立计算的,上游Stage可能只有4个Task,下游Stage可能有200个Task,因为下游分区数完全由shuffle那一步的分区器决定,和上游分区数没有必然关系。
3. 窄依赖的并行度:父RDD分区数如何直接约束下游Task
3.1 常见窄依赖算子的分区数映射:map、filter、union、coalesce
窄依赖算子是流水线式传播分区结构的,不同算子对分区数的影响并不完全相同。我整理了一张常用对照表:
| 算子 | 依赖类型 | 子RDD分区数变化 | 对并行度的影响 |
|---|---|---|---|
| map / flatMap / filter | OneToOneDependency | 不变,等于父RDD分区数 | 并行度不变 |
| mapPartitions / mapPartitionsWithIndex | OneToOneDependency | 不变 | 并行度不变 |
| union | RangeDependency | 所有父RDD分区数之和 | 并行度变大,但数据量也随之变大 |
| zipPartitions | OneToOneDependency或其他 | 与父分区数较小者一致 | 并行度由各个父RDD共同约束 |
| coalesce(n, shuffle=false) | 窄依赖 | 减少到n | Task数量减少,但父分区只归并到某一个子分区 |
这里最容易看走眼的是coalesce。很多人以为coalesce就是单纯变一下分区数,是个轻量操作,但实际上它有两种形态。coalesce(n, shuffle=false)不触发shuffle,属于窄依赖,父RDD的每个分区只会被归入一个子分区,数据不会跨Executor重分布;coalesce(n, shuffle=true)则等价于repartition(n),是宽依赖。
窄依赖对并行度的影响很直接:子RDD的分区数被父RDD锁死,你想通过map阶段加分区是加不上去的。如果你的并行度瓶颈在map、filter这类算子阶段,唯一的办法是在这之前用repartition或coalesce调整分区结构,或者直接在读取源数据时就多切几个分区。
3.2 coalesce与repartition:看似相似,并行度语义完全不同
repartition的源码实现其实就是在调用coalesce,但它强制要求shuffle:
def repartition(numPartitions: Int): RDD[T] = withScope { coalesce(numPartitions, shuffle = true) }关键差异在shuffle这个布尔值上。coalesce(4, shuffle=false)是窄依赖,它只是简单地把多个父分区合并成4个子分区,数据分布可能不均衡,比如连续的一段分区会被合并到一个子分区里,某些Executor就会拿到偏多的数据。repartition(4)走的是ShuffleDependency,所有数据经过哈希重分布,每个子分区的数据量相对均匀,但要付出shuffle的落盘和传输代价。
实际场景里怎么选?如果数据经过filter后明显变少,想缩减分区数,用coalesce就够了,避免一次无谓的shuffle。如果下游并行度不足,或者数据凌乱需要均匀打散,用repartition。我在生产中见过不少人不管三七二十一都用repartition,把本来可以窄依赖处理的阶段硬生生切出一个Stage,白白增加一轮shuffle,这种开销在小文件场景下尤其明显。
3.3 窄依赖阶段容易忽略的性能点
窄依赖区间虽然没有shuffle,但并行度不合理同样会放大问题。最典型的一个坑是:在一个被cache的RDD上反复做窄依赖操作。比如某baseRDD分区数只有4,你把它cache住,后面十几个job都在它上面跑map、filter,那么这十几个job的Stage任务数永远只有4。即使集群有几百个核,每个job都在白白浪费算力。这种情况下,正确的做法是在cache之前先repartition(合适的并行度),让后续所有下游作业都受益。
另一个点是JVM层面的串扰。窄依赖的Task在同一个Executor上跑时,如果单分区数据量过大,Task内部的内存压力会传导到GC。并行度太小会导致每个Task处理的数据块过大,看似CPU在忙,实际大量时间耗在GC上。判断方法很直接:在Spark UI里看某个Stage的GC Time是不是占大头。遇到这种情况,把分区数上调,让每个Task的数据量降下来,GC自然会缓解。
4. 宽依赖的并行度:Shuffle下游分区数与Partitioner的决策逻辑
4.1 下游分区数由Partitioner决定:HashPartitioner源码解读
宽依赖最核心的代码点是Partitioner。ShuffleDependency构造时就收了一个partitioner参数,下游RDD的分区数完全由它决定。看HashPartitioner的实现:
class HashPartitioner(partitions: Int) extends Partitioner { require(partitions >= 0, s"Number of partitions ($partitions) cannot be negative.") def numPartitions: Int = partitions def getPartition(key: Any): Int = key match { case null => 0 case _ => Utils.nonNegativeMod(key.hashCode(), numPartitions) } }numPartitions就是shuffle之后的分区数。每个key先取hashCode,再做非负取模,落到某个下游分区。整个流程里,上游RDD有多少分区、多少Task,对下游分区数都没有影响,下游分区数是你在创建ShuffleDependency时给定的。
RangePartitioner则是另一种策略,它会对key做采样后排序,生成一组有序的边界,保证相邻key尽量落在同一个分区,适合排序类的需求。源码里的核心是sketch和determineBounds两个方法,前者对分区的key抽样,后者算出边界数组。用RangePartitioner时,下游分区数同样是构造参数指定的。
4.2 Shuffle两侧的并行度为什么要分开看
一个宽依赖切出的两个Stage,并行度需要分开评估。上游Stage的Task数量等于父RDD分区数,每个Task负责把本分区的数据按分区器写成本地shuffle文件;下游Stage的Task数量等于partitioner.numPartitions,每个Task从所有上游节点的文件中拉取属于自己那部分的数据。
举个例子更容易懂。父RDD有10个分区,reduceByKey时指定分区器为100,那上游Stage有10个Task做shuffle write,下游Stage有100个Task做shuffle read。数据量不变,但并行度从10跳到了100。反过来,如果父RDD有1000个分区,你指定repartition(10),那上游Stage有1000个Task,每个Task要写10份数据文件,下游Stage只有10个Task,每个Task要拉1000份数据文件。
这就解释了为什么宽依赖两侧的并行度要分开调。有些时候问题不在上游,而在下游分区数太少导致每个Task拉取的数据量太大;有些时候刚好相反,上游Task太多而每个分区的数据量太小,shuffle write的调度开销和文件数量反而拖慢速度。你盯着Spark UI看,要分别看两个Stage的Task数量是否各自合理,而不是想当然地认为并行度是“一条线传下来”的。
4.3 上下游并行度不匹配会引发什么问题
并行度不匹配带来的问题,在生产环境里通常以三种面貌出现。
第一种是下游分区数远小于上游,导致单个下游Task要拉取海量数据。这种情况最容易表现为OOM或长GC。我在某离线作业里见过上游几千个分区数据汇总后,下游只有10个分区,每个Task拉几百GB数据,跑了一个多小时都没结束,最后靠把shuffle分区数调到几百才解决。宽依赖这里有个经验值:shuffle后单分区数据量尽量控制在128MB到256MB之间,分区数不是越大越好,但也不能小到让单个Task吃不下。
第二种是下游分区数过多,产生了大量小Task。每个Task都有调度、序列化、反序列化的固定开销,Task数量一旦上万,光调度的CPU开销就很可观,还会生成大量小文件影响后续读取。这种情况常见于对每个key调groupByKey却给了超大分区数。
第三种是上下游之间的数据倾斜。上游某几个分区数据量特别大,对应的shuffle输出文件也特别大,下游拉取时这几个Task就成了长尾。这时候单纯调并行度解决不了根本问题,要做的是对key加盐、两阶段聚合,或者换成RangePartitioner让数据分布更均匀。并行度调整只是第一步,不是万能药。
5. 并行度参数选型与实操调优:从默认值到案例
5.1 并行度相关参数速查与默认值解读
平时调并行度打交道最多的就是几个参数,我先列表说明,免得大家一头扎进代码里找不到方向:
| 参数 | 默认值 | 位置 | 作用 |
|---|---|---|---|
| spark.default.parallelism | local模式取本机核数,集群模式取executor总核数 | SparkConf | 未显式指定分区数时,parallelize、reduceByKey等操作的默认分区数 |
| spark.sql.shuffle.partitions | 200 | SparkConf | SQL和DataFrame操作中shuffle阶段默认分区数 |
| spark.executor.cores | 1(YARN下通常手动配置) | SparkConf | 单个Executor可并行执行Task的核数 |
| spark.executor.instances | 动态分配时自动 | SparkConf | Executor个数,乘以cores可估算集群总并行能力 |
| spark.files.maxPartitionBytes | 128MB | SparkConf | 读取文件时单个分区目标大小,影响初始分区数 |
| spark.sql.files.minPartitionNum | 根据文件大小动态 | SparkConf | 读取文件时的最小分区数下限,调大可提高初始并行度 |
spark.default.parallelism是很多人容易忽略的。如果你用textFile时不指定minPartitions,或者直接parallelize不指定切片数,Spark会退回这个默认值。集群模式下它默认是executor总核数,如果你的executor实例数不多,这个值可能相当小,初始读取阶段的分区数就会特别少。所以我一般建议在提交任务时显式配置spark.default.parallelism,避免依赖环境自动推断。
5.2 并行度估算思路:读入、流水线、shuffle三段式
我的调优习惯是分三段算并行度,每段独立评估。
读取阶段,看源数据总量和目标分区大小。maxPartitionBytes是128MB,意思是每个分区尽量不超过128MB数据。如果输入是1TB,理想情况下大约分成8000个分区左右,但还要考虑文件个数,每个文件至少拆出一个分区,所以小文件多的时候分区数会被拉高。
窄依赖流水线阶段,沿用上游分区数,不需要单独调整。你在这个阶段加了repartition,等于强行插入一次shuffle,要确认真的有收益才值得做。比如后续有cache复用或者要均匀打散数据,可以做;纯粹为了让“并行度好看”,没必要。
shuffle阶段,看两个点:上游输出数据总量和期望的单分区大小。比如上游Stage输出的shuffle数据一共200GB,你希望每个下游Task处理256MB,那分区数可以取800附近。如果你同时考虑集群总核数,可以再取一个“核数乘以2到3”的参考值,两者取最大。核心原则是让每个Task的工作量均匀且可控,而不是固定套某个数。
5.3 一次离线任务的并行度调优实录
分享一个我实际处理过的模拟场景。某离线数据清洗作业,每天处理几十GB的日志文件,任务经常跑到半小时以上。打开Spark UI后,第一个Stage的Task数量只有2,可是集群给了20个Executor、每Executor 4核,总并行能力是80。显然,读取阶段就没把并行度提起来。查代码发现读文件时没指定minPartitions,spark.default.parallelism又没配,默认值远低于集群能力。
处理方式分三步走。第一步,读取时显式指定minPartitions,并把spark.files.maxPartitionBytes调小一点,让初始分区数从2提高到接近80。改完之后第一个Stage从2个Task变成了70多个Task,任务总时间立刻缩短了一大截。第二步,后面接了一个大表join,发现shuffle后的Stage还是200个分区,但单分区数据量超过1GB,引起几个Task频繁GC。把spark.sql.shuffle.partitions调到600后,每个Task的数据量降到了合理区间,GC时间明显回落。第三步,最后输出时发现生成了600个小文件,对下游查询不友好,于是在写之前用coalesce把分区数压到和目标文件数一致,这里用的窄依赖,省掉了最后一次shuffle的代价。
整个改动加起来不到十行配置和几个算子调整,任务从半小时跑到七八分钟。核心思路就是一句话:每个阶段单独看分区数是否合理,不要图省事“一把梭”式地全局repartition。
6. 并行度排查实录:常见问题、误区和Spark UI观察点
6.1 Spark UI上怎么看并行度是否合理
先说说我在排查并行度问题时固定会看的几个位置。最关键是Stage页面里的Total Tasks,它就是该Stage实际生成的Task数量,等于这个Stage的并行度上限。旁边还有个Completed Tasks,跑完就一致了。如果Total Tasks远小于集群总核数,比如集群80核,Task只有10个,那就说明这个Stage的并行度不够,资源在空转。
再往下看Shuffle Write和Shuffle Read两个指标。Shuffle Write的总量可以帮你估算上游输出的数据规模,Shuffle Read的总量可以验证下游拉取的数据量。如果某个Task的Shuffle Read特别大,而其他Task很小,这不是并行度问题,是数据倾斜,需要单独处理。还要顺带看一眼GC Time,如果某个Stage的GC时间占总运行时间的比例很高,并行度大概率偏小,每个Task处理的数据块太大了。
Executor页面也值得留意。看CPU利用率柱状图,如果长期只有百分之二三十,而代码里又找不到其他瓶颈,多半就是分区数设低了。反过来,如果CPU长期接近100%但任务还是慢,那瓶颈可能不在并行度,而在单分区的计算复杂度,这时候盲加分区没有意义。
6.2 并行度不足和数据倾斜的区分技巧
并行度不足和数据倾斜都会表现为任务慢,但判断方法完全不同,而且经常被混为一谈。并行度不足的特征是:Task总数少,且每个Task耗时都比较长,但耗时分布相对均匀,所有Task一起慢。数据倾斜的特征则是:Task总数并不少,大多数Task很快跑完,只有一两个Task消耗了绝大部分时间,Shuffle Read量也不均衡。
我曾经遇到过一起“伪数据倾斜”的案例。任务里有一个Stage只有8个Task,打开UI一看,8个Task的耗时都特别长,其中有1个更突出。乍看像倾斜,但细看Shuffle Read,8个Task的数据量其实差不多,只是整体并行度太低,每个Task的任务量都过大,那个突出的Task只是数据量稍微多了一点。把分区数调大后,8个Task变成了100多个,突出点自然消失了,整体耗时也掉了下来。所以判断倾斜前,先确认并行度本身是否正常,否则很容易把并行度问题误诊成倾斜,白白折腾加盐和两阶段聚合。
6.3 我踩过的几个并行度误区
第一个误区:分区数越大越好。刚开始调优的时候我也犯过这个错,把repartition的分区数调到几千,以为并行度越高越快。结果Task数量爆炸,每个Task只处理几MB数据,调度开销和序列化开销占了主导,反而比原来更慢。并行度要匹配集群算力和数据规模,不是单纯追求数字。
第二个误区:map阶段加repartition就能解决shuffle之后的问题。repartition只能改变当前Stage及其下游窄依赖阶段的并行度,shuffle后的并行度由下一步的ShuffleDependency分区器决定。你在map阶段repartition成100个分区,后面reduceByKey如果不指定分区数,Spark仍会按spark.default.parallelism生成新的shuffle分区,前面调整的效果会被覆盖掉。
第三个误区:只调spark.sql.shuffle.partitions,不关注RDD API里的分区算子。很多SQL任务调完shuffle.partitions就好使了,但RDD代码里用到groupByKey、reduceByKey时,如果不显式传分区数,它仍然走默认并行度,和SQL参数完全是两套机制。所以排查并行度问题时,先分清楚你的计算是RDD还是Dataset API,再找准对应的参数入口。
最后分享一个我一直在用的检查习惯:改完并行度相关配置后,不要只看总耗时,先回Stage页面确认每个Stage的Task数量有没有按预期变化,再看Shuffle Read的总量是否合理。一次调优如果Task数量确实上去了但耗时没降,说明瓶颈不在并行度;如果Task数量都没变,那说明你改的参数根本没作用到这条计算链路上。并行度这个东西,调没调对,UI上十几秒就能看出来,别凭感觉下结论。