MapTask 工作机制:从读数据到写磁盘,一个 Map 任务的完整一生
一个 Map 任务,从开始到结束,走五个阶段:
Read → Map → Collect → Spill → Combine
- Read:把文件读进来,转成键值对
- Map:跑你的业务逻辑
- Collect:把结果写进内存缓冲区
- Spill:缓冲区满了,写到磁盘
- Combine:把磁盘上的小文件合并成大文件
flowchart TD
A[InputSplit 数据分片] -->|Read 阶段| B[RecordReader 解析为 ]
B -->|Map 阶段| C[用户自定义 map 函数处理为 ]
C -->|Collect 阶段| D[写入内存缓冲区]
D -->|Spill 阶段| E[分区 排序 溢写至磁盘]
E -->|Combine 阶段| F[局部聚合合并溢写文件]
F --> G[输出最终中间结果文件]
Map 阶段:跑你的业务逻辑
map() 函数接收一个键值对,输出零个或多个键值对。
WordCount 的例子最直观:
1 | protected void map(LongWritable key, Text value, Context context) { |
输入一行 "hello world",输出 <"hello", 1> 和 <"world", 1>。
这个阶段是 CPU 密集的,但通常不是瓶颈。 真正的瓶颈在后面——数据怎么写出。
Collect 阶段:写到内存缓冲区
Map 输出的每一条键值对,不能直接写磁盘(太慢了)。先写到内存缓冲区,攒够了再一起写磁盘。
- 缓冲区默认 100MB
- 达到 80% 时触发溢写
缓冲区是环形的——Map 一直在写,溢写线程在另一个方向读,互不干扰。
关键设计:缓冲区的存在,让 Map 的写操作几乎是内存速度。 没有缓冲区,每条记录写一次磁盘,Map 任务会慢 10 倍以上。
Spill 阶段:缓冲区满了,写到磁盘
缓冲区达到阈值(80%),触发溢写。这一步做三件事:
1. 分区
按 Key 决定这条数据去哪个 Reduce。默认按 Key 的哈希值取模:
1 | reduce分区 = hash(key) % reduce数量 |
2. 排序
每个分区内部按 Key 排序。为什么排序? Reduce 端拿到数据后要按 Key 分组,如果数据已经是排好序的,分组就快很多。
3. 溢写
排好序的数据写到磁盘上的一个临时文件(spill_{N}.out)。
一个 Map 任务可能产生多个溢写文件——缓冲区每满 80% 就溢写一次。
Combine 阶段:合并 + 局部聚合
Map 任务结束的时候,磁盘上可能有好几个溢写文件。Combine 阶段把它们合并成一个最终文件。
合并的时候做两件事:
1. 归并排序
多个溢写文件按分区和 Key 做多路归并,合成一个大文件。每个分区内部仍然有序。
2. Combiner(如果配置了)
Combiner 是 Map 端的”迷你 Reduce”——在合并的时候,把相同 Key 的 Value 先聚合一下。
1 | 合并前:<"hello",1>, <"hello",1>, <"world",1> |
Combiner 的价值: 减少了数据量,溢写文件变小,后续 Shuffle 传输也变少。
Combiner 的约束: 必须满足结合律——求和、计数可以用,平均值不行。
最终产出:
file.out:合并后的数据文件,包含所有分区file.index:索引文件,记录每个分区在file.out中的起始位置
Reduce 端就是通过索引文件快速定位到属于自己分区的数据位置,然后拉取。
五个阶段的数据流
1 | 读数据 → 跑 Map → 写内存(100MB) → 满80%溢写磁盘 → 多个溢写文件合并为一个 |
关键理解:Map 任务不是”算完就完了”,它一边算一边把数据整理好,留给 Reduce。
调优的几个关键参数
| 参数 | 默认值 | 调什么 |
|---|---|---|
mapreduce.task.io.sort.mb |
100MB | 缓冲区大小。大一点,溢写次数少一点 |
mapreduce.map.sort.spill.percent |
0.8(80%) | 溢写阈值。调大到 0.85 能多用点缓冲 |
mapreduce.task.io.sort.factor |
10 | 合并时一次合并多少个文件。调大到 20 减少合并轮次 |
mapreduce.map.memory.mb |
1024MB | Map 任务总内存。缓冲区调大了这里要跟着调 |
mapreduce.map.output.compress |
false | Map 输出压缩。开 Snappy,减少溢写和传输 |
最常见的调优操作:
- 缓冲区 100MB → 256MB(溢写次数减少 60%)
- 开启 Snappy 压缩(Shuffle 数据量减少 50%)
- 启用 Combiner(Map 输出数据量进一步减少)
几个容易踩的坑
1. 内存溢出(OOM)
缓冲区调太大,超过了 MapTask 总内存(mapreduce.map.memory.mb),就会 OOM。
解决: 调缓冲区的时候,同步调大 MapTask 内存。比如缓冲区 256MB,MapTask 内存至少给 1.5GB。
2. 数据倾斜
某个 Key 的数据特别多,它所在的分区数据量巨大,溢写和合并都慢。
解决: 自定义分区器,让大 Key 的数据能分散到多个 Reduce。或者先做一次预处理,把大 Key 打散。
3. 小文件导致 Map 任务太多
10000 个 1MB 文件 → 10000 个 Map 任务 → 调度开销巨大。
解决: 用 CombineTextInputFormat 把小文件合并成大的 InputSplit,减少 Map 任务数。