MapReduce 编程模型
概述
MapReduce 是 Hadoop 的编程模型:把"大规模数据处理"抽象成 Map(映射)与 Reduce(归约)两个阶段,框架自动处理并行化、调度与容错。虽然现代离线计算已被 Spark 取代,但 MapReduce 的思想——分而治之、移动计算、中间数据排序——仍是理解整个大数据计算体系的基石。
一、核心思想
1.1 一句话总结
把对海量数据的处理,拆成可以并行执行的 Map 任务和汇总结果的 Reduce 任务,中间由框架完成排序与分组。
1.2 五阶段总览
输入数据
│
▼
1. Split(输入分片) 文件切分为分片,每个分片一个 Map 任务
│
▼
2. Map(映射) 逐条处理,输出 <key, value> 中间结果
│
▼
3. Shuffle(洗牌) 分区、排序、合并、归并——最核心也最耗时
│
▼
4. Reduce(归约) 对每组 key 聚合计算,输出结果
│
▼
5. 输出 写入 HDFS(或指定输出格式)二、五个阶段的细节
2.1 Split:输入分片
- 输入文件按 块大小(默认 128MB)切分为逻辑分片
- 每个分片对应一个 Map 任务(map task),并行执行
- 分片优先分配给数据所在节点(数据本地性),减少网络传输
文件 1GB = 8 个 128MB 分片 → 8 个 Map 任务并行处理2.2 Map:映射阶段
Map 函数逐条处理输入,输出键值对:
java
// 示例:WordCount 的 Map
public class TokenizerMapper extends Mapper<Object, Text, Text, IntWritable> {
private final static IntWritable one = new IntWritable(1);
private Text word = new Text();
@Override
protected void map(Object key, Text value, Context context)
throws IOException, InterruptedException {
// 按空格分词,输出 <单词, 1>
StringTokenizer itr = new StringTokenizer(value.toString());
while (itr.hasMoreTokens()) {
word.set(itr.nextToken());
context.write(word, one);
}
}
}Map 输出的中间结果先写内存缓冲区(默认 100MB),达到阈值后溢写(spill)到本地磁盘。
2.3 Shuffle:洗牌阶段(核心)
Shuffle 横跨 Map 与 Reduce 两侧,是 MapReduce 性能的关键:
Map 端 Reduce 端
┌────────────────┐ ┌────────────────┐
│ Map 输出缓冲 │ │ 拉取(Fetch) │
│ → 分区 │ │ → 合并(Merge)│
│ → 排序 │ │ → 归并排序 │
│ → 溢写(合并) │ │ → 分组 │
└────────────────┘ └────────────────┘| 环节 | 说明 |
|---|---|
| 分区(Partition) | 按 key 哈希决定发往哪个 Reduce |
| 排序(Sort) | 每个分区内按 key 排序 |
| 合并(Combiner) | Map 端局部聚合,减少传输量 |
| 拉取(Fetch) | Reduce 从各 Map 节点拉取自己分区的数据 |
| 归并(Merge) | 多路归并排序,形成有序分组 |
2.4 Reduce:归约阶段
Reduce 函数对同一 key 的所有 value 做聚合:
java
// WordCount 的 Reduce
public class IntSumReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
@Override
protected void reduce(Text key, Iterable<IntWritable> values, Context context)
throws IOException, InterruptedException {
int sum = 0;
for (IntWritable val : values) {
sum += val.get(); // 对同一单词的计数求和
}
context.write(key, new IntWritable(sum));
}
}三、Combiner 与 Partitioner
3.1 Combiner:Map 端预聚合
Combiner 是在 Map 端执行的局部 Reduce,显著减少网络传输:
无 Combiner:每个 Map 输出大量 <单词,1>
有 Combiner:Map 端先汇总("hello" 出现 1000 次 → 输出 1 条)
传输量从 1000 条降为 1 条java
job.setCombinerClass(IntSumReducer.class); // 复用 Reducer 类即可注意:Combiner 必须在 Reduce 里"可交换可结合"才能安全复用(求和、计数、最大值可以;平均值不可以)。
3.2 Partitioner:决定数据去往哪个 Reduce
默认按 key 哈希 % Reduce 数 分区,可自定义控制数据分布:
java
// 自定义分区:按首字母分桶
public class FirstLetterPartitioner extends Partitioner<Text, IntWritable> {
@Override
public int getPartition(Text key, IntWritable value, int numPartitions) {
char c = key.toString().charAt(0);
return (c - 'a') % numPartitions;
}
}java
job.setPartitionerClass(FirstLetterPartitioner.class);
job.setNumReduceTasks(26); // 每个字母一个 Reduce3.3 数据倾斜的根源
如果 Partitioner 分布不均(如大量 key 落在同一分区),某些 Reduce 任务负载极高,其余空闲——这就是数据倾斜的常见来源。缓解手段:
- 增加随机前缀打散 key(需要最终结果可接受)
- 采用两阶段聚合(先局部聚合再全局)
- 对热点 key 单独处理
四、WordCount 完整作业
java
public class WordCount {
public static void main(String[] args) throws Exception {
Configuration conf = new Configuration();
Job job = Job.getInstance(conf, "word count");
job.setJarByClass(WordCount.class);
// 设置 Map / Reduce 类
job.setMapperClass(TokenizerMapper.class);
job.setReducerClass(IntSumReducer.class);
// 设置中间输出类型
job.setMapOutputKeyClass(Text.class);
job.setMapOutputValueClass(IntWritable.class);
// 设置最终输出类型
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(IntWritable.class);
// 输入输出路径
FileInputFormat.addInputPath(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));
// 提交并等待完成
System.exit(job.waitForCompletion(true) ? 0 : 1);
}
}bash
# 打包并提交
hadoop jar wordcount.jar WordCount /input /output五、关键调优参数
5.1 Map 端
| 参数 | 默认值 | 说明 |
|---|---|---|
mapreduce.task.io.sort.mb | 100 | Map 输出缓冲大小 |
mapreduce.map.sort.spill.percent | 0.8 | 溢写阈值 |
mapreduce.map.combine.minspills | 3 | 溢写次数达到后启用 Combiner |
mapreduce.map.memory.mb | 1024 | Map 容器内存 |
5.2 Reduce 端
| 参数 | 默认值 | 说明 |
|---|---|---|
mapreduce.reduce.shuffle.parallelcopies | 5 | 并行拉取线程数 |
mapreduce.reduce.input.buffer.percent | 0.0 | 内存缓冲占堆比例 |
mapreduce.reduce.memory.mb | 1024 | Reduce 容器内存 |
5.3 通用
| 参数 | 说明 |
|---|---|
mapreduce.job.reduces | Reduce 数量(默认 1,需手动设置) |
mapreduce.map.memory.mb / reduce.memory.mb | 容器内存,需与 YARN 配置匹配 |
mapreduce.map.cpu.vcores | Map 容器 CPU |
六、常见问题排查
| 现象 | 原因 | 处理 |
|---|---|---|
| Reduce 永远等不满 | Reduce 数过少或 Map 输出倾斜 | 调大 Reduce 数、均衡分区 |
| 大量小文件 | 输入分片过碎 | 合并小文件或用 CombineFileInputFormat |
| OOM | 容器内存不足 | 调大 map/reduce 内存参数 |
| 任务本地性差 | 分片与节点不匹配 | 检查机架感知与数据分布 |
| 中间数据爆炸 | 无 Combiner 或 Map 输出过多 | 加 Combiner、精简输出 |
七、MapReduce 与 Spark 的关系
| 维度 | MapReduce | Spark |
|---|---|---|
| 中间结果 | 落盘 | 内存优先 |
| 迭代计算 | 每轮重写磁盘 | 内存复用,快数十倍 |
| 延迟 | 分钟级 | 秒级 |
| 编程模型 | Map/Reduce 两阶段 | RDD/DataFrame 更灵活 |
MapReduce 的意义在于定义了分布式计算的范式:Shuffle、分区、排序、容错、本地性这些概念,Spark/Flink 中依然存在。理解了 MapReduce 的瓶颈(磁盘 IO),也就理解了 Spark 为何成功。
八、小结
| 阶段 | 关键点 |
|---|---|
| Split | 分片决定并行度与本地性 |
| Map | 逐条转换,输出键值对 |
| Shuffle | 分区 + 排序 + 合并 + 拉取 + 归并 |
| Reduce | 同 key 聚合输出 |
| 优化 | Combiner 减传输、Partitioner 均负载、内存参数匹配 |
参考链接: