0%

ReduceTask工作机制

ReduceTask 工作机制:拉数据、排序、分组、计算——一个 Reduce 任务的完整一生

Map 阶段干完活了,输出了一堆中间文件,散落在各个节点上。Reduce 的任务就是:把这些数据拉过来,按 Key 分组,然后执行你的业务逻辑。

ReduceTask 走四个阶段:

Copy → Merge → Sort → Reduce

  • Copy:从各个 Map 节点把属于自己的数据拉过来
  • Merge:把拉过来的数据合并、排序
  • Sort:最终按 Key 排好,分好组
  • Reduce:跑你的业务逻辑,输出结果
flowchart TD  
    A[Map 输出文件] -->|Copy 阶段| B[ReduceTask 拉取分区数据]  
    B -->|Merge 阶段| C[内存/磁盘合并数据]  
    C -->|Sort 阶段| D[按 Key 分组排序]  
    D -->|Reduce 阶段| E[执行 reduce 函数输出结果]  
    E --> F[写入 HDFS 最终结果]

Copy 阶段:从 Map 节点拉数据

Reduce 任务启动后,第一件事就是去各个 Map 节点把属于自己分区的数据拉过来

怎么知道哪些数据是自己的?

Map 输出的时候已经分好区了。Reduce 任务编号是 0,就拉分区 0 的数据;编号是 1,就拉分区 1 的数据。

并行拉取:

Reduce 同时开多个线程去拉数据,默认 5 个线程并行。集群网络好的时候可以调大这个数,加快拉取速度。

数据放哪?

拉回来的数据先放内存缓冲区(占 Reduce 内存的 70%)。缓冲区满了就往磁盘溢写。

Merge 阶段:合并 + 排序

拉过来的数据是碎片化的——来自不同的 Map 节点,每个 Map 的输出在自己的节点上已经是排好序的,但多个文件之间还没合并。

Merge 阶段做的就是把多个有序文件合并成一个有序文件

内存合并:

内存缓冲区里的数据到了一定阈值(默认 66%),就触发一次内存合并——按 Key 排序,然后溢写到磁盘。

磁盘合并:

磁盘上的溢写文件多了(默认 10 个),就触发一次磁盘合并——多个小文件归并排序成一个大文件。

合并的意义: Reduce 最终只需要处理一个有序文件,而不是几十上百个碎片文件。合并轮次越少,磁盘 I/O 越少。

如果配置了 Combiner: 合并的时候做局部聚合,进一步减少数据量。

Sort 阶段:最终排序 + 分组

所有数据合并完毕后,最后一个有序文件已经按 Key 全局有序了。Sort 阶段做最后一件事:分组

分组就是把相同 Key 的所有 Value 放在一起,形成一个迭代器,传给 reduce() 函数。

1
2
3
4
5
6
排序后的数据:
<"apple",1>, <"apple",1>, <"apple",1>, <"banana",1>, <"banana",1>

分组后:
Key="apple"[1, 1, 1] → 传给 reduce()
Key="banana"[1, 1] → 传给 reduce()

默认分组规则: 使用 Key 的 compareTo() 方法,相等的 Key 分到一组。

自定义分组: 如果你想按 Key 的某一部分分组(比如按订单 ID 分组,忽略时间字段),可以写自定义 GroupingComparator

1
2
3
4
5
6
7
8
9
10
11
12
public class OrderGroupingComparator extends WritableComparator {  
protected OrderGroupingComparator() {
super(OrderKey.class, true);
}

@Override
public int compare(WritableComparable a, WritableComparable b) {
OrderKey k1 = (OrderKey) a;
OrderKey k2 = (OrderKey) b;
return k1.getOrderId().compareTo(k2.getOrderId()); // 仅按订单 ID 分组
}
}

Reduce 阶段:执行业务逻辑

数据整理好了,终于轮到 reduce() 函数上场。

1
2
3
4
5
6
7
protected void reduce(Text key, Iterable<IntWritable> values, Context context) {
int sum = 0;
for (IntWritable val : values) {
sum += val.get();
}
context.write(key, new IntWritable(sum));
}

每个 Key 调用一次 reduce(),处理完输出结果到 HDFS。

这个阶段的性能通常不是瓶颈——真正的重活都在前面的 Copy 和 Merge(网络传输 + 磁盘 I/O)。

四阶段的数据流

1
2
3
4
5
6
7
8
9
Map 输出文件(散落在各节点)

Copy:并行拉取到 Reduce 节点(网络传输)

Merge:内存合并 + 磁盘合并(排序 + 归并)

Sort:最终全局有序 + 分组

Reduce:执行业务逻辑,输出结果

关键理解:Reduce 任务 80% 的时间花在 Copy 和 Merge 上,不是在跑 reduce() 函数。

调优的关键参数

参数 默认值 调什么
拉取阶段
mapreduce.reduce.shuffle.parallelcopies 5 并行拉取线程数。网络好时调到 10-15
mapreduce.reduce.shuffle.input.buffer.percent 0.7(70%) 内存缓冲区占比。调到 0.8 减少溢写
合并阶段
mapreduce.reduce.shuffle.merge.percent 0.66(66%) 内存合并触发阈值。调大到 0.8 多用内存
mapreduce.task.io.sort.factor 10 一次合并多少个文件。调到 20-30 减少合并轮次
总内存
mapreduce.reduce.memory.mb 1024MB Reduce 总内存。缓冲区调大了这里要跟着调

几个容易踩的坑

1. 数据倾斜——一个 Key 的数据特别多

某个 Key 的数据量特别大,它所在的 Reduce 任务处理时间远超其他 Reduce。

表现: 大部分 Reduce 任务已经完成,一两个任务还在跑,进度卡在 90%+ 不动。

解决思路:

  • 给热点 Key 加随机后缀,分散到多个 Reduce
  • 用两阶段聚合:Map 端先聚合一次,Reduce 端再聚一次
  • 自定义分区器,让大 Key 的数据能分散

2. 内存溢出(OOM)

Reduce 内存不够用了。

表现: 作业失败,日志里有 OutOfMemoryError

解决:

  • 调大 mapreduce.reduce.memory.mb
  • 调小 mapreduce.reduce.shuffle.input.buffer.percent(少给缓冲区点内存)
  • 开启 Map 输出压缩,减少 Shuffle 数据量

3. Copy 阶段太慢

网络传输是瓶颈。

解决:

  • 调大并行拉取线程数(parallelcopies
  • 开启 Map 输出压缩(Snappy),减少网络传输量
  • 检查机架感知配置,优先拉取同机架的数据

4. 磁盘 I/O 过高

Merge 阶段频繁读写磁盘。

解决:

  • 调大内存缓冲区占比(input.buffer.percent
  • 调大合并因子(sort.factor),减少合并轮次
  • 开启数据压缩,减少磁盘读写量

Reduce 数量和分区数的关系

1
分区数 = ReduceTask 数

每个分区对应一个 Reduce 任务。所以:

  • Reduce 数量 = 输出文件数量
  • Reduce 太多 → 输出小文件多,NameNode 压力大
  • Reduce 太少 → 单点处理压力大,并行度不够

经验值: Reduce 数量 = 集群可用核心数 × 0.9 左右。核心 100 个,Reduce 设 90 个。

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