考察点
这题区分「背过八股」和「真调过 Spark 性能」的人。面试官想听:HashShuffle 为什么被废弃(小文件爆炸的量化理解)、SortShuffleManager 里 bypass 和普通 sort 分支的触发条件、以及 Tungsten 到底优化了什么(很多人只会念「钨丝计划」四个字)。追问常往 shuffle 参数调优(buffer、spill、压缩)、shuffle read 的归并过程、外部 shuffle service 这些方向走。
参考答案
为什么 shuffle 是 Spark 最重的环节
shuffle 是宽依赖的数据重分配过程:map 端把每条记录按分区器算目标 partition,写到本地;reduce 端从所有 map 端拉取属于自己的数据并归并。它贵在三处——全量网络 IO、磁盘 spill、以及排序/归并的 CPU 开销。一次大 shuffle 的耗时经常占到整个 job 的一半以上,所以 Spark 历次大版本都在打磨它。
HashShuffle 时代(Spark 1.2 之前)
最早的 HashShuffleManager 逻辑简单粗暴:map 端每个 task 为每个 reduce partition 各写一个文件。M 个 map task、R 个 reduce partition,就产生 M×R 个小文件。
问题很直接:10 万 map × 1 万 reduce = 10 亿个文件。每个文件都要打开/关闭句柄、做 buffer 管理,文件系统元数据压力爆炸,而且大量随机写让小 IO 拖垮磁盘。后来加了 ConsolidateShuffle 的优化(同一个 executor 上复用文件组),缓解但没根治。
SortShuffleManager(1.2 起,2.0 后成为唯一)
思路对齐 MapReduce:每个 map task 只写一个 data 文件,加一个 index 文件记录各 partition 段落的偏移。文件数从 M×R 降到 2M,reduce 端按 index 区间去各 map 端拉数据。它内部有三个分支,按条件自动选择:
1. BypassMergeSortShuffleWriter(免排序版)
触发条件:reduce partition 数 ≤ spark.shuffle.sort.bypassMergeThreshold(默认 200),且 map 端没有聚合语义(比如不是 reduceByKey 这种带 combine 的算子)。逻辑是退化成 hash 式的临时文件写法(每个 partition 一个临时文件),最后归并成一个文件 + 索引。因为 partition 少,文件数可控,还省掉了排序开销——这就是阈值默认 200 的原因。
2. SortShuffleWriter(普通排序版)
map 端用 PartitionedPairBuffer(内存)或外部排序把记录按 (partitionId, key) 排序,满了就 spill 到磁盘,最后归并所有 spill 文件成一个有序文件。需要排序语义的算子(sortByKey)和 map 端预聚合(reduceByKey 的 combine)都走这条路。
3. UnsafeShuffleWriter(Tungsten 序列化版)
这是 Tungsten 项目的一部分,触发条件较苛刻:序列化器支持重定位(Kryo 默认支持)、不需要 map 端聚合、partition 数 < 2^24。它直接对序列化后的二进制字节做排序——比的是 byte 数组前缀,不做反序列化,不建 Java 对象。配合 Tungsten 的堆外内存管理,Java 对象开销和 GC 压力大幅下降,cache locality 也更好。
Tungsten 到底优化了什么
Tungsten 不是一个 shuffle writer 那么简单,它是一整套物理执行优化,shuffle 只是受益者之一:
- 二进制内存布局:数据以二进制的 UnsafeRow 形式存放在堆外内存(sun.misc.Unsafe 管理),绕开 JVM 对象头(一个空 Java 对象都有 16 字节对象头)和 GC。
- 基于二进制比较的排序:String/数组类型的排序直接比较序列化后的字节,语义上等价于反序列化后比较,但省掉了对象构造。spark.sql 的执行路径全面走这套机制。
- Whole-stage codegen:把多个算子的代码融合编译成一个 Java 函数,消除虚函数调用和迭代器开销,数据尽量留在寄存器里流动。这对 shuffle 前后的 map 端计算提速明显。
shuffle 的另一半:read 端与实用调优
map 端写完只是上半场。reduce 端从 MapOutputTracker 拿到各 map 输出位置,由 BlockStoreShuffleReader 拉取、归并。几个生产上真正有用的参数:
| 参数 | 默认 | 作用 |
|---|---|---|
| spark.shuffle.compress | true | shuffle 写盘压缩,省 IO 费 CPU |
| spark.shuffle.file.buffer | 32k | 写文件缓冲,机械盘可调大 |
| spark.reducer.maxSizeInFlight | 48m | reduce 端单次拉取上限 |
| spark.shuffle.io.maxRetries | 3 | 拉取失败重试,配合 retryWait |
| spark.shuffle.service.enabled | false | 外部 shuffle service,动态资源下必开 |
外部 shuffle service 值得单独说:executor 被回收(动态分配)后它的 shuffle 文件还在 NodeManager 的 ESS 进程手里,reduce 端照常能拉。开动态资源(spark.dynamicAllocation.enabled)的集群必须开 ESS,否则 task 会因为 fetch failed 大面积重试。
可能的追问
- shuffle 过程中 spill 是怎么回事? map 端排序缓冲(默认初始 5MB,可增长)装不下时把已排序的部分写磁盘成一个 spill 文件,最后把所有 spill 文件做归并排序输出。spill 次数多说明内存给少了,可以调 spark.shuffle.memoryFraction 相关配置或增大 executor 内存。
- FetchFailed 怎么排查? 常见三个原因:map 端 executor OOM/超时被杀、ESS 没开或异常、网络拥塞。应对是开 ESS、调大 shuffle.io 重试次数、检查 map 端 GC 日志。FetchFailed 会触发 stage 重算,偶发可接受,频发必须查根因。
- Spark 3.x 之后 shuffle 还有什么新东西? AQE 能动态合并/拆分 shuffle partition、动态切换 join 策略;3.2 起 push-based shuffle 把 map 输出推给 ESS 合并成大块,减少随机读,大集群下收益明显。