专业书籍精读 · DDIA · 第 10 章
Designing Data-Intensive Applications · Ch 10 · Martin Kleppmann · 2017
你打开淘宝看到的「猜你喜欢」、打开抖音刷到的推荐、在 Google 搜到的结果——这些背后都有一类不慌不忙、夜里悄悄跑的计算:把成千上万人一整天的行为攒起来,一次性算出「谁可能喜欢什么」「哪个词对应哪些网页」。这一章讲的就是这种攒一大批、一次性算完的活儿,叫批处理,以及它最经典的招式——MapReduce。
批处理就像洗一大桶衣服:不是来一件洗一件(那是「在线服务」,你点一下要马上有反应),而是攒够一整桶,一次性丢进洗衣机——不追求单件多快,追求「一桶下来总共省时省水」。它不在乎你等一秒还是等一小时,只在乎整批算完花多少时间、每小时能洗多少。
难点在数据太大了——大到一台机器的硬盘装不下、一个 CPU 一辈子也算不完。你得把活儿拆给上千台机器一起干。可一旦人多,新麻烦就来了:怎么把活公平地分下去?算到一半有台机器死机了怎么办?算完的半成品放哪?这些「协调一大群机器」的脏活,才是真正让人头疼的地方。
MapReduce 的核心点子,其实就是「分头数,再归堆汇总」。想象几百个人一起统计一座图书馆里每个作者有几本书:第一步(Map)——每人分一摞书,把每本书写成一张小卡片「作者名 → 1」;第二步(归堆)——把所有卡片按作者名排好队、堆到一起,同一个作者的卡片自然就挨在一块了;第三步(Reduce)——每个作者一堆卡片,数一数就是他的书数。「按名字排队、把同名的凑一堆」这一步是整台机器的心脏——正是它让分散在上千台机器上的同类数据,最终能碰头汇总。
靠这套「分头数再汇总」,工程师能用一大堆便宜机器啃下海量数据:建出搜索引擎的索引、算出推荐名单、训练模型。更妙的是它不怕出错——输入数据只读不改,某台机器算砸了,换台机器把这块重算一遍就行,连人写错了代码都能改完重跑、不留后遗症。一句诚实的代价:它天生「慢性子」,要等一整批都算完才出结果,所以只适合能等的离线活儿,实时性的场景得靠后面讲的「流处理」。
批处理 = 攒一大批数据、一次性算完,图的是吞吐(每小时啃多少)而非快。招式是 MapReduce:分头数(Map)→ 按名字排队把同类凑堆 → 汇总(Reduce),靠「排序凑堆」把上千台机器的数据碰到一起;输入只读、算错就重算,所以特别皮实。
想进到具体机制、join 策略和示意图? → 切到精读版
这一章讲批处理(batch processing):把海量数据当作不可变的输入、成批读进来、算出派生结果(搜索索引、推荐、报表),追求的是吞吐量而非响应时间。核心是把 Unix 的老智慧(小工具 + 管道)搬到上千台机器上——MapReduce:你只写 map 和 reduce 两个函数,框架靠一次分布式排序把同类数据凑到一起。一句要你记住的反直觉:你以为批处理的灵魂是「算得快」,其实是「输入不可变 + 确定性」带来的那份『算错了大不了重来』的从容。
map(从每条记录里挑出「键→值」对)和 reduce(把同一个键的所有值汇总);框架负责把它们铺到上千台机器上跑。(用户ID, 一次点击)。MapReduce 全程在搬运键值对。用户ID 对上号。本章属全书 Part III「派生数据(Derived Data)」的开篇。Part II 讲的是怎么在多台机器上正确地存与读(复制、分区、事务、共识);Part III 掉头问怎么从已有数据批量算出新数据。作者先把数据系统分成三类——在线服务、批处理、流处理;本章讲批处理(离线、有界数据),下一章(Ch11)讲流处理(无界、实时),二者共用一套「不可变数据 + 派生」的世界观。对应现实里的 Hadoop / Spark 生态、数据仓库 ETL、离线特征与模型训练。
假设你要为一个上亿用户的网站做两件事:① 每天扫一遍 10 TB 的访问日志,算出「每个网页被访问了多少次、哪些是热门」;② 把用户行为喂进推荐系统,算出「猜你喜欢」。这类活的共同点是:数据大到单机存不下、算不完,但又允许「今晚跑、明早出」——不追求实时。
难点从来不是那两行统计逻辑,而是怎么驾驭一大群机器:如何把 10 TB 切开、公平地分给上千台机器并行啃;某台机器算到一半宕了,怎么不让整批白跑;中间半成品放哪、算完的结果怎么安全地全量替换旧结果。不解决这些,你要么被单机容量卡死,要么写一堆脆弱的、一台机器崩就全盘皆输的手工脚本。这一章给的,正是一套把「协调机群、容错、数据流」这些脏活标准化的通用骨架。
别急着上分布式。作者先带你看一个 Unix 单机例子:从一份网站日志里揪出访问量最高的 5 个网址,一行命令就够:awk '{print $7}' access.log | sort | uniq -c | sort -rn | head。awk 抽出网址那一列,sort 把它们排成队(相同的挨到一起),uniq -c 数出每个出现几次,再 sort -rn 按次数倒排、head 取前几名。
这条管道藏着三条影响深远的设计原则(Doug McIlroy 提出的 Unix 哲学):①每个程序只做一件事、做好(sort 只管排序,但排得极好——数据比内存大时它自动用磁盘做归并排序);②一个程序的输出能当另一个的输入(用管道 | 串起来);③统一接口——大家都读写纯文本行流,所以任意两个工具都能拼。MapReduce 几乎是这套哲学的分布式翻版:把「统一接口=文本流、组合=管道」换成「统一接口=HDFS 文件、组合=作业链」。
MapReduce 让你只写两个纯函数,剩下的机群调度、容错、数据搬运全由框架包办:
map,你从中挑出并吐出若干 (键, 值) 对。例:读一行日志,吐出 (网址, 1)。reduce,把该键名下的一串值交给你汇总。例:(网址, [1,1,1,…]) → 求和 → (网址, 总次数)。关键要看懂:「按键排序、把同键聚拢」这一步是整个模型的心脏。Mapper 把输出按目标 Reducer 分区(通常 hash(键) mod R,R 是 Reducer 数)、在本地排好序落盘;每个 Reducer 再把属于自己的那些分片从各 Mapper 拉过来、做归并——这样散落在上千台机器上的同一个键,最终一定在同一个 Reducer 处碰头。数据规模能有多大?开源的 Hadoop 曾在数千台机器上排序 1 PB 级的数据,正是这套机制在撑。
单个 MapReduce 作业只能表达「一进一出」的简单变换,真实任务(如推荐系统)往往要把几十个作业串成工作流——前一个的输出目录当后一个的输入,由 Airflow、Oozie 之类的调度器编排。Google 内部一个复杂流程曾串起约 50 个乃至更多的 MapReduce 作业。
推荐、分析几乎都要 join:手上有一份「用户活动日志」(谁点了什么,海量),要拼上「用户资料表」(用户的年龄地区,相对小),才能算「各年龄段爱看什么」。在批处理里,你不能像在线数据库那样逐条去远程查资料表——上亿次网络往返会慢到天荒地老、还把数据库打垮。于是有了几套截然不同的 join 策略:
Reduce 端 join(sort-merge join,排序归并连接):把两份数据都用 用户ID 当键喂进 Map,框架一洗牌,同一个用户的活动记录和资料记录自然聚到同一个 Reducer,在那里一拼即可。好处是不对输入做任何假设、通用;代价是两份数据都得完整排序、洗牌,开销大。
Map 端 join:如果小表(资料表)小到能塞进内存,就不必洗牌——广播哈希连接(broadcast hash join)把整张小表加载成哈希表、复制到每个 Mapper,Mapper 一边流式扫大表、一边查表拼上,全程无 Reduce。若两份数据已按相同方式分区,则用分区哈希连接(partitioned hash join),每个 Mapper 只需装载对应的那一片小表。Map 端 join 快得多,但前提苛刻(要么一边够小,要么两边分区对齐)。
热点键的坑:若某个用户是超级大 V、活动记录奇多,负责该键的那个 Reducer 会被压垮、拖慢整批(数据倾斜)。对策是「分片 / 倾斜 join」:把热点键的记录随机拆散到多个 Reducer 分头算,另一边数据复制多份配合——这与 Ch6 分区那章处理热点的思路一脉相承。
批处理任务的产出通常是:建搜索索引(如为全站文档生成 Lucene 倒排索引文件)、批量构建键值存储(把算好的推荐结果、机器学习模型打包成只读文件,整体加载进线上服务库,如早期 LinkedIn 的 Voldemort 只读存储)。这里藏着本章最重要的思想,它直接呼应了 Unix 哲学:
一句话:把「计算逻辑」和「读写线路」彻底分开,让数据流清清爽爽——正是这份克制,换来了批处理系统标志性的皮实与可维护。
MapReduce 一统江湖多年,但有个硬伤:每个作业都把中间结果完整写回 HDFS(还带副本),下一个作业再读回来。一个几十级的工作流,就要反复「落盘—读盘」几十趟;而且后一步必须干等前一步整个跑完才能起步。这叫中间状态的物化(materialization)——稳,但慢、且浪费。
数据流引擎(dataflow engines:Spark、Tez、Flink)把整个工作流看成一张算子图(DAG)一次性提交,于是能:不必每步都排序、不必落 HDFS 而是流水线式(pipelining)把数据直接递给下一个算子(或只留在内存),算子一有输入就开工。容错怎么办?不再靠「处处存副本」,而是靠谱系重算——记住每份中间数据是怎么算出来的(Spark 的 RDD 血缘),丢了就照着重算(前提是算子确定性)。官方基准里,迭代式作业(同一份数据反复算很多轮,如机器学习、图算法)靠把数据留在内存,能比 MapReduce 快一个数量级。
还有两类专门化:迭代 / 图计算用 Pregel 模型(又称批量同步并行 BSP)——每个顶点给邻居发消息、一轮轮迭代到收敛,天生适合 PageRank、最短路径这类要反复扫图的算法。高层 API——没人愿意天天手写 map/reduce,于是 Hive、Pig、Spark SQL、Flink Table 让你写声明式的类 SQL,由查询优化器自动选 join 算法、排执行计划(这正呼应了 Ch2「声明式胜过命令式」的主题)。
批处理的选型,本质是在通用性、速度、容错代价、上手门槛之间做取舍。三张表把最常考、最实用的对比拎出来。
表 1 · MapReduce vs 数据流引擎(Spark / Tez / Flink)
| MapReduce | 数据流引擎 | |
|---|---|---|
| 中间结果 | 物化到 HDFS(带副本),反复读写盘 | 流水线 / 留内存,尽量不落盘 |
| 执行单位 | 一堆独立作业,后步干等前步整个跑完 | 整张算子图一次提交,算子有输入即开工 |
| 容错 | 靠中间结果的副本,宕了从落盘处续 | 靠谱系重算(需算子确定性),或按需 checkpoint |
| 速度 | 基线;多级工作流被落盘拖慢 | 迭代 / 多级作业常快 一个数量级 |
| 皮实度 | 极稳——超大作业中途挂也能续 | 重算链太长时代价高,超大作业需 checkpoint 兜底 |
| 适用 | 超大、超长、容错优先的批作业 | 迭代式(ML / 图)、交互式、追求速度 |
表 2 · 三种 join 策略怎么选
| 策略 | 怎么做 | 前提 | 代价 / 何时选 |
|---|---|---|---|
| Reduce 端 sort-merge | 两边都以 join 键当键喂进 Map,洗牌后同键聚到一个 Reducer 再拼 | 无前提,最通用 | 两边全排序 + 洗牌,最慢;两边都很大、无法预分区时选它 |
| Map 端 广播哈希 | 小表整张读进每个 Mapper 内存哈希表,流式扫大表查表即拼 | 一边小到能进内存 | 无洗牌、最快;一大一小的 join 首选 |
| Map 端 分区哈希 | 两边按相同方式分区,Mapper 只装对应那片小表 | 两边已按 join 键同样分区 | 省内存又免洗牌;数据管线里两表本就对齐分区时选它 |
表 3 · 三类数据系统:批处理放在哪
| 在线服务 online | 批处理 batch | 流处理 stream(Ch11) | |
|---|---|---|---|
| 输入 | 请求随到随处理 | 有界:固定的一大批 | 无界:永不结束的事件流 |
| 首要指标 | 响应时间(p99) | 吞吐量 | 吞吐 + 端到端延迟 |
| 结果新鲜度 | 实时 | 「今晚跑、明早出」 | 秒级 / 亚秒级 |
| 典型系统 | Web 服务、OLTP 数据库 | Hadoop MapReduce、Spark | Kafka + Flink / Storm |
批处理是整个大数据生态的地基。Google 的 MapReduce + GFS(2003–2004 两篇论文)直接催生了开源 Hadoop(HDFS + MapReduce),让「用一堆廉价机器啃 PB 级数据」从 Google 的独门秘技变成人人可用的基础设施。此后 Hive(SQL-on-Hadoop)、Spark(数据流引擎)、Flink 层层演进;今天的数据仓库 ETL、离线特征工程、推荐 / 广告的离线训练管线,骨子里都是这一章的思想。面试里,「MapReduce 的 shuffle 在做什么」「reduce-side 和 map-side join 怎么选」「为什么 Spark 比 MapReduce 快」几乎是数据岗必考题——答案全在本章。
map/reduce 两个函数、把并行化 / 容错 / 数据分发全交给框架,就能在上千台廉价机器上处理 TB–PB 级数据——奠定了整个批处理范式。J. Dean & S. Ghemawat《MapReduce》, OSDI 2004 ↗① 一句话:批处理 = 把有界、不可变的海量数据成批读入、算出派生结果,图吞吐不图快。
② 三类系统:在线服务(盯响应时间)、批处理(盯吞吐)、流处理(无界实时)——本章讲批处理,下一章讲流。
③ 思想根在 Unix 哲学:小工具做好一件事 + 管道组合 + 统一接口;MapReduce 是它的分布式翻版。
④ MapReduce = Map(逐条吐键值对)→ 按键排序洗牌(发动机)→ Reduce(同键汇总);同一个键必在同一 Reducer 碰头。
⑤ Join 三招:reduce-side sort-merge(通用最慢)、map-side 广播哈希(一大一小、最快)、分区哈希(两边同分区);热点键靠分片 join 化解。
⑥ 灵魂在不可变输入 + 确定性算子 + 无副作用 → 失败重算、可回滚、人为容错,皆源于此。
⑦ 超越 MapReduce:数据流引擎(Spark/Tez/Flink)把中间态从物化落盘改为流水线 / 内存 + 谱系重算,迭代作业快一个数量级;另有 Pregel(图)与高层声明式 API(Hive/Spark SQL)。
⑧ 落地:Google MapReduce/GFS → Hadoop → Hive/Spark/Flink,撑起整个数据仓库与离线管线;面试高频考点。