Lecture 5: MapReduce 与 Hadoop —— 大规模分布式数据处理(MapReduce and Hadoop)
第四部分:通信原语与分布式数据结构
这一部分回答:多个进程之间如何高效、可靠、可扩展地交换信息与状态? 从「重新执行代替恢复」的 MapReduce,到概率性的 Gossip,到故障检测与成员管理, 到结构化与非结构化的 P2P,再到把这一切组装成真实键值存储的 Cassandra。
Lecture 5: MapReduce 与 Hadoop —— 大规模分布式数据处理(MapReduce and Hadoop)
讲义对应:CS 425 FA2026 Lecture 3(
L3.FA26.txt,57 页);补充:CS 425 FA2025 Lecture 4(L4.FA25.txt,39 页) 教材对应:Coulouris 5th Ed. Ch. 21(Google 案例研究:GFS / Chubby / Bigtable / MapReduce 这一组系统) 阅读材料:cs425_data/papers/mapreduce-osdi04.pdf—— Dean & Ghemawat, MapReduce: Simplified Data Processing on Large Clusters, OSDI 2004(重点 §3.1 执行流程、§3.2 master 数据结构、§3.3 容错、§3.4 局部性、§3.5 任务粒度、§3.6 备份任务、§4.1–4.6 精化、§5 性能)
5.1 概述
本讲要解决的核心问题是:当数据量大到单机无论如何都处理不完时,如何让一个只会写顺序代码的程序员也能用上几千台机器,而且完全不用操心并行化、数据分发和机器故障? MapReduce 给出的答案是:把计算限制在一个极窄的编程模型里——用户只写 map 与 reduce 两个纯函数,其余四件事(并行化 map、把中间数据搬到 reduce、并行化 reduce、为四个阶段实现存储)全部由框架负责。
这套设计在分布式系统史上之所以重要,是因为它第一次把“用重新执行(re-execution)代替状态恢复”这一容错思想做成了工业级范式:worker 挂了不要紧,把它做过的 map 任务在别的机器上重跑一遍即可;不需要日志、不需要检查点、不需要副本协调。代价是要求用户函数确定性且可重复执行。这条思想随后被 Spark 的 RDD lineage(详见 Lecture 24)与流处理系统的重放语义(详见 Lecture 23)直接继承,因此本讲是整门课”容错”主线的一个枢纽。
5.2 核心概念与分布式机制图解
5.2.1 动机:把”并行化 + 分布 + 容错”这对程序员隐藏(Motivation)
定义与目的:MapReduce 是一个批处理(batch processing)编程模型与运行时系统:用户提供
map/reduce两个函数,系统在成百上千台商用 PC 上自动并行执行,并容忍机器故障。直观解释(”它是什么?”):设想你要统计一座巨型图书馆里每个词出现了多少次。一个人从早读到晚要读三百年。你会怎么做?找来一千个读者,把书架按位置切成一千份,每人负责一份,各人先在自己手上把”词 → 次数”归并成一张小表(这一步就是 combiner),然后把所有写着同一个词的小纸条投进同一个信箱(这一步就是 partition + shuffle),最后每个信箱由一个人汇总(这一步就是 reduce)。你作为馆长只需要规定两件事:“每个人怎么把一页书变成一张小表”(map)和“每个信箱怎么把小纸条汇总成一个数”(reduce)。至于谁读哪个书架、纸条怎么送、有人生病了怎么办——那是图书馆管理员(框架)的事。
这个类比里最关键的一点是:划分解法的自由度被限制住了,正因为限制得足够窄,系统才能全自动地做并行化与容错。
机制图解:MapReduce 的分工是一条清晰的分界线。
┌──────────────────── 用户视角(Externally)────────────────────┐
│ 1. 写一个很短的 Map 程序 + 一个很短的 Reduce 程序 │
│ 2. 指定并行度:M(map 任务数)、R(reduce 任务数)与分区函数 │
│ 3. 提交 job, 等结果 │
│ 4. 几乎不需要懂并行/分布式编程 │
└───────────────────────────┬────────────────────────────────────┘
│ 这条线就是抽象边界
┌───────────────────────────┴────────────────────────────────────┐
│ 1. 并行化 Map —— 每个 map 任务互相独立, 天然可并行 │
│ 2. Map → Reduce 的数据搬运 (shuffle) │
│ 3. 并行化 Reduce —— 每个 reduce 任务互相独立 │
│ 4. 为四个阶段实现存储: │
│ map 输入 ← 分布式文件系统 (GFS / HDFS) │
│ map 输出 → map 节点的**本地磁盘** │
│ reduce 输入 ← 多个远程 map 节点的本地磁盘 │
│ reduce 输出 → 分布式文件系统, 每个 reduce 任务一个文件 │
│ ★ 并保证 map 阶段与 reduce 阶段之间的屏障 (barrier) │
└────────────────────────────────────────────────────────────────┘
- 关键假设与系统模型:目标环境是”几百到几千台通过交换式以太网连接的商用 PC”(论文 §3)。具体地:(1) 机器是双核 x86 + 2–4GB 内存;(2) 网络是 100Mbps 或 1Gbps,但总对分带宽远低于此;(3) 机器数量巨大,机器故障是常态而非异常(Google 2004 年 8 月的统计:平均每个 job 有 1.2 台 worker 死亡);(4) 存储是直连的廉价 IDE 磁盘,由 GFS/HDFS 用副本提供可靠性。
5.2.2 从函数式原语到分布式执行模型(map / reduce)
定义与目的:
map与reduce两个词直接借自 Lisp 等函数式语言。目的是在”表达力”与”可自动并行化”之间取一个可证明安全的平衡点。直观解释:”map” 就是逐条独立加工:给一串记录,对每条记录独立地做同一件事,记录之间互不影响,因此任意划分都能并行。”reduce” 就是把同一个 key 下的一批记录捏成一个:它关心的是”集合”而不是”个体”,因此必须先把同一 key 的记录凑到一起才能算。
Lisp: (map square '(1 2 3 4)) => (1 4 9 16) 每条记录独立处理
(reduce + '(1 4 9 16)) => 30
即 (+ 16 (+ 9 (+ 4 1))) 成批处理
MapReduce 的类型签名(论文 §2.2):
map (k1, v1) -> list(k2, v2)
reduce (k2, list(v2)) -> list(k3, v3) [论文写作 list(v2), 含义相同]
其中: 输入键值域 (k1,v1) 与中间键值域 (k2,v2) **可以不同**;
中间键值域 (k2,v2) 与输出键值域 (k3,v3) **必须相同**
—— 因为 reduce 的输出可以再喂给下一个 MapReduce 作业
- 机制图解:WordCount(词频统计)是最小可用的示例。输入的一”条记录”是
<文件名, 文件内容>(Hadoop 的 text 输入格式下实际是<行偏移量, 行内容>):
输入 <filename, file text> map 输出 (中间 KV)
┌──────────────────────────────────┐ ┌────────────────────────┐
│ Welcome Everyone │ │ (Welcome, 1) │
│ Hello Everyone │ ───▶ │ (Everyone, 1) │
└──────────────────────────────────┘ │ (Hello, 1) │
│ (Everyone, 1) │
MAP TASK 1 └────────────────────────┘
每条记录调用一次 map
中间 KV 按 key 分给某个 reduce reduce 输出
┌────────────────────────┐ ┌──────────────────────┐ ┌──────────────┐
│ (Welcome, 1) │ │ reduce#1: │ │ (Everyone,2) │
│ (Everyone, 1) │ ───▶ │ Everyone -> 2 │──▶│ (Hello, 1) │
│ (Hello, 1) │ │ (Welcome, 1) │ │ (Welcome, 1) │
│ (Everyone, 1) │ │ reduce#2: │ └──────────────┘
└────────────────────────┘ │ Hello -> 1 │
└──────────────────────┘
- 关键假设与系统模型:
map必须对单条记录无副作用且只依赖该记录;reduce只依赖(key, 该 key 的全部 value)。二者都定义为确定性函数(见 5.3.4 对非确定性的讨论)。框架不保证map/reduce在同一台机器上执行,也不保证各 key 的处理顺序——这一点直接决定了后面所有正确性论证的形态。
5.2.3 角色与系统架构(Master / Worker / Client / DFS / Combiner / Partitioner)
- 定义与目的:一次 MapReduce 执行由六个角色组成,它们共同把”两个纯函数”变成”一次分布式计算”。
| 角色 | Hadoop 1.x 名称 | Hadoop 2.x(YARN)名称 | 职责 |
|---|---|---|---|
| Client | Client | Client | 提交 job、指定 M/R 与分区函数、等待结果;master 挂掉时由它决定是否重试 |
| Master | JobTracker | ApplicationMaster(每 job 一个) | 保存全部任务状态(idle / in-progress / completed)、分配任务、ping worker、记录中间文件位置表、触发 backup task |
| Worker | TaskTracker | NodeManager + Container | 实际执行 map / reduce 任务,把中间结果写本地磁盘并上报位置 |
| DFS | HDFS | HDFS | 存放输入 split 与最终输出;提供 3 副本容错;为局部性调度提供 block 位置信息 |
| Combiner | Combiner | Combiner | 可选。在 map 端本地对同一 partition 内同一 key 的 value 做预聚合,削减 shuffle 数据量 |
| Partitioner | Partitioner | Partitioner | 决定中间 key 去哪个 reduce:默认 hash(key) mod R |
直观解释:Master 像工头,手上只有一块白板:每个任务现在什么状态、归谁做、做完了的 map 把中间文件放在哪。它从不搬运中间数据,只传播”元数据”(文件位置与大小),真正的数据是 reduce worker 直接向 map worker 拉的。这样做的好处是 master 几乎不会成为带宽瓶颈。
机制图解:同一批机器同时扮演”存储节点”与”计算节点”,这是局部性优化的物理基础。
┌──────────────────────── 一个集群 (同一批机器) ────────────────────────┐
│ │
│ Client ──① 提交 job──▶ Master (Resource Manager + per-job AM) │
│ │ │
│ ② 分配 map/reduce │ ③ 回报中间文件位置表 (只是元数据!) │
│ ┌───────────────┼────────────────┐ │
│ ▼ ▼ ▼ │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ Node A │ │ Node B │ │ Node C │ │
│ │ DataNode A │ │ DataNode B │ │ DataNode C │ ← 存 DFS 块 │
│ │ NodeManager A│ │ NodeManager B│ │ NodeManager C│ │
│ │ │ │ │ │ │ │
│ │ map#0 ✅ │ │ map#1 ✅ │ │ map#2 ✅ │ ← 跑 map 任务 │
│ │ reduce#0 ✅ │ │ reduce#1 ✅ │ │ │ ← 跑 reduce │
│ │ [本地磁盘] │ │ [本地磁盘] │ │ [本地磁盘] │ ← 中间数据在这 │
│ └──────┬───────┘ └──────┬───────┘ └──────┬───────┘ │
│ │ ④ shuffle: reduce 端 HTTP fetch │ │
│ └────────────────┼────────────────┘ │
│ ▼ │
│ ⑤ 输出文件写回 DFS (3 副本) │
└───────────────────────────────────────────────────────────────────────┘
关键: "本地写, 远程读"(local write, remote read)。
1/N 的数据不需要走网络 —— 这一份就是 map 端本地读到的输入。
- 关键假设与系统模型:
- 单 master:论文的实现里 master 是单点,不存在 master 选举,故障即整个 job 失败。
- master 保存 O(M·R) 状态:对每个已完成的 map 任务,master 记录它产生的 R 个中间文件区域的位置与大小;每对 (map, reduce) 约占 1 字节内存。
- 中间数据只写本地磁盘:不复制、不落 DFS。这是”已完成 map 任务在机器故障后必须重做”的直接原因,也是整个容错设计的枢纽假设。
5.2.4 完整执行流程(Execution Overview)
定义与目的:论文 §3.1 把一次 MapReduce 执行拆成 7 步;理解这 7 步的顺序与数据落在哪块盘上,就理解了 MapReduce 的全部工程权衡。
输入切分:输入文件被切成 M 个 split(论文 16MB–64MB 一个,可由用户调参),每个 split 对应一个 map 任务。GFS/HDFS 的 block 大小是 64MB(论文口径)/ 128MB(Hadoop 2.x 默认),split 通常与 block 对齐,但必须能在任意位置切断而不是切断一行——text 输入格式保证 split 边界落在行边界上。reduce 侧把中间 key 空间切成 R 份,分区函数与 R 都由用户指定。
机制图解:下图是完整的 MapReduce 数据流,”在哪台机器上执行”已直接标注在图中。
HDFS / GFS: file.dat => block0 | block1 | block2 | block3 | ...
block 默认 128MB (论文与讲义口径 64MB); 每块 3 副本: 2 个同机架 + 1 个异机架
split0 split1 split2 split3 <- M = split 数 = map 任务数
| | | |
+---------------+ +---------------+ +---------------+ +---------------+
| W0: map#0 | | W1: map#1 | | W2: map#2 | | W3: map#3 | <- 优先 data-local 调度
| MAP(k1,v1) | | MAP(k1,v1) | | MAP(k1,v1) | | MAP(k1,v1) |
| ->list(k2,v2) | | ->list(k2,v2) | | ->list(k2,v2) | | ->list(k2,v2) |
| 写入内存缓冲 | | 写入内存缓冲 | | 写入内存缓冲 | | 写入内存缓冲 |
+---------------+ +---------------+ +---------------+ +---------------+
| | | |
+-----------------+-----------------+-----------------+
|
| (1) 缓冲达到阈值 -> spill: 按 PARTITION(k2) = hash(k2) mod R 分桶
| 桶内按 k2 快排 (所以 map 端输出有序) -> 可选 COMBINE 预聚合
| 每个 map 任务写出 R 个分区文件到**本地磁盘**
v
+==============================================+
| (2) SHUFFLE: reduce 端按 Master 给出的位置表 |
| 主动 HTTP 拉取(fetch) 属于自己那个分区 |
+==============================================+
|
v
reduce#0 reduce#1 reduce#2 <- R 个 reduce 任务
| | |
+-----------------+ +-----------------+ +-----------------+
| 拉取 p=0 分片 | | 拉取 p=1 分片 | | 拉取 p=2 分片 |
| 归并 M 份 | | 归并 M 份 | | 归并 M 份 |
| 按 k2 排序 | | 按 k2 排序 | | 按 k2 排序 |
| 按 key 分组 | | 按 key 分组 | | 按 key 分组 |
| REDUCE(k2,v*) | | REDUCE(k2,v*) | | REDUCE(k2,v*) |
| 写 tmp 文件 | | 写 tmp 文件 | | 写 tmp 文件 |
+-----------------+ +-----------------+ +-----------------+
+---------------------+---------------------+
|
v (3) 原子 rename: tmp_00000 -> part-r-00000 等最终输出文件
(要么整个 job 的输出可见, 要么不可见)
+==============================================+
| HDFS 最终输出: part-r-00000 ... 共 R 个文件 |
+==============================================+
- 逐步解说(对应上图编号):
- 切分并启动:用户程序里的 MapReduce 库把输入文件切成 M 份(每份 16–64MB),然后在集群上启动这份程序的许多副本。
- 选出 master:其中一份副本是特殊的——master,其余的 worker 由 master 分配工作。共有 M 个 map 任务与 R 个 reduce 任务。master 挑选空闲 worker,分别派发 map 或 reduce 任务。
- 执行 map:被派到 map 任务的 worker 读取对应的输入 split,解析出键值对,逐条交给用户
Map函数;产生的中间 KV 先缓存在内存里。 - 溢写本地磁盘:缓冲区的键值对周期性地按分区函数写成 R 个区域落到本地磁盘。这些区域的位置被回报给 master,master 负责把它们转发给 reduce worker。
- shuffle + 排序:reduce worker 收到 master 关于位置的通告后,用 RPC 从各 map worker 的本地磁盘读取缓冲数据。读完全部中间数据后,按中间 key 排序,使同一个 key 的所有出现聚在一起。如果中间数据大到装不进内存,就使用外部排序(external sort)。
- 执行 reduce:reduce worker 遍历已排序的中间数据,每遇到一个不同的中间 key,就把该 key 与对应的 value 集合交给用户
Reduce函数;Reduce的输出追加到该 reduce 分区的最终输出文件。 - 唤醒用户程序:所有 map 与 reduce 任务都完成后,master 唤醒用户程序,MapReduce 调用返回。
为什么必须保证 map 与 reduce 之间的屏障(barrier):因为 reduce 要处理的是某一个 key 的全部 value,而”全部”意味着必须等所有 map 都产出完毕才知道是否齐全。如果允许 reduce 提前开始,那么一个尚未完成的 map 之后产出的 value 就可能被漏掉,reduce 得到的是部分结果——这是静默的错误结果,比崩溃更危险。这就是讲义反复强调”Ensure that no Reduce starts before all Maps are finished”的原因。
- 关键假设与系统模型:输入数据在整个 job 期间不变(否则重执行的 map 可能读到不同的数据,串行等价性被破坏)。map 输出不复制。reduce 输出直接落全局 DFS,因此天然冗余。
5.2.5 一个具体例子:WordCount 端到端走一遍
定义与目的:把抽象流程落到真实字节上,是理解 shuffle 与 combiner 收益的最快方式。
例子设定:输入两个 split,R = 2,为便于手工验算,本处用玩具分区函数
P(word) = len(word) mod 2(真实系统用稳定哈希,见 5.2.7)。map输出(word, 1),combine/reduce都是求和。
输入 split0 = [ "Welcome Everyone", "Hello Everyone" ]
输入 split1 = [ "Hello Hadoop", "Welcome Hadoop" ]
【map 阶段】map 任务逐条记录调用 MAP(k1,v1) = for each word w: emit(w, 1)
map#0 原始输出 : (Welcome,1) (Everyone,1) (Hello,1) (Everyone,1) 共 4 条
map#1 原始输出 : (Hello,1) (Hadoop,1) (Welcome,1) (Hadoop,1) 共 4 条
【本机按 partition 分桶】P(Welcome)=7%2=1, P(Everyone)=8%2=0,
P(Hello)=5%2=1, P(Hadoop)=6%2=0
map#0: p=0 桶 [(Everyone,1),(Everyone,1)] p=1 桶 [(Welcome,1),(Hello,1)]
map#1: p=0 桶 [(Hadoop,1),(Hadoop,1)] p=1 桶 [(Hello,1),(Welcome,1)]
【combiner 本地预聚合】桶内按 key 排序后求和
map#0: p=0 => (Everyone,2) p=1 => (Welcome,1) (Hello,1)
map#1: p=0 => (Hadoop,2) p=1 => (Hello,1) (Welcome,1)
★ shuffle 传输量从 8 条降到 4 条 (本例 2x; 真实数据上常见 10x 以上, 见 5.4.2)
【shuffle + 归并排序】reduce#0 拉取所有 map 的 p=0 分片, reduce#1 拉取所有 p=1 分片
reduce#0 的输入 (归并排序后) : (Everyone,2) (Hadoop,2)
reduce#1 的输入 (归并排序后) : (Hello,1) (Hello,1) (Welcome,1) (Welcome,1)
【reduce 阶段】按 key 分组后调用 REDUCE(k2, list(v2)) = sum(list)
reduce#0: Everyone -> 2 ; Hadoop -> 2
reduce#1: Hello -> 2 ; Welcome -> 2
【最终输出】写入 DFS, 共 R = 2 个文件
part-r-00000 : Everyone 2 \n Hadoop 2
part-r-00001 : Hello 2 \n Welcome 2
【正确性校验】串行统计: Welcome 2, Everyone 2, Hello 2, Hadoop 2 —— 完全一致 ✅
- 关键假设与系统模型:这个例子同时也说明了为什么 reduce 的输出是 R 个文件而不是 1 个:每个 reduce 任务独立写自己的文件,框架不做合并。用户通常直接把这 R 个文件当作下一个 MapReduce 作业的输入,或者交给另一个能处理多文件输入的应用。
5.2.6 Shuffle 与 Sort:为什么它不是”可选优化”
定义与目的:Shuffle 指”把 map 的输出按分区搬到对应 reduce 节点”这一整个阶段(分区 + 排序 + 网络传输 + 归并);Sort 是它的组成部分。
直观解释:reduce 的语义是
REDUCE(k2, list(v2))——它要的是”这个 key 的全部 value 的一张完整清单“。而 map 任务是多台机器并行跑的,同一个 key 的 value 零散地分布在 M 台机器各自的本地磁盘上。要把它们凑成”一张清单”,就必须:(1) 按 key 决定去向(partition),(2) 在每台机器内部把同 key 的记录排到一起(sort),(3) 到了 reduce 端把 M 份有序序列归并成一份有序序列(merge sort)。排序不是为了让输出好看,而是为了让”同一个 key 的所有 value 连续出现”,从而 reduce worker 可以流式处理——读到一个 key 就把它的 value 全部收齐、调用一次REDUCE、然后丢进输出。机制图解:map 端与 reduce 端用的是两种不同的排序算法。
map 端: 内存缓冲 -> 分桶 -> 每个桶内部 快排(quicksort) -> 写盘
|-- 分区 p=0 的溢写文件: (apple,1)(apple,1)(banana,1)... 按 key 有序
|-- 分区 p=1 的溢写文件: (cat,1)(dog,1)(dog,1)... 按 key 有序
|-- 分区 p=2 的溢写文件: (egg,1)... 按 key 有序
(桶数 = R; 若缓冲区溢出多次, 每个桶会有多个溢写文件, 最后做一次归并)
reduce 端: 从 M 个 map 节点 fetch 到 M 个有序 run
run_0: (apple,1)(banana,1)(cat,1)
run_1: (apple,1)(cat,1)(cat,1)
run_2: (banana,1)(cat,1)
│
▼ 归并排序 (merge sort) —— 外排, 内存装不下就多路归并 + 落盘
(apple,1)(apple,1)(banana,1)(banana,1)(cat,1)(cat,1)(cat,1)(cat,1)
│
▼ 按 key 分组 (group by key)
apple -> [1,1] banana -> [1,1] cat -> [1,1,1,1]
│
▼ 每个 key 调用一次 REDUCE
(apple,2) (banana,2) (cat,4)
有序性的额外红利:论文 §4.2 给出顺序保证——在给定分区内,中间键值对按 key 递增顺序被处理。这让”输出本身就排好序”成为免费功能(Distributed Sort 就是靠它实现的),也让输出文件支持按 key 的随机访问查找。
关键假设与系统模型:排序键是中间 key
k2,不是 value。因此 reduce 拿到的list(v2)内部顺序是不确定的(取决于各 map 的完成顺序与 fetch 顺序)——用户的REDUCE绝不能依赖这个顺序,否则就变成了非确定性函数。
5.2.7 Combiner 与 Partitioner
- 定义与目的:
- Partitioner:决定中间 key
k2去哪个 reduce,即映射 $P: K_2 \to \{0,1,\dots,R-1\}$,默认 $P(k) = \text{hash}(k) \bmod R$。 - Combiner:一个可选的、运行在 map 端的本地归约。论文 §4.3 指出:每个 map 任务产生的中间 key 常有大量重复(词频服从 Zipf 分布,
<the, 1>会出现成百上千次),如果能在本地先合并,就能大幅减少需要过网络的数据。
- Partitioner:决定中间 key
直观解释:combiner 就像”在把纸条投进信箱之前,先在自己桌上把写着同一个词的小纸条摞成一摞,只写一张总数上去”。代价是只有当”摞纸条”这个操作与最终汇总可以互换时才允许这么做——这就是结合律与交换律的要求。
- 机制图解:
不用 combiner: 用 combiner:
map#0 输出 10000 条 (the,1) map#0 本地: the -> 10000, 输出 1 条
│ │
├── 全部走网络 ──▶ reduce#p ├── 只有 1 条走网络 ──▶ reduce#p
│ │
map#1 输出 8000 条 (the,1) map#1 本地: the -> 8000, 输出 1 条
│ │
└── 全部走网络 ──▶ reduce#p └── 只有 1 条走网络 ──▶ reduce#p
reduce 端: 10000 + 8000 = 18000
reduce 端: 求和 18000 次加法 结果完全相同, 但网络流量降了 9000 倍
Partitioner 的重要性:默认哈希分区”倾向于给出相当均衡的分区”(论文 §4.1),但有时必须自定义。例如输出 key 是 URL、而用户希望同一个 host 的所有条目落到同一个输出文件,就应使用
hash(Hostname(urlkey)) mod R。Sort benchmark 则必须用范围分区(range partitioning)而非哈希——因为哈希会打乱 key 的全局顺序,而”全局有序”正是排序作业的语义本身(讲义明确指出”can’t use hashing!”),并且切分点应当根据数据分布来定,才能让各 reduce 负载均衡(论文的做法是先跑一个采样 MapReduce 收集 key 分布,据此计算切分点)。Combiner 与 Reducer 的区别(含正确性要求):
| 维度 | Combiner(合并器) | Reducer(归约器) |
|---|---|---|
| 运行位置 | 每个 map worker 上,本地 | 每个 reduce worker 上 |
| 运行时机 | map 输出溢写本地磁盘之前;多个溢写文件归并时可能再跑一次 | 全部 map 完成后,shuffle 拉取并归并排序之后 |
| 输入范围 | 一个 map 任务内、同一个 partition 内、同一个 key 的 value | 一个 reduce 任务内、某个 key 的 value 的全局全集(来自全部 M 个 map) |
| 输出去向 | 本地磁盘的中间文件(随后被 shuffle 拉走) | 最终输出文件(写入全局 DFS) |
| 调用次数 | 不确定:由溢写次数与归并策略决定,可能一次也不调用,也可能调用多次 → 必须允许被重复执行且不改变结果 | 每个 key 恰好一次 |
| 对结果的约束 | 不得改变结果。必须满足结合律与交换律,且 REDUCE(k,V) 必须可表达为 REDUCE(k, {COMBINE(k,V₁), COMBINE(k,V₂), …}) | 它本身就是语义的定义,无额外约束 |
| 典型实现 | 通常直接复用 reduce 函数的代码(论文 §4.3:”typically the same code is used”);要求中间 value 类型与输出 value 类型兼容(如 sum 的输入输出都是计数) | 用户自定义 |
| 非法例子 | mean(平均)、median(中位数)、count distinct —— 都不满足结合律 | 任意确定性函数皆可 |
| 主要收益 | 减少 shuffle 网络流量与 map 端磁盘写入;缓解 reduce 端的 key 倾斜 | —— |
5.2.8 容错:本讲的核心(Fault Tolerance)
定义与目的:MapReduce 的设计前提是”机器故障是常态”。它的容错哲学只有一句话:用”重新执行”代替”恢复状态”。没有日志回放、没有检查点恢复(master 除外)、没有数据副本协调,只有”把这个任务再做一遍”。
直观解释:这就像考试时发现某位同学的答题卡被咖啡泼了。传统数据库的做法是”从备份里恢复这张卡”(要先有备份、要有恢复协议);MapReduce 的做法是”让他再答一遍”(因为答题过程是确定性的,答案必然一样)。后者简单得多,而且正确性论证只需要一句话:输入没变 + 函数确定 ⇒ 重跑的结果与第一次相同。代价是要求计算是幂等可重放的——这就是为什么 MapReduce 要求用户函数确定性。
机制图解:master 与 worker 的完整交互时序(含一次 worker 故障的检测与恢复)。
Client Master (JobTracker / AM) Worker W0 Worker W1
| | | |
|--submit(job)------->| | |
Master: create M map tasks + R reduce tasks, all state = idle
| |--assign(map#0, split0)----->| |
| |--assign(map#1, split1)------|------------------------->|
W0/W1: read HDFS split, call map(), buffer intermediate KV in memory
| |--ping(progress=40%)--------<| |
| |--ping(progress=30%)---------|-------------------------<|
| |--done(map#0, tmp[0..R-1])--<| |
Master: state[map#0] <- completed; record R partition file locations
| | | |
X W1 misses several heartbeat periods -> Master marks W1 as failed
map#1 output sits on W1 local disk (unreachable) -> reset to idle
| |--reassign(map#1)----------->| |
| |--done(map#1, tmp[0..R-1])--<| |
== BARRIER: no reduce starts before all M map tasks complete ==
| |--assign(reduce#0, p=0)------|------------------------->|
W1: HTTP-fetch p=0 shards from every map worker, merge-sort by key
| |--done(reduce#0)-------------|-------------------------<|
Master: atomic rename tmp -> part-r-00000, final output becomes visible
|--return(ok)--------<| | |
- Worker 故障(Worker Failure):master 周期性地 ping 每个 worker。若某个 worker 在一段时间内没有响应,master 就把它标记为 failed。
- 该 worker 上已经完成的所有 map 任务被重置回初始的 idle 状态,从而可以在其他 worker 上重新调度。为什么已完成的任务也要重做?因为它的输出存放在那台故障机器的本地磁盘上,已经不可访问了——reduce 任务再也取不到这份数据。这是”map 输出只写本地磁盘”这一设计假设的直接后果。
- 该 worker 上正在进行的 map 或 reduce 任务同样被重置为 idle 并重新调度。
- 已完成的 reduce 任务不需要重做,因为它们的输出存放在全局文件系统中,与 reducer 机器的生死无关。
- reduce worker 的缓存失效:当 map 任务先由 worker A 执行、后来因 A 故障而由 worker B 重新执行时,所有正在跑 reduce 的 worker 都会收到通知;尚未从 A 读取数据的 reduce 任务改从 B 读取。
- 大规模故障的韧性:论文记录了一次真实事件——集群维护导致每次约 80 台机器同时不可达、持续数分钟,MapReduce master 只是重新执行了这些机器做过的计算并继续前进,最终完成了作业。
- 最新 Hadoop 的分层实现:YARN 中 NodeManager 向 ResourceManager 心跳;服务器故障时 RM 等待心跳超时,然后通知所有受影响的 ApplicationMaster,由 AM 采取动作(在其他机器上重启任务)。任务级故障则由 NM 自己跟踪:若某个任务在运行中失败,标记为 idle 并重启它。
Master 故障(Master Failure):论文的做法很直接——让 master 周期性地把上述数据结构写成检查点;如果 master 挂掉,可以从最近一次检查点启动一个新的副本。但论文同时指出:因为只有一个 master,它挂掉的概率本身很低,所以当前实现选择直接中止 MapReduce 计算,由 client 检查到这一情况后自行决定是否重试整个作业。——这是 MapReduce 最大的可用性弱点。工业界后来的改进路径是:Hadoop 1.x 的 JobTracker 也是单点;Hadoop 2.x 把它拆成 ResourceManager(全局,用 checkpoint + 备用 RM 做 active-standby)+ 每作业一个 ApplicationMaster(失败时由 RM 在新容器里重启,并与其仍在运行的任务重新同步)。这与 Chubby/ZooKeeper 的思路一致,详见后续关于协调服务的章节。
- 原子提交(Atomic Commit)与输出语义:这是保证”要么整个 job 的输出可见,要么不可见”的关键机制。
- 每个正在进行的任务把输出写到私有的临时文件里:reduce 任务产生 1 个临时文件,map 任务产生 R 个临时文件(每个 reduce 任务一个)。
- map 任务完成时,worker 向 master 发消息,消息里带上这 R 个临时文件的名字。如果 master 收到的是一个已经完成过的 map 任务的完成消息,它直接忽略(因为备份任务或重执行可能产生重复上报)。否则它把 R 个文件名记入 master 的数据结构。
- reduce 任务完成时,reduce worker 把它的临时输出文件原子地 rename 成最终输出文件。如果同一个 reduce 任务在多台机器上各执行了一次,就会有多次针对同一个最终文件名的 rename。这里依赖底层文件系统提供的原子 rename 操作来保证:最终文件系统的状态只包含某一次 reduce 执行产生的数据。
- 为什么这就够了?因为两次执行的结果相同(用户函数确定性,且输入相同),所以”最后一次 rename 覆盖前一次”不会让用户看到不一致的数据。反过来说,如果用户函数不确定,原子 rename 就只能保证”不撕裂”,不能保证”内容正确”。
- 故障类型 → 检测方式 → 恢复动作汇总:
| 故障类型 | 检测方式 | 恢复动作 | 依据与代价 |
|---|---|---|---|
| Map worker 崩溃(crash-stop) | master 周期性 ping,超时未响应即标记 failed(Hadoop 默认 10 分钟量级;本讲代码用 0.6s 便于观察) | 该 worker 上所有 map 任务(含已完成的)重置为 idle 并重新调度;已完成的 map 的中间文件位置作废;通知所有 reduce worker 改从新的位置读取 | map 输出在故障机本地磁盘,不可访问,必须重做(论文 §3.3) |
| Reduce worker 崩溃 | 同上 | 仅把 in-progress 的 reduce 重置为 idle;已完成的 reduce 不重做 | 输出已通过原子 rename 落在全局 DFS |
| Reduce 任务被误判为失败(假阳性,机器其实还活着) | 心跳超时 | 允许同一 reduce 任务存在两个并发副本;两次都对同一最终文件做 rename | 安全性由原子 rename 保证:最终文件是”某一次”执行的结果;两次内容相同(确定性函数) |
| Straggler(慢节点:坏磁盘、CPU/内存/网络竞争、缓存被关) | master 记录每个任务的进度百分比,发现某个任务长期大幅落后于同阶段其他任务 | 在另一台空闲机器上启动 backup task(speculative execution,推测执行);先完成的那个副本胜出,另一个副本可以被 kill | 论文实测:sort 程序关闭 backup 机制后总时间延长 44% |
| Master(JobTracker / AM)崩溃 | 论文:不检测,job 直接失败,由 client 重试;YARN:RM 通过 AM 心跳超时检测 | 论文:整个 job 重跑;YARN:RM 在新的容器里重启 AM,AM 再与它原本还在跑的任务重新对账 | 论文明确承认 master 是单点 |
| ResourceManager 崩溃 | 备用 RM 收不到心跳 / 选主超时 | 备用 RM 从旧检查点恢复并接管(active-standby) | YARN 的高可用方案;与 Lecture 15/16 的选主与协调机制同源 |
| 坏记录导致用户函数确定性崩溃 | 每个 worker 安装信号处理器捕获 SIGSEGV/SIGBUS;调用用户函数前把参数序号存入全局变量,崩溃时信号处理器发出带序号的 “last gasp” UDP 包给 master | 当 master 看到同一条记录失败超过一次,就在下次重新执行该任务时跳过这条记录 | 这是可选模式(论文 §4.6);它会破坏”输出 = 串行输出”的严格语义,只能用于统计类分析 |
| 主副本与备份副本同时完成 | 收到第二个完成消息时 | master 发现该任务已是 completed,直接忽略重复消息 | 论文 §3.3 的明确规定;也是 5.4.1 代码中验证的行为 |
- 关键假设与系统模型:crash-stop(进程死掉后不再产生任何行为,不发送错误消息、不产生拜占庭行为);心跳超时是故障检测器,因此存在假阳性(把慢的判成死的)与假阴性(死的没及时判出来),恢复机制必须对假阳性安全——这正是原子 rename 存在的原因。网络被假定为可靠但可能很慢:消息可能延迟任意长,但不会损坏。
5.2.9 Straggler 与备份任务(Backup Tasks)
定义与目的:Straggler 指”完成最后几个 map 或 reduce 任务时,耗时异常长的一台机器”。它没有故障,只是慢。备份任务机制用来消除它造成的尾部延迟。
直观解释:木桶效应——整个 job 的完成时间由最慢的那一个任务决定。如果 1000 个任务里 999 个在 10 分钟内做完了,最后一个却要 30 分钟,那么整个 job 就是 30 分钟,另外 999 台机器陪跑。备份任务的思路是让一条备用的流水线同时做同一个零件的加工,谁先做好就用谁的——本质上是把”延迟”问题转化为”吞吐”问题,用冗余计算买时间。
Straggler 的成因(论文 §3.6):坏磁盘导致频繁的可纠正错误,读性能从 30MB/s 掉到 1MB/s;集群调度系统在该机器上又排了别的任务,导致 CPU、内存、本地磁盘或网络带宽竞争;论文还遇到过一个真实 bug——机器初始化代码使 CPU 缓存被禁用,受影响机器上的计算速度慢了一百倍以上。讲义补充的成因同样包括:坏磁盘、网络带宽、CPU、内存。
机制图解:备份任务对总完成时间的影响(数值取自论文 §5.3 与 §5.4 的 sort benchmark)。
(a) backup task DISABLED -- total 1283 s (44% slower than enabled)
|t=0 |200s |400s |600s |800s |1000s |1283s
| |
W0 ###########################################.......................
W1 ###########################################.......................
... ###########################################.......................
Wk-1 ##################################################################
Wk-2 ##################################################################
^
every machine except the 5 stragglers is IDLE from ~t=830s on;
the last 5 reduce tasks run alone until t=1283s.
(b) backup task ENABLED -- total 891 s
|t=0 |200s |400s |600s |800s|891s
| |
W0 ###########################################.......................
W1 ###########################################.......................
... ###########################################.......................
Wk-1 ################################===========.......................
Wm ................................###########.......................
legend: # useful work = duplicated work later discarded . idle
Wk-1 = the straggler; Wm = idle machine picked to run the backup copy at t~620s;
the first replica to finish wins, the other one is killed / its output discarded.
- 论文给出的实测结论:
- 排序程序在禁用备份任务的配置下,960 秒时除 5 个 reduce 任务外全部完成,这最后几个 straggler 又跑了 300 秒才结束;总时间从 891 秒涨到 1283 秒,增加 44%。
- 备份机制经过调优,通常只会让作业消耗的计算资源增加不超过几个百分点。
- 这正是”为什么它能降低几个数量级的尾部延迟”的答案:尾部是极少数任务的长尾,复制它们的成本很小(几个百分点),但消除的延迟极大(44% 甚至更多)。
为什么备份任务是安全的:因为”一个任务被认为完成”的判据是它的第一个副本完成。另一个副本的结果随后被 master 忽略(论文明确规定:收到已完成任务的完成消息就丢弃)。若被忽略的副本是 reduce,它那次多余的 rename 也因原子性 + 结果相同而无害。讲义进一步追问:”重启的 reduce 任务怎么拿到 map 的输入?”答案是 master 保存着
loc[m][r]位置表,它会增量地把新位置推送给正在进行 reduce 的 worker。- 关键假设与系统模型:备份任务假设存在空闲资源(因此它在作业接近尾声、大部分 worker 已经空闲时才最有效;论文正是”当 MapReduce 作业接近完成时,master 调度剩余 in-progress 任务的备份执行”);它也假设任务是幂等可重放的(同样是确定性函数的要求)。
5.2.10 局部性优化(Locality Optimization)
定义与目的:网络带宽是稀缺资源(论文 §3.4 的第一句话)。局部性优化利用”输入数据就存在集群机器的本地磁盘上”这一事实,把计算搬到数据旁边,而不是把数据搬到计算旁边。
直观解释:”与其把一座山的矿石运到冶炼厂,不如把冶炼炉搬到矿山。” 在 MapReduce 规模下,这一点非常关键:grep 实验里输入读取速率峰值超过 30GB/s,而集群的网络总带宽只有 100–200Gbps(约 12–25GB/s)——如果所有输入都要过网络,读取速率根本不可能达到 30GB/s。
机制图解:三级调度优先级。
级别 1 DATA-LOCAL 把 map 任务调度到**存有该 split 副本的机器**上
└─ 延迟 ~0, 输入读取完全不走网络 ← 绝大多数情况应命中这一级
级别 2 RACK-LOCAL 退而求其次, 调度到与该副本**同一个机架**的机器上
└─ 走机架内交换, 通常有较高带宽、较低延迟
级别 3 OFF-RACK 都不行就随便找一台
└─ 数据必须跨机架传输, 带宽代价最高
副本放置策略(3 副本): 第 1 个副本在本机(写者所在节点)
第 2 个副本在**另一个机架**
第 3 个副本在同一机架的另一节点
-> 讲义口径: "2 个在同一机架, 1 个在不同机架"
-> 这样安排是为了**降低写文件时的跨机架带宽**; 同时保证机架级容错
-> 使用更多机架不会影响 MapReduce 的整体调度性能
论文的实测验证:在集群中大比例的 worker 上跑大型 MapReduce 作业时,绝大多数输入数据都是本地读取的,不消耗网络带宽。sort benchmark 的图中,输入速率明显高于 shuffle 速率与输出速率,正是因为局部性优化让大部分数据绕开了相对带宽受限的网络。
关键假设与系统模型:假设 DFS 的 block 位置信息对 master 可见(GFS/HDFS 的 NameNode 保存 block → DataNode 映射);假设集群是层次化拓扑(机架 + 交换机);假设机器同时是存储节点与计算节点。
5.2.11 Hadoop 生态:HDFS、YARN 与应用层
- HDFS(Hadoop Distributed File System):
- NameNode:保存整个文件系统的命名空间(目录树、文件 → block 列表、block → DataNode 位置映射)。它不存数据本身,所有元数据常驻内存,并通过 edit log + fsimage 持久化。它是单点(Hadoop 3.x 引入多 NameNode 改善)。
- DataNode:实际存放 block 数据;定期向 NameNode 发送心跳与 block report;客户端直接与 DataNode 传输数据(NameNode 只给位置,不经过数据)。
- block 与副本:默认 block 128MB(Hadoop 2.x+;论文与讲义口径为 64MB),每块 3 副本,跨机架放置(2 + 1)。
- 与 MapReduce 的配合:InputFormat 按 block 边界切分输入(text 格式保证在行边界断开),把 split → block 位置交给 master,master 据此做 data-local 调度。这就是 5.2.10 的物理基础。
- 注意:HDFS 是”一次写、多次读”的批处理文件系统,不适合随机小写;需要随机读写的场景用 HBase。
- YARN(Yet Another Resource Negotiator):Hadoop 2.x 起的调度层,作用是把资源管理与计算框架解耦。MapReduce 从此只是”跑在 YARN 上的一个应用”,Spark、Tez、Flink 也是。
┌─────────────────────────────────────────────────────────────┐
│ Resource Manager (RM, 全局唯一, 有备 RM) │
│ ├── Scheduler (Capacity Scheduler / Fair Scheduler) │
│ └── 把每台服务器看作一组逻辑 **Container** │
│ Container = 固定 CPU + 固定内存 │
│ (类似 Linux cgroups, 但更轻量) │
└───────┬─────────────────────────────────────┬───────────────┘
│ ① 我要 container │ ② container 完成
│ │
┌──────────┴──────────┐ ┌────────────┴──────────┐
│ Node A │ │ Node B │
│ Node Manager A │ │ Node Manager B │
│ ┌────────────────┐ │ │ ┌─────────────────┐ │
│ │ ApplicationMgr1│ │ │ │ ApplicationMgr2 │ │
│ └────────────────┘ │ │ ├─────────────────┤ │
│ │ │ │ Task (App2) │ │
└─────────────────────┘ │ └─────────────────┘ │
└───────────────────────┘
一个 job 拿到 container 的四步:
1. AM -> RM : "我需要一个 container"
2. RM 的 Capacity Scheduler 分配 -> "Node B 上有一个"
3. AM -> NM(B) : "请在这个 container 里启动我的 task"
4. task 运行, 结束后 container 归还 RM
- 三个主要组件:Global Resource Manager (RM) 负责调度;per-server Node Manager (NM) 是每台机器上的守护进程,跟踪本机的 container 与 task;per-application ApplicationMaster (AM) 负责与 RM/NM 协商 container、并检测本作业的任务失败。
心跳的妙用:心跳消息同时捎带(piggyback)container 请求,避免为每个请求单独发消息。
- 应用层生态:
- Hive:把 SQL 翻译成 MapReduce/Tez/Spark 作业,面向数据仓库式的批量分析。它补上了 MapReduce 缺失的 schema、catalog 与统计信息。
- Pig:Pig Latin 数据流语言(LOAD / FILTER / GROUP / JOIN / FOREACH / STORE),编译成一串 MapReduce 作业。
- HBase:基于 HDFS 的列式、支持随机读写的分布式表存储(Google Bigtable 的开源实现)。MapReduce 可以直接以 HBase 表为输入/输出。
- Spark:用内存中的 RDD + lineage(血统)做容错,对迭代式与交互式负载快得多(详见 Lecture 24)。
- Hadoop 3.x 相对 2.x 的改进(讲义列举):用 Docker 取代 container;用纠删码(erasure coding)取代 3 副本(降低存储开销,但恢复时需要网络重建);多 NameNode 解决命名解析单点;GPU 支持(面向机器学习);节点内磁盘均衡(应对改造过的磁盘);在队列间抢占之外增加队列内抢占。
5.2.12 MapReduce 与 SQL 的关系及 MapReduce 的局限
- MapReduce ≈ 一个表达能力很弱的关系查询引擎:
| SQL 概念 | MapReduce 对应物 |
|---|---|
SELECT(投影/变换) | map |
WHERE(过滤) | map 中不 emit 即可 |
GROUP BY + 聚合函数 | partition + shuffle + reduce |
ORDER BY | 范围分区 + reduce 端归并排序(不能用哈希分区) |
JOIN | map-side join(小表广播)或 reduce-side join(按 join key 分区) |
| 视图/多步查询 | 串联多个 MapReduce 作业 |
- MapReduce 相对 SQL 的五个根本性局限:
- 单遍(single pass):一个 MapReduce 作业只做一次 map + 一次 reduce,没有迭代。复杂的多步分析必须串成一条 MapReduce 作业链,而每一个作业都要落盘再读回。
- 无 schema、无索引、无统计信息:优化器没有任何可依据的元数据。而 MapReduce 干脆没有优化器——M、R、combiner、分区函数全由用户手工指定。
- 迭代弱(weak iteration):PageRank、K-means、梯度下降这类算法需要反复扫描同一份数据,每轮都要”从 HDFS 读 → shuffle → 写回 HDFS”。磁盘 I/O 主宰了运行时间,而 CPU 大部分时间在等 I/O。
- 高延迟 / 交互式查询不可行:作业启动开销很大——grep 实验总耗时 150 秒里约有 60 秒是启动开销(程序分发到所有 worker、与 GFS 交互打开 1000 个输入文件、获取局部性信息)。这种量级的延迟无法支持交互式查询。
- 不适合小数据:框架的固定开销(进程启动、调度、屏障同步)与数据量无关,几 MB 的数据用 MapReduce 反而比单机 Python 慢几个数量级。
- 两条演进路线:(1) 保留 MapReduce 的接口,换掉执行引擎——Tez / Spark 用 DAG 表达多阶段计算,中间结果留在内存/SSD,不强制落 HDFS;(2) 直接提高抽象层次——Hive / Impala / Presto 让用户写 SQL。二者都继承了 MapReduce 的核心遗产:shuffle、分区、数据局部性、”重算代替恢复”的容错范式。
5.3 算法伪代码与正确性分析
算法 5.3.1:MapReduce 整体执行协议(Master 状态机 + Worker 循环)
假设与系统模型
- 进程模型:1 个 master(单点,不参与计算),$n$ 个 worker,均可作为 map worker 或 reduce worker。
- 故障模型:crash-stop。worker 可能在任意时刻停止工作且不发出任何进一步的消息;master 与输入数据不发生故障(master 故障的处理见 5.2.8,本算法不覆盖)。
- 通道假设:消息通道可靠但可能任意延迟(不丢失、不重复、不损坏、不重排到违反因果的程度——实际上本算法对重排也不敏感,因为用
attempt代数做幂等过滤)。 - 故障检测器:基于超时的完美性近似——master 周期性 ping,超时未响应即判定 failed。检测器是不完美的(存在假阳性与假阴性),因此恢复动作必须对假阳性安全。
- 参数:$M$ 个 map 任务(split 数)、$R$ 个 reduce 任务(分区数),$f$ 为故障 worker 数(无上界,但 $f$ 有限)。
- 用户函数假设:
MAP与REDUCE是确定性函数,且在有限时间内终止。
伪代码
─────────────────────────── Master 端 ───────────────────────────
常量: TIMEOUT (心跳超时阈值), STRAGGLE_THRESHOLD
状态变量:
T_map[m] ∈ {IDLE, IN_PROGRESS, COMPLETED} for m = 0..M-1
T_red[r] ∈ {IDLE, IN_PROGRESS, COMPLETED} for r = 0..R-1
owner[t] 当前执行任务 t 的 worker (仅 t 非 IDLE 时有效)
attempt[t] 任务 t 的执行代数 (generation), 每次重新调度 +1
loc[m][r] map m 产生的第 r 个分区文件的位置与大小
hb[w] worker w 最近一次心跳的时间戳
alive[w] worker w 是否存活
busy[w] worker w 正在执行的任务 (或 nil)
start[t] 任务 t 本次执行的开始时间 (用于 straggler 判定)
backup[t] 任务 t 是否已经启动了备份副本
初始化:
for m in 0..M-1: T_map[m] ← IDLE; attempt[m] ← 0
for r in 0..R-1: T_red[r] ← IDLE; attempt[r] ← 0
for w in workers: alive[w] ← true; hb[w] ← now(); busy[w] ← nil
主循环 (反复执行直到返回):
── A. 故障检测 ──────────────────────────────────────────────
for each w in workers:
if alive[w] = true ∧ now() − hb[w] > TIMEOUT:
alive[w] ← false
upon WorkerFailed(w) // 见下面的处理例程
── B. 终止条件 ──────────────────────────────────────────────
if ∀m: T_map[m] = COMPLETED ∧ ∀r: T_red[r] = COMPLETED:
return SUCCESS // 唤醒用户程序, 返回 R 个输出文件
── C. 任务分配 (考虑局部性) ─────────────────────────────────
for each w with alive[w] = true ∧ busy[w] = nil:
if ∃ m: T_map[m] = IDLE: // 阶段 1: map
m ← 优先选择满足 "数据本地性" 的 IDLE map 任务 // data-local > rack-local > 任意
attempt[m] ← attempt[m] + 1
T_map[m] ← IN_PROGRESS; owner[m] ← w; busy[w] ← m; start[m] ← now()
send ⟨RUN_MAP, m, attempt[m], input_path(m)⟩ to w
else if (∀m: T_map[m] = COMPLETED) ∧ (∃ r: T_red[r] = IDLE): // 阶段 2: reduce
r ← 任取一个满足 T_red[r] = IDLE 的 r // ★ 屏障: 只有全部 map 完成才进这里
attempt[r] ← attempt[r] + 1
T_red[r] ← IN_PROGRESS; owner[r] ← w; busy[w] ← r
send ⟨RUN_REDUCE, r, attempt[r], { (m, loc[m][r]) : m = 0..M−1 }⟩ to w
── D. Straggler 检测与备份任务 ───────────────────────────────
for each task t with state[t] = IN_PROGRESS:
if now() − start[t] > STRAGGLE_THRESHOLD ∧ backup[t] = false
∧ ∃ an idle alive worker w′:
backup[t] ← true
// ★ 注意: attempt[t] 不变。两个副本共享同一个代数, 因此谁先提交谁赢
send ⟨RUN_MAP, t, attempt[t], input_path(t)⟩ to w′ // t 为 map 任务时
busy[w′] ← t; start_copy[t] ← now()
事件处理:
upon receive ⟨MAP_DONE, m, a, files[m][0..R−1], counters⟩ from w:
busy[w] ← nil
if a ≠ attempt[m] ∨ T_map[m] = COMPLETED: return // 过期结果 / 重复完成 -> 丢弃
T_map[m] ← COMPLETED; owner[m] ← w; loc[m][*] ← files[m][*]
push loc[m][*] to every worker currently running a reduce // 增量推送位置
upon receive ⟨REDUCE_DONE, r, a⟩ from w:
busy[w] ← nil
if a ≠ attempt[r] ∨ T_red[r] = COMPLETED: return // 同上: 幂等过滤
T_red[r] ← COMPLETED // 最终输出已由 worker 端原子 rename
upon receive ⟨PING, w, progress, counters⟩ from w:
hb[w] ← now() // 心跳捎带进度与计数器
upon WorkerFailed(w):
b ← busy[w]; busy[w] ← nil
if b ≠ nil: state[b] ← IDLE // 未完成的任务回到 idle
for each m with owner[m] = w ∧ T_map[m] = COMPLETED: // ★ 已完成的 map 也必须重做
T_map[m] ← IDLE; loc[m][*] ← ⊥; backup[m] ← false // 因为输出在 w 的本地磁盘上
for each r with T_red[r] = IN_PROGRESS ∧ owner[r] = w:
T_red[r] ← IDLE // reduce 的临时文件作废
for each r with T_red[r] = COMPLETED: (保持不变) // ★ 输出在全局 DFS, 不必重做
if 有 map 任务被重置为 IDLE:
for each r with T_red[r] = IN_PROGRESS: // 屏障被打破
T_red[r] ← IDLE // 正在跑的 reduce 也必须重来
─────────────────────────── Worker 端 ───────────────────────────
状态: my_id, 当前任务的进度 progress ∈ [0,1]
初始化: 启动一个心跳线程, 每 HB_INTERVAL 发送一次 PING
loop: send ⟨PING, my_id, progress, counters⟩ to master; sleep(HB_INTERVAL)
主循环:
upon receive ⟨RUN_MAP, m, a, path⟩:
(k1,v1)* ← InputFormat(path) // 按行切分; 每个 k1 是行偏移量
buffer ← ∅
for each (k1,v1) in (k1,v1)*:
for each (k2,v2) in MAP(k1,v1): buffer.append((k2,v2))
progress ← 已处理字节 / 总字节
if |buffer| > BUFFER_LIMIT: Spill()
Spill(); MergeSpills() // 产生 R 个分区文件 (m, 0..R−1)
send ⟨MAP_DONE, m, a, file_names(m,0..R−1), counters⟩ to master
function Spill():
for p in 0..R−1:
partition[p] ← { (k2,v2) ∈ buffer : PARTITION(k2) = p }
for p in 0..R−1:
sort(partition[p]) by k2 // 快排 -> 该分区内 key 有序
if COMBINER 已定义:
partition[p] ← { COMBINE(k2, values) : (k2, values) = group_by_key(partition[p]) }
write partition[p] to local_disk as file (m, p)
buffer ← ∅
upon receive ⟨RUN_REDUCE, r, a, locations⟩:
fetched ← ∅
for m in 0..M−1:
fetched ← fetched ∪ HTTP_fetch(locations[m][r]) // 远程读 map worker 的本地磁盘
progress ← m / M
merged ← MergeSort(fetched) // 归并 M 个有序 run; 内存不够则外排
for each (k2, list(v2)) in GroupByKey(merged): // merged 已按 key 有序, 可流式分组
for each (k3,v3) in REDUCE(k2, list(v2)):
append (k3,v3) to tmp_file(r)
AtomicRename(tmp_file(r), final_path(r)) // ★ 原子提交, 完成后最终输出才可见
send ⟨REDUCE_DONE, r, a⟩ to master
算法逻辑解说(用一个具体数值例子走一遍)
- 设 $M=8, R=3$,4 个 worker W0–W3,某个 worker W1 在完成第 1 个 map 后崩溃。
- $t_0$:master 把 map#0–#3 派给 W0–W3。每个 map 任务耗时约 0.3 个时间单位。
- $t_1$:W1 完成 map#1,发出
<MAP_DONE, 1, attempt=1, files>。master 把T_map[1] ← COMPLETED、owner[1] ← W1、loc[1][*]记下,并把位置推送给(此时还不存在的)reduce worker。随后 W1 进程消失。 - $t_2$:master 在若干轮循环里继续派活。它把 map#5 派给了 W1(因为它还认为 W1 活着)——这个任务永远不会完成,
busy[W1] = ('map', 5)。 - $t_3 = t_2 + \text{TIMEOUT}$:
now() − hb[W1] > TIMEOUT,master 判定 W1 failed。于是:busy[W1] = ('map',5)→T_map[5] ← IDLE(未完成的任务回收);owner[1] = W1 ∧ T_map[1] = COMPLETED→T_map[1] ← IDLE,loc[1][*] ← ⊥(已完成的 map 也必须重做);- 因为发生了第 2 步,任何
IN_PROGRESS的 reduce 也被回退为 IDLE($M=8$ 时此刻通常还没有 reduce 在跑)。
- $t_4$:map#1 与 map#5 被重新派给空闲的 W0(
attempt分别变为 2 和 2)。 - $t_5$:一个迟到的结果到达——W1 在死前已把 map#1 的结果发出,但更可能的情形是:备份任务或重执行导致同一个 map 出现了两份结果。master 检查
a ≠ attempt[m]或T_map[m] = COMPLETED,两者命中其一就丢弃。代码中的实际日志是:[master] 丢弃 map#5 的过期结果(gen=1, 当前=2)与[master] map#2 已由另一个副本提交, 忽略 W2 的迟到结果(输出去重)。 - $t_6$:所有 map 任务 COMPLETED,屏障跨过,reduce#0–#2 被派发,各自 fetch 自己的那个分区、归并排序、逐 key 调用
REDUCE、原子 rename。 - $t_7$:全部 reduce COMPLETED,master 返回 SUCCESS。
正确性论证
安全性 S1(同一个 map 任务至多有一个被采纳的结果):master 只在 a = attempt[m] ∧ T_map[m] ≠ COMPLETED 时才把 MAP_DONE 提交进状态。attempt 每次重新调度(包括因故障重置后的重调度)都自增,因此任何来自旧代数的消息都被 a ≠ attempt[m] 过滤。对于备份任务,attempt 不变,此时由 T_map[m] = COMPLETED 这一条件过滤重复消息。两条过滤条件合起来覆盖了”重执行”与”备份执行”两种重复来源,所以每个 map 任务的状态至多一次从 IN_PROGRESS 转移到 COMPLETED。reduce 同理。这是幂等性的核心,也是 master 不需要区分”这条消息是不是第一次发来的”的原因。
安全性 S2(不会在有 map 未完成时运行 reduce):分配 reduce 的分支有一个前置条件 ∀m: T_map[m] = COMPLETED。而任何 map 被重置为 IDLE 时(WorkerFailed 处理例程的第 2 步),所有 IN_PROGRESS 的 reduce 都会在同一例程内被回退为 IDLE。因此”存在 IDLE 的 map”与”存在 IN_PROGRESS 的 reduce”这两个条件在 master 的任意可达状态里互斥。这保证了 reduce 永远不会读到不完整的中间数据集合。依赖的假设:master 是单点、状态更新是原子的(单线程主循环)。
安全性 S3(输出等于某次串行执行的输出):见 5.3.4 的定理。
活性 L1(只要还有存活 worker,就总有任务被派发):若存在 IDLE 任务且存在空闲的存活 worker,分配分支一定会在本轮循环里派发至少一个任务。因此”任务停滞”只可能发生在两种情形:(a) 所有任务都已完成(正常终止),或 (b) 所有 worker 都已被占用或全部死亡。情形 (b) 由 L2 排除。
活性 L2(在公平性假设下最终完成):设故障发生次数有限(这是”crash-stop 且机器最终被修复/替换”的公平性假设)。每次 worker 故障最多把 $M+R$ 个任务回退为 IDLE,故障次数有限 ⇒ 回退总次数有限 ⇒ 存在一个时刻 $t^*$,此后不再有任务被回退。此时每个任务被派发给某个 worker 后,由”用户函数有限时间终止”这一假设,会在有限时间内完成,且期间不会再被回退。因此所有任务在有限时间内 COMPLETED,master 在 B 步返回 SUCCESS。依赖的假设:(i) 至少有一个 worker 最终不再崩溃(否则算法会无限重试——这是活性的损失,不是安全性的损失:它不会给出错误答案,只是不给答案);(ii) 输入数据在 job 期间不变。
安全性 S4(不会返回部分结果):master 只在全部 R 个 reduce 都 COMPLETED 时才唤醒用户程序。”全部 COMPLETED”意味着每个 reduce 的 REDUCE_DONE 都通过了幂等过滤,即每个 reduce 的最终输出文件都已经被原子 rename。因此用户看到的输出是完整的 R 个文件,不存在”只写了一半”的中间状态。反向地,若 job 因 master 失败或超时而中止,用户程序根本没有被唤醒——“要么完整可见,要么完全不可见” ——这正是原子提交要保证的。
复杂度
- 消息复杂度:心跳 $O(f \cdot n \cdot \frac{T_{\text{job}}}{\text{HB\_INTERVAL}})$($n$ 为 worker 数,与 job 时长成正比);任务派发 $O(M + R + \text{reassignments})$;完成上报 $O(M + R)$;位置推送 $O(R)$ 每个完成的 map 任务,总计 $O(M \cdot R)$。注意心跳消息与作业时长成正比,这是大规模长作业中的主要消息开销。
- master 空间复杂度:$O(M + R)$ 的任务状态 + $O(M \cdot R)$ 的中间文件位置表(每对约 1 字节,论文 §3.5)。这是 $M$、$R$ 不能无限增大的根本原因。
- master 时间复杂度:每次调度决策 $O(M + R)$;总调度决策数 $O(M + R + \text{reassignments})$。
- 失败恢复代价:单次 worker 故障最多导致 $O(M + R)$ 个任务重做;若故障机器上已完成 $k$ 个 map,则额外重做 $k$ 个 map 任务的全部输入。
算法 5.3.2:WordCount 的 map / combiner / reduce
假设与系统模型
- 输入为文本文件,InputFormat 把每一行解析为
<行在文件中的字节偏移量, 行内容>;输入 key 是什么并不重要(WordCount 的map忽略它)。 - 词的分割规则
tokenize由用户定义(论文的 C++ 示例与 Hadoop 示例都用”按空白字符切分”)。 - 计数用整数,不会溢出(若会溢出,
sum就不再满足结合律——整数加法在有限精度下会失去结合律,这也是 combiner 合法性的一个隐含前提)。 - 分区函数为 $P$,reduce 数 $R$。确定性函数假设:
tokenize必须是纯函数(不能依赖当前 locale 的随机状态或外部可变状态)。
伪代码
map(k1, v1):
// k1 = 行偏移量 (使用中可忽略), v1 = 行内容
for each word w in tokenize(v1):
EmitIntermediate(w, 1)
combine(k2, list(v2)): // 可选; 与 reduce 同代码, 但只在 map 端本地运行
s ← 0
for each v in list(v2): s ← s + v
EmitIntermediate(k2, s) // 输出回到中间文件, 之后被 shuffle 拉走
reduce(k2, list(v2)):
s ← 0
for each v in list(v2): s ← s + v
Emit(k2, s) // 输出写到最终输出文件
算法逻辑解说(以 5.2.5 的例子走一遍)
- 输入 split0 =
["Welcome Everyone", "Hello Everyone"],tokenize按空白切分。 map被调用 2 次(每行一次);第一次 emit(Welcome,1) (Everyone,1),第二次 emit(Hello,1) (Everyone,1)。合计 4 条中间记录。- 分区(玩具函数
len(w) mod 2):Everyone→0,Welcome→1,Hello→1。 combine在每个分区内、每个 key 上被调用:p=0 桶得到(Everyone, [1,1])→ emit(Everyone,2);p=1 桶得到(Welcome,[1])与(Hello,[1])→ emit(Welcome,1) (Hello,1)。中间记录从 4 条降到 3 条。- 两个 map 任务的产出被 shuffle 到对应 reduce:reduce#0 收到
(Everyone,2) (Hadoop,2);reduce#1 收到(Hello,1)(Hello,1)(Welcome,1)(Welcome,1)。 reduce对Everyone调用一次得到 2;对Hello调用一次,输入列表是[1,1],得到 2。每个 key 恰好调用一次reduce。- 最终输出两个文件,内容与串行统计完全一致。
正确性论证
安全性(无 combiner 时):设输入行的多重集为 $L$。map 按定义把每行 $l$ 映射为多重集 $\text{MAP}(l) = \{(w,1) : w \in \text{tokenize}(l)\}$。分区的正确性(见算法 5.3.3)保证:对任意词 $w$,$\bigcup_l \text{MAP}(l)$ 中所有 key 为 $w$ 的记录都进入唯一的 reduce 任务 $P(w)$。该 reduce 任务的输入中,$w$ 的 value 列表长度恰为 $w$ 在 $L$ 中出现的总次数(因为 map 对每个出现 emit 一个 1,而分区不丢失、不重复记录)。REDUCE 对该列表求和,得到 $w$ 的词频。依赖的假设:分区不丢记录(通道可靠)、每个 map 任务的输出被完整 shuffle(reduce 拉取全部 M 个 map 的分区文件)。
安全性(有 combiner 时)——结合律与交换律论证:设某个 reduce 任务的输入多重集为 $V$,被 map 端的 combiner 分成了 $k$ 个子多重集 $V_1, \dots, V_k$(每个子集来自某个 map 任务的某次 spill 或归并)。无 combiner 时计算的是 $f(V)$,其中 $f$ 是 REDUCE 在 value 多重集上的语义(这里 $f(V) = \sum_{v \in V} v$)。有 combiner 时计算的是 $f(\{f(V_1), \dots, f(V_k)\})$。
- 结合律给出 $f(V_1 \uplus V_2) = f(\{f(V_1), f(V_2)\})$($\uplus$ 表示多重集并)。对 $k$ 归纳:$f(V_1 \uplus \cdots \uplus V_k) = f(\{f(V_1), \dots, f(V_k)\})$。
- 交换律保证 $f$ 对元素的次序不敏感,因此 spill 顺序、map 完成顺序、tie-break 顺序都不影响结果。
- 整数加法满足结合律与交换律(在无溢出的前提下),因此 combiner 合法,结果与无 combiner 时完全相同。∎
- 反例(必须点明):$f = $ 平均值 不满足结合律。$\text{mean}(1,2,3) = 2$,而 $\text{mean}(\text{mean}(1,2), 3) = \text{mean}(1.5, 3) = 2.25$。直觉上的原因是:局部均值丢失了”这个局部有多少个样本”这一权重信息,而各 split 中该 key 的样本数并不相等。中位数、
count distinct同理不合法。
活性:map 对每行调用一次,行数有限;combine/reduce 对每个 key 的 value 列表做一次线性求和,列表有限。故所有函数在有限时间内终止,结合算法 5.3.1 的活性论证,job 最终完成。
复杂度
map单次调用时间 $O(\vert v_1\vert )$(行长度);一次 map 任务总时间 $O(\vert split\vert + \vert K_{\text{local}}\vert \log \vert K_{\text{local}}\vert )$,其中 $K_{\text{local}}$ 是该 split 内不同 key(词)的数量——排序/分组的开销。- 中间输出规模:无 combiner 时 $\vert V\vert = \sum_l \vert \text{tokenize}(l)\vert $(即总词数 $N$);有 combiner 时 $\vert V^{\prime}\vert = \sum_{m=1}^{M} k_m \le \min(N,\; M \cdot K)$,其中 $k_m$ 是 split $m$ 中不同词数,$K$ 是全局不同词数。对于 Zipf 分布的文本,$\sum_m k_m \ll N$,这正是 combiner 的巨大收益所在。
reduce端空间:$O(\vert V_r\vert )$(若超过内存则外排,$O(\vert V_r\vert \log \vert V_r\vert )$ 时间、$O(\vert V_r\vert )$ 磁盘)。- 时间:$T_{\text{map}} = O(N)$,$T_{\text{sort}} = O(\vert V\vert \log \vert V\vert )$ 分布在 $M$ 个 map 与 $R$ 个 reduce 上。
算法 5.3.3:Hash Partitioner 与 shuffle 的正确性
假设与系统模型
- 中间键集合 $K_2$;$R$ 个 reduce 任务;分区函数 $P(k) = \text{hash}(k) \bmod R$。
- 核心假设(三条,缺一不可): (i) 所有 map 任务使用完全相同的 $R$ 与 $P$(由 job 配置统一给出,不由 worker 自行决定); (ii)
hash是一个纯函数:对同一个 key,在任何机器、任何时刻、任何进程中都返回相同的值。特别注意:Python 的hash()对字符串默认带随机化种子(PYTHONHASHSEED),在不同进程里返回值不同,因此绝不能用它做分区函数;必须使用hashlib.md5之类的稳定哈希(本讲代码正是这么做的); (iii) 分区器无状态:$P(k)$ 只依赖 $k$,不依赖调用次数、不依赖已经处理了哪些 key。 - 中间数据在 job 期间不被修改;网络可靠。
伪代码
function PARTITION(k2): // 由用户在 job 提交时指定, 所有 map 任务一致
return hash_stable(k2) mod R // hash_stable 必须是确定性的纯函数
// map 端 (每个 spill 时执行)
for each (k2, v2) in buffer:
p ← PARTITION(k2)
append (k2, v2) to partition_file[p] // 同一 map 任务的同一个 p 写入同一个文件
sort(partition_file[p]) by k2
write partition_file[p] to local disk; report (size, path) to master
// reduce 端 (任务 r)
for m in 0..M−1:
data[m] ← HTTP_fetch(loc[m][r]) // 只拉分区号为 r 的那一份
merged ← MergeSort(data[0], ..., data[M−1])
for each (k2, list(v2)) in GroupByKey(merged):
emit REDUCE(k2, list(v2))
算法逻辑解说
- 一个 map 任务产生 $R$ 个分区文件(而不是 1 个),就是为了让每个 reduce 只拉”属于自己的那一份”。假设 $R=3$,某个 map 任务处理了 10000 条中间记录,它们被拆成 3 份分别落到
(m,0)、(m,1)、(m,2)。 - 一个 reduce 任务 $r$ 需要拉 $M$ 个文件(每个 map 一个),做 $M$ 路归并排序。若 $M = 15000$、$R = 4000$(论文 sort benchmark 的配置),则总共存在 $M \cdot R = 6 \times 10^7$ 个”潜在的分区文件”,master 为每个已完成的 map 记录 4000 个位置三元组——这就是 $O(M \cdot R)$ 状态的来源(占内存约 $M \cdot R$ 字节,即 60MB 量级,论文说”常数因子很小”)。
正确性论证
安全性(同一个 key 的全部 value 必然被送到同一个 reduce): 设 key $k$ 在 map 任务 $m_1$ 与 $m_2$ 上都产生了中间记录。
- $m_1$ 计算 $P(k) = \text{hash\_stable}(k) \bmod R = p_1$,把记录写入它自己的第 $p_1$ 个分区文件。
- $m_2$ 调用同一个纯函数(假设 ii)、以同样的 $R$(假设 i),得到 $p_2 = \text{hash\_stable}(k) \bmod R$。由于
hash_stable是确定性函数,$\text{hash\_stable}(k)$ 的取值与调用者无关,故 $p_2 = p_1 =: p$。$m_2$ 把记录写入它自己的第 $p$ 个分区文件。 - reduce 任务 $r$ 的输入定义为”所有 $m$ 的第 $r$ 个分区文件的并”(见 reduce 端伪代码:只 fetch
loc[m][r])。若 $p = r$,则两条记录都进入 $r$ 的输入;若 $p \ne r$,则两条都不进入。 - 因此对任意 key $k$,其全部中间 value 恰好出现在唯一一个 reduce 任务 $r = P(k)$ 的输入中,且不会出现在任何其他 reduce 的输入中。
- 由 5.3.2 的安全性论证,$r$ 的
REDUCE看到的是 $k$ 的完整 value 多重集。∎- 三条假设各自的作用:假设 (i) 保证不同 map 不会用不同的 $R$(例如某个 worker 误用了默认值);假设 (ii) 保证同一个 key 的哈希值跨机器一致(用 Python
hash()就会违反这一条,导致同一 key 被送到两个不同 reduce,结果静默错误——每个 reduce 各算出一个部分计数);假设 (iii) 保证 partition 决策只依赖 key。
- 三条假设各自的作用:假设 (i) 保证不同 map 不会用不同的 $R$(例如某个 worker 误用了默认值);假设 (ii) 保证同一个 key 的哈希值跨机器一致(用 Python
活性与负载均衡(性能,而非正确性):
- 设 $n_k$ 为 key $k$ 的 value 个数($k$ 的”重量”),reduce $r$ 的负载 $L_r = \sum_{k: P(k)=r} n_k$,总负载 $\vert V\vert = \sum_k n_k$。
- 上界下界(分区函数的能力极限):令 $L^* = \max_k n_k$ 为最重单个 key 的重量。无论 $P$ 怎么选,承载 $k^$ 的那个 reduce 满足 $L_r \ge L^$;同时由鸽巢原理存在 $r$ 使 $L_r \ge \vert V\vert / R$。因此 \(\max_r L_r \;\ge\; \max\!\left(\frac{\vert V\vert }{R},\; L^*\right).\) 当 $L^* \gg \vert V\vert /R$ 时(重尾分布,例如词频的 Zipf 分布中 “the” 的计数远超平均值),任何分区函数都无法让负载均衡。这是分区函数的能力上界,不是实现缺陷。解决办法不是换 hash,而是换策略:用 combiner 把 $n_k$ 压成”每个 map 一个数”($n_k$ 从可能的 $10^6$ 降到 $M$),或者对超重 key 加盐(salted key)拆成多个 reduce 再跑第二轮聚合。
- 哈希均匀性:若
hash近似把 $K_2$ 均匀映射到 $R$ 个桶(好的混合函数能消除 key 自身的结构,例如 URL 的公共前缀),则 $\mathbb{E}[L_r] = \vert V\vert / R$,偏差为 $O(\sqrt{\vert V\vert /R})$ 量级;用 MD5/SHA 的最后几个字节取模(分子通常取质数)是常见做法。 - 范围分区的适用场景:当输出必须全局有序(Sort 作业)时,哈希分区会打乱 key 的全局顺序,因此必须改用范围分区。此时切分点必须根据数据的实际分布来确定才能均衡——论文的做法是先用一个 MapReduce 预扫描采样 key 分布,据此计算 split point,再跑正式的排序 pass。若用等宽切分点,重尾分布会让某个 reduce 承担绝大部分数据。
复杂度
- 分区计算:每个中间记录 $O(1)$ 次哈希(
hash_stable(k2)的实际开销取决于实现,MD5 约 $O(\vert k_2\vert )$)。 - 空间:每个 map 任务本地 $O(R)$ 个文件句柄与缓冲;master $O(M \cdot R)$ 位置元数据。
- 通信量:无 combiner 时 shuffle 记录数 $= \vert V\vert $;有 combiner 时 $= \sum_{m=1}^{M} k_m \le \min(\vert V\vert ,\; M \cdot K)$;连接数(reduce 向 map 发起的 fetch 次数)最多 $O(M \cdot R)$,这也是为什么 $R$ 太大会让小文件与连接数爆炸。
5.3.4 全局正确性:crash-stop 下的完成性与串行等价性
假设与系统模型
- 故障模型:worker 为 crash-stop,故障数任意但有限(公平性:存在某个时刻之后所有存活 worker 不再崩溃);master 与输入数据不发生故障(这是本讲结论的边界,master 故障的后果已由 5.2.8 单独讨论:整个 job 失败并由 client 重试)。
- 用户函数:
MAP与REDUCE均为确定性函数(同一输入必得同一输出),且在有限时间内终止。 - 通道可靠;底层文件系统提供原子 rename。
定理 5.1(完成性 / Liveness) 在上述假设下,若至少有一个 worker 永不故障,则 job 最终完成,master 返回 SUCCESS;否则 master 可能永远重试(不终止),但绝不会返回错误的或部分的输出。
证明:
- 由算法 5.3.1 的 WorkerFailed 例程,每次 worker 故障最多把 $M + R$ 个任务的状态从 {IN_PROGRESS, COMPLETED} 位移回 IDLE。
- 故障次数有限 ⇒ 被回退的任务实例总数有限 ⇒ 存在时刻 $t^*$,此后不再有任何任务被回退。
- 由算法 5.3.1 的活性论证 L1,只要存在 IDLE 任务与空闲的存活 worker,分配分支每轮必然派发任务;空闲 worker 只会因为”被分配任务”而变忙,而每个任务在有限时间内完成(用户函数终止性假设),故任务不会永久积压。
- 因此 $t^*$ 之后每个 IDLE 任务最终被派发并在有限时间内 COMPLETED;所有任务在有限时间内全部 COMPLETED。
- master 的终止条件要求”全部 map 且全部 reduce 都 COMPLETED”,此时它在 B 步返回 SUCCESS,唤醒用户程序。∎ 若连一个永久存活的 worker 都不存在,则第 3 步的假设被破坏,算法可能无限重试。这是活性(终止性)的损失,不是安全性的损失:master 从未返回 SUCCESS,用户程序从未被唤醒,因此不会有人读到错误结果。工程实现通常再加一个 deadline/失败计数上限,把”无限重试”变成”显式报告 FAILURE”——本讲代码中的
deadline就是这一机制。
定理 5.2(串行等价性 / Safety,顺序无关性) 若 MAP 与 REDUCE 确定,则无论发生多少次 worker 故障、多少个备份任务被执行,最终 R 个输出文件的并集,等于对同一输入的一次无故障串行执行所产生的输出(在 key 到 value 的映射意义上完全相等)。
证明:
- 每个 map 任务恰好有一个被采纳的结果。由算法 5.3.1 的幂等过滤(
a = attempt[m] ∧ T_map[m] ≠ COMPLETED),每个 map 任务最多提交一次;由定理 5.1,它最终会提交。记 map 任务 $m$ 被采纳的那次执行为 $e(m)$。注意 $e(m)$ 可能是原始执行、重执行或备份执行中的任意一次。 - 被采纳的是哪一次执行无关紧要。$e(m)$ 的输入是固定的输入 split(假设输入不变),
MAP是确定性函数,故 $e(m)$ 产生的中间 KV 多重集 $I(m)$ 与”是哪一次执行”无关。这是”重算代替恢复”能够成立的全部理由——它把”选哪一次结果”这个难题消解掉了。 - 分区唯一性:由算法 5.3.3,key $k$ 的全部 value 都进入唯一一个 reduce 任务 $r = P(k)$。
- 每个 reduce 任务的输入是确定的:reduce 任务 $r$ 的提交输入为 $\biguplus_{m=0}^{M-1} \big(I(m) \,\big\vert \,_{P(\cdot)=r}\big)$。由第 2 步,每个 $I(m)$ 是确定的多重集,故这个并也是确定的(与 shuffle 的到达顺序、map 的完成顺序无关)。
- 每个 reduce 任务的输出是确定的:
REDUCE是确定性函数,输入(连同 key 的分组方式)确定 ⇒ 输出确定。记 reduce 任务 $r$ 的输出为 $O(r)$。 - 与串行执行的对应:一次串行执行的定义是:计算所有 $I(m)$,合并成 $\biguplus_m I(m)$,按 key 分组,对每个 key 调用一次
REDUCE,把所有输出收集起来。而 $\{O(r)\}_{r=0}^{R-1}$ 正是把同一个分组过程按 $P(k)$ 切成 $R$ 块分别计算的结果(不同 key 之间互不影响,因为REDUCE的输入只包含同一个 key 的 value),拼起来与串行结果完全相同。∎ - 原子提交保证可见性:每个任务先写私有临时文件,reduce 完成时才原子 rename 到最终文件名。若同一 reduce 被执行多次,会有多次 rename,但由底层文件系统的原子性,最终文件内容是恰好某一次执行的结果;由第 5 步这些执行的结果相同,所以用户观察到的输出完整且一致。reduce 输出不需要在 reducer 机器故障后重做,正是因为 rename 后的文件已经在全局 DFS 上,不随机器消失。∎
非确定性函数下的弱化语义(必须点明) 论文给出了精确的弱化表述:当 MAP 或 REDUCE 非确定时,某个 reduce 任务 $R_1$ 的输出等价于某一次串行执行中 $R_1$ 的输出;但另一个 reduce 任务 $R_2$ 的输出可能对应另一次不同的串行执行。
- 原因:设 map 任务 $M$ 先后被执行了两次(一次因故障重做,或一次是 backup)。$e(R_1)$ 可能读了 $M$ 的第 1 次执行的输出,而 $e(R_2)$ 读了第 2 次执行的输出。由于 $M$ 非确定,两次输出不同,于是 $R_1$ 与 $R_2$ 各自”内部自洽”,但拼在一起不对应任何一次串行执行。
- 现实中的反例:
map里 emit 随机数或time.time();reduce里依赖 value 列表的顺序取”第一个元素”当代表(shuffle 顺序不确定);reduce里写外部文件/发网络请求(副作用);用未初始化的全局计数器。 - 结论:MapReduce 的”顺序无关性 / 串行等价”不是框架自动给你的,而是”框架 + 用户承诺函数确定性”这一契约共同保证的。 忘记这一点,会写出偶发错误的作业——而且是那种只在有机器故障时(也就是生产环境里)才出现的错误。
复杂度小结
- 总任务数 $M + R$;最坏情况下每个任务被重执行的次数无上界(但有限);备份任务使计算资源使用增加”不超过几个百分点”(论文 §3.6)。
- 端到端时间与空间的上界见 5.5。
5.4 代码示例与分布式实现
5.4.1 多进程 MapReduce 框架模拟器(含故障注入与 backup task)
这份代码用 multiprocessing 在一台机器上模拟一个完整的 MapReduce 集群:一个 master 进程 + 多个 worker 进程,覆盖 split → map → buffer → partition/sort → combiner → shuffle → merge → reduce → 原子提交的完整链路,并注入两种真实世界的异常:worker 崩溃与 straggler。
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""Lecture 5 代码 1: 多进程 MapReduce 模拟器 (master + map/reduce workers)
特性: split/partition/combiner/shuffle/归并排序/心跳超时重调度/backup task
"""
import multiprocessing as mp
import queue
import threading
import hashlib
import random
import time
import os
from collections import defaultdict, Counter
# ---------- 用户自定义函数 (UDF): 只有这三个函数与具体作业有关 ----------
def udf_map(key, value):
return [(w, 1) for w in value.split()]
def udf_combine(key, values):
return [(key, sum(values))]
def udf_reduce(key, values):
return [(key, sum(values))]
def partition_of(key, R):
return int(hashlib.md5(str(key).encode()).hexdigest(), 16) % R
# ---------- Worker 进程 ----------
def worker_main(wid, in_q, result_q, hb, cfg):
stop = threading.Event()
def beat():
while not stop.is_set():
hb[wid] = time.time()
time.sleep(cfg['hb_interval'])
threading.Thread(target=beat, daemon=True).start()
done = 0
while True:
try:
task = in_q.get(timeout=0.2)
except queue.Empty:
continue
if task is None:
break
kind, tid, gen, payload = task
time.sleep(cfg['task_delay']) # 模拟真实计算耗时
if cfg['slow'].get(wid):
time.sleep(cfg['slow_delay']) # 注入 straggler
if kind == 'map':
buckets = defaultdict(list)
for i, line in enumerate(payload):
for k, v in udf_map(i, line):
buckets[partition_of(k, cfg['R'])].append((k, v))
raw = sum(len(v) for v in buckets.values())
parts = {}
for p, kvs in buckets.items():
if cfg['use_combiner']: # combiner: 本地预聚合
grp = defaultdict(list)
for k, v in kvs:
grp[k].append(v)
out = []
for k in sorted(grp):
out.extend(udf_combine(k, grp[k]))
parts[p] = out
else:
parts[p] = sorted(kvs)
result_q.put(('map_done', tid, gen, wid, parts, raw))
else:
grp = defaultdict(list)
for k, v in payload: # payload 已在 master 端归并排序
grp[k].append(v)
out = []
for k in sorted(grp):
out.extend(udf_reduce(k, grp[k]))
result_q.put(('reduce_done', tid, gen, wid, out, 0))
done += 1
if cfg['crash_worker'] == wid and done == cfg['crash_after']:
time.sleep(0.2) # 先让完成消息抵达 master
stop.set()
os._exit(1) # crash-stop
# ---------- Master 进程 ----------
def run_job(lines, M, R, nw, cfg):
mgr = mp.Manager()
hb = mgr.dict()
result_q = mp.Queue()
qs = [mp.Queue() for _ in range(nw)]
procs = []
for w in range(nw):
p = mp.Process(target=worker_main, args=(w, qs[w], result_q, hb, cfg), daemon=True)
p.start()
procs.append(p)
hb[w] = time.time()
chunk = max(1, -(-len(lines) // M))
splits = [lines[i * chunk:(i + 1) * chunk] for i in range(M)]
split_loc = [m % nw for m in range(M)] # 模拟 HDFS 副本所在机器
mstate = {m: 'idle' for m in range(M)}
rstate = {r: 'idle' for r in range(R)}
mgen, rgen = defaultdict(int), defaultdict(int)
mowner, inter, final, mstart, backed = {}, {}, {}, {}, set()
busy = {w: None for w in range(nw)}
alive = {w: True for w in range(nw)}
log, raw_total, shuffled = [], 0, 0
t0 = time.time()
elapsed = 0.0
def assign(w, kind, tid, payload, bump=True, note=''):
if kind == 'map':
if bump:
mgen[tid] += 1
mstart[tid] = time.time()
mstate[tid] = 'in_progress'
busy[w] = ('map', tid, mgen[tid])
qs[w].put(('map', tid, mgen[tid], payload))
else:
rgen[tid] += 1
rstate[tid] = 'in_progress'
busy[w] = ('reduce', tid, rgen[tid])
qs[w].put(('reduce', tid, rgen[tid], payload))
if note:
log.append('[master] ' + note)
while True:
now = time.time()
if now - t0 > cfg['deadline']:
log.append('[master] 超出 deadline, 放弃')
break
# (1) 收集 worker 上报
while True:
try:
kind, tid, gen, wid, payload, extra = result_q.get_nowait()
except queue.Empty:
break
busy[wid] = None
if kind == 'map_done':
if gen != mgen[tid]:
log.append('[master] 丢弃 map#%d 的过期结果(gen=%d, 当前=%d)' % (tid, gen, mgen[tid]))
continue
if mstate[tid] == 'completed':
log.append('[master] 忽略 map#%d 的重复完成消息(来自 W%d): 备份/重试输出去重' % (tid, wid))
continue
mstate[tid] = 'completed'
mowner[tid] = wid
inter[tid] = payload
raw_total += extra
shuffled += sum(len(v) for v in payload.values())
log.append('[master] map#%d 提交成功(W%d), shuffle 记录 %d 条'
% (tid, wid, sum(len(v) for v in payload.values())))
else:
if gen != rgen[tid]:
log.append('[master] 丢弃 reduce#%d 的过期结果(gen=%d, 当前=%d)' % (tid, gen, rgen[tid]))
continue
if rstate[tid] == 'completed':
log.append('[master] 忽略 reduce#%d 的重复完成消息(来自 W%d)' % (tid, wid))
continue
rstate[tid] = 'completed'
final[tid] = payload
log.append('[master] reduce#%d 提交成功(W%d), 输出 %d 条' % (tid, wid, len(payload)))
# (2) 心跳超时 -> 判定 worker 失败
for w in range(nw):
if alive[w] and now - hb.get(w, 0) > cfg['timeout']:
alive[w] = False
log.append('[master] W%d 心跳超时(>%.1fs) -> 判定失败' % (w, cfg['timeout']))
b = busy[w]
busy[w] = None
if b:
kind, tid, _ = b
(mstate if kind == 'map' else rstate)[tid] = 'idle'
log.append('[master] 回收 W%d 上未完成的 %s#%d 为 idle' % (w, kind, tid))
lost = [m for m, o in mowner.items() if o == w]
for m in lost:
mstate[m] = 'idle'
backed.discard(m)
inter.pop(m, None)
del mowner[m]
log.append('[master] map#%d 的中间输出在 W%d 本地磁盘, 不可访问 -> 必须重新执行' % (m, w))
if lost:
for r in range(R):
if rstate[r] == 'in_progress':
rstate[r] = 'idle'
log.append('[master] map 屏障被打破, reduce#%d 回退为 idle' % r)
# (3) 调度: map 阶段 -> 屏障 -> reduce 阶段
maps_done = all(v == 'completed' for v in mstate.values())
idle = [w for w in range(nw) if alive[w] and busy[w] is None]
if not maps_done:
for m in [m for m in range(M) if mstate[m] == 'idle']:
w = split_loc[m]
if w in idle:
idle.remove(w)
assign(w, 'map', m, splits[m], note='map#%d -> W%d (data-local)' % (m, w))
for m in [m for m in range(M) if mstate[m] == 'idle']:
if not idle:
break
w = idle.pop(0)
assign(w, 'map', m, splits[m], note='map#%d -> W%d (rack-local/off-rack)' % (m, w))
else:
for r in [r for r in range(R) if rstate[r] == 'idle']:
if not idle:
break
payload = []
for m in range(M):
payload.extend(inter[m].get(r, []))
payload.sort() # reduce 端归并排序
w = idle.pop(0)
assign(w, 'reduce', r, payload,
note='reduce#%d 的输入已归并排序(%d 条) -> W%d' % (r, len(payload), w))
# (4) straggler -> backup task
if cfg['speculative']:
idle = [w for w in range(nw) if alive[w] and busy[w] is None]
for m in range(M):
if (idle and mstate[m] == 'in_progress' and m not in backed
and now - mstart.get(m, now) > cfg['straggle']):
backed.add(m)
w = idle.pop(0)
assign(w, 'map', m, splits[m], bump=False,
note='map#%d 进度落后 %.2fs -> 在 W%d 启动 backup task'
% (m, now - mstart[m], w))
if maps_done and all(v == 'completed' for v in rstate.values()):
elapsed = time.time() - t0
break
time.sleep(0.01)
# (5) 收尾: 收集在途的迟到结果. 真实 Hadoop 会直接 kill 落败的 task attempt
t_grace = time.time()
while time.time() - t_grace < 0.6:
try:
kind, tid, gen, wid, payload, extra = result_q.get_nowait()
except queue.Empty:
time.sleep(0.02)
continue
busy[wid] = None
if kind == 'map_done' and mstate[tid] == 'completed':
log.append('[master] map#%d 已由另一个副本提交, 忽略 W%d 的迟到结果(输出去重)' % (tid, wid))
elif kind == 'reduce_done' and rstate[tid] == 'completed':
log.append('[master] reduce#%d 已提交, 忽略 W%d 的迟到结果' % (tid, wid))
else:
log.append('[master] 作业已结束, 丢弃 %s#%d 的过期结果(gen=%d)' % (kind, tid, gen))
for q in qs:
q.put(None)
for p in procs:
p.join(timeout=1.0)
if p.is_alive():
p.terminate()
mgr.shutdown()
result = sorted([kv for r in range(R) for kv in final.get(r, [])])
return {'result': result, 'log': log, 'elapsed': elapsed,
'raw_map_records': raw_total, 'shuffle_records': shuffled,
'splits': splits}
def make_corpus(n_lines, seed=2026):
rng = random.Random(seed)
vocab = ['the', 'quick', 'brown', 'fox', 'jumps', 'over', 'lazy', 'dog', 'map',
'reduce', 'shuffle', 'hadoop', 'zookeeper', 'partition', 'combiner',
'backup', 'straggler', 'locality', 'dfs', 'trace']
weights = [20.0 / (i + 1) for i in range(len(vocab))] # 近似 Zipf
return [' '.join(rng.choices(vocab, weights=weights)[0] for _ in range(rng.randint(8, 16)))
for _ in range(n_lines)]
def serial_wordcount(lines):
c = Counter()
for ln in lines:
for w in ln.split():
c[w] += 1
return sorted(c.items())
if __name__ == '__main__':
rng = random.Random(2026)
lines = make_corpus(400)
M, R, NW = 8, 3, 4
base = {'hb_interval': 0.05, 'task_delay': 0.30, 'timeout': 0.6, 'deadline': 60.0,
'R': R, 'use_combiner': True, 'slow': {}, 'slow_delay': 0.0,
'crash_worker': None, 'crash_after': 0, 'straggle': 0.6, 'speculative': True}
truth = serial_wordcount(lines)
print('单机串行 word count: %d 个不同单词, 总词数 %d' % (len(truth), sum(c for _, c in truth)))
# 场景 A: 正常执行 + combiner
a = run_job(lines, M, R, NW, dict(base))
print('\n===== 场景 A: 正常执行 (M=%d R=%d workers=%d combiner=on) =====' % (M, R, NW))
print('split#0 的输入(前 2 行):', a['splits'][0][:2])
buckets = Counter()
for ln in a['splits'][0]:
for w in ln.split():
buckets[partition_of(w, R)] += 1
print('map#0 的 partition 桶大小(每个桶 = 一个 reduce 的输入):',
[buckets[p] for p in range(R)])
print('耗时 %.2fs | map 原始输出 %d 条 | shuffle 传输 %d 条 (压缩比 %.1fx)'
% (a['elapsed'], a['raw_map_records'], a['shuffle_records'],
a['raw_map_records'] / max(1, a['shuffle_records'])))
print('最终结果前 4 项:', a['result'][:4])
assert a['result'] == truth, '场景 A 结果与串行不一致'
print('断言通过: 场景 A 结果 == 单机串行结果')
# 场景 B: 随机挑一个 map worker, 让它在完成第 1 个 map 后崩溃
victim = rng.randrange(1, NW)
b = run_job(lines, M, R, NW, dict(base, crash_worker=victim, crash_after=1))
print('\n===== 场景 B: W%d 在完成 1 个 map 后崩溃 (不再心跳) =====' % victim)
for line in b['log']:
if '失败' in line or '重新执行' in line or '丢弃' in line or '忽略' in line \
or '回退' in line or '回收' in line:
print(' ', line)
print('耗时 %.2fs' % b['elapsed'])
assert b['result'] == truth, '场景 B 结果与串行不一致'
print('断言通过: 场景 B 结果 == 单机串行结果 (worker 崩溃未破坏正确性)')
# 场景 C: straggler, 分别关闭/开启 backup task
slow = dict(base, slow={2: True}, slow_delay=1.2)
c1 = run_job(lines, M, R, NW, dict(slow, speculative=False))
c2 = run_job(lines, M, R, NW, dict(slow, speculative=True))
print('\n===== 场景 C: W2 是 straggler (每个任务多花 1.2s) =====')
print('无 backup task: 耗时 %.2fs' % c1['elapsed'])
print('有 backup task: 耗时 %.2fs (尾部延迟降低 %.0f%%)'
% (c2['elapsed'], 100 * (c1['elapsed'] - c2['elapsed']) / c1['elapsed']))
for line in c2['log']:
if 'backup' in line or '忽略' in line:
print(' ', line)
assert c1['result'] == truth and c2['result'] == truth
print('断言通过: 两种配置结果都 == 单机串行结果')
实际运行输出(固定种子,可完全复现):
单机串行 word count: 20 个不同单词, 总词数 4790
===== 场景 A: 正常执行 (M=8 R=3 workers=4 combiner=on) =====
split#0 的输入(前 2 行): ['quick straggler over straggler partition partition jumps brown fox',
'dog jumps the jumps the shuffle brown shuffle lazy over the the dfs backup reduce']
map#0 的 partition 桶大小(每个桶 = 一个 reduce 的输入): [229, 121, 255]
耗时 0.93s | map 原始输出 4790 条 | shuffle 传输 160 条 (压缩比 29.9x)
最终结果前 4 项: [('backup', 78), ('brown', 449), ('combiner', 88), ('dfs', 74)]
断言通过: 场景 A 结果 == 单机串行结果
===== 场景 B: W1 在完成 1 个 map 后崩溃 (不再心跳) =====
[master] W1 心跳超时(>0.6s) -> 判定失败
[master] 回收 W1 上未完成的 map#5 为 idle
[master] map#1 的中间输出在 W1 本地磁盘, 不可访问 -> 必须重新执行
[master] 丢弃 map#5 的过期结果(gen=1, 当前=2)
耗时 1.67s
断言通过: 场景 B 结果 == 单机串行结果 (worker 崩溃未破坏正确性)
===== 场景 C: W2 是 straggler (每个任务多花 1.2s) =====
无 backup task: 耗时 3.01s
有 backup task: 耗时 1.24s (尾部延迟降低 59%)
[master] map#2 进度落后 0.62s -> 在 W1 启动 backup task
[master] map#2 已由另一个副本提交, 忽略 W2 的迟到结果(输出去重)
断言通过: 两种配置结果都 == 单机串行结果
【代码做什么?】
udf_map/udf_combine/udf_reduce是唯一与作业相关的三个函数,其余全是框架——这正是 MapReduce 抽象边界的体现。run_job把输入行切成 M 个 split(连续分块,模拟 HDFS block 边界),并记录split_loc[m] = m % nw,模拟”这个 split 的副本在哪台机器上”。- 每个 worker 是一个独立进程,带一个心跳线程,每 0.05 秒把
hb[wid] = time.time()写进mp.Manager().dict()——这就是跨进程共享的心跳表。 - master 主循环每轮做四件事:(1) 非阻塞地排空结果队列;(2) 检查心跳超时并执行故障恢复;(3) 调度(先按
split_loc做 data-local 分配,再补 rack/off-rack);(4) straggler 检测与 backup 任务的启动。 - map 任务在 worker 端:逐行调用
udf_map,把中间 KV 按partition_of分桶 → 桶内按 key 排序 → 可选 combiner → 作为一个整体上报给 master。 - reduce 任务:master 端把 M 个 map 的第 r 个分区合并、排序后作为 payload 下发(模拟 shuffle + 归并),worker 端按 key 分组调用
udf_reduce。 - 三个场景分别验证:正常执行、worker 崩溃后正确性不受影响、backup task 把尾延迟从 3.01s 降到 1.24s。
- 最后用
assert把分布式结果与serial_wordcount的串行结果逐项比较。
【分布式机制透视】
- 进程与通道:
mp.Process是 worker,mp.Queue是单向消息通道(master → worker 的任务队列、worker → master 的结果队列)。每个 worker 有自己的任务队列,因此 master 可以点对点指派任务(这正是实现”把任务派给哪台机器”和”data-local 调度”的前提)。 - 心跳与故障检测:心跳线程独立于计算线程,因此计算卡住不会导致心跳停止——这一点很重要,它区分了”慢(straggler)”与”死(failed)”。注入崩溃用的
os._exit(1)会同时杀死两个线程,于是心跳自然停止,master 通过 0.6 秒的超时判定它 failed。 - 状态维护:master 是唯一的状态持有者——
mstate/rstate(任务状态)、mowner(谁做的)、inter(中间数据,模拟 map worker 的本地磁盘)、mgen/rgen(执行代数,用于幂等过滤)。 - 并发与时序:
mgen[tid]是解决”过期消息”的关键。当 map#5 从 gen=1 被重新调度为 gen=2 后,那个在已经死掉的 W1 队列里排队、或者由某个慢 worker 迟到的 gen=1 结果到达时,master 用gen != mgen[tid]一句话就能丢弃它——这就是分布式系统中的”逻辑版本号”模式,与后面会学到的 Lamport 时钟/向量时钟解决”消息乱序”是同一类思想。 - 与原型的对应物:
hb字典 ↔ TaskTracker 心跳表;inter↔ map worker 的本地磁盘(因此 W1 失败时它必须被清空);final↔ HDFS 上的最终输出(因此不受 worker 故障影响);split_loc↔ NameNode 的 block→DataNode 映射;timeout↔mapreduce.task.timeout。 - 真实性的取舍:真实 Hadoop 中 reduce 是主动去 map 节点 fetch 的(pull),本模拟为了简化让 master 把数据下发(push)。此外本模拟中 map worker 失败后,数据还有一份副本在 master 内存里(
inter),代码显式地inter.pop(m)才模拟出”数据丢失”。这两点都是为了在同一台机器上可观察而做的简化,不影响被验证的协议逻辑。
【与理论的对应】
- 主循环的 (1) 段对应算法 5.3.1 的
upon receive ⟨MAP_DONE, ...⟩/⟨REDUCE_DONE, ...⟩,其中gen != mgen[tid]对应伪代码里的a ≠ attempt[m],mstate[tid] == 'completed'对应T_map[m] = COMPLETED——这两条就是算法 5.3.1 中”幂等过滤”的实现,代码中的日志丢弃 map#5 的过期结果(gen=1, 当前=2)与已由另一个副本提交, 忽略 ... 迟到结果分别验证了它们。 - 主循环的 (2) 段对应伪代码的
upon WorkerFailed(w):lost循环对应”已完成的 map 必须重做”(安全性论证 S1 与定理 5.2 第 2 步的工程前提),rstate[r] = 'idle'对应”屏障被打破时正在跑的 reduce 也要回退”(安全性 S2)。 - 主循环的 (3) 段对应伪代码的
── C. 任务分配,其中if not maps_done:与else:的分叉正是屏障的实现(安全性 S2)。 - 主循环的 (4) 段对应伪代码的
── D. Straggler 检测与备份任务;backup[t]对应backed集合;bump=False对应”备份执行不改变attempt“,于是两个副本共享代数、由COMPLETED条件决定谁胜出——这正是论文”先完成者胜”的实现方式。 - 场景 A 的
assert a['result'] == truth验证定理 5.2;场景 B 验证”任意多个 worker 崩溃(只要 master 与输入仍在)结果不变”;场景 C 验证论文 §5.4 的实测结论(关闭 backup 后总时间显著变长)。
5.4.2 Combiner 正确性验证与 shuffle 流量测量
这份代码用一个单进程的 shuffle 模拟器,在同一个作业语义下对比三种配置的 shuffle 流量,并验证 combiner 的合法性条件。
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""Lecture 5 代码 2: Combiner 的正确性验证与收益测量
同一个 job 语义下对比 (a) 无 combiner (b) 有 combiner (c) 用错 combiner
统计 shuffle 阶段真正被发送的记录条数与字节数。
"""
import hashlib
import random
from collections import defaultdict
R = 4 # reduce 任务数 = 分区数
M = 8 # map 任务数 = split 数
def partition_of(key, R):
return int(hashlib.md5(key.encode()).hexdigest(), 16) % R
def make_splits(n_splits=M, lines_per_split=500, seed=7):
"""模拟代理服务器日志: 少数 host 出现极多次 (Zipf 分布)."""
rng = random.Random(seed)
hosts = ['h%03d' % i for i in range(80)]
weights = [1.0 / (i + 1) for i in range(80)]
splits = []
for _ in range(n_splits):
splits.append([(rng.choices(hosts, weights=weights)[0], rng.randint(1, 5000))
for _ in range(lines_per_split)])
return splits
def sum_fn(vals): # 可结合 + 可交换 -> 可做 combiner
return sum(vals)
def mean_fn(vals): # 不可结合 -> 不可做 combiner
return sum(vals) / len(vals)
def run_job(splits, R, combine_fn, reduce_fn):
"""返回 (shuffle 记录数, shuffle 字节数, 最终结果 dict)."""
recs_sent = bytes_sent = 0
buckets = {p: defaultdict(list) for p in range(R)}
for split in splits:
local = defaultdict(list) # 一个 map task 的本地缓冲
for k, v in split:
local[k].append(v)
for k, vs in local.items(): # combiner 在这里生效
emitted = [(k, combine_fn(vs))] if combine_fn else [(k, v) for v in vs]
for kk, vv in emitted:
buckets[partition_of(kk, R)][kk].append(vv) # shuffle
recs_sent += 1
bytes_sent += len(kk) + 8 # key 的 ASCII 字节 + 8 字节数值
result = {}
for p in range(R):
for k, vs in buckets[p].items(): # reduce
result[k] = reduce_fn(vs)
return recs_sent, bytes_sent, result
if __name__ == '__main__':
random.seed(2026)
splits = make_splits()
total_records = sum(len(s) for s in splits)
print('输入: %d 个 split, 共 %d 条 <host, bytes> 记录' % (M, total_records))
# --- 作业 1: reduce 求总和(sum 可结合可交换, combiner 合法) ---
r0, b0, ref = run_job(splits, R, None, sum_fn)
r1, b1, got = run_job(splits, R, sum_fn, sum_fn)
print('\n--- 作业 1: 每 host 的字节总和 (sum) ---')
print('无 combiner : shuffle %6d 条, %8d 字节' % (r0, b0))
print('有 combiner : shuffle %6d 条, %8d 字节' % (r1, b1))
print('记录压缩比 %.1fx, 字节压缩比 %.1fx' % (r0 / r1, b0 / b1))
assert got == ref, 'combiner 改变了 sum 的语义!'
print('断言通过: 有/无 combiner 的最终结果完全相同 (%d 个 key)' % len(ref))
# --- 作业 2: reduce 求平均(mean 不可结合, combiner 非法) ---
_, _, m_ref = run_job(splits, R, None, mean_fn)
_, _, m_bad = run_job(splits, R, mean_fn, mean_fn)
diff = [k for k in m_ref if abs(m_ref[k] - m_bad[k]) > 1e-9]
print('\n--- 作业 2: 每 host 的平均响应字节 (mean) ---')
print('无 combiner : %d 个 key, 例如 %s = %.2f' % (len(m_ref), 'h000', m_ref['h000']))
print('有 combiner : 其中 %d 个 key 的结果被改变, 例如 h000 = %.2f' % (len(diff), m_bad['h000']))
assert diff, 'mean 居然可以用作 combiner? 与理论不符'
print('断言通过: 用 mean 做 combiner 得到错误结果 -> 违反结合律的函数不能当 combiner')
print(' 原因: map#i 的局部均值没有携带"该局部有多少个样本"这一权重信息,')
print(' 而各 split 里 h000 的样本数并不相等, 均值不可结合。')
实际运行输出:
输入: 8 个 split, 共 4000 条 <host, bytes> 记录
--- 作业 1: 每 host 的字节总和 (sum) ---
无 combiner : shuffle 4000 条, 48000 字节
有 combiner : shuffle 572 条, 6864 字节
记录压缩比 7.0x, 字节压缩比 7.0x
断言通过: 有/无 combiner 的最终结果完全相同 (80 个 key)
--- 作业 2: 每 host 的平均响应字节 (mean) ---
无 combiner : 80 个 key, 例如 h000 = 2505.32
有 combiner : 其中 80 个 key 的结果被改变, 例如 h000 = 2502.81
断言通过: 用 mean 做 combiner 得到错误结果 -> 违反结合律的函数不能当 combiner
原因: map#i 的局部均值没有携带"该局部有多少个样本"这一权重信息,
而各 split 里 h000 的样本数并不相等, 均值不可结合。
【代码做什么?】
make_splits生成 M=8 个 split,每个 500 条<host, bytes>记录;host 按 Zipf 分布抽取,模拟真实日志的数据倾斜。run_job是最小化的 MapReduce 引擎:每个 split 是一个 map 任务,本地按 key 聚成local字典(这一步是 map 的输出缓冲),然后是否调用 combiner 由combine_fn是否为None决定。- 每条被 emit 的记录都记账:
recs_sent += 1、bytes_sent += len(kk) + 8——这就是 shuffle 阶段真正过网络的数据量。 - 记录按
partition_of(MD5 稳定哈希)进入 R=4 个分区,每个分区各自做 reduce。 - 作业 1 用
sum:验证有/无 combiner 的结果完全相同,并打印流量压缩比。 - 作业 2 用
mean:验证有/无 combiner 的结果必然不同——这是对 5.3.2 中结合律论证的反例验证。
【分布式机制透视】
- map 端的本地归约:
local = defaultdict(list)对应 Hadoop 的 map 输出缓冲区;combine_fn(vs)对应溢写时在每个 partition 内做的 combiner。真实系统中这个缓冲区是固定大小的环形缓冲(默认 100MB,80% 触发溢写),因此 combiner 会被调用多次——本代码为了简化只调用一次,但正确性要求(允许被调用任意次)是一样的。 - shuffle 流量:真实系统中 shuffle 走 HTTP,reduce 端主动 fetch;本代码用”往
buckets[p]里 append”来代表”跨网络的发送”,用recs_sent/bytes_sent直接量化流量。 - 数据倾斜:Zipf 分布让
h000的记录数远超其他 host。这正是 combiner 收益最大的场景(7.0x),也是哈希分区在负载均衡上无能为力的场景(见 5.3.3 的 $\max(\vert V\vert /R, L^*)$ 下界)。
【与理论的对应】
- 作业 1 的
assert got == ref直接验证算法 5.3.2 中”有 combiner 时安全性”的归纳论证:$f(V_1 \uplus \cdots \uplus V_k) = f(\{f(V_1),\dots,f(V_k)\})$。这里 $k = M = 8$(每个 map 一个局部归约结果)。 - 作业 2 的
assert diff验证反面:mean 不满足结合律,所以combine_fn = mean_fn是非法的,结果必然错。它同时说明了 combiner 正确性检查的实操判据——问自己”f(f(V1), f(V2), ...)是否等于f(V1 ∪ V2 ∪ ...)“。 - 压缩比 7.0x 与 5.3.3 中的公式 $S_{\text{with combine}} = \sum_{m=1}^{M} k_m$ 吻合:本例中每个 split 贡献的不同 host 数约为 $572/8 \approx 71.5$ 个,而每个 split 有 500 条记录。
5.5 性能与可扩展性分析
5.5.1 总时间模型
把一次 MapReduce 作业的墙钟时间拆开,可以粗略写成
\[T_{\text{job}} \;\approx\; T_{\text{startup}} \;+\; T_{\text{map}} \;+\; T_{\text{shuffle}} \;+\; T_{\text{reduce}} \;+\; T_{\text{tail}}\]- 启动开销 $T_{\text{startup}}$:把用户程序分发到所有 worker、打开输入文件、获取局部性信息。论文实测 grep 作业总耗时约 150 秒,其中约 60 秒是启动开销——这是一个不可忽略的常数项,也是 MapReduce 不适合小作业与交互式查询的直接原因。
- map 阶段:$T_{\text{map}} \approx \left\lceil \frac{M}{W} \right\rceil \cdot t_{\text{map}}$,其中 $W$ 是可用 worker 数、$t_{\text{map}}$ 是单个 map 任务耗时。要做到线性加速,前提是 $M \gg W$(见 5.5.3 的任务粒度)。
- shuffle 阶段:$T_{\text{shuffle}} \approx \dfrac{S}{B_{\text{net}}}$,其中 $S$ 是需要跨网络的中间数据量,$B_{\text{net}}$ 是集群可用带宽。shuffle 可以与后续 map 任务重叠(论文 sort 实验显示 “shuffling starts as soon as the first map task completes”),因此它是流水化的,不是严格串行的一段时间。
- reduce 阶段:$T_{\text{reduce}} \approx \left\lceil \frac{R}{W} \right\rceil \cdot t_{\text{reduce}}$,其中 $t_{\text{reduce}}$ 包含拉取 + 归并排序 + 逐 key 调用 REDUCE + 写输出。
- 尾部 $T_{\text{tail}}$:由 straggler 造成。这一项是 backup task 专门对付的。它的量级可以非常大:论文实测 sort 作业关闭 backup 后,最后 5 个 reduce 任务多花了 300 秒,使总时间从 891 秒涨到 1283 秒。
5.5.2 Shuffle 通信量:上界与 combiner 的实际效果
| 配置 | shuffle 记录数 $S$ | 说明 |
|---|---|---|
| 无 combiner | $S = \lvert V \rvert$(全部中间 KV 记录) | 每条 map 输出记录都要过网络;$V$ 由数据决定,与 $M$、$R$ 无关 |
| 有 combiner | $S = \sum_{m=1}^{M} k_m \le \min(\lvert V \rvert,\; M \cdot K)$ | $k_m$ = split $m$ 中不同 key 数,$K$ = 全局不同 key 数 |
| fetch 连接数(最坏) | $O(M \cdot R)$ | 每个 reduce 要向每个 map 拉一次;小文件与连接数随 $M \cdot R$ 增长 |
| master 元数据 | $O(M \cdot R)$ 字节 | 每对 (map, reduce) 约 1 字节 |
| 实测(5.4.2 代码,Zipf 输入) | $4000 \to 572$ 条,压缩 7.0x | 输入倾斜越重,combiner 收益越大 |
| 实测(5.4.1 代码,20 词表) | $4790 \to 160$ 条,压缩 29.9x | 词表越小(不同 key 越少),收益越大 |
注意 $S$ 的两个极端:当每个 map 的中间结果几乎全是不同 key 时($k_m \approx \vert V_m\vert $),combiner 几乎没用;当 key 高度重复时(Zipf),combiner 能把 $S$ 从 $\vert V\vert $ 压到 $M \cdot K$。因此 combiner 的收益完全取决于中间 key 的重复度,而不是取决于数据总量。
5.5.3 M 与 R 的取值权衡
论文 §3.5 的核心结论是:理想情况下 M 与 R 都应当远大于 worker 机器数。让每个 worker 做很多个任务的好处是:(1) 更好的动态负载均衡(快的机器自然拿到更多任务);(2) 加速故障恢复——一台机器上已经完成的许多 map 任务可以分散到其余所有机器上重做。
| 维度 | $M$(map 任务数) | $R$(reduce 任务数) |
|---|---|---|
| 由什么决定 | 由输入切分决定:$M = \lceil \lvert D \rvert / \text{split size} \rceil$,与 HDFS block 对齐(论文 16–64MB,HDFS 2.x 默认 128MB) | 用户指定;每个 reduce 产生一个独立输出文件 |
| 变大的好处 | 更好的动态负载均衡;故障恢复更快(已完成 map 可分散重做);单任务输入更小,更容易做 data-local | 更高的 reduce 并行度;单个输出文件更小,下游处理更灵活 |
| 变大的代价 | master 需要 $O(M+R)$ 次调度决策;每个 map 任务有自己的启动开销(进程/JVM 启动、输入打开) | 输出文件数爆炸;master 内存状态 $O(M \cdot R)$;每个 map 任务要维护 $R$ 个分区文件句柄与缓冲 |
| 经验取值 | 让单个任务处理 16–64MB 输入(这样局部性优化最有效);论文实际用 $M = 200{,}000$,配 2000 台机器 | $R$ 取”预期使用机器数的小倍数“;论文实际用 $R = 5{,}000$ |
| 极端情况 | $M=1$:完全没有 map 并行度,局部性优化失效 | $R=1$:完全没有 reduce 并行度,但输出是单文件、无分区不均衡问题(grep 实验就是 $M=15000, R=1$) |
| 硬性上界 | master 的调度决策数 $O(M+R)$ | master 的内存 $O(M \cdot R)$(每对约 1 字节);$R$ 还常被用户因”输出文件数”而主动限制 |
5.5.4 论文实测数据
| 实验 | 配置 | 结果 |
|---|---|---|
| Grep | $10^{10}$ 条 100 字节记录(约 1TB),$M=15000$,$R=1$,约 1800 台机器 | 峰值扫描速率 > 30 GB/s(1764 个 worker 时);约 80 秒时 map 速率归零;总耗时约 150 秒,其中约 60 秒为启动开销 |
| Sort | $10^{10}$ 条 100 字节记录(约 1TB),$M=15000$,$R=4000$,输出 2 副本共 2TB | 输入速率峰值约 13 GB/s(低于 grep,因为 map 要花一半时间和 I/O 写本地中间结果);shuffle 从第一个 map 完成就开始,约 600 秒结束;输出写入约 850 秒结束;总耗时 891 秒(当时 TeraSort 公开最好成绩为 1057 秒) |
| Sort,关闭 backup task | 同上 | 960 秒时还有 5 个 reduce 未完成,它们又跑了 300 秒;总耗时 1283 秒,比开启时慢 44% |
| Sort,故意杀掉 200 / 1746 个 worker 进程 | 同上,作业开始数分钟后杀进程(机器本身仍正常,调度器立刻重启了新 worker 进程) | 图上出现负的输入速率(因为已完成的 map 工作随机器消失而需要重做);重执行很快完成;总耗时 933 秒,仅比正常执行慢 5% |
| Google 2004 年 8 月生产统计 | 29,423 个作业 | 平均作业时长 634 秒;消耗 79,186 机器·天;读入 3,288 TB 输入;产生 758 TB 中间数据;写出 193 TB 输出;平均每作业 157 台 worker;平均每作业 1.2 台 worker 死亡;平均 3,351 个 map 任务、55 个 reduce 任务;共 395 种 map 实现、269 种 reduce 实现、426 种组合 |
从这些数字里能读出的设计结论:
- “平均每作业 1.2 台 worker 死亡”——故障不是异常,是日常。这解释了为什么容错必须是框架的第一等公民,而不是可选项。
- “输入 3288TB → 中间数据 758TB → 输出 193TB”——输入远大于输出,说明 MapReduce 的主要负载是”从海量数据里提炼少量信息”(grep 类)与”重排数据表示”(sort 类);中间数据 758TB 体现了 shuffle 的巨大开销。
- 杀 200 台机器只慢 5%——这是”重算代替恢复”的效率证明:因为 $M$ 远大于机器数,被重做的 map 任务可以立即分散到大量空闲机器上并行重跑。
- 关闭 backup 慢 44%——尾部延迟可以主宰总时间。
5.5.5 MapReduce 不适合的场景
| 场景 | 为什么不适合 | 量化原因 | 后来靠什么解决 |
|---|---|---|---|
| 迭代式算法(PageRank、K-means、梯度下降) | 每轮迭代都要重新从 HDFS 读、shuffle、写回 HDFS;磁盘 I/O 主导,CPU 空转 | 每轮至少 2 次全量落盘 + 1 次全量 shuffle;10 轮迭代 = 10 个 MR 作业 | Spark:RDD 常驻内存 + lineage 重算(Lecture 24) |
| 交互式查询 | 启动开销是几十秒量级 | grep 实验 150 秒中有约 60 秒是启动开销 | Hive/Impala/Presto(MPP 执行引擎)、Spark SQL |
| 小数据 | 框架固定开销与数据量无关 | 几 MB 数据用 MR 比单进程 Python 慢几个数量级 | 直接单机处理,或 Spark local mode |
| 低延迟 / 流式处理 | 屏障(barrier)要求所有 map 完成才能开始 reduce;批处理天生高延迟 | 屏障意味着 $T \ge \max_m t_{\text{map}}(m)$ | 流处理系统(Lecture 23):无全局屏障、连续算子 |
| 多阶段复杂查询 | 每个 MR 作业都要落盘再读回 | $k$ 个阶段的查询 = $k$ 次全量落盘 | Tez/Spark 的 DAG 执行引擎 |
| 需要强事务/更新语义 | MR 输出是”一次写、多次读”的批产物,不支持就地更新 | —— | HBase / 数据库(OLTP 场景) |
这些局限如何催生了 Spark(Lecture 24):Spark 保留了 MapReduce 的三件核心遗产——shuffle、分区、用重算做容错——但做了三处关键改造:(1) 把中间结果留在内存(RDD 是分布式的、不可变的分区集合),消除迭代式负载的磁盘往返;(2) 把多个阶段串成 DAG(一个 job 可以有多个窄依赖/宽依赖阶段,不必每个阶段都落 HDFS);(3) 把”重算”的粒度从”任务”细化到”分区 + lineage”:RDD 记住它是怎么从父 RDD 算出来的(lineage),丢失的分区可以只重算它自己那一支依赖链,而不必重跑整个阶段的任务。从”重算任务”到”重算分区”,正是 MapReduce 容错思想在 Spark 里的自然演进。
5.6 关键要点
- 设计哲学:用”重新执行”代替”恢复状态”。 MapReduce 不维护任何可恢复的中间状态,只要求用户函数确定性,于是”重算”与”第一次算”结果相同,容错逻辑退化成一句”把它再做一遍”。简单,且正确性容易论证——这是分布式系统容错的经典范式。
- 抽象边界的划法决定了系统的成败。 用户只写
map与reduce两个纯函数;框架包办并行化 map、shuffle、并行化 reduce、四阶段存储与屏障。表达的狭窄换来的是自动并行化与透明容错,而不是能力的损失——因为绝大多数批处理作业都能塞进这个形状。 - 屏障是整个模型的骨架。 “所有 map 完成才能开始 reduce”保证了 reduce 看到的每个 key 的 value 集合是完整的。一旦有 map 被重做,屏障必须重新建立——这也是”map 故障会导致正在跑的 reduce 一并回退”的原因。
- 数据放在哪块盘上,决定了容错怎么做。 map 输出在 map 机器的本地磁盘(所以已完成的 map 必须重做)↔ reduce 输出在全局 DFS(所以已完成的 reduce 不必重做)。这一条不对称性是 MapReduce 容错规则的唯一来源。
- 原子提交把”随机故障”变成”全有或全无”。 每个任务先写私有临时文件,完成时原子 rename。重执行导致的多次 rename 因为”结果相同 + rename 原子”而无害,用户永远看不到半个作业的输出。
- 尾部延迟必须用冗余买时间。 最慢的任务决定总时间,因此备份任务(推测执行)用几个百分点的额外资源换掉 44% 的尾部延迟。这个权衡之所以划算,是因为 straggler 是极少数而空闲资源是多数。
- 确定性是用户与框架之间的契约,不是附加条款。 框架能保证”顺序无关、等价于串行执行”,前提是
map/reduce确定性。非确定函数下语义退化为”每个 reduce 各自等价于某次串行执行,但彼此可能对应不同的执行”——这是只在故障时才暴露的、极难调试的一类 bug。
5.7 常见陷阱与注意事项
- 误以为”已完成的 map 任务不会重做”。
- 为什么错:map 的输出写在 map worker 的本地磁盘上,那台机器一挂,这份数据就彻底不可访问了。若不重做,reduce 会少拉到一个 map 的全部数据,从而得到静默错误的(偏小的)聚合结果——这比崩溃危险得多。
- 正确做法:worker 被判 failed 时,把它负责的所有 map 任务(无论 IDLE / IN_PROGRESS / COMPLETED)都重置为 idle;同时让所有仍在运行的 reduce 重新拉取数据(论文明确要求通知所有 reduce worker 改从新的 map worker 读取)。
- 把 combiner 当成”随便写都行”的优化。
- 为什么错:combiner 会在任意次数、任意分组的子集上被调用(每次溢写、每次归并都可能调),因此它必须满足结合律与交换律。用
mean、median或count distinct当 combiner,结果会悄悄出错(见 5.4.2 的作业 2:80 个 key 全部被改错)。 - 正确做法:先用判据自检——
REDUCE(k, V1 ∪ V2 ∪ …)是否等于REDUCE(k, {COMBINE(k,V1), COMBINE(k,V2), …})?若是,就可以直接把 reduce 函数复用为 combiner(这也是论文推荐的常见做法);若否(如均值),就老老实实不加 combiner,或者在 combiner 里额外携带样本数(sum, count)并让 reduce 合并这两个量。
- 为什么错:combiner 会在任意次数、任意分组的子集上被调用(每次溢写、每次归并都可能调),因此它必须满足结合律与交换律。用
- 用 Python 内置的
hash()做分区函数。- 为什么错:CPython 对
str的hash()默认启用随机化种子(PYTHONHASHSEED),因此同一个字符串在不同进程里的哈希值不同。于是一个 key 在 map 任务 A 里被分到 p=0、在 map 任务 B 里被分到 p=2,同一个 key 的 value 被拆到两个 reduce,各自算出一个部分结果。这种错误没有任何报错,只是结果偏小。 - 正确做法:使用稳定的哈希,例如
int(hashlib.md5(k.encode()).hexdigest(), 16) % R。真实 Hadoop 用HashPartitioner(基于 key 的hashCode(),对 Java 而言是稳定且跨 JVM 一致的)。
- 为什么错:CPython 对
- 在
reduce里依赖 value 列表的顺序。- 为什么错:
list(v2)的顺序由各 map 任务的完成顺序、shuffle 的到达顺序决定,不确定;而且重执行/备份任务会让顺序发生变化。依赖顺序的 reduce(例如”取第一个元素当代表”、”输出时保留输入顺序”)就变成了非确定性函数,从而丧失串行等价性保证。 - 正确做法:把 reduce 写成一个对列表顺序不敏感的纯函数(求和、取最大、拼接但先排序等)。若确实需要顺序,就在 reduce 内显式排序,但要注意这会增加内存与时间开销。
- 为什么错:
- 以为 master 有高可用。
- 为什么错:论文的实现里 master 是单点:它的数据结构(任务状态、中间文件位置表)只存在内存里,master 一挂,整个作业就中止,用户只能重跑。Hadoop 1.x 的 JobTracker 同样如此。
- 正确做法:理解这一代系统的边界——master 靠”client 重试”而不是”自动接管”来容错。Hadoop 2.x/YARN 把它拆成全局 RM(checkpoint + 备用 RM)+ 每作业一个 AM,才真正改善了这一问题。
- 把 M 和 R 设得越大越好。
- 为什么错:master 必须做 $O(M + R)$ 次调度决策并保存 $O(M \cdot R)$ 的状态;每个 map 任务的启动开销是固定的;每个 reduce 任务产生一个独立输出文件。$M$、$R$ 太大时,调度与元数据开销会超过并行化带来的收益,还会产生海量小文件。
- 正确做法:按论文 §3.5 的经验值——让每个 map 任务处理 16–64MB 输入(保证局部性),让 $R$ 是预期机器数的小倍数。
- 在 map 里产生副作用(写文件、发请求、累加全局计数器)。
- 为什么错:map 任务可能被执行多次(故障重做、备份任务),副作用会被执行多次;而且备份任务的输出最终会被丢弃,副作用却已经发生了。论文 §4.5 明确要求用户自己保证副作用的原子性与幂等性。
- 正确做法:把副作用写成”先写临时文件、完成后原子 rename”;把计数工作交给框架提供的 counter 机制(master 会在聚合时消除重复执行带来的重复计数)。
- 认为 reduce 可以”看到全局信息”。
- 为什么错:一个 reduce 任务只能看到属于它自己那个分区的 key;不同 reduce 任务之间没有任何通信,也不能共享全局变量。所以”求出出现次数最多的那个词”这种需要跨 reduce 比较的问题是一个 MR 作业做不到的(讲义中的练习正是问这个)。
- 正确做法:串联第二个 MapReduce 作业——第一个作业输出
<word, count>,第二个作业的 map 把所有记录重新 key 成同一个 key(例如1),只用一个 reduce;该 reduce 就能看到全部<word, count>并选出最大值。或者用自定义 partitioner 保证全局唯一 reduce。
- 忘记输入数据在作业期间不能变。
- 为什么错:如果输入 split 在中途被修改(例如上游作业正在写同一个目录),那么被重执行的 map 任务会读到不同的数据,串行等价性被打破。
- 正确做法:把 map 的输入当作不可变的;上游作业必须完全结束并原子提交后,下游才能开始读取。
5.8 思考题(带答案)
Q1(计算题) 某集群有 2000 台机器,每个 worker 同时只跑一个任务。输入数据 1TB,HDFS block 为 128MB。作业设 $R = 4000$。(a)应该设 $M$ 为多少?简述理由。(b)若每个 map 任务处理一条 split 平均耗时 20 秒,每个 reduce 任务平均耗时 60 秒,忽略启动开销、shuffle 与故障,估算总时间。(c)若某个 map worker 在执行了 5 个 map 任务后崩溃,master 需要重做多少个 map 任务?为什么?
答:
- (a)按 HDFS block 对齐,$M = 1\text{TB} / 128\text{MB} = 1024 \times 1024 / 128 = 8192$ 个 map 任务。若按论文经验(每任务 16–64MB)则 $M$ 可以取 16384–65536。关键是要让 $M \gg$ 机器数(2000),以取得动态负载均衡与快速故障恢复。
- (b)$T_{\text{map}} \approx \lceil M / W \rceil \cdot 20 = \lceil 8192/2000 \rceil \times 20 = 5 \times 20 = 100$ 秒。$T_{\text{reduce}} \approx \lceil R / W \rceil \cdot 60 = \lceil 4000/2000 \rceil \times 60 = 2 \times 60 = 120$ 秒。由于屏障的存在,reduce 不能在 map 全部完成前开始,故 $T \approx 100 + 120 = 220$ 秒(实际中 shuffle 能与 map 尾部重叠,因此会略优于这个上界)。
- (c)6 个:该 worker 上已完成的 5 个 map 任务因为输出只在它的本地磁盘上、已不可访问,必须全部重做;再加上它崩溃时正在执行的那 1 个。注意它做过的 reduce 任务(如果有)不需要重做,因为 reduce 的输出已经通过原子 rename 落到全局 DFS 上了。
Q2(”直观但错误的想法”辨析题) 有同学提出:”既然网络是瓶颈,那我们可以把 map 的输出直接写进 HDFS,这样 worker 挂了数据也还在,就不用重做已完成的 map 任务了。这个优化显然更好。” 请指出这个想法错在哪里。
答:错在它用错误的方式解决了一个已经被解决的问题,而且代价很高:
- 它破坏了局部性优化的收益。 map 输出是中间数据,体量往往与输入同量级(论文实测:3288TB 输入产生 758TB 中间数据)。把这些数据写进 HDFS 意味着一次完整的网络传输 + 3 副本复制(3 倍网络与磁盘写入)。而把中间数据写本地磁盘,之后只有真正需要它的那个 reduce 任务去拉取一次(1 次网络传输,且只传一份)。在 sort benchmark 中,正是因为输入和中间结果尽量走本地磁盘,输入速率才能达到 13–30 GB/s 而网络不至于被打爆。
- 它并不能消除重做。 即使中间数据在 HDFS 里,一旦某台 map worker 崩溃、它的 block 副本只剩 2 个,HDFS 会立刻启动副本重建——这同样是一次完整的网络传输,只是把成本从”重算”挪到了”复制”,而且复制量可能更大(重算只重算丢失的那几个任务,复制要维持 3 副本)。
- 重算本来就不贵。 论文实测杀掉 200/1746 个 worker 进程,总时间只增加 5%。原因正是 $M$ 远大于机器数:被重做的 map 任务可以立刻分散到大量空闲机器上并行重跑。
- 真正的设计哲学是”重算比恢复便宜”:本地写、远程读,加上”用重新执行代替状态恢复”,整体上是最优的工程权衡。如果想减少重算量,正确的方向是增大 $M$(把重算粒度变细)而不是把中间数据搬进 DFS。
Q3(推演题) 词频统计作业中,某个 map 任务的输入 split 是 "a b a a b",$R = 2$,分区函数 $P(w) = \text{len}(w) \bmod 2$(len("a")=1,len("b")=1)。另一位 reduce worker 处理的 split 是 "b c c"。(a)分别给出两个 map 任务在有 combiner 与无 combiner 时产生的 shuffle 记录内容。(b)验证最终结果与串行统计一致。(c)如果分区函数改成 P(w) = len(w) mod 3 而 $R$ 仍然是 2,会发生什么?
答:
- (a)$P(a) = 1 \bmod 2 = 1$,$P(b) = 1$,$P(c) = 1$。所有 key 都落到 p=1(这是这个玩具分区函数的一个真实缺陷:所有单字母词的哈希值相同)。
- map#0(
"a b a a b"):无 combiner 时输出 5 条:(a,1)(b,1)(a,1)(a,1)(b,1);有 combiner 时先在本桶内按 key 分组求和,输出 2 条:(a,3)、(b,2)。 - map#1(
"b c c"):无 combiner 时输出 3 条:(b,1)(c,1)(c,1);有 combiner 时输出 2 条:(b,1)、(c,2)。 - shuffle 总量从 8 条降到 4 条。
- map#0(
- (b)reduce#1 的输入(归并排序后)为
(a,1)(a,1)(a,1)(b,1)(b,1)(b,1)(c,1)(c,1);按 key 分组后a→[1,1,1]、b→[1,1,1]、c→[1,1];求和得a→3, b→3, c→2。串行统计:a出现 3 次、b出现 3 次、c出现 2 次。两者完全一致 ✅。有 combiner 时 reduce 看到的输入是(a,3)(b,2)(b,1)(c,2),分组后a→[3]、b→[2,1]、c→[2],求和仍为3, 3, 2✅——这正是结合律与交换律在起作用。reduce#0 的输入为空,输出空文件。 - (c)分区函数与 $R$ 的组合是合法的:$P: K \to \{0,1\}$,$R=2$ 仍然是”值域与桶数匹配”。改成
mod 3而 $R=2$ 也不会破坏正确性——因为所有 map 任务用的都是同一个 $P$ 与同一个 $R$,同一个 key 依然会稳定地落到同一个 reduce(P(a)=1, P(b)=1, P(c)=0)。真正会破坏正确性的错误是:不同 map 任务用了不同的 $R$ 或不同的 hash 函数(例如某个 worker 用了 Python 内置hash(),其值随进程变化),那样同一个 key 会被送到两个不同的 reduce,各自算出部分结果。此外要注意:mod 3的值域是 $\{0,1,2\}$ 而只有 2 个 reduce,桶 2 永远不会被使用——这是负载均衡问题(reduce#1 承担全部负载)而不是正确性问题。
Q4(应用题) 一个作业需要计算”每个用户的会话平均时长”,其中 map 输出 <user, 会话时长>,reduce 计算平均值。(a)能否用 mean 作为 combiner?为什么?(b)如果不能,如何改造使得既有 combiner 的收益又保持正确?(c)如果某个用户有 5 亿条会话记录,而其他用户只有几百条,会出现什么问题?如何缓解?
答:
- (a)不能。均值不满足结合律:
mean(mean(1,2), 3) = mean(1.5, 3) = 2.25 ≠ mean(1,2,3) = 2。根本原因是局部均值丢失了”样本数”这一权重信息,而各 map 任务里该用户的会话数并不相等。若强行使用,每个用户的结果都会被改错(这正是 5.4.2 作业 2 的结论)。 - (b)改造 value 的类型,使其携带足够信息:让 map 输出
<user, (sum, count)>,combiner 做逐分量相加((s1,c1) + (s2,c2) = (s1+s2, c1+c2)),reduce 端先合并所有(sum, count)再算sum / count。逐分量加法满足结合律与交换律,因此 combiner 合法,且 reduce 得到的仍是全局的 sum 与 count,结果与串行完全一致。这是一个通用技巧:把不可结合的聚合改写为”可结合的充分统计量(sufficient statistic)”的聚合。 - (c)数据倾斜(data skew)。由 5.3.3 的下界 $\max_r L_r \ge \max(\vert V\vert /R, L^*)$:承载这个超级用户的 reduce 任务的负载至少是 5 亿条,而其他 reduce 可能只有几千条——一个 reduce 跑了几个小时,其余全部空闲,整个作业被它拖住(这正是 MapReduce 版的 straggler)。
- 缓解手段:① combiner:
(sum, count)的可结合性使得每个 map 只发出 1 条记录,把该 reduce 的输入从 5 亿条压到 $M$ 条($M$ 可能只是几万)——这一步通常就足够了;② 若单 key 仍然过重($M$ 本身很大),用加盐(salting):map 输出<user#suffix, (sum,count)>,其中suffix = hash(...) mod K,用 K 个 reduce 并行做第一轮部分聚合,再跑第二个 MapReduce 作业做第二轮最终聚合;③ 极端情况下改用范围分区 + 采样,把重 key 单独拆出来处理。
- 缓解手段:① combiner:
