MapReduce简介
MapReduce深度解析:分而治之,把大数据拆成小任务
处理 TB 级的数据,单机跑不动。MapReduce 的思路很简单:把大任务拆成小任务,分到多台机器上并行跑,最后把结果汇总。
这个思路来自函数式编程的 Map 和 Reduce——Map 负责”拆”和”算”,Reduce 负责”合”。这套模型成了 Hadoop 时代的计算标准。
MapReduce 在干什么?分三阶段
一个 MapReduce 作业,数据走三个阶段:
输入 → Map(拆分+计算)→ Shuffle(排序+传输)→ Reduce(汇总)→ 输出
Map 阶段: 把输入数据拆成多个分片,每个分片交给一个 Map 任务处理,输出一堆键值对 (key, value)。
Shuffle 阶段: 把 Map 输出的键值对按 key 分组、排序,然后传给 Reduce。
Reduce 阶段: 每个 Reduce 任务拿到一组相同 key 的所有 value,聚合计算,输出最终结果。
完整流程:一条数据怎么走
第 1 步:输入分片(Input Split)
输入文件被切成多个逻辑分片(InputSplit),每个分片对应一个 Map 任务。默认按块大小切(128MB),所以一个 1GB 文件大约产生 8 个 Map 任务。
第 2 步:Map 处理
每个 Map 任务读取自己的分片,逐行处理,调用 map() 函数输出键值对。
比如 WordCount,输入 "hello world" → 输出 <"hello", 1> 和 <"world", 1>。
第 3 步:Shuffle——最复杂的环节
Shuffle 是数据从 Map 到 Reduce 的传输过程,包含三个关键动作:
- 分区(Partition):决定每个键值对去哪个 Reduce。默认按 key 的哈希值取模。
- 排序(Sort):Map 端输出按 key 排序,Reduce 端拉取后再次合并排序。
- 传输(Transfer):Reduce 从所有 Map 任务拉取属于自己分区的数据。
第 4 步:Reduce 处理
每个 Reduce 任务收到一组 (key, values),调用 reduce() 函数做聚合,输出结果到 HDFS。
数据处理流程

Shuffle 机制:MapReduce 的核心
Shuffle 是 MapReduce 中最复杂且关键的环节,负责 数据的分区、传输和排序,直接影响性能。
Map 端 Shuffle
flowchart TD
A[Map输出] --> B[内存缓冲区]
B --> C{达到80%阈值?}
C -->|是| D[溢写磁盘]
C -->|否| B
D --> E[分区+排序]
E --> F[多个溢写文件]
F --> G[合并为最终文件]
- 内存缓冲区:默认 100MB(可通过
io.sort.mb配置); - 分区规则:默认使用
HashPartitioner,根据 key 的哈希值决定 Reduce 分区; - 排序优化:溢写文件在合并时会进行 归并排序,确保最终文件按 key 有序。
Reduce 端 Shuffle
flowchart TD
A[Map输出文件] --> B[Reduce拉取数据]
B --> C[内存合并]
C --> D{内存不足?}
D -->|是| E[溢写磁盘]
D -->|否| F[最终合并]
E --> F
F --> G[按key分组]
G --> H[调用reduce]
- 数据拉取:ReduceTask 通过 HTTP 并行拉取多个 MapTask 的输出;
- 内存合并:拉取的数据先在内存中合并,超过阈值则溢写到磁盘;
- 最终合并:所有数据拉取完成后,磁盘文件和内存数据再次合并并排序。
一个完整的例子:WordCount

下面是经典 WordCount 案例的完整实现:
提供了五个可编程组件,分别是InputFormat、Mapper、Partitioner、Reducer和OutputFormat
注意新API是在org.apache.hadoop.mapreduce包下
依赖
<dependencies>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-common</artifactId>
<version>3.3.0</version>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-client</artifactId>
<version>3.3.0</version>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-hdfs</artifactId>
<version>3.3.0</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-assembly-plugin</artifactId>
<configuration>
<descriptorRefs>
<descriptorRef>jar-with-dependencies</descriptorRef>
</descriptorRefs>
<archive>
<manifest>
<addClasspath>true</addClasspath>
<mainClass>com.zhanghe.study.mapreduce.wordcount.WordCountDriver</mainClass> <!-- 你的主类名 -->
</manifest>
</archive>
</configuration>
<executions>
<execution>
<id>make-assembly</id>
<phase>package</phase>
<goals>
<goal>single</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
代码展示
/**
* Map阶段
* @author zh
* @date 2021/3/27 22:21
*/
// 泛型中的含义
//1.输入数据key的类型
//2.输入数据value的类型
//3.输出数据key的类型
//4.输出数据value的类型
public class WordCountMapper extends Mapper<LongWritable, Text, Text,IntWritable> {
// 输出的key
Text text = new Text();
// 输出的value
IntWritable intWritable = new IntWritable(1);
// 每个kv对都会进入该方法
// 用来处理业务逻辑
protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
// 读取一行的数据
String line = value.toString();
// 使用空格切割单词
String[] words = line.split(" ");
// 遍历单词
for(String word : words){
text.set(word);
context.write(text,intWritable);
}
}
}
/**
* Reduce阶段
* @author zh
* @date 2021/3/27 22:35
*/
// 泛型的含义
//1.map阶段key的类型
//2.map阶段value的类型
//3.输出数据key的类型
//4.输出数据value的类型
public class WordCountReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
IntWritable value = new IntWritable();
// 用来处理map阶段的结果
protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException {
int sum = 0;
// 累加求和
for(IntWritable intWritable : values){
sum+=intWritable.get();
}
value.set(sum);
context.write(key,value);
}
}
public class WordCountDriver {
public static void main(String[] args) throws IOException, ClassNotFoundException, InterruptedException {
Configuration conf = new Configuration();
// 获取job
Job job = Job.getInstance(conf);
// 设置jar的存储位置
job.setJarByClass(WordCountDriver.class);
// 设置map和reduce
job.setMapperClass(WordCountMapper.class);
job.setReducerClass(WordCountReducer.class);
// 设置map的输出类型
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]));
// job提交,并等待完成
job.waitForCompletion(true);
}
}
Shuffle 调优:MapReduce 的性能瓶颈
Shuffle 是 MapReduce 最耗时的环节——数据要排序、要网络传输、要落盘。调优主要围绕这几件事:
1. 增大 Map 端内存缓冲区
默认 100MB,达到 80% 阈值时溢写磁盘。增大缓冲区能减少溢写次数:
<property>
<name>mapreduce.task.io.sort.mb</name>
<value>256</value>
</property>
2. 开启 Map 端输出压缩
Map 输出压缩后,网络传输量减少:
conf.set("mapreduce.map.output.compress", "true");
conf.set("mapreduce.map.output.compress.codec",
"org.apache.hadoop.io.compress.SnappyCodec");
3. 启用 Combiner——Map 端局部聚合
Combiner 是 Map 端的”小 Reduce”,在数据发往 Reduce 之前先做一次局部聚合,减少 Shuffle 数据量。
job.setCombinerClass(WordCountReducer.class);
Combiner 必须满足结合律——求和、计数、最大值可以用,平均值不行。
4. 合理设置 Reduce 数量
Reduce 数 = 集群可用核心数 × 0.9 左右。太少则单点瓶颈,太多则小文件过多。