0%

Shuffle机制

MapReduce Shuffle 深度解析:数据从 Map 到 Reduce 的完整旅程

MapReduce 里有一句老话:“Map 和 Reduce 的代码好写,Shuffle 的调优最难。”

Shuffle 是 Map 输出到 Reduce 输入之间的所有数据处理——分区、排序、溢写、拉取、合并、分组。数据要在网络和磁盘之间来回倒腾,任何一个环节没配好,作业就跑得慢。

Shuffle 是什么?为什么它那么重要?

Shuffle 这个词翻译过来是”洗牌”,在 MapReduce 里就是 把 Map 的输出重新整理,让相同 Key 的数据聚到同一个 Reduce 去处理。

它干三件事:

  1. 分区:决定每条 Map 输出去哪个 Reduce
  2. 排序:让 Reduce 拿到的数据按 Key 有序
  3. 传输:把数据从 Map 节点搬到 Reduce 节点

整个 MapReduce 作业的时间,Shuffle 经常占 30%-50%。 所以 Shuffle 调优是 MapReduce 性能优化的核心。

Map 端 Shuffle:写数据、排序、溢写、合并

Map 任务一边生产数据,一边把数据准备好交给 Reduce。

第 1 步:写内存缓冲区

Map 输出的数据先写到内存缓冲区,默认 100MB。数据在这里待着,攒够了一批发到磁盘。

第 2 步:分区

每条 Map 输出按 Key 计算去哪个 Reduce。默认用 Key 的哈希值取模:

1
reduce分区 = hash(key) % reduce数量

如果默认分区不均匀(比如某个 Key 特别多),可以自己写 Partitioner 定制分区逻辑。

  • 默认分区器**:HashPartitioner,通过 Key 的哈希值取模分区:

    1
    2
    3
    public int getPartition(K key, V value, int numReduceTasks) {  
    return (key.hashCode() & Integer.MAX_VALUE) % numReduceTasks;
    }
  • 自定义分区:当默认分区规则不满足需求(如按业务字段分区)时,可通过继承Partitioner重写分区逻辑:

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    public class CustomPartitioner extends Partitioner<Text, IntWritable> {  
    @Override
    public int getPartition(Text key, IntWritable value, int numReduceTasks) {
    // 按 Key 前缀分区(如 "order_001" 分到分区 0,"user_001" 分到分区 1)
    if (key.toString().startsWith("order")) {
    return 0 % numReduceTasks;
    } else {
    return 1 % numReduceTasks;
    }
    }
    }

    需在 Driver 中配置:

    1
    job.setPartitionerClass(CustomPartitioner.class);

第 3 步:排序

缓冲区里的数据按 Key 排序。每个分区内部都是有序的。

1
2
3
4
5
6
7
8
9
10
11
12
public class OrderKey implements WritableComparable<OrderKey> {  
private String id;
private double price;

@Override
public int compareTo(OrderKey o) {
// 先按 ID 升序,再按价格降序
int cmp = this.id.compareTo(o.id);
if (cmp != 0) return cmp;
return Double.compare(o.price, this.price); // 价格降序
}
}

第 4 步:溢写(Spill)

缓冲区满了(默认 80%)就写到磁盘上,形成一个溢写文件。一次 Map 任务可能产生多个溢写文件。

第 5 步:合并

Map 任务结束时,把所有溢写文件按分区、按 Key 归并排序,合并成一个大文件。这样就只有一个最终输出文件,Reduce 拉取的时候方便。

如果启用了 Combiner: 在排序之后、溢写之前,先做一次局部聚合。相同 Key 的 Value 先合并一下,减少溢写的数据量。

Reduce 端 Shuffle:拉数据、合并、分组

Reduce 任务从各个 Map 节点把属于自己的数据拉回来,整理好交给 reduce() 函数。

第 1 步:拉取数据(Fetch)

Reduce 向所有 Map 任务发起 HTTP 请求,拉取属于自己分区的数据。默认同时拉 5 个 Map。

什么时候开始拉?Map 完成一定比例(默认 5%)就开始,不用等所有 Map 完成。

第 2 步:内存合并

拉回来的数据先放内存缓冲区(占 Reduce 内存的 70%)。这里边收边做归并排序。

第 3 步:溢写磁盘

内存不够了(默认到 66%),就写到磁盘上,形成溢写文件。

第 4 步:最终合并

所有数据拉完后,把内存里的和磁盘上的所有数据做多路归并排序,合并成一个有序的大文件。

第 5 步:分组

把相同 Key 的 Value 放在一起,形成一个迭代器,传给 reduce() 函数处理。

  • 默认分组器**:WritableComparator,通过 Key 的 compareTo 方法判断是否为同一组;

  • 自定义分组:当需要按 Key 的部分字段分组时(如按订单 ID 分组,忽略其他字段),可继承WritableComparator重写分组逻辑:

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    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;
    // 仅按订单 ID 分组,忽略价格字段
    return k1.getId().compareTo(k2.getId());
    }
    }

    需在 Driver 中配置:

    1
    job.setGroupingComparatorClass(OrderGroupingComparator.class);

总体流程

flowchart TD  
    subgraph Map端Shuffle
        A[Map 输出] --> B[内存缓冲区]  
        B --> C{达到溢写阈值?}  
        C -->|是| D[分区+排序+溢写磁盘]  
        C -->|否| B  
        D --> E[多个溢写文件]  
        E --> F[归并排序合并为一个文件]  
    end  
    subgraph Reduce端Shuffle  
        F --> G[Reduce 拉取数据]  
        G --> H[内存合并]  
        H --> I{内存不足?}  
        I -->|是| J[溢写磁盘]  
        I -->|否| K[最终归并排序]  
        J --> K  
        K --> L[按 Key 分组]  
    end  
    L --> M[Reduce 处理]

核心配置参数

参数 默认值 调优建议
Map端缓冲区
mapreduce.task.io.sort.mb 100MB 内存充足时调到 200-300MB,减少溢写次数
mapreduce.map.sort.spill.percent 0.8(80%) 调到 0.85-0.9 可以多用缓冲区空间
Reduce端拉取
mapreduce.reduce.shuffle.parallelcopies 5 集群网络好时调到 10-15,加速拉取
mapreduce.reduce.shuffle.input.buffer.percent 0.7(70%) 调到 0.8 减少磁盘溢写
合并
mapreduce.task.io.sort.factor 10 调到 20-30,减少合并轮次(但要内存够)
压缩
mapreduce.map.output.compress false 开启 Snappy 压缩,减少网络传输

实际经验:io.sort.mb 是最常见的优化。默认 100MB 太小,大作业溢写频繁,磁盘 I/O 暴涨。调到 256MB 或 512MB,溢写次数大幅减少。

Shuffle 常见问题与优化策略

1. 减少数据量

  • 启 Combiner:Map 端提前聚合,减少传给 Reduce 的数据
  • 开启压缩:Map 输出用 Snappy 压缩,网络传输量减少 50%-70%

2. 减少磁盘 I/O

  • 调大 io.sort.mb:减少溢写次数
  • 调大 sort.factor:减少合并轮次

3. 减少网络传输

  • 调大 Reduce 拉取并行度
  • 提高数据本地化率:让 Map 任务跑在数据所在节点

4. 处理数据倾斜

某个 Key 的数据特别多,导致一个 Reduce 任务特别慢。解决思路:

  • 给热点 Key 加随机后缀,分散到多个 Reduce
  • 用自定义分区器,让大 Key 的数据不是全去一个 Reduce
  • 先做一次局部聚合,减少热点 Key 的数据量

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