深入 Hadoop 13 - MapReduce 编程模型与分而治之
第一、第二阶段把 HDFS 和 YARN 讲完了。从这一篇起进入 MapReduce——Hadoop 三大子系统的最后一个,也是历史最久的一个。
MapReduce 常被介绍成"Hadoop 的计算引擎"。这个描述对应了功能但漏了 MapReduce 的本质。准确的说法是:MapReduce 是把"分而治之"这个通用算法思路工程化为可运行框架的产物——把大输入切成 InputSplit 给 Map 并行处理,用 Partitioner 把 Map 输出按 key 分发到 Reduce,Reduce 聚合后写 HDFS。三阶段切分(Map / Shuffle / Reduce)和三个用户钩子(Mapper / Combiner / Reducer)是 MapReduce 全部抽象的核心。
本篇只抓一个问题:MapReduce 为什么把计算切成 Map → Shuffle → Reduce 三阶段、Combiner 在什么场景下能省 Shuffle 流量、Partitioner 与 Reduce 个数的关系。
三个阶段的本质
一次 MapReduce 作业的数据流可以压成一句话:
1 | |
三个阶段各自承担不同的工程目标:
Map 阶段把输入切片转成 key-value 对。Mapper 接收 InputSplit 里的一条记录(默认是一行文本),输出 0 到 N 个 key-value 对。Map 阶段是并行执行的——每个 InputSplit 由一个独立的 Map Task 处理,多个 Map Task 在不同节点上同时跑。Map 的并行度等于 InputSplit 数量,由输入文件大小和 block 大小决定(默认 128MB 一个 split)。
Shuffle 阶段把 Map 输出按 key 重新分发到 Reduce。Map 的输出在本地磁盘上按 Partition 排序,Reduce Task 通过 HTTP 从所有 Map Task 拉取自己负责的 Partition。Shuffle 是 MapReduce 数据交换最昂贵的阶段——网络传输 + 磁盘 I/O 双重开销。Shuffle 的细节下一篇展开。
Reduce 阶段把同一个 key 的所有 value 聚合成最终输出。Reducer 接收 (key, [v1, v2, v3, …]),按用户定义的逻辑处理,输出 0 到 N 个 key-value 对到 HDFS。Reduce 的并行度等于 Reduce Task 数量,由用户在作业配置里指定(mapreduce.job.reduces)。
三阶段切分对应"分而治之"的工程化:Map 切输入(并行处理),Shuffle 重组(按 key 分组),Reduce 聚合(产出结果)。这种切分让 MapReduce 可以处理任意"按键聚合"的工作负载——word count、group by、join、排序等。
Mapper 与 InputSplit 的关系
Mapper 是用户实现的 Map 逻辑。InputSplit 是 MapReduce 框架对输入的切分单元。两者关系是"一个 InputSplit 对应一个 Mapper 实例"。
1 | |
这段 Mapper 代码处理一行文本,把每个单词输出为 (word, 1)。MapReduce 框架对每条输入记录调用一次 map 方法。
InputSplit 的切分由 InputFormat 决定。默认的 TextInputFormat 按 HDFS block 边界切——一个 128MB block 一个 InputSplit。所以一个 1GB 文件切成 8 个 InputSplit,对应 8 个 Map Task 并行处理。
注意 InputSplit 不完全等于 HDFS block。InputFormat 可以按逻辑边界切分(例如按行边界、按 CSV 记录边界),不严格按物理 block 边界。但默认 TextInputFormat 让 InputSplit 与 HDFS block 一一对应,实现"数据本地化"——Map Task 在持有该 block 的 DataNode 上执行,避免网络读。
输入处理的关键 hook:
1 | |
InputFormat 是 MapReduce 处理不同输入格式(文本、Sequence File、Parquet、ORC、自定义格式)的扩展点。
Reducer 与 Partitioner 的关系
Reducer 是用户实现的 Reduce 逻辑。Partitioner 决定 Map 输出的每个 key-value 对发给哪个 Reduce Task。
1 | |
这段 Reducer 代码接收 (word, [1, 1, 1, …]),把所有 1 求和,输出 (word, total)。
Partitioner 决定 Map 输出的去向。默认的 Partitioner 是 HashPartitioner:
1 | |
这个公式保证同一个 key 的所有 key-value 对发到同一个 Reduce Task。numReduceTasks 由 mapreduce.job.reduces 配置决定。
如果 numReduceTasks 配置为 0,作业没有 Reduce 阶段,Map 输出直接写 HDFS。这种"Map-only"作业适合不需要聚合的场景——例如数据清洗、格式转换。
Partitioner 的设计决定数据分布的均匀性。如果某个 key 的数据量远大于其他 key(数据倾斜),所有这些数据集中到一个 Reduce Task,导致该 Reduce 慢,整个作业延迟被拖。生产环境常见的"热点 key"问题就是 Partitioner 选型不当导致的。
应对数据倾斜的常见 Partitioner:
1 | |
TotalOrderPartitioner 是 Hadoop 自带的特殊 Partitioner,用于 TeraSort 等需要全局排序的作业。它先对输入采样,建立 key 范围分割点,让 Reduce Task 0 处理最小 key、Reduce Task 1 处理次小、依此类推。
Combiner:Map 端的 mini Reduce
Combiner 是可选的 Map 端本地聚合。它在 Map Task 内部对 Map 输出做一次本地 Reduce,减少 Shuffle 传输的数据量。
1 | |
Combiner 让 Shuffle 数据量从 O(原始记录数) 降到 O(distinct key 数)。在 word count 这种"key 集中"的场景,Shuffle 流量可以减少 10-100 倍。
Combiner 的限制是它必须满足结合律和交换律——Combiner 的输出类型必须等于 Mapper 的输出类型(因为 Combiner 在 Reducer 之前跑,Reducer 接收的输入必须和 Mapper 输出类型一致)。所以求和、最大值、最小值这类操作可以用 Combiner,求平均值不能用(平均值不可结合)。
实际作业配置 Combiner:
1 | |
word count 这种场景 Combiner 和 Reducer 用同一个类——都是把 (key, [v1, v2, …]) 求和。Combiner 在 Map 端先求一次和,Reducer 在 Reduce 端再求总和。
Combiner 不是强制运行的——MapReduce 框架根据内存缓冲状态决定是否运行 Combiner。如果 Map 输出在 Spill 之前内存缓冲没攒够,Combiner 可能不运行。所以 Combiner 是优化不是正确性保证——作业必须假设 Combiner 可能不运行仍然得到正确结果。
Map 与 Reduce 的并行度
Map 并行度 = InputSplit 数量,由输入文件大小决定。1TB 输入默认切成 8192 个 Map Task。这个并行度用户通常不需要调整——InputFormat 决定。
Reduce 并行度 = mapreduce.job.reduces 配置值,用户决定。Reduce 数过少会让单 Reduce 处理过多数据,延迟长;过多会让每个 Reduce 处理很少数据,调度开销大。生产经验:每个 Reduce 处理 1-5 GB 数据比较合适。1TB 输入建议 Reduce 数 200-1000。
1 | |
特殊场景:
Map-only 作业:mapreduce.job.reduces = 0,没有 Reduce 阶段,Map 输出直接写 HDFS。适合数据转换、过滤。
强聚合作业:Reduce 数少(10-50),让 Reduce Task 数量小但每个处理大量数据。适合数据量小但聚合逻辑复杂的场景。
排序作业:Reduce 数等于 TeraSort 配置(通常按集群规模选)。TotalOrderPartitioner 保证 Reduce 之间全序。
MapReduce 的并行度是手动调优的关键维度。Spark 的 Adaptive Query Execution(AQE)可以运行时动态调整 Shuffle 分区数,MapReduce 没有这种自适应能力——Reduce 数在作业提交时固定,运行中不能改。
实验:观察 MapReduce 作业结构
提交一个 word count 作业(Hadoop 自带示例):
1 | |
作业运行时打开 ResourceManager Web UI,找到这个作业,点进 ApplicationMaster Web UI(MapReduce 自带):
1 | |
观察几个细节:
Map Task 数 = 输入文件 block 数。如果 /input 下文件总大小 1GB(切成 8 个 128MB block),Map Task 数 = 8。
Reduce Task 数默认 = 1。这个默认值通常需要用户主动配置提高(mapreduce.job.reduces)。
Reduce Task 在 Map Task 完成一定比例后才启动(默认 mapreduce.job.reduce.slowstart.completedmaps = 0.05,即 5% Map 完成后启动 Reduce)。这是 MapReduce 的优化——让 Reduce Task 提前启动,一边拉 Shuffle 数据一边等其他 Map 完成。
每个 Task 点进去可以看到 counters——Map 输入记录数、Map 输出记录数、Spill 记录数、Shuffle 字节数、Reduce 输入记录数、Reduce 输出记录数。这些 counter 是性能调优的关键观测点。
命令行层面,mapred job -status <jobId> 展示作业整体进度:
1 | |
输出 Map/Reduce 完成度、counter、最近 task 失败等。
模式提炼
MapReduce 编程模型体现的设计模式:
1 | |
这个模式不只是 MapReduce。Spark 的 Map/Reduce 算子(map / reduceByKey)是同一抽象,差异在 Spark 内存优先、MapReduce 落盘。Flink 的 keyBy / window 算子也用类似分区模型。Presto / Hive 的 SQL 执行最终也落到 Map-Reduce-like 的 stage 切分。
数据库领域的对应物是并行查询执行。Greenplum / Vertica 的并行查询把 SQL 拆成多个 slice 在不同节点上执行,slice 之间通过网络重分发数据,本质与 MapReduce 相同。差异在数据库的"slice"由优化器自动决定,MapReduce 的"Map/Reduce"由用户代码决定。
这种"分而治之 + 数据流编排"是分布式数据处理的标准范式。后续的 Spark、Flink、Beam 都在这个范式内做改进(内存优先 / 流处理 / 统一编程模型),没有跳出"Map → Shuffle → Reduce"的三段式。
工程迁移表
| MapReduce 概念 | Spark | Flink | Hive | 并行数据库 |
|---|---|---|---|---|
| InputSplit | Partition(RDD) | Source split | InputFormat | table partition |
| Mapper | map 算子 | MapFunction | Mapper operator | query operator |
| Partitioner | partitionBy | keyBy | reducer | distribution |
| Combiner | map-side aggregation | combiner | map-side aggregation | partial aggregation |
| Reducer | reduceByKey / groupBy | reduce / window | reducer | final aggregation |
| Shuffle(落盘) | Shuffle(落盘 / 内存) | Network + state | Shuffle | redistribute |
| OutputFormat | OutputFormat / sink | sink | OutputFormat | table write |
注意 Spark 这一列。Spark 的 map / reduceByKey 算子与 MapReduce 的 Mapper / Reducer 抽象几乎对应,差异在 Spark 的 Shuffle 可以走内存(窄依赖不需要落盘)、MapReduce 强制落盘。这是为什么 Spark 比 MapReduce 快——不是算法差异,是中间数据放置策略不同。
常见误解
误解一:“MapReduce 必须有 Reduce 阶段”。Map-only 作业是合法的——mapreduce.job.reduces = 0 时 Map 输出直接写 HDFS。适合数据转换、过滤、格式转换等不需要聚合的场景。
误解二:“Combiner 总是让作业变快”。Combiner 只在 Map 输出 key 集中时显著减少 Shuffle 流量。如果 Map 输出的 key 全部不同(例如 UUID),Combiner 没什么可合并的,反而增加本地计算开销。Combiner 适合"求和 / 最大值 / 最小值"等可结合操作。
误解三:“Reduce 数越多并行度越高越好”。Reduce 数过多让每个 Reduce 处理很少数据,调度开销(启动 Reduce Task、Shuffle 握手、OutputFormat 初始化)超过实际计算时间。生产经验:每个 Reduce 处理 1-5 GB 数据最划算。
误解四:“MapReduce 不适合流处理”。MapReduce 设计目标就是批处理,不是流处理。流处理场景应该用 Flink / Spark Streaming / Storm。MapReduce 处理流数据靠"微批"——把流切成小批次,每个批次跑一个 MapReduce。这是早期 Spark Streaming 的思路,但延迟在分钟级,不如原生流处理。
误解五:“MapReduce 已经被 Spark 完全取代”。MapReduce 在大规模 ETL 场景仍然广泛使用——Hive on MapReduce、Pig on MapReduce 在很多公司仍是主力。MapReduce 的优势是稳定性(运行了 20 年)和可预测性(性能模型简单)。Spark 在交互式分析、机器学习场景取代了 MapReduce,但批处理 ETL 仍然是 MapReduce 的领地。
练习
-
提交一个 word count 作业(
hadoop jar hadoop-mapreduce-examples-*.jar wordcount /input /output),观察 Map Task 数 = 输入文件 block 数。修改输入为 1 个大文件 vs 多个小文件,对比 Map Task 数变化。 -
配置
mapreduce.job.reduces = 10,重新提交 word count,观察 Reduce Task 数变化。然后实现一个自定义 Partitioner,让所有以元音字母开头的单词发到 Reduce-0,其他发到 Reduce-1,验证 Partitioner 是否生效。 -
在
apache/hadoop源码里找到Mapper.java、Reducer.java、Partitioner.java、InputFormat.java(hadoop-mapreduce-client-core 模块),观察四个核心抽象的接口定义。 -
思考题:如果让 MapReduce 支持任意 DAG(不是固定 Map-Reduce 两阶段),需要怎么改造?这是 Tez 和 Spark 做的改造,思考它们的做法与 MapReduce 的差别。
系列导航
| 序号 | 主题 | 状态 |
|---|---|---|
| 00-06 | HDFS 存储层(00 导读 / 01-06 HDFS 各主题) | 第一阶段(已完成) |
| 07-12 | YARN 资源管理层(07 架构 / 08 资源 / 09 调度 / 10 生命周期 / 11 HA / 12 Timeline) | 第二阶段(已完成) |
| 13 | MapReduce 编程模型:分而治之的工程化表达 | 本篇 |
| 14 | Shuffle 全流程:Map 端、Reduce 端与磁盘 I/O 的代价 | 下一篇 |
| 15 | MRv2 on YARN:ApplicationMaster 与 Task Attempt | |
| 16-19 | HA、安全与 Common 基础设施(RPC / 序列化 / 安全 / 监控) | 第四阶段,待开始 |
| 20-22 | 演进、生态与对比(Hadoop 生态 / 对象存储 + K8s / 设计遗产) | 第五阶段,待开始 |
参考资料
- Jeffrey Dean, Sanjay Ghemawat. MapReduce: Simplified Data Processing on Large Clusters. OSDI 2004.(MapReduce 设计原型论文,三阶段切分和钩子式扩展的源头)
- Apache Hadoop 官方文档:MapReduce Tutorial. https://hadoop.apache.org/docs/current/hadoop-mapreduce-client/hadoop-mapreduce-client-core/MapReduceTutorial.html
- Apache Hadoop 源码:
Mapper.java、Reducer.java、Partitioner.java、InputFormat.java. https://github.com/apache/hadoop - Tom White. Hadoop: The Definitive Guide. O’Reilly, 4th Edition 2015. Chapter 6 “MapReduce 工作机制”、Chapter 7 “MapReduce 类型与格式” 详细描述了三阶段切分和 InputFormat/OutputFormat/Partitioner 等扩展点。
- Jimmy Lin, Chris Dyer. Data-Intensive Text Processing with MapReduce. Morgan & Claypool, 2010.(MapReduce 算法设计教材,讲解 Combiner、in-mapper combining 等优化模式)
