上一篇讲了 MapReduce 编程模型。本篇展开 MapReduce 性能代价最重的环节——Shuffle。

Shuffle 常被介绍成"Map 到 Reduce 之间的数据传输"。这个描述对应了功能但完全没解释 Shuffle 为什么慢。准确的说法是:Shuffle 由 Map 端的 Spill-Sort-Merge、Reduce 端的 Fetch-Merge 和网络 HTTP 传输共同组成,涉及磁盘、网络、内存缓冲和排序算法。在许多 Reduce-heavy 作业里,Shuffle 会成为最主要的延迟来源;调优 Shuffle 往往就是调优 MapReduce 的核心部分。

本篇只抓一个问题:Map 端 Shuffle 怎么从内存缓冲流到本地磁盘分区文件、Reduce 端怎么从所有 Map 拉取并合并、CombineFileInputFormat 在什么场景下能减少 Map 数。

Map 端 Shuffle 的五个步骤

Map Task 输出的 key-value 对不是直接发给 Reduce,而是经过一个六步骤的本地流水线:

1
2
3
4
5
6
1. map() 输出进入环形缓冲区(默认 100MB)
2. 缓冲区使用率达到阈值(默认 80%)时触发 Spill
3. Spill 时按 Partition 排序,每个 Partition 内按 key 排序
4.(可选)对每个 Partition 跑 Combiner
5. Spill 输出为本地磁盘文件
6. Map 完成时合并所有 Spill 文件为单个分区输出文件

详细展开每一步:

环形缓冲区(Circular Buffer)是 Map 输出的临时存放地。MapReduce 用一个固定大小的 byte 数组(默认 mapreduce.task.io.sort.mb = 100MB)作为缓冲,map() 输出的 key-value 对序列化后写入这个数组。数组是环形——写到尾部后绕回头部继续写。

缓冲区同时维护一个索引数组,记录每个 Partition 的起始位置和数据长度。索引数组让 Spill 时能快速按 Partition 切分。

Spill 在缓冲区使用率达到 mapreduce.map.sort.spill.percent(默认 0.80)时触发。Spill 不是阻塞 map() 的——Spill 在后台线程做,map() 继续往缓冲区另一半写。如果缓冲区另一半也满了(map() 比 Spill 快),map() 阻塞等待 Spill 完成。

Spill 的核心动作是排序。Spill 线程把缓冲区里的数据按 Partition 排序(先按 Partition 编号分桶,每个桶内按 key 排序)。排序使用快速排序,时间复杂度 O(n log n)。这就是为什么 MapReduce 输出最终是按 key 排序的——不是用户要求,是 Spill 排序的副产物。

排序完成后,Spill 线程把每个 Partition 的数据写入本地磁盘的临时文件(spillN.out,N 是 Spill 序号)。如果配置了 Combiner,Spill 在写盘前会对每个 Partition 跑一次 Combiner。如果配置了压缩(mapreduce.map.output.compress=true),Spill 写盘前会压缩数据。

Map Task 完成时,所有 Spill 文件被合并成单个分区输出文件(output 文件)。合并使用归并排序——同时打开多个 Spill 文件(默认 mapreduce.task.io.sort.factor = 10 个),按 key 归并。如果 Spill 文件数超过 factor,合并是分多轮做的(每轮合并 10 个文件,最终合并到 1 个)。

合并完成后,Map Task 把单个 output 文件的元数据(每个 Partition 的偏移、长度)发给 ApplicationMaster,AM 把这些信息转给 Reduce Task。Reduce Task 用这些元数据从 Map Task HTTP 拉取自己负责的 Partition。

Reduce 端 Shuffle 的三个步骤

Reduce Task 启动后开始从所有 Map Task 拉取自己负责的 Partition。这个流程可以拆成三步:

1
2
3
1. Fetch:并发 HTTP 请求所有 Map Task,拉取 Partition 数据
2. Merge:归并排序所有拉取的数据,按 key 全局有序
3. Group:把同一个 key 的所有 value 收集成 Iterable 传给 reduce()

Fetch 是 Reduce 端最重的操作。Reduce Task 默认用 mapreduce.reduce.shuffle.parallelcopies(默认 5)个线程并发拉取。每次拉取是 HTTP 请求到 NodeManager 上的 ShuffleHandler 辅助服务;mapreduce.shuffle.port 在 Hadoop 3.4.1 的默认值是 13562。

Reduce Task 怎么知道哪些 Map Task 持有自己需要的 Partition?通过 ApplicationMaster。AM 收集所有 Map Task 的完成状态(output 文件元数据 + HTTP 地址),Reduce Task 周期性向 AM 询问 “已完成的 Map 列表”。Reduce Task 拿到列表后立即开始拉取,不等所有 Map 完成——这是为什么 Reduce Task 可以与 Map Task 重叠执行(前面讲的 slowstart.completedmaps)。

Merge 是把多个有序的小文件合并成有序的大文件。Reduce Task 维护一个 in-memory buffer(mapreduce.reduce.shuffle.input.buffer.percent,默认 0.70 × JVM heap),Fetch 来的数据先进内存。内存满后写本地磁盘。

最终 reduce() 调用时,Reduce Task 把所有数据(内存 + 磁盘)按 key 归并排序,按 key 分组形成 Iterable。reduce() 处理完一个 key 后立即处理下一个,不需要等所有数据齐——这是流式归并,内存里同时持有的只是当前正在处理的 key 的所有 value。

Shuffle 的网络代价

Shuffle 涉及大量网络传输。一个 100 个 Map Task + 100 个 Reduce Task 的作业,逻辑上最多形成 100 × 100 = 10000 份 map-output 分片拉取。如果每个分片平均 100MB,总网络流量约 10000 × 100MB = 1,000,000MB,按十进制约 1TB,不是 10TB。若要达到 10TB,要么分片均值接近 1GB,要么 Map/Reduce 数继续放大。

MapReduce 的 Shuffle 是 “all-to-all” 模式——每个 Reduce 都要从每个 Map 拉数据。这与 Spark 的 Shuffle 类似,但与 Flink 的 streaming shuffle 不同(Flink 是 key-based 路由,不是 all-to-all)。

网络 Shuffle 的优化方向:

数据压缩:Map 输出在 Spill 时压缩(Snappy / LZO / Zstd),网络传输的字节数减少 50-90%。代价是 CPU 开销。

数据本地化:如果 Map 输出的 Partition 与 Reduce 在同一节点,Shuffle 走 loopback,避开网络。这是 MapReduce 的默认优化——Reduce Task 优先从本地 Map Task 拉取。但 Reduce Task 数通常远少于 Map Task 数,每个 Reduce 必须从大部分 Map 拉取,本地化命中率有限。

Combiner 减少 Map 输出量:前面讲过,Combiner 在 Map 端本地聚合,减少 Spill 数据量。Spill 减少 = Shuffle 网络流量减少。

CombineFileInputFormat 合并小文件:如果输入是大量小文件,普通文件输入格式容易产生过多 Map Task,进而放大 Shuffle 连接数。CombineFileInputFormat 把多个小文件合并成一个 CombineFileSplit,减少 Map 数。

CombineFileInputFormat:减少 Map 数

FileInputFormat 是文件输入格式的抽象基类,默认按 block 大小和文件可切分性生成 FileSplit;对于大量远小于 block 的小文件,通常会形成接近"每文件一个 split"的 Map 数。一个目录里有 10000 个 1MB 小文件,常见配置下会创建接近 10000 个 Map Task。每个 Map Task 启动开销(JVM 启动 + 初始化)几秒,累计开销会被放大。

CombineFileInputFormat 把多个小文件合并成 CombineFileSplit,每个 split 可以达到 block 大小(128MB)。10000 个 1MB 文件用 CombineFileInputFormat 切成约 80 个 split,Map 数从 10000 降到 80,启动开销减少约 100 倍。

1
2
3
4
job.setInputFormatClass(CombineTextInputFormat.class);
CombineFileInputFormat.setMaxInputSplitSize(job, 128 * 1024 * 1024); // 128MB
CombineFileInputFormat.setMinInputSplitSizeNode(job, 10 * 1024 * 1024); // 节点级最小 10MB
CombineFileInputFormat.setMinInputSplitSizeRack(job, 40 * 1024 * 1024); // 机架级最小 40MB

这里不能直接把 CombineFileInputFormat.class 塞给 setInputFormatClass:Hadoop 3.4.1 Javadoc 里 CombineFileInputFormat 是抽象类,真实作业要使用 CombineTextInputFormatCombineSequenceFileInputFormat,或自定义继承它并实现 RecordReader 的具体类。这三个参数控制 CombineFileSplit 的边界。最大 split 大小决定合并上限,节点级最小决定同节点内合并的下限,机架级最小决定跨节点合并的下限。

CombineFileInputFormat 还考虑数据本地化——优先合并同节点的文件,再合并同机架的文件。这让 Map Task 的数据本地化率(在持有数据的节点上跑的比例)仍然保持高。

小文件问题在 Hadoop 生态里普遍存在——日志归档、爬虫结果、传感器数据等都产生大量小文件。CombineFileInputFormat 是处理这类输入的标准工具,生产环境的小文件作业几乎都用它。

Shuffle Handler 的 HTTP 协议

Reduce Task 与 Map 输出所在节点之间的 Shuffle 走 HTTP。MapReduce 在 NodeManager 中通过 mapreduce_shuffle auxiliary service 启动 org.apache.hadoop.mapred.ShuffleHandler,由它监听端口并向 Reduce 端提供已完成 Map 的中间输出。

Shuffle Handler 的实现基于 Netty,默认端口由 mapreduce.shuffle.port 控制,Hadoop 3.4.1 源码常量 DEFAULT_SHUFFLE_PORT 为 13562。每次拉取是一个 HTTP GET 请求,URL 形如:

1
http://map-task-host:13562/mapOutput?job=job_1721900000000_0001&reduce=5&map=attempt_1721900000000_0001_m_000003

这个请求让 Map Task 把 job_…0001 作业里、Map Task attempt_…m_000003 输出的、Partition 5 的数据返回给 Reduce。

Shuffle Handler 的关键性能参数:

1
2
3
4
5
6
7
8
9
mapreduce.shuffle.max.threads(默认 0)
- NodeManager ShuffleHandler 的最大处理线程数
- 0 表示使用 Netty 事件循环的默认线程模型,不要写成每个 Map Task 的线程数

mapreduce.shuffle.max.connections(默认 0)
- NodeManager ShuffleHandler 最大 TCP 连接数,0 表示不按该项额外限制

mapreduce.shuffle.connection-keep-alive.timeout(默认 5s)
- keep-alive 连接超时

这些参数在大集群很重要。一个 NodeManager 上可能有多个已完成 Map 输出同时被许多 Reduce 拉取,ShuffleHandler 配置不当会让该节点的网络或连接队列先成为瓶颈。

Shuffle Handler 默认不加密。如果集群启用了安全(第十八篇),Shuffle Handler 配置 SSL/TLS 加密,但加密会增加 CPU 开销和延迟。

实验:观察 Shuffle 性能

实验状态:UNVERIFIED_RUNTIME。下面列的是需要在真实 Hadoop 3.4.1 集群上采集的 counters 与判断方法,本轮没有执行作业,也不把示例数值写成已验证结果。

提交一个 MapReduce 作业后,在 ApplicationMaster Web UI 的 Counters 标签页可以看到大量 Shuffle 相关 counter:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
Map-Reduce Framework Counter (per Task):
Map input records: 1000000
Map output records: 5000000
Map output bytes: 500000000
Combine input records: 5000000
Combine output records: 500000 ← Combiner 输入 5M,输出 500K,节省 90%
Spilled records: 500000 ← Spill 记录数 = Combiner 输出
BytesWritten (local disk): 50000000 ← 本地磁盘写入字节

Shuffle Metrics (per Reduce Task):
Combine input records: 0 (Reduce 端也有 Combine,通常为 0)
Combine output records: 0
Reduce input records: 500000 ← Reduce 处理的记录数
Reduce input groups: 50000 ← distinct key 数
Reduce output records: 50000
Reduce shuffle bytes: 60000000 ← 网络拉取字节
Reduce shuffle records: 500000
Failed shuffle: 0 ← 拉取失败次数(应该为 0)
Merged map outputs: 100 ← 合并了多少 Map 输出

这些 counter 是 Shuffle 性能调优的关键观测点。几个判断标准:

Spilled records / Map output records 接近 1 是正常。如果远大于 1(例如 3 倍),说明 Spill 太频繁,需要增大 mapreduce.task.io.sort.mb

Combine output records / Combine input records 应该是减少比例(例如 10%)。如果接近 100%,说明 Combiner 没起到合并作用(key 全部 distinct),Combiner 应该去掉。

Failed shuffle 必须接近 0。如果非 0,说明 Reduce 拉取失败重试,可能是网络问题或 Map Task 的 Shuffle Handler 配置不足。

Reduce shuffle bytes 是 Shuffle 网络流量。配合 Map output bytes 看压缩比——开启 Map 输出压缩后,shuffle bytes 应该显著小于 map output bytes。

通过修改 mapreduce.task.io.sort.mbmapreduce.task.io.sort.factormapreduce.reduce.shuffle.parallelcopies 等参数,重新跑作业,对比 counter 变化,可以直观看到 Shuffle 调优的效果。

模式提炼

MapReduce Shuffle 体现的设计模式:

1
2
3
4
5
6
7
8
模式:环形缓冲 + Spill 排序 + 网络拉取 + 归并合并

- 数据生产端用固定大小缓冲,满了就落盘(避免 OOM)
- 落盘前按 Partition + key 排序(让下游合并简单)
- 数据消费端并发拉取所有生产端数据
- 归并合并保持流式(不全量入内存)
- 中间走网络 HTTP(标准化协议)
- 调优靠调整缓冲大小、并发数、合并因子

这个模式不只是 MapReduce。Spark Shuffle 几乎是同构——SortShuffleManager 的内部数据流与 MapReduce 几乎相同。Flink 的 Batch Shuffle 也用类似流程(差异在 Flink 支持流式 shuffle)。Presto / Hive 的 Exchange operator 也是同一思路——上游落盘排序,下游拉取合并。

数据库的 parallel hash join 也是类似模式——build 端按 hash bucket 分桶落盘,probe 端按相同 hash 拉。差别在数据库按 hash 分桶,MapReduce 按 key 排序(hash 是 sort 的特殊情况)。

工程迁移表

MapReduce Shuffle 概念 Spark Shuffle Flink Shuffle Hive Exchange DB Parallel Join
环形缓冲 Serializer buffer network buffer sort buffer work_mem
Spill 排序 SortShuffleWriter sort-shuffle external sort external sort
Partition 文件 data file + index partitioned file bucket file hash bucket
Shuffle HTTP netty / transport netty jetty proprietary
Combiner map-side aggregation combiner map-side aggregation partial aggregation
Merge memory + disk merge chained merge merge join merge join
CombineFileInputFormat CombineFileInputFormat file source Combine input partition prune

注意 Spark 这一列。Spark 2.x 起默认用 SortShuffleManager,与 MapReduce 的 Shuffle 流程几乎相同。Spark 的优势是 Stage 内的窄依赖不需要 Shuffle(窄依赖在内存中 pipeline),只有宽依赖(Shuffle 边界)才走 Shuffle——这与 MapReduce 的"每个作业强制 Map-Reduce 两阶段"形成对比,是 Spark 性能优势的核心来源。

常见误解

误解一:“Shuffle 在内存里做很快”。MapReduce Shuffle 强制落盘——Map 端 Spill 落本地磁盘,Reduce 端 Fetch 后内存满也落盘。这是 MapReduce 设计的容错要求(落盘才能让 Task 失败重试不需要重跑 Shuffle)。Spark 在 Stage 内可以不落盘(窄依赖),但 Shuffle 边界仍然落盘。

误解二:“Shuffle 网络是 MapReduce 性能瓶颈”。网络确实是瓶颈之一,但磁盘 I/O 也是。Spill + Merge 涉及多次本地磁盘读写,磁盘带宽不够时 Spill 排队,整个 Map Task 变慢。生产环境 MapReduce 集群通常用 SSD 加速 Shuffle——相比 HDD,SSD 让 Shuffle 吞吐提升数倍。

误解三:“CombineFileInputFormat 总是好的”。CombineFileInputFormat 适合小文件场景。如果输入本来是大文件(128MB+),CombineFileInputFormat 不会改变 InputSplit 数。错误配置 CombineFileInputFormat 的 maxInputSplitSize(例如设为 1GB)可能让单个 InputSplit 跨多个 block,破坏数据本地化。

误解四:“Combiner 是 Reducer 的简化版本”。Combiner 和 Reducer 用同一个类(在 word count 场景),但它们的运行时机和输入不同。Combiner 在 Map 端 Spill 前跑,输入是单个 Map Task 的输出。Reducer 在 Reduce 端跑,输入是所有 Map Task 的对应 Partition 拉取后合并的结果。Combiner 是本地聚合优化,Reducer 是最终聚合。

误解五:“Reduce 数越多越好(并行度高)”。Reduce 数过多让每个 Reduce 拉取的 Partition 很小(例如 1MB),但每次 HTTP 请求的固定开销(连接握手、Shuffle Handler 处理)几乎与拉取大小无关。1MB 的拉取可能 50% 时间是 HTTP 开销。每个 Reduce 处理 1-5 GB 才能让网络有效利用。

练习

  1. 提交一个 word count 作业,在 ApplicationMaster Web UI 的 Counters 标签页观察 Map output records、Spilled records、Reduce shuffle bytes。修改 mapreduce.task.io.sort.mb 从 100MB 改到 512MB,重新跑,对比 Spilled records 变化。

  2. apache/hadoop 源码里找到 MapTask.javaReduceTask.javaShuffle.javaShuffleHandler.java(hadoop-mapreduce-client-core / hadoop-mapreduce-client-shuffle 模块),观察 Spill、Fetch、Shuffle Handler 的核心实现。

  3. 准备一个目录里放 1000 个小文件(每个 1MB),分别用 FileInputFormat 和 CombineFileInputFormat 提交 MapReduce 作业,对比 Map Task 数和作业总执行时间。

  4. 思考题:如果让 MapReduce Shuffle 走 RDMA(Remote Direct Memory Access)替代 HTTP,Shuffle 性能会怎么变化?哪些环节是 RDMA 也无法加速的?

系列导航

序号 主题 状态
00 导读:节点总会失败
01 HDFS 架构与三层切分
02 文件写入路径与流水线
03 文件读取路径与副本选择
04 NameNode 内存模型与启动恢复
05 HDFS HA 与脑裂防御
06 HDFS 3.x 演进与纠删码
07 YARN 架构与三方契约
08 YARN 资源模型、Container 与 NodeLabel
09 YARN 调度器对比:FIFO、Capacity、Fair
10 YARN 应用程序生命周期
11 YARN HA 与 Federation
12 YARN Timeline Service v2
13 MapReduce 编程模型与分而治之 上一篇
14 MapReduce Shuffle 全流程 本篇
15 MRv2 on YARN ApplicationMaster 与 Task Attempt 下一篇
16 Hadoop RPC 协议栈
17 序列化与压缩
18 Hadoop 安全 Kerberos Token ProxyUser
19 监控与运维 Metrics JMX 日志聚合
20 Hadoop 生态 Hive HBase Pig
21 Hadoop 与对象存储 Kubernetes 演进对比
22 Hadoop 设计遗产从 GFS MapReduce 到云原生

参考资料