
摘要Shuffle 是 Spark 作业中最昂贵的操作也是性能瓶颈的高发区。Spark 2.0 起 SortShuffleManager 成为唯一的默认实现通过 ExternalSorter 分区排序 溢写归并 Index File 精确索引三大机制彻底解决了 HashShuffle 的 M×R 文件爆炸问题。本文从 SortShuffleManager 的三种 ShuffleHandle 选择策略、ExternalSorter 溢写与归并流程、Index File 索引机制、Tungsten 堆外排序优化、HashShuffle → SortShuffle 演进对比五个维度配合 2 张架构图 源码追踪完整拆解 SortShuffle 原理。关键词SortShuffleManager, ExternalSorter, Index File, BypassMerge, Tungsten, HashShuffle, Shuffle 优化一、开篇Shuffle 为何是 Spark 的「阿喀琉斯之踵」凡是涉及reduceByKey/groupByKey/join/sortBy等宽依赖算子都会触发一次 Shuffle。Shuffle 涉及磁盘 I/O、网络传输和序列化反序列化是 Spark 性能调优的核心战场。Spark Shuffle 演进时间线 Spark 0.x-1.x: HashShuffleManager (M×R 文件已弃用) Spark 1.2: SortShuffleManager 引入可选 Spark 2.0: SortShuffleManager 成为唯一默认 Spark 2.x: Tungsten-sort (UnsafeShuffleWriter)二、SortShuffle 架构全景2.1 三种 ShuffleHandle 选择策略// SortShuffleManager.registerShuffle() 决策逻辑if(依赖不需要 map-side combine不要求排序分区数bypassMergeThreshold){→ BypassMergeSortShuffleHandle(无排序每个分区独立文件最后合并为一个 data file)}elseif(记录可序列化不需要 aggregation不需要排序分区数16777216){→ SerializedShuffleHandle(UnsafeShuffleWriter — Tungsten 堆外排序)}else{→ BaseShuffleHandle(SortShuffleWriterExternalSorter 通用路径)}2.2 ExternalSorter 溢写与归并ExternalSorter.insertAll(records): ① 数据先写入内存中的 PartitionedAppendOnlyMap ② 内存超过阈值 → spill() 溢写磁盘 - 按 (partitionId, key) 排序 - 写入临时文件受 spark.shuffle.spill.compress 控制 ③ 所有数据处理完毕 → 多路归并 - 将所有溢写文件 内存中剩余数据 - 按 (partitionId, key) 归并排序 ④ 输出到最终 data file index file2.3 Index File — SortShuffle 的关键创新Index File 结构每分区 16 字节 [p0_offset(8B)|p0_length(8B)|p1_offset|p1_length|...] 价值 ├── Reduce 端可精确定位每个分区的数据块 ├── 无需全量扫描 data file ├── 支持并行拉取不同分区 └── Netty 零拷贝传输sendfile三、BypassMerge 机制当分区数较少≤200且不需要排序时BypassMerge 直接为每个分区写入独立临时文件最后合并为一个 data file 一个 index file。省去了排序开销适合groupByKey/join等不需要 map-side combine 的场景。spark.shuffle.sort.bypassMergeThreshold200 # 默认值 # 如果期望的分区数 200使用 sort-based 路径 # 如果较小且无排序需求走 bypass 路径更快四、Tungsten-Sort二进制堆外排序当满足 SerializedShuffleHandle 条件时SortShuffleManager 启用UnsafeShuffleWriter。核心机制将序列化后的二进制数据放在堆外内存中直接按分区 ID 排序数据指针——全程不反序列化零GC压力。Tungsten-sort 启用条件 ├── spark.shuffle.managersort ├── 无 map-side aggregation ├── 无排序要求 ├── 分区数 16777216 (2^24) ├── 序列化器支持对象重定位Kryo / Java └── spark.sql.tungsten.enabletrue默认五、HashShuffle → SortShuffle 演进核心对比HashShuffle 缺陷SortShuffle 解决方案M×R 个文件2M 个文件dataindex无排序Reduce 全量扫描分区排序 Index File 精确定位无溢写内存可能 OOMExternalSorter 自动溢写归并文件句柄耗尽BypassMerge 合并为单文件无堆外优化Tungsten UnsafeShuffleWriter六、核心源码追踪// SortShuffleManager.write() — Shuffle Write 入口overridedefgetWriter[K,V](handle:ShuffleHandle,...):ShuffleWriter[K,V]{handlematch{caseunsafe:SerializedShuffleHandle[Kunchecked,Vunchecked]newUnsafeShuffleWriter(env,blockManager,...)casebypass:BypassMergeSortShuffleHandle[Kunchecked,Vunchecked]newBypassMergeSortShuffleWriter(blockManager,...)caseother:BaseShuffleHandle[Kunchecked,Vunchecked,_]newSortShuffleWriter(shuffleBlockResolver,handle,mapId,context)}}// SortShuffleWriter.write() → ExternalSortersorternewExternalSorter(context,aggregator,ordering,serializer)sorter.insertAll(records)// ① 插入排序溢写sorter.writePartitionedFile(outputWriter)// ② 输出 dataindex七、关键配置速查# SortShuffle 核心参数 spark.shuffle.managersort # 默认 spark.shuffle.sort.bypassMergeThreshold200 # bypass 阈值 spark.shuffle.spill.compresstrue # 溢写压缩 spark.shuffle.compresstrue # shuffle 输出压缩 spark.shuffle.file.buffer32k # 写缓冲 spark.reducer.maxSizeInFlight48m # reduce 拉取缓冲 spark.shuffle.sort.io.plugin.class # 自定义 IO 插件 # Tungsten spark.sql.tungsten.enabledtrue # 默认启用八、总结SortShuffleManagerSpark 2.0 起唯一默认实现三种 ShuffleHandle 覆盖所有场景。BypassMerge 处理小分区直写Serialized 启用 Tungsten 堆外排序Base 走 ExternalSorter 通用路径。ExternalSorter内存中 PartitionedAppendOnlyMap → 超阈值自动溢写 → 多路归并排序 → 输出 data index 文件。Index File16B/分区实现精准分区定位。演进路径HashShuffleM×R 文件爆炸→ SortShuffle2M 文件 排序 Index→ Tungsten-Sort堆外二进制排序零GC。作者starzy博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践