Hadoop Shuffle 流程详解:从 Map 到 Reduce 的数据流转(红色字体为个人经验补充,AI基础上做了修改和补充)
1. Hadoop Shuffle 流程概述
Shuffle 是 Hadoop MapReduce 框架中连接 Map 和 Reduce 阶段的关键过程。它负责将 Map 任务输出的中间结果进行分区、排序、合并,并传输给对应的 Reduce 任务。理解 Shuffle 流程对于优化 MapReduce 作业性能至关重要。
2. 完整 Shuffle 流程图
以下是 Hadoop Shuffle 的完整流程图,展示了从 Map 输出到 Reduce 输入的完整数据流转过程:
flowchart TD
subgraph Map端
A["Map Task 执行"] --> B["Map 输出 key-value 对"]
B --> C["写入环形缓冲区 Circular Buffer"]
C --> D{"缓冲区是否达到阈值?"}
D -- 是 --> E["启动溢出线程 Spill Thread"]
D -- 否 --> C
E --> F["分区 Partition 按 key 哈希到不同分区"]
F --> G["分区内排序 Sort 按 key 排序"]
G --> H["写入本地磁盘 Spill 文件"]
H --> I{"所有 Map 输出完成?"}
I -- 否 --> C
I -- 是 --> J["合并 Merge 合并所有 Spill 文件"]
J --> K["生成最终 Map 输出文件 每个分区一个文件"]
end
subgraph Reduce端
K --> L["Reduce Task 拉取数据 HTTP GET"]
L --> M["合并 Map 输出 内存/磁盘合并"]
M --> N["排序 Sort 合并后再次排序"]
N --> O["分组 Group 相同 key 的 values 合并"]
O --> P["Reduce 函数输入"]
P --> Q["Reduce 执行"]
Q --> R["Reduce 输出 最终结果"]
end
subgraph 关键参数影响
S["io.sort.mb 缓冲区大小"] --> C
T["io.sort.spill.percent 溢出阈值"] --> D
U["io.sort.factor 合并因子"] --> J
V["mapreduce.task.io.sort.factor Reduce 端合并因子"] --> M
W["mapreduce.reduce.shuffle.parallelcopies 并行拷贝数"] --> L
end
3. Map 端 Shuffle 流程(输入是根据split进行顺序读入)
3.1 Map 输出与环形缓冲区
Map 任务执行用户定义的 map 函数,生成 key-value 对作为中间输出。这些输出不会直接写入磁盘,而是先写入一个环形缓冲区(Circular Buffer)。
- 数据流向:Map 函数 → 内存缓冲区
- 关键参数:
io.sort.mb(默认 100MB)控制缓冲区大小/新版本 mapreduce.task.io.sort.mb
3.2 溢出(Spill)过程
当缓冲区使用量达到阈值(由 io.sort.spill.percent / mapreduce.map.sort.spill.percent控制,默认 0.8)时,会启动溢出线程:
- 分区(Partition):根据 key 的哈希值决定数据属于哪个 Reduce 分区(分区根据输入或者默认reduce多少进行分 hash(key)/reduce数;
- 排序(Sort):在
每个分区内按 key 进行分区 ,key排序(快速排序/归并 看数据量大小,小的用快速,hadoop自有算法),在内存排序 - 写入磁盘:生成临时的 Spill 文件到本地磁盘
- 锁:溢出时候会有个锁,比如80%溢出,那么1-80%位置会被标记锁住,另外二十继续写;
3.3 合并(Merge)与最终输出
Map 任务完成后,将所有 Spill 文件合并成一个分区且排序的文件:
- 合并因子:
io.sort.factor(默认 10)控制一次合并的文件数 - 最终输出:每个 Map 任务生成一个文件(跟文件块没任何关系就是一个大文件),包含所有分区数据
- 索引文件:同时生成索引文件,记录每个分区在文件中的偏移量
3.34 思考点
1,这里有几个设计的思考点,为什么map的阶段不能边溢出边合并;而是要等待全部溢出完再合s并;答:查阅一些材料和AI作答;map端的io压力大,因为有个本地读;reduce就没有,reduce压力主要在网络io;所以他们根据架构开发和整体取舍就一直这么处理;
2,为什么合并是串行而不是并行;答:我理解还是io问题;并行的话 无上限;后面hadoop版本调整为64做一组归并;
3,map数量根据什么分配? 分片;正常跟块一个大小;大小不一 增加网络消耗;
4. Reduce 端 Shuffle 流程
4.1 数据拉取(Copy Phase)
Reduce 任务启动后,通过 HTTP 从各个 Map 任务的输出节点拉取属于自己的分区数据:
- 并行拷贝:
mapreduce.reduce.shuffle.parallelcopies(默认 5)控制同时拉取的线程数 - 内存缓冲区:拉取的数据先放入内存缓冲区,大小由
mapreduce.reduce.shuffle.input.buffer.percent控制
4.2 合并与排序(Merge & Sort Phase)
Reduce 端需要进行多轮合并:
- 内存合并:当内存中的数据达到阈值时,在内存中合并后溢写到磁盘(这里也涉及排序,但是按key)
- 磁盘合并:磁盘上的文件数量达到阈值(
mapreduce.task.io.sort.factor)时进行合并 - 最终排序:所有 Map 数据拉取完成后,进行最终合并排序,生成一个完全排序的文件
4.3 分组(Group)与 Reduce 输入
排序后的文件按 key 分组,相同 key 的 values 被合并成一个迭代器,作为 Reduce 函数的输入:
- 分组比较器:
mapreduce.job.output.group.comparator.class控制分组逻辑 - Reduce 输入:每个 key 及其对应的 values 迭代器传递给 reduce 函数
4.4 思考点
1,并发拷贝是5:答:经验,但可以调,看shuffle处理io跟得上吗;
二,分区数,根据什么确认?提交AM的时候的redcuce数量参数确认;唯一的;但是遇到order by 就强制1;
5. 关键参数及其影响
| 参数 | 默认值 | 影响流程 | 调优建议 |
|---|---|---|---|
io.sort.mb |
100MB | Map 端缓冲区大小 | 增大可减少溢出次数,但占用更多内存 |
io.sort.spill.percent |
0.8 | 溢出触发阈值 | 降低可提前溢出,减少内存压力 |
io.sort.factor |
10 | Map 端合并因子 | 增大可减少合并轮次,但需要更多内存 |
mapreduce.task.io.sort.factor |
10 | Reduce 端合并因子 | 影响磁盘文件合并效率 |
mapreduce.reduce.shuffle.parallelcopies |
5 | Reduce 并行拷贝数 | 增大可加快数据拉取速度 |
mapreduce.reduce.shuffle.input.buffer.percent |
0.7 | Reduce 内存缓冲区占比 | 增大可减少磁盘 I/O |
mapreduce.reduce.merge.inmem.threshold |
1000 | 内存合并阈值 | 控制内存合并的触发条件 |
上面的reduce 有俩参数控制内存合并,一个是缓冲一个是内存合并,这俩都会溢出;先到先出;按官网的说法是为了优化io;
6. 数据输入输出总结
6.1 Map 端数据流
- 输入:HDFS 数据块(InputSplit)
- 处理:map(key1, value1) → list(key2, value2)
- 输出:分区且排序的中间文件(每个 Map 一个)
6.2 Reduce 端数据流
- 输入:所有 Map 任务中对应分区的数据
- 处理:reduce(key2, [value2, value2, ...]) → (key3, value3)
- 输出:最终结果写入 HDFS
7. 性能优化建议
- 减少数据量:使用 Combiner 减少 Map 输出数据量(这个过程是在map端溢出之前,排序之后,按key)
- 合理分区:自定义 Partitioner 避免数据倾斜(常规是加盐最好)
- 压缩中间数据:启用 Map 输出压缩(
mapreduce.map.output.compress) - 调整缓冲区:根据集群内存情况调整缓冲区大小
- 监控 Spill 次数:通过计数器监控,避免过多磁盘 I/O
8. 总结
Hadoop Shuffle 是 MapReduce 框架中最复杂且最耗时的环节,涉及内存管理、磁盘 I/O、网络传输和排序算法。通过理解完整的 Shuffle 流程、关键参数的影响以及优化策略,可以显著提升 MapReduce 作业的性能。在实际应用中,需要根据数据特征和集群资源进行针对性调优。
更多推荐




所有评论(0)