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

Shuffle 常被介绍成"Map 到 Reduce 之间的数据传输"。这个描述对应了功能但完全没解释 Shuffle 为什么慢。准确的说法是:Shuffle 是 Map 端的 Spill-Sort-Merge 流水线、Reduce 端的 Fetch-Merge 流水线、中间的网络 HTTP 传输共同组成的复杂流程,涉及磁盘 I/O、网络 I/O、内存缓冲、排序算法等多个层。Shuffle 是 MapReduce 作业 50%+ 延迟的来源,调优 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 数组(默认 io.sort.mb = 100MB)作为缓冲,map() 输出的 key-value 对序列化后写入这个数组。数组是环形——写到尾部后绕回头部继续写。

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

Spill 在缓冲区使用率达到 io.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 文件(默认 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 请求到 Map Task 的 Shuffle 服务(Netty HTTP server,默认端口 8080+随机)。

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 的作业,Shuffle 网络连接数是 100 × 100 = 10000 个 HTTP 请求。如果每个 Partition 平均 100MB,总网络流量是 10TB。

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,Map 数过多导致 Shuffle 网络连接数爆炸。CombineFileInputFormat 把多个小文件合并成一个 InputSplit,减少 Map 数。

CombineFileInputFormat:减少 Map 数

默认的 FileInputFormat 按每个文件单独切 InputSplit。一个目录里有 10000 个 1MB 小文件,会创建 10000 个 Map Task。每个 Map Task 启动开销(JVM 启动 + 初始化)几秒,10000 个 Map Task 的启动开销是几小时。

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

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

这三个参数控制 InputSplit 的边界。最大 split 大小决定合并上限,节点级最小决定同节点内合并的下限,机架级最小决定跨节点合并的下限。

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

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

Shuffle Handler 的 HTTP 协议

Reduce Task 与 Map Task 之间的 Shuffle 走 HTTP。MapReduce 框架在每个 Map Task 内嵌一个 HTTP server(Shuffle Handler),监听一个端口接受 Reduce 的拉取请求。

Shuffle Handler 的实现(Hadoop 2.x+)基于 Netty,支持 NIO 和 keep-alive。每次拉取是一个 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 = 2 × CPU 核数)
- 每个 Map Task 的 Shuffle HTTP server 最大并发线程数
- 限制 Map Task 能同时服务多少 Reduce 请求

mapreduce.shuffle.max.connections(默认 0 = 2 × max.threads)
- 单个 Map Task 最大 TCP 连接数

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

这些参数在大集群很重要。一个 Map Task 可能被几百个 Reduce Task 同时请求,Shuffle Handler 配置不当会让 Map Task 网卡打满,其他 Reduce 拉不动。

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

实验:观察 Shuffle 性能

提交一个 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 太频繁,需要增大 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。

通过修改 io.sort.mbio.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。修改 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-06 HDFS 存储层(00 导读 / 01-06 HDFS 各主题) 第一阶段(已完成)
07-12 YARN 资源管理层 第二阶段(已完成)
13 MapReduce 编程模型:分而治之的工程化表达 上一篇
14 Shuffle 全流程:Map 端、Reduce 端与磁盘 I/O 的代价 本篇
15 MRv2 on YARN:ApplicationMaster 与 Task Attempt 下一篇
16-19 HA、安全与 Common 基础设施 第四阶段,待开始
20-22 演进、生态与对比 第五阶段,待开始

参考资料