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。

数据处理流程

MapReduce流程

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的流程分析

下面是经典 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对都会进入该方法
    // 用来处理业务逻辑
    @Override
    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阶段的结果
    @Override
    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 左右。太少则单点瓶颈,太多则小文件过多。