0%

MapTask工作机制

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
2
3
4
5
protected void map(LongWritable key, Text value, Context context) {
for (String word : value.toString().split(" ")) {
context.write(new Text(word), new IntWritable(1));
}
}

输入一行 "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
2
合并前:<"hello",1>, <"hello",1>, <"world",1>
合并后(Combiner 生效):<"hello",2>, <"world",1>

Combiner 的价值: 减少了数据量,溢写文件变小,后续 Shuffle 传输也变少。

Combiner 的约束: 必须满足结合律——求和、计数可以用,平均值不行。

最终产出:

  • file.out:合并后的数据文件,包含所有分区
  • file.index:索引文件,记录每个分区在 file.out 中的起始位置

Reduce 端就是通过索引文件快速定位到属于自己分区的数据位置,然后拉取。

五个阶段的数据流

1
2
3
4
5
6
7
8
读数据 → 跑 Map → 写内存(100MB) → 满80%溢写磁盘 → 多个溢写文件合并为一个
↑ ↑
CPU密集 内存速度


最终文件(file.out + file.index)

Reduce 通过索引拉取自己的分区数据

关键理解: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,减少溢写和传输

最常见的调优操作:

  1. 缓冲区 100MB → 256MB(溢写次数减少 60%)
  2. 开启 Snappy 压缩(Shuffle 数据量减少 50%)
  3. 启用 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 任务数。

欢迎关注我的其它发布渠道