跳转到内容
← 返回算法
分布式算法计算机科学 · 分布式系统 · 大数据处理26 分钟阅读

MapReduce

MapReduce

2004 年,谷歌工程师杰弗里·迪恩(Jeffrey Dean)和桑贾伊·格马瓦特(Sanjay Ghemawat)发表了《MapReduce:大型集群上的简化数据处理》。这篇论文描述了 Google 内部用于处理 PB 级数据的编程模型,并在发表后迅速成为大数据时代的基础范式,催生了 Hadoop 生态系统,影响了此…

MapReduce分布式计算大数据并行计算Google

2004 年,谷歌工程师杰弗里·迪恩(Jeffrey Dean)和桑贾伊·格马瓦特(Sanjay Ghemawat)发表了《MapReduce:大型集群上的简化数据处理》。这篇论文描述了 Google 内部用于处理 PB 级数据的编程模型,并在发表后迅速成为大数据时代的基础范式,催生了 Hadoop 生态系统,影响了此后十年的数据工程实践。

背景:集群计算的工程噩梦

2000 年代初,搜索引擎需要处理数以亿计的网页、日志、索引。单台服务器无法胜任,但分布式计算面临的工程问题令人头疼:如何分发任务?如何处理机器故障?如何聚合结果?

迪恩和格马瓦特注意到,Google 内部大量分布式计算代码都在重复相同的模式:将输入分成块、并行处理每块、聚合结果。他们将这个模式抽象为 MapReduce。

核心抽象:两个函数

MapReduce 的编程模型极简:用户只需实现两个函数:

Map 函数

Map(k1,v1)[(k2,v2)]\text{Map}(k_1, v_1) \to [(k_2, v_2)]

接受一个键值对,输出一组中间键值对。这里的关键约束是:一次 Map 调用只能看见一条输入,看不到别的文档、也看不到别的 Map 的输出。正是这个约束让框架可以随意把 Map 任务分派到任何机器、随意重试。

Reduce 函数

Reduce(k2,[v2])[(k3,v3)]\text{Reduce}(k_2, [v_2]) \to [(k_3, v_3)]

接受一个中间键和该键对应的所有值的列表,输出最终结果。同样地,一次 Reduce 调用只能看见一个键,看不到别的键——所以不同的键可以并行归约。

经典例子——单词计数(Word Count)

Map(document_id, text):
  for each word w in text:
    emit(w, 1)

Reduce(word, [count1, count2, ...]): emit(word, sum(counts)) ```

将这两个函数提交给 MapReduce 框架,框架自动处理:数据分发、并行执行、故障恢复、结果聚合。

执行流程

输入文件 → 分片(Split)→ Map 阶段
                              ↓
                          Shuffle(按 key 分组,跨网络传输)
                              ↓
                          Reduce 阶段 → 输出文件
```

详细步骤

  1. 输入分片:输入文件被分成 16–64 MB 的块(Chunk)
  2. Map 任务分配:调度器将每个 Chunk 分配给一台空闲 Worker,Worker 在本地磁盘运行 Map 函数
  3. Shuffle 阶段:中间键值对按键的哈希值分配到不同的 Reduce Worker,通过网络传输(这是最昂贵的阶段)
  4. 排序:每个 Reduce Worker 对收到的中间键值对按键排序
  5. Reduce 任务执行:对每个不同的键,调用 Reduce 函数
  6. 输出:写入分布式文件系统(如 GFS/HDFS)

手算走查:三篇文档的 word count 全流程

"框架自动处理一切"这句话最容易糊过去的地方,是中间那步 shuffle 到底搬了什么。用三篇极短的文档把每一个键值对都写出来:

text
d1: the cat sat on the mat
d2: the dog sat on the log
d3: the cat and the dog
```

Map 阶段(三个 Map 任务可以在三台机器上同时跑,互不知情):

Map 任务输出的中间键值对对数
Map(d1)(the,1) (cat,1) (sat,1) (on,1) (the,1) (mat,1)6
Map(d2)(the,1) (dog,1) (sat,1) (on,1) (the,1) (log,1)6
Map(d3)(the,1) (cat,1) (and,1) (the,1) (dog,1)5

中间数据总量:17 对。输入是 17 个词,输出就是 17 对——(w, 1) 这种写法把中间数据量和输入 token 数绑死了,这一点马上会成为问题。

Shuffle 阶段。设 $R = 2$ 个 Reducer,分区函数用"首字母在字母表中的序号 mod 2\bmod\ 2"(真实实现用 hash(key) mod R,这里换成一个能手算的规则):

首字母序号mod 2\bmod\ 2去哪个 Reducer
thet = 200R0
dogd = 40R0
logl = 120R0
anda = 11R1
catc = 31R1
matm = 131R1
ono = 151R1
sats = 191R1

每个 Reducer 收到什么(收到后先按键排序,所以下面是字典序):

Reducer排序后的分组收到的对数输出
R0dog → [1,1] · log → [1] · the → [1,1,1,1,1,1]9dog 2 log 1 the 6
R1and → [1] · cat → [1,1] · mat → [1] · on → [1,1] · sat → [1,1]8and 1 cat 2 mat 1 on 2 sat 2

最终结果:the 6, cat 2, sat 2, on 2, dog 2, mat 1, log 1, and 1,加起来正好 17,对上了。

这两张表里最该注意的不是结果,是分配的不均。R0 只拿到 3 个键,却收到 9 对;R1 拿到 5 个键,只收到 8 对。更极端的是单个键:the 一个词就占了 17 对中的 6 对,35% 的中间数据挤在一个键上。而一个键只能由一个 Reducer 处理——这是 MapReduce 模型的硬约束,也是下一节所有麻烦的根源。

combiner 与热键:中间数据是真正的瓶颈

上面的流程里,(the,1) 被独立地跨网络传了 6 次。在真实语料里这个数字是几十万、几千万——英语里 the 是最高频的词,一篇文档出现几十次是常态。

combiner(合并器)的作用就是先在本地把它们折起来。它是一个在 Map 任务所在机器上、对该任务的输出先跑一遍的"局部 Reduce":

Map 任务无 combiner有 combiner省下
Map(d1)6 对(the,2) (cat,1) (sat,1) (on,1) (mat,1) = 5 对1
Map(d2)6 对(the,2) (dog,1) (sat,1) (on,1) (log,1) = 5 对1
Map(d3)5 对(the,2) (cat,1) (and,1) (dog,1) = 4 对1
合计17 对14 对3 对(18%)

三篇六个词的文档上只省 18%,看着不值当。但把文档换成一篇一千词的网页,(the,1) 从几十对折成一对——combiner 的收益随单个 Map 输入里键的重复度增长,而这个重复度在真实数据上通常很高。省下的不只是网络流量,还有 Map 端和 Reduce 端两次落盘与排序的量。

combiner 不是白用的,它有一个严格的前提:归约函数必须满足结合律与交换律summaxmincount 都满足,先局部加还是最后一起加结果一样;median(中位数)和平均值不满足——把三组数各取中位数再取中位数,和把所有数放一起取中位数是两个答案。这条前提和"能不能做两阶段聚合"是同一条。

热键(hot key)是 combiner 治不了的病。 combiner 只能压缩单个 Map 任务内部的重复,压不掉"全世界的 the 最终都要汇到同一个 Reducer"这件事。当一个键的数据量大到单台机器装不下或算不完,整个作业就被这一个 Reducer 拖住——其他几百个 Reducer 早已空闲,进度条卡在 99%。这在真实数据上极常见:电商里的爆款 SKU、社交网络里的名人账号、日志里的 NULL 或默认值。

标准解法是加盐 + 两阶段聚合

text
第一阶段:把热键随机拆开
  Map:    emit("the#" + random(0, 63), 1)     // 一个键变 64 个
  Reduce: 得到 64 个部分和,如 the#0→912, the#1→887, ...

第二阶段:把盐去掉再聚一次 Map: emit("the", 部分和) // 只有 64 条输入 Reduce: 912 + 887 + ... = 最终结果 ```

代价是多一轮完整的 MapReduce(又一次落盘 + shuffle)。该记住的是这个交换:MapReduce 用"一个键一个 Reducer"换来了极简的编程模型和无脑的容错,代价是数据分布只要不均匀,模型就会把不均匀原封不动地传导成性能问题;而修补手段(加盐)本质上是手工把并行度加回来

故障容忍:重试即可

MapReduce 最重要的工程贡献之一是其简单而有效的故障处理机制

  • Worker 故障:Master 定期 ping Worker。若 Worker 失联,其正在进行的 Map/Reduce 任务被标记为失败,重新分配给其他 Worker。Map 任务的中间结果也需要重新计算(因为存在本地磁盘上)。
  • Master 故障:Master 将状态定期写入检查点文件,故障后从最近的检查点恢复。
  • 任务慢尾(Straggler):某些 Worker 由于硬件原因极慢,Master 在任务接近完成时,为仍在进行的任务启动"备份执行(Backup Execution)",取先完成的结果。

这种设计依赖的关键假设:Map 和 Reduce 函数是纯函数(相同输入 → 相同输出),可以安全地重试

现场:5 个慢任务让整个作业慢了 44%

论文第 5.4 节做了一个只改一个开关的对照实验,是全篇最有说服力的数字。

实验用的是 sort 基准:约 1800 台机器(每台两颗 2 GHz Xeon、4 GB 内存、两块 160 GB IDE 硬盘、千兆网),排序 101010^{10} 条 100 字节记录,约 1 TB。

配置完成时间尾部发生了什么
开启备份任务891 秒正常收尾
关闭备份任务1283 秒(慢 44%)第 960 秒时只剩 5 个 reduce 任务未完成,这 5 个又拖了 300 秒

把这两行读明白:整个集群 1800 台机器、上千个任务,被 5 个慢任务多拖了 300 秒——占总时长的四分之一。这 5 台机器没有坏,只是慢:可能是磁盘有坏道导致读取从 30 MB/s 掉到 1 MB/s,可能是被其他作业抢走了 CPU 和缓存,可能是机器配置里某个错误关掉了处理器缓存。论文把这类任务叫 straggler(长尾任务)

"备份任务(backup task)"的做法简单到近乎粗暴:作业接近完成时,Master 给所有仍在进行的任务再起一份副本,谁先完成算谁的,另一份直接杀掉。它只在收尾阶段启用,所以额外开销很小——论文写明这套机制经过调优后,"通常把该操作使用的计算资源提高不超过百分之几"。

这个设计能成立,靠的还是上一节那个假设:任务是纯函数,跑两遍结果一样,重复计算不会造成副作用。反过来说,一旦 Map/Reduce 函数带副作用(往外部数据库写、给用户发邮件、消耗一次性令牌),备份任务就会造成重复执行——这是把 MapReduce 用错的经典方式之一。

该记住的是这条推理链:容错和抗长尾在 MapReduce 里是同一个机制的两个用途。因为任务可以无副作用地重跑,所以机器挂了可以重试;也正因为可以重跑,所以机器慢了可以并行再跑一份。函数式约束不是审美偏好,它是整套容错策略的前提。

现场:2004 年 8 月,Google 到底跑了多少 MapReduce

论文附了一张 2004 年 8 月的实际用量表。这些数字比任何架构图都更能说明 MapReduce 当时的地位:

指标数值
作业总数29 423
消耗的机器·天79 186
读入数据3 288 TB
产生的中间数据758 TB
写出数据193 TB
平均作业完成时间634 秒
平均每作业的 worker 机器数157
平均每作业的 map 任务数3 351
平均每作业失效的 worker 数1.2

几个可以自己算出来的推论:

  • 29 423 个作业除以 31 天,约 950 个作业/天,一小时 40 个。MapReduce 那时已经不是特殊工具,而是日常基础设施。
  • 79 186 机器·天除以 31 天,相当于约 2550 台机器全天候被 MapReduce 占着
  • 中间数据 758 TB 是输入 3288 TB 的 23%。这 758 TB 全部要写本地磁盘、跨网络传给 Reducer、再排序落盘一次——这就是上一节 combiner 想压缩的那个量,也是后来 Spark 主攻的那个环节。
  • 最后一行最值得停下来看:一个平均只跑 634 秒(约 10.5 分钟)、用 157 台机器的作业,期间平均有 1.2 台 worker 失效。折算下来单机在这十分钟里失效的概率约 1.2/1570.76%1.2/157 \approx 0.76\%

把最后这个数字放在 2004 年的语境里:Google 用的是廉价 PC 硬件,规模一上去,"机器会挂"从异常事件变成了统计上的必然。在这个故障率下,"框架自动重试"不是保守设计,而是唯一可行的设计——如果让用户代码自己处理"我的 157 台机器里平均有一台会在作业中途消失",没有人写得出正确的分布式程序。MapReduce 真正卖的不是 map 和 reduce 这两个函数,是"你可以假装机器不会坏"这个幻觉,而它靠纯函数约束把这个幻觉做成了真的。

历史影响

MapReduce 论文发表后,开源社区基于其思想实现了 Hadoop(道格·卡廷 Doug Cutting 等人,2006),成为大数据时代最重要的开源项目之一:

  • Hadoop MapReduce:MapReduce 的开源实现
  • HDFS(Hadoop Distributed File System):基于 GFS 的分布式文件系统
  • Hive:将 SQL 查询编译为 MapReduce 作业
  • Pig:数据流语言,编译为 MapReduce

雅虎、Facebook、百度等公司曾在数万台服务器规模的 Hadoop 集群上运行 MapReduce 作业。

MapReduce 的衰落与遗产

2012 年之后,MapReduce 逐渐被更高效的框架取代:

Apache Spark(2012):将中间数据保存在内存中(而非磁盘),迭代计算速度比 Hadoop MapReduce 快 10–100 倍。Spark 的 RDD(弹性分布式数据集)保留了 MapReduce 的函数式思想,但更灵活。

Google Cloud Dataflow(2014 年 Google I/O 发布)与后来的 Apache Beam:将批处理和流处理统一为同一编程模型,进一步抽象了计算框架。

MapReduce 本身已不再是首选,但其思想——将数据处理分解为无状态的转换(Map)和聚合(Reduce),由框架处理并行化和容错——渗透进了所有后继系统。

转折:Google 自己先放弃了它

外界还在大建 Hadoop 集群的时候,Google 内部已经在往外走。两个时间点值得记住。

2010 年,索引流水线换掉了 MapReduce。 Google 的 Caffeine 索引系统改用 Percolator(Peng 与 Dabek,OSDI 2010)做增量处理:不再"攒够一批网页就全量重建一次索引",而是每抓到一个页面就用分布式事务就地更新。论文报告的效果是——在处理同样文档吞吐量的前提下,中位索引延迟降低了 100 倍以上,搜索结果里文档的平均"年龄"减半。

MapReduce 为什么干不了这件事?不是实现不好,是抽象本身的性质:它的成本模型是"批量摊薄"——一次全量扫描的固定开销(分片、调度、shuffle、落盘)只有摊到几十亿个文档上才划算。要更新一个文档,就得重跑整批。批处理框架的最小工作单位是一批,这不是可以优化掉的常数。

2014 年,Google 公开说了出来。 在 Google I/O 的主题演讲上,负责技术基础设施的高级副总裁 Urs Hölzle 说:"我们已经不太用 MapReduce 了。"他给的理由不是"太慢",而是"简单作业它很好,但一旦你开始搭流水线就变得太笨重,而如今一切都是分析流水线"。同一场会议上 Google 发布了 Cloud Dataflow。

这句话点出的是 MapReduce 一个被低估的问题:它是两步的模型,而真实工作是多步的。图算法、迭代优化、多次连接的 ETL,都要拆成一串 MapReduce 作业首尾相接,每一步之间都强制落盘一次。抽象没有错,只是粒度选在了错的层级——它把"一次 map + 一次 reduce"当成了原语,而真正该当原语的是"一个由若干转换组成的有向图"。

对照:Spark 赢在哪里,以及不是赢在哪里

Spark 的 RDD(Zaharia 等,NSDI 2012)换掉的是容错方式。MapReduce 靠"把中间结果写到磁盘上"来容错——数据在磁盘上,任务挂了重读即可。RDD 换成了血缘(lineage):只记住"这个分区是从哪些父分区、经过哪些确定性变换得来的",丢了就照着谱系重算。中间结果因此不必落盘,多步流水线可以留在内存里连着跑。

两个排序纪录把差距量成了具体数字:

年份系统数据量机器数耗时
2008Hadoop MapReduce(Yahoo,Owen O'Malley)1 TB(101010^{10} 条记录)910209 秒(前纪录 297 秒)
2014Hadoop MapReduce(前纪录保持者)100 TB210072 分钟
2014Spark(Databricks,Daytona GraySort 冠军)100 TB20623 分钟

机器少 10 倍、时间短 3 倍,合起来是 30 倍的资源效率差。

但这里有个必须说清的地方:这次比赛的排序全程在磁盘(HDFS)上完成,没有使用 Spark 的内存缓存。所以那 30 倍不能简单归给"内存计算"这个流行说法——它主要来自 shuffle 实现的重写(基于 Netty 的网络传输、更少的拷贝)和调度开销的下降。

该记住的是:"Spark 比 MapReduce 快 10–100 倍"这句流传最广的话,在最严格的对照实验里其实站不住——真正的差别是中间数据要不要落盘由程序决定,而不是由框架强制,以及一个作业能不能表达成多步流水线而非多个独立作业。把性能差异归因到"内存 vs 磁盘"是一种偷懒的解释,它让人误以为买更多内存就能解决问题。

代价与争议

磁盘 I/O 瓶颈:每个 MapReduce 阶段的输出都写入磁盘,Shuffle 阶段大量网络传输+磁盘写入,对迭代算法(机器学习)极度低效。Spark 等内存计算框架因此兴起。

编程模型有限:不是所有算法都能自然地表达为 Map-Reduce 两步。图算法、迭代优化等需要多个 MapReduce 轮次,效率低,代码复杂。

投机执行的争议:备份执行(Straggler 处理)会消耗额外资源,在资源受限的集群上可能加剧竞争而非改善性能。

"大数据"的反思:MapReduce 的设计假设数据规模需要数百台机器,但 McSherry 等人("Scalability! But at what COST?", USENIX HotOS 2015)指出,许多 MapReduce/Spark 集群作业在单台笔记本电脑上用简单的单线程代码能更快完成——分布式系统的开销常被低估。

跨域连接

  • 函数式编程:combiner 可用、reduce 可任意分组、结果与执行顺序无关,这三件事的充要条件是同一条——归约操作满足结合律且有单位元,也就是构成幺半群。求和、计数、取最值满足,中位数与平均值不满足:分组取中位数再取中位数不等于整体中位数。能不能加 combiner 不是调优选项,是代数性质决定的。
  • 统计学:可分解的统计量才好分布式计算。均值要拆成和与计数两个可加量,方差还要再带上平方和才能两阶段聚合,而分位数没有有限维的可加充分统计量,只能用近似草图换取可合并性。这条界限解释了大数据系统里反复出现的现象:方差便宜,分位数贵且只有近似值。
  • 不平等经济学:作业完成时间由最慢的那个归约任务决定,该看的是分布的尾部而非平均。用平均任务时长评估集群效率,与用人均收入评估福利犯的是同一个统计错误——都把高度偏斜的分布压成一个数,掩盖了真正起作用的那一端。备份执行之所以有效,正是因为它针对尾部而非均值。
  • 集合预报与决策:集合预报天然尴尬并行,每个成员独立积分,因此能吃满上万核。但并行度不等于信息量:集合的价值来自成员之间的差异,成员再多而扰动方式单一,不确定性估计也不会变好。这与大数据处理的教训同构——加机器只能缩短时间,改变不了数据里究竟有多少可提取的结构。
  • 计算社会科学:把"用了集群"当成"做了大数据",是学科采纳新工具时的典型误判。已有的系统性对照表明,相当多的分布式作业在单机上用朴素单线程代码反而更快——分布式带来的调度、序列化与网络开销常被低估。合理的基线是先把单机做到极致,再用它判断分布式是否真的赚回了开销。

参考文献

  • Dean, J. & Ghemawat, S. MapReduce: Simplified Data Processing on Large Clusters. OSDI 2004, pp. 137–150.(第 5.3–5.4 节的 sort 基准与备份任务对照;第 6 节 2004 年 8 月用量表)
  • Ghemawat, S., Gobioff, H. & Leung, S.-T. The Google File System. SOSP 2003, pp. 29–43.
  • Peng, D. & Dabek, F. Large-scale Incremental Processing Using Distributed Transactions and Notifications. OSDI 2010.(Percolator 与 Caffeine 索引系统)
  • Zaharia, M. et al. Resilient Distributed Datasets: A Fault-Tolerant Abstraction for In-Memory Cluster Computing. NSDI 2012.(Spark 论文)
  • McSherry, F., Isard, M. & Murray, D. G. Scalability! But at what COST? HotOS XV, USENIX 2015.
  • O'Malley, O. TeraByte Sort on Apache Hadoop. Sort Benchmark 技术报告,2008 年 5 月.(910 节点、1 TB、209 秒)
  • Dean, J. & Barroso, L. A. "The Tail at Scale." Communications of the ACM 56(2), 74–80, 2013.(长尾延迟的一般性分析,备份任务思想的推广)

延伸阅读

  • Google Cloud Dataflow 与 Apache Beam 官方文档 — https://beam.apache.org/documentation/ (批流统一模型的编程接口)
  • White, T. Hadoop: The Definitive Guide. 4th ed. O'Reilly, 2015.(Hadoop MapReduce 的配置项、shuffle 调优与 combiner 用法)
  • Kleppmann, M. Designing Data-Intensive Applications. O'Reilly, 2017.(第 10 章批处理,对 MapReduce 的设计取舍有一份很好的复盘)
  • Databricks 技术博客:Spark the fastest open source engine for sorting a petabyte,2014 年 —— https://www.databricks.com/blog/2014/10/10/spark-petabyte-sort.html