1.1 什么是 MapReduce 本节摘要:MapReduce 是 Google 于 2004 年公开发表、随后由 Hadoop 开源落地的分布式计算模型。它把"分而治之"抽象成 map 与 reduce 两个用户函数,把并行、调度、容错、数据分发全部收编为框架职责。本节从定义出发,还原它诞生时的工程命题,拆开它的两函数契约,最后用一个可以在本机复算的迷你例子把"模型"落到可执行代码上。 一句话定义与它的三层含义 MapReduce 是一种用于大规模数据集并行计算的编程模型:用户只需编写 map 和 reduce 两个函数,框架负责把作业切分、分发到成百上千台机器上执行,并处理失败重试与机器间的数据传输。这句话可以再剥成三层: 第一层,它是模型而非产品。
本节摘要:MapReduce 是 Google 于 2004 年公开发表、随后由 Hadoop 开源落地的分布式计算模型。它把"分而治之"抽象成 map 与 reduce 两个用户函数,把并行、调度、容错、数据分发全部收编为框架职责。本节从定义出发,还原它诞生时的工程命题,拆开它的两函数契约,最后用一个可以在本机复算的迷你例子把"模型"落到可执行代码上。
MapReduce 是一种用于大规模数据集并行计算的编程模型:用户只需编写 map 和 reduce 两个函数,框架负责把作业切分、分发到成百上千台机器上执行,并处理失败重试与机器间的数据传输。这句话可以再剥成三层:
第一层,它是模型而非产品。map 与 reduce 借自 Lisp 等函数式语言的列表操作:map 对每条记录独立变换,reduce 把同键的值归并。模型不绑定任何实现——Google 的内部 C++ 实现、Hadoop 的 Java 实现、乃至今天在教材里用单机模拟的教学实现,跑的都是同一套语义。
第二层,它是面向普通工程师的分布式系统封装。2004 年之前,写一个并行的日志统计程序,工程师要自己处理数据划分、任务分发、机器故障、慢机器拖尾、结果汇总,代码里业务逻辑常常只占三成。MapReduce 的革命性在于:这七成"分布式脏活"被收进框架,业务代码收敛到两个函数签名里。
第三层,它的运行单位是作业。一个 MapReduce 作业把输入数据切分成若干分片,每个分片由一个 Map 任务处理;Map 的中间输出按键分区后,被汇聚到若干 Reduce 任务做最终归并。作业、任务、分片、分区这四个词,构成了后面所有章节的词汇表。
| 层次 | 关注点 | 谁负责 |
|---|---|---|
| 模型 | map 与 reduce 的输入输出契约 | 用户 |
| 框架 | 切分、调度、Shuffle、容错、计数 | 实现 |
| 作业 | 一次计算从提交到产出的生命周期 | 两者协作 |
理解一个技术"为什么出现在那个时间点",比背它的定义更有用。2003 到 2004 年,Google 面对的现实是:网页规模冲到数十亿页,为了给搜索引擎建倒排索引,要把全网页面解析、提词、统计词频、再按词聚合一轮。这个计算量单机无解,只能横跨上千台廉价机器。
廉价机器意味着每月必有若干台损坏,于是任何并行程序都必须回答三个问题:某台机器中途挂了,它已完成的部分怎么办?某台机器特别慢,整个作业要不要等它?数据分布在各台机器的磁盘上,怎么把它们搬到一起做聚合?在 MapReduce 之前,每个并行程序各自 answering,答案还常常是错的。
Google 的观察很朴素:内部大量批处理计算——倒排索引、排序、词频统计、链接图分析——都能套进同一个壳子:"对每条记录独立映射,再按键归并"。既然模式重复出现,就值得把这个壳子做成基础设施。2004 年的论文报告了彼时的成绩:单次索引构建从数千行手写并行代码,缩到几百行 MapReduce 代码;上千节点的作业自动挺过机器故障。次年 Hadoop 把这套思想开源成 Java 实现,MapReduce 从一家公司的内部武器变成行业公共底座。
论文前:并行逻辑 + 容错 + 数据搬运 + 业务 约 3000 行 论文后:map 函数 + reduce 函数 + 作业配置 约 300 行 省下的 90% 代码,全部长进了框架
把模型压缩到最小,就是两个函数签名(以 Hadoop Java API 为例):
// map:输入一对,输出零到多对中间键值 protected void map(KEYIN key, VALUEIN value, Context context) { // 对 value 做切分、过滤、变换 context.write(中间键, 中间值); } // reduce:输入一个键和它的全部值,输出零到多对最终结果 protected void reduce(KEYIN key, Iterable<VALUEIN> values, Context context) { // 遍历 values 做求和、最值、拼接 context.write(key, 最终值); }
契约的关键约束只有两条:map 对每条输入必须互相独立(不依赖处理顺序、不依赖其他记录),reduce 只通过键拿到它该看的数据(不同键之间不通信)。只要业务能接受这两条,框架就能把 map 随便复制到任何节点、把 reduce 的输入提前排好序送上门。反过来,一旦你在 map 里偷偷记住了"上一条记录",或是在 reduce 里想访问另一个键的结果,范式就破了——这也是第 7 章讨论其局限的伏笔。
在两函数之间,框架自动补齐了一段用户不写却真实存在的逻辑:map 的输出按中间键分区,同键数据跨节点聚到同一个 reduce,且送达时按键有序。这段"数据迁徙"就是 Shuffle,第三幕的主角。
图中实线是数据流向,虚线标注职责归属:用户只拥有两个函数,其余整条流水线都姓"框架"。
光说不练假把式。下面用原生 Java MapReduce 写一个"按部门汇总工资"——它是 WordCount 的近亲,但带了一点业务味。输入是若干行"姓名,部门,工资":
赵敏,研发,12000 钱磊,研发,15000 孙蕙,市场,9000 李昂,市场,11000 周婷,研发,13000
map 阶段把每行拆开,以部门为键、工资为值原样发出;reduce 阶段把同部门的工资累加:
public class DeptSalarySum { public static class DeptMapper extends Mapper<LongWritable, Text, Text, LongWritable> { private Text dept = new Text(); private LongWritable salary = new LongWritable(); @Override protected void map(LongWritable offset, Text line, Context context) { String[] cols = line.toString().split(","); if (cols.length != 3) return; // 脏行直接跳过 dept.set(cols[1]); salary.set(Long.parseLong(cols[2])); context.write(dept, salary); } } public static class SumReducer extends Reducer<Text, LongWritable, Text, LongWritable> { private LongWritable result = new LongWritable(); @Override protected void reduce(Text dept, Iterable<LongWritable> salaries, Context context) throws IOException, InterruptedException { long sum = 0; for (LongWritable s : salaries) sum += s.get(); result.set(sum); context.write(dept, result); } } public static void main(String[] args) throws Exception { Job job = Job.getInstance(); job.setJarByClass(DeptSalarySum.class); job.setMapperClass(DeptMapper.class); job.setReducerClass(SumReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(LongWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }
提交运行:
hadoop fs -put emp.csv /input/dept/ hadoop jar dept-sum.jar DeptSalarySum /input/dept /output/dept hadoop fs -cat /output/dept/part-r-00000
输出(五条记录、两个部门,可手工验算):
市场 20000 研发 40000
值得在脑内推演一遍框架在这五行数据上做了什么:假设输入被切成两个分片,前两行与后三行各归一个 Map 任务。第一个 Map 发出(研发,12000)(研发,15000),第二个发出(市场,9000)(市场,11000)(研发,13000);框架按"部门"分区,市场组进一个 Reduce,研发组进另一个 Reduce,且各自按键排好序送达。于是 reduce 函数拿到的 values 是有序列表,求和逻辑水到渠成。数据怎么切、怎么搬、怎么排序,用户代码里一个字没提——这就是模型的意义。
初学者常把 MapReduce 与三个近邻混为一谈,这里划清边界:
| 概念 | 关系 | 一句话区分 |
|---|---|---|
| HDFS | 共生不隶属 | MapReduce 是计算,HDFS 是存储,前者读后者的块 |
| Hadoop | 实现与其上的生态 | MapReduce 是 Hadoop 最早的计算引擎之一 |
| Spark | 后来者 | 同一范式思想(宽依赖处发生 Shuffle),引擎换成内存优先 |
还有一对易混词:数据并行与任务并行。MapReduce 是数据并行——所有 Map 任务跑的是同一段代码,只是各管一份数据;它不是把一个算法的不同步骤分给不同机器。认清这一点,就能解释为什么矩阵乘法这种需要全局通信的算法在 MapReduce 里写起来格外别扭,而日志统计类任务写起来行云流水。
下一节把镜头拉远:这个两函数模型作为"范式",究竟强在哪、贵在哪。