MapReduce Shuffle 深度解析:数据从 Map 到 Reduce 的完整旅程
MapReduce 里有一句老话:“Map 和 Reduce 的代码好写,Shuffle 的调优最难。”
Shuffle 是 Map 输出到 Reduce 输入之间的所有数据处理——分区、排序、溢写、拉取、合并、分组。数据要在网络和磁盘之间来回倒腾,任何一个环节没配好,作业就跑得慢。
Shuffle 是什么?为什么它那么重要?
Shuffle 这个词翻译过来是”洗牌”,在 MapReduce 里就是 把 Map 的输出重新整理,让相同 Key 的数据聚到同一个 Reduce 去处理。
它干三件事:
- 分区:决定每条 Map 输出去哪个 Reduce
- 排序:让 Reduce 拿到的数据按 Key 有序
- 传输:把数据从 Map 节点搬到 Reduce 节点
整个 MapReduce 作业的时间,Shuffle 经常占 30%-50%。 所以 Shuffle 调优是 MapReduce 性能优化的核心。
Map 端 Shuffle:写数据、排序、溢写、合并
Map 任务一边生产数据,一边把数据准备好交给 Reduce。
第 1 步:写内存缓冲区
Map 输出的数据先写到内存缓冲区,默认 100MB。数据在这里待着,攒够了一批发到磁盘。
第 2 步:分区
每条 Map 输出按 Key 计算去哪个 Reduce。默认用 Key 的哈希值取模:
reduce分区 = hash(key) % reduce数量
如果默认分区不均匀(比如某个 Key 特别多),可以自己写 Partitioner 定制分区逻辑。
默认分区器**:HashPartitioner,通过 Key 的哈希值取模分区:
public int getPartition(K key, V value, int numReduceTasks) {
return (key.hashCode() & Integer.MAX_VALUE) % numReduceTasks;
}
自定义分区:当默认分区规则不满足需求(如按业务字段分区)时,可通过继承Partitioner重写分区逻辑:
public class CustomPartitioner extends Partitioner<Text, IntWritable> {
@Override
public int getPartition(Text key, IntWritable value, int numReduceTasks) {
if (key.toString().startsWith("order")) {
return 0 % numReduceTasks;
} else {
return 1 % numReduceTasks;
}
}
}
需在 Driver 中配置:
job.setPartitionerClass(CustomPartitioner.class)
第 3 步:排序
缓冲区里的数据按 Key 排序。每个分区内部都是有序的。
public class OrderKey implements WritableComparable<OrderKey> {
private String id;
private double price;
@Override
public int compareTo(OrderKey o) {
int cmp = this.id.compareTo(o.id);
if (cmp != 0) return cmp;
return Double.compare(o.price, this.price);
}
}