Spark 中“Shuffle 有几种”要分两个角度回答:
- 按 RDD 依赖分类:宽依赖产生 Shuffle;
- 按 Shuffle 的实现方式分类:Hash Shuffle、Sort Shuffle 等。
实际开发中通常重点理解第二种。
一、按计算依赖分类
1. Narrow Dependency:窄依赖
父 RDD 的一个分区只被一个子 RDD 分区使用,不需要跨节点重新分发数据。
常见算子:
map filter flatMap mapPartitions例如:
rdd.map(x=>x*2)数据可以在当前分区内直接处理,不发生 Shuffle。
2. Shuffle Dependency:宽依赖
父 RDD 的一个分区可能被多个子 RDD 分区使用,需要根据 Key 重新分区和网络传输。
常见算子:
reduceByKey groupByKey join distinct sortByKey repartition例如:
rdd.reduceByKey(_+_)相同 Key 的数据可能来自不同节点,需要把它们发送到同一个下游分区,这就是 Shuffle。
二、按实现方式分类
1. Hash Shuffle
早期 Spark 主要使用 Hash Shuffle。
Map Task 会根据目标分区和 Key 进行 Hash,然后把数据写入不同的临时文件。
例如有:
M 个 Map Task R 个 Reduce 分区可能产生接近:
M × R个文件。
特点
- 实现简单;
- 根据 Hash 分区;
- 文件数量可能非常多;
- 容易造成小文件和文件句柄压力;
- 现代 Spark 中已经不是主流实现。
2. Sort Shuffle
现代 Spark 默认主要使用 Sort Shuffle。
Map Task 先把数据写入内存缓冲区,内存不足时溢写到磁盘;最后将多个溢写文件归并,并按照分区信息组织输出。
Reduce Task 再读取对应分区的数据。
特点
- 文件数量相对可控;
- 支持内存溢写;
- 通常比旧 Hash Shuffle 更适合大规模数据;
- 是现代 Spark 的主流 Shuffle 实现。
执行过程大致是:
Map Task -> 内存缓冲 -> Spill 到磁盘 -> 多个 Spill 文件归并 -> 生成 Shuffle 输出 -> Reduce Task 拉取三、Sort Shuffle 内部的几种写法
严格来说,Sort Shuffle 还可以根据数据规模和配置走不同路径。
1. Unsafe Shuffle Writer
针对 UnsafeRow 等二进制格式优化,利用 Tungsten 内存管理和高效排序。
通常适用于 Spark SQL、DataFrame、Dataset 场景。
特点:
- 二进制处理;
- 减少对象创建;
- 内存利用率较高;
- 性能通常较好。
2. Serialized Shuffle Writer
数据序列化后在内存中排序,减少 Java 对象开销。
3. BypassMergeSortShuffleWriter
当满足特定条件时使用,例如:
- 没有 map-side combine;
- 下游分区数不超过
spark.shuffle.sort.bypassMergeThreshold,默认通常是 200。
它会分别写入各个分区文件,最后再合并这些文件。
优点:
- 不需要对记录进行排序;
- 某些场景下更快。
缺点:
- 仍然可能产生较多临时文件;
- 不支持 map-side combine。
4. SerializedShuffleWriter
用于序列化数据的 Shuffle 写入路径,通常在不需要特殊排序逻辑时使用。
实际使用哪一种,由 Spark 的执行计划、数据格式、是否需要聚合以及配置共同决定。
四、Shuffle 的两个重要阶段
一次 Shuffle 通常包含两个阶段。
Map 阶段
Map Task:
- 读取上游数据;
- 根据分区器计算目标分区;
- 可能执行 map-side combine;
- 写入内存或磁盘;
- 生成 Shuffle 文件;
- 向 Driver 汇报输出位置。
Reduce 阶段
Reduce Task:
- 从各个 Executor 拉取自己负责的分区数据;
- 进行归并、聚合或排序;
- 输出最终结果。
可以理解为:
Map 端写 Reduce 端拉五、哪些操作会触发 Shuffle?
通常会触发
rdd.groupByKey()rdd.reduceByKey(_+_)rdd.aggregateByKey(...)rdd.join(other)rdd.distinct()rdd.sortByKey()rdd.repartition(n)DataFrame / SQL 中常见的 Shuffle 来源:
GROUPBYJOINORDERBYDISTINCTDISTRIBUTEBYREPARTITION例如:
df.groupBy("user_id").count()需要把相同user_id的数据发送到同一个分区,通常会产生 Shuffle。
通常不会触发
map filter flatMap mapValues filterByRangeDataFrame 中常见的窄依赖操作:
df.select(...)df.filter(...)df.withColumn(...)但要注意:具体是否产生 Shuffle,最终应以物理执行计划为准。
六、reduceByKey和groupByKey的区别
两者都会 Shuffle,但效率通常不同。
groupByKey
rdd.groupByKey()先把同一个 Key 的所有 Value 发送到一起,再聚合。
reduceByKey
rdd.reduceByKey(_+_)可以在 Map 端先进行局部聚合,也就是 map-side combine,减少网络传输量。
例如原始数据:
(a, 1) (a, 1) (a, 1)Map 端可以先变成:
(a, 3)再发送到 Reduce 端。
所以:
需要聚合时,通常优先使用
reduceByKey、aggregateByKey或combineByKey,而不是先groupByKey。
七、DataFrame 中常见的 Shuffle 类型
Spark SQL 中经常看到以下分区方式:
1. HashPartitioning
按照 Key 的 Hash 值分区:
hash(key) % numPartitions常用于:
GROUP BY;- 等值 Join;
dropDuplicates;repartition("key")。
2. RangePartitioning
按照范围分区,常用于排序相关操作:
df.repartitionByRange("order_id")或执行全局排序时使用。
3. RoundRobinPartitioning
轮询分发数据,常见于:
df.repartition(10)不指定 Key 时,通常按照轮询方式重新分区。
八、Broadcast Join 不一定需要 Shuffle
如果一张表很小,可以广播到各个 Executor:
frompyspark.sql.functionsimportbroadcast result=large_df.join(broadcast(small_df),"user_id")这种情况下:
- 小表被广播;
- 大表通常不需要为了 Join Key 做全量 Shuffle;
- 可以避免大规模网络重分区。
但广播表需要能放进 Executor 内存,不能盲目使用。
九、如何查看是否发生 Shuffle?
DataFrame / SQL:
df.explain("formatted")重点观察是否出现:
Exchange例如:
Exchange hashpartitioning(user_id, 200)通常表示发生了 Shuffle。
RDD 可以查看依赖:
rdd.toDebugString如果出现:
ShuffledRDD说明存在 Shuffle 依赖。
Spark UI 中也可以查看:
- Shuffle Read;
- Shuffle Write;
- Fetch Wait Time;
- Records Read / Write;
- Spill Memory;
- Spill Disk。
总结
如果按最常用的实现方式回答,Spark Shuffle 主要可以说:
1. Hash Shuffle:历史实现 2. Sort Shuffle:现代主流实现如果进一步细分 Sort Shuffle 的写入器,还包括:
- Unsafe Shuffle Writer - Serialized Shuffle Writer - BypassMergeSortShuffleWriter而从 Spark 计算模型看,真正决定是否发生 Shuffle 的核心是:
窄依赖:不需要跨分区重组 宽依赖:需要跨分区重组,会产生 Shuffle在实际排查性能时,最值得关注的是:
- 是否发生了不必要的
Exchange; - Shuffle 分区数是否合理;
- 是否存在数据倾斜;
- 是否产生了磁盘 Spill;
- 是否可以使用 map-side combine;
- 是否可以使用 Broadcast Join。