Spark Shuffle 机制
概述
Shuffle 是 Spark 中最昂贵的操作:数据跨节点重分布、落盘、序列化。理解 Shuffle 的演进与调优,是 Spark 性能优化的核心。本文讲清 Hash/Base Sort/Tungsten 三种 Shuffle 的实现与差异、Shuffle 读写流程与调优。
一、Shuffle 的本质
1.1 什么时候发生
| 算子 | 说明 |
|---|---|
| reduceByKey / groupByKey | 按键重分区 |
| join | 按 key 关联 |
| repartition | 改变分区数 |
| distinct | 去重分组 |
| sortByKey | 全局排序 |
1.2 代价分析
Map 端:写出分区文件(序列化、可能溢写磁盘)
网络:数据跨节点传输
Reduce 端:拉取、归并、聚合Shuffle 涉及磁盘 IO + 网络 IO + 排序/聚合 CPU,是性能瓶颈的主要来源。
二、Shuffle 整体流程
Map Task(ShuffleMapTask)
├── ShuffleMapOutputWriter:按键分区写出
├── 内存缓冲(spark.shuffle.file.buffer)
├── 超过阈值溢写磁盘,多次溢写合并
└── 产出分区文件 + 索引文件
Reduce Task(ResultTask)
├── 根据 map 输出位置拉取数据(可聚合多个块)
├── 内存缓冲(spark.reducer.maxSizeInFlight)
└── 归并排序 / 聚合后计算三、Shuffle 演进史
3.1 Hash Shuffle(Spark 1.2 前)
| 特点 | 说明 |
|---|---|
| 每个 Map Task 为每个 Reduce 分区写一个文件 | 文件数 = MapTask × 分区数 |
| 小文件爆炸 | 1000 Map × 1000 分区 = 100 万文件 |
| 优化版 | Consolidate 合并文件,同 Executor 复用 |
缺点:文件句柄与 IO 开销巨大,已废弃。
3.2 Sort Shuffle(Base Sort,Spark 1.2+ 默认)
| 特点 | 说明 |
|---|---|
| 每个 Map Task 只写一个数据文件 + 一个索引文件 | 文件数 = MapTask × 2 |
| 内存中按分区排序 | 支持溢出合并 |
| 稳定性好 | 有额外排序开销 |
Hash: M1 → 分区0文件、分区1文件、分区2文件...
Sort: M1 → 单文件(内部按分区号排序)+ 索引3.3 Tungsten Shuffle(Shuffle 2.0 起默认)
| 特点 | 说明 |
|---|---|
| 基于 Tungsten 二进制内存 | 避免对象序列化 |
| 直接内存写入分区 | 减少 GC |
| 偏移量索引 | 快速定位 |
| 优化小文件合并 | shuffle.mergeLocations |
演进总结:文件数从 M×R 降到 M×2,内存与 GC 持续优化,网络与 IO 开销最小化。
四、Shuffle 读优化
4.1 网络请求合并
| 参数 | 默认 | 说明 |
|---|---|---|
spark.reducer.maxSizeInFlight | 48m | 每 Reduce 同时拉取的数据量 |
spark.reducer.maxReqsInFlight | 20 | 并发请求数 |
spark.reducer.maxBlocksInFlightPerAddress | 5 | 每节点并发块数 |
4.2 排序与聚合
| 机制 | 说明 |
|---|---|
| 按分区聚合 | 同类 key 合并在一个分区 |
| 内存聚合 | spark.shuffle.aggregate.bufferSize 聚合缓冲区 |
| 归并排序 | 多个溢写文件合并成有序流 |
五、Shuffle 调优参数表
| 参数 | 默认 | 作用 |
|---|---|---|
spark.shuffle.file.buffer | 32k | Map 端写出缓冲 |
spark.shuffle.memoryFraction | 0.2(旧版) | Shuffle 内存占比(新版由统一内存管理接管) |
spark.shuffle.spill.compress | true | 溢写压缩 |
spark.shuffle.compress | true | 网络传输压缩 |
spark.shuffle.io.maxRetries | 3 | 拉取失败重试 |
spark.shuffle.io.retryWait | 5s | 重试间隔 |
spark.shuffle.sort.bypassMergeThreshold | 200 | 分区数小于此值时跳过合并排序 |
spark.shuffle.service.enabled | true | 外部 Shuffle 服务(动态分配时必开) |
5.1 bypassMergeThreshold
| 分区数 | 行为 |
|---|---|
| < 200 | 不做排序合并,直接写多个分区文件(类 Hash,更快) |
| > 200 | 走完整 Sort Shuffle |
5.2 压缩与序列化
| 手段 | 收益 |
|---|---|
| 压缩(lz4/zstd) | 网络/磁盘 IO 减 50%+,略耗 CPU |
| Kryo 序列化 | 比 Java 序列化快 10 倍、更小 |
scala
sparkConf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")六、Shuffle 数据倾斜与排查
6.1 倾斜表现
| 症状 | 原因 |
|---|---|
| 个别 Reduce Task 极慢 | key 分布不均,某 key 数据量巨大 |
| 某 Executor 内存溢出 | 单个 key 聚合数据超内存 |
6.2 常用手段
| 手段 | 说明 |
|---|---|
| 加盐 | 对热点 key 加随机前缀拆散,聚合后去盐 |
| 两阶段聚合 | 本地聚合 + 全局聚合 |
| 提高并行度 | 增加分区数,分散热点 |
| 过滤极端 key | 单独处理后 union |
| map 端预聚合 | 用 reduceByKey 而非 groupByKey |
七、Shuffle 调优决策清单
优先检查:
1. 是否真的需要 Shuffle(能否 map 端聚合/广播 join)
2. 分区数是否合理(2-3 倍 Executor 核数)
3. 压缩与 Kryo 是否开启
4. 是否出现数据倾斜
避免:
groupByKey(无预聚合)→ 换 reduceByKey
大表 join 大表 → 换 broadcast join 或分桶常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 小文件爆炸(Shuffle 后) | 提高分区数与动态合并,控制输出分区 |
| 拉取失败重试超时 | 网络问题或 Executor 频繁重启,调大重试参数 |
| Shuffle 写磁盘过慢 | 调大 file.buffer、开启压缩、用 SSD |
| OOM during shuffle | 减少并发拉取、提高并行度、治倾斜 |