Lecture 5: MapReduce 与 Hadoop —— 大规模分布式数据处理(MapReduce and Hadoop)

目录 · ← l4 · l6 →

第四部分:通信原语与分布式数据结构

这一部分回答:多个进程之间如何高效、可靠、可扩展地交换信息与状态? 从「重新执行代替恢复」的 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 给出的答案是:把计算限制在一个极窄的编程模型里——用户只写 mapreduce 两个纯函数,其余四件事(并行化 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)

  • 定义与目的mapreduce 两个词直接借自 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)名称职责
ClientClientClient提交 job、指定 M/R 与分区函数、等待结果;master 挂掉时由它决定是否重试
MasterJobTrackerApplicationMaster(每 job 一个)保存全部任务状态(idle / in-progress / completed)、分配任务、ping worker、记录中间文件位置表、触发 backup task
WorkerTaskTrackerNodeManager + Container实际执行 map / reduce 任务,把中间结果写本地磁盘并上报位置
DFSHDFSHDFS存放输入 split 与最终输出;提供 3 副本容错;为局部性调度提供 block 位置信息
CombinerCombinerCombiner可选。在 map 端本地对同一 partition 内同一 key 的 value 做预聚合,削减 shuffle 数据量
PartitionerPartitionerPartitioner决定中间 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 个文件    |
                                 +==============================================+
  • 逐步解说(对应上图编号)
    1. 切分并启动:用户程序里的 MapReduce 库把输入文件切成 M 份(每份 16–64MB),然后在集群上启动这份程序的许多副本。
    2. 选出 master:其中一份副本是特殊的——master,其余的 worker 由 master 分配工作。共有 M 个 map 任务与 R 个 reduce 任务。master 挑选空闲 worker,分别派发 map 或 reduce 任务。
    3. 执行 map:被派到 map 任务的 worker 读取对应的输入 split,解析出键值对,逐条交给用户 Map 函数;产生的中间 KV 先缓存在内存里
    4. 溢写本地磁盘:缓冲区的键值对周期性地按分区函数写成 R 个区域落到本地磁盘。这些区域的位置被回报给 master,master 负责把它们转发给 reduce worker。
    5. shuffle + 排序:reduce worker 收到 master 关于位置的通告后,用 RPC 从各 map worker 的本地磁盘读取缓冲数据。读完全部中间数据后,按中间 key 排序,使同一个 key 的所有出现聚在一起。如果中间数据大到装不进内存,就使用外部排序(external sort)
    6. 执行 reduce:reduce worker 遍历已排序的中间数据,每遇到一个不同的中间 key,就把该 key 与对应的 value 集合交给用户 Reduce 函数;Reduce 的输出追加到该 reduce 分区的最终输出文件
    7. 唤醒用户程序:所有 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> 会出现成百上千次),如果能在本地先合并,就能大幅减少需要过网络的数据。
  • 直观解释: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 taskspeculative 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 端归并排序(不能用哈希分区)
JOINmap-side join(小表广播)或 reduce-side join(按 join key 分区)
视图/多步查询串联多个 MapReduce 作业
  • MapReduce 相对 SQL 的五个根本性局限
    1. 单遍(single pass):一个 MapReduce 作业只做一次 map + 一次 reduce,没有迭代。复杂的多步分析必须串成一条 MapReduce 作业链,而每一个作业都要落盘再读回
    2. 无 schema、无索引、无统计信息:优化器没有任何可依据的元数据。而 MapReduce 干脆没有优化器——M、R、combiner、分区函数全由用户手工指定
    3. 迭代弱(weak iteration):PageRank、K-means、梯度下降这类算法需要反复扫描同一份数据,每轮都要”从 HDFS 读 → shuffle → 写回 HDFS”。磁盘 I/O 主宰了运行时间,而 CPU 大部分时间在等 I/O。
    4. 高延迟 / 交互式查询不可行:作业启动开销很大——grep 实验总耗时 150 秒里约有 60 秒是启动开销(程序分发到所有 worker、与 GFS 交互打开 1000 个输入文件、获取局部性信息)。这种量级的延迟无法支持交互式查询。
    5. 不适合小数据:框架的固定开销(进程启动、调度、屏障同步)与数据量无关,几 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$ 有限)。
  • 用户函数假设MAPREDUCE确定性函数,且在有限时间内终止。

伪代码

─────────────────────────── 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] ← COMPLETEDowner[1] ← W1loc[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。于是:
    1. busy[W1] = ('map',5)T_map[5] ← IDLE(未完成的任务回收);
    2. owner[1] = W1 ∧ T_map[1] = COMPLETEDT_map[1] ← IDLEloc[1][*] ← ⊥已完成的 map 也必须重做);
    3. 因为发生了第 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→0Welcome→1Hello→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)
  • reduceEveryone 调用一次得到 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$ 上都产生了中间记录。

  1. $m_1$ 计算 $P(k) = \text{hash\_stable}(k) \bmod R = p_1$,把记录写入它自己的第 $p_1$ 个分区文件。
  2. $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$ 个分区文件。
  3. reduce 任务 $r$ 的输入定义为”所有 $m$ 的第 $r$ 个分区文件的并”(见 reduce 端伪代码:只 fetch loc[m][r])。若 $p = r$,则两条记录都进入 $r$ 的输入;若 $p \ne r$,则两条都不进入。
  4. 因此对任意 key $k$,其全部中间 value 恰好出现在唯一一个 reduce 任务 $r = P(k)$ 的输入中,且不会出现在任何其他 reduce 的输入中。
  5. 由 5.3.2 的安全性论证,$r$ 的 REDUCE 看到的是 $k$ 的完整 value 多重集。∎
    • 三条假设各自的作用:假设 (i) 保证不同 map 不会用不同的 $R$(例如某个 worker 误用了默认值);假设 (ii) 保证同一个 key 的哈希值跨机器一致(用 Python hash() 就会违反这一条,导致同一 key 被送到两个不同 reduce,结果静默错误——每个 reduce 各算出一个部分计数);假设 (iii) 保证 partition 决策只依赖 key。

活性与负载均衡(性能,而非正确性)

  • 设 $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 重试)。
  • 用户函数MAPREDUCE 均为确定性函数(同一输入必得同一输出),且在有限时间内终止。
  • 通道可靠;底层文件系统提供原子 rename

定理 5.1(完成性 / Liveness) 在上述假设下,若至少有一个 worker 永不故障,则 job 最终完成,master 返回 SUCCESS;否则 master 可能永远重试(不终止),但绝不会返回错误的或部分的输出

证明

  1. 由算法 5.3.1 的 WorkerFailed 例程,每次 worker 故障最多把 $M + R$ 个任务的状态从 {IN_PROGRESS, COMPLETED} 位移回 IDLE。
  2. 故障次数有限 ⇒ 被回退的任务实例总数有限 ⇒ 存在时刻 $t^*$,此后不再有任何任务被回退。
  3. 由算法 5.3.1 的活性论证 L1,只要存在 IDLE 任务与空闲的存活 worker,分配分支每轮必然派发任务;空闲 worker 只会因为”被分配任务”而变忙,而每个任务在有限时间内完成(用户函数终止性假设),故任务不会永久积压。
  4. 因此 $t^*$ 之后每个 IDLE 任务最终被派发并在有限时间内 COMPLETED;所有任务在有限时间内全部 COMPLETED。
  5. master 的终止条件要求”全部 map 且全部 reduce 都 COMPLETED”,此时它在 B 步返回 SUCCESS,唤醒用户程序。∎ 若连一个永久存活的 worker 都不存在,则第 3 步的假设被破坏,算法可能无限重试。这是活性(终止性)的损失,不是安全性的损失:master 从未返回 SUCCESS,用户程序从未被唤醒,因此不会有人读到错误结果。工程实现通常再加一个 deadline/失败计数上限,把”无限重试”变成”显式报告 FAILURE”——本讲代码中的 deadline 就是这一机制。

定理 5.2(串行等价性 / Safety,顺序无关性)MAPREDUCE 确定,则无论发生多少次 worker 故障、多少个备份任务被执行,最终 R 个输出文件的并集,等于对同一输入的一次无故障串行执行所产生的输出(在 key 到 value 的映射意义上完全相等)。

证明

  1. 每个 map 任务恰好有一个被采纳的结果。由算法 5.3.1 的幂等过滤(a = attempt[m] ∧ T_map[m] ≠ COMPLETED),每个 map 任务最多提交一次;由定理 5.1,它最终会提交。记 map 任务 $m$ 被采纳的那次执行为 $e(m)$。注意 $e(m)$ 可能是原始执行、重执行或备份执行中的任意一次
  2. 被采纳的是哪一次执行无关紧要。$e(m)$ 的输入是固定的输入 split(假设输入不变),MAP 是确定性函数,故 $e(m)$ 产生的中间 KV 多重集 $I(m)$ 与”是哪一次执行”无关。这是”重算代替恢复”能够成立的全部理由——它把”选哪一次结果”这个难题消解掉了
  3. 分区唯一性:由算法 5.3.3,key $k$ 的全部 value 都进入唯一一个 reduce 任务 $r = P(k)$。
  4. 每个 reduce 任务的输入是确定的:reduce 任务 $r$ 的提交输入为 $\biguplus_{m=0}^{M-1} \big(I(m) \,\big\vert \,_{P(\cdot)=r}\big)$。由第 2 步,每个 $I(m)$ 是确定的多重集,故这个并也是确定的(与 shuffle 的到达顺序、map 的完成顺序无关)。
  5. 每个 reduce 任务的输出是确定的REDUCE 是确定性函数,输入(连同 key 的分组方式)确定 ⇒ 输出确定。记 reduce 任务 $r$ 的输出为 $O(r)$。
  6. 与串行执行的对应:一次串行执行的定义是:计算所有 $I(m)$,合并成 $\biguplus_m I(m)$,按 key 分组,对每个 key 调用一次 REDUCE,把所有输出收集起来。而 $\{O(r)\}_{r=0}^{R-1}$ 正是把同一个分组过程按 $P(k)$ 切成 $R$ 块分别计算的结果(不同 key 之间互不影响,因为 REDUCE 的输入只包含同一个 key 的 value),拼起来与串行结果完全相同。∎
  7. 原子提交保证可见性:每个任务先写私有临时文件,reduce 完成时才原子 rename 到最终文件名。若同一 reduce 被执行多次,会有多次 rename,但由底层文件系统的原子性,最终文件内容是恰好某一次执行的结果;由第 5 步这些执行的结果相同,所以用户观察到的输出完整且一致。reduce 输出不需要在 reducer 机器故障后重做,正是因为 rename 后的文件已经在全局 DFS 上,不随机器消失。∎

非确定性函数下的弱化语义(必须点明) 论文给出了精确的弱化表述:当 MAPREDUCE 非确定时,某个 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 的迟到结果(输出去重)
断言通过: 两种配置结果都 == 单机串行结果

【代码做什么?】

  1. udf_map / udf_combine / udf_reduce唯一与作业相关的三个函数,其余全是框架——这正是 MapReduce 抽象边界的体现。
  2. run_job 把输入行切成 M 个 split(连续分块,模拟 HDFS block 边界),并记录 split_loc[m] = m % nw,模拟”这个 split 的副本在哪台机器上”。
  3. 每个 worker 是一个独立进程,带一个心跳线程,每 0.05 秒把 hb[wid] = time.time() 写进 mp.Manager().dict()——这就是跨进程共享的心跳表
  4. master 主循环每轮做四件事:(1) 非阻塞地排空结果队列;(2) 检查心跳超时并执行故障恢复;(3) 调度(先按 split_loc 做 data-local 分配,再补 rack/off-rack);(4) straggler 检测与 backup 任务的启动。
  5. map 任务在 worker 端:逐行调用 udf_map,把中间 KV 按 partition_of 分桶 → 桶内按 key 排序 → 可选 combiner → 作为一个整体上报给 master。
  6. reduce 任务:master 端把 M 个 map 的第 r 个分区合并、排序后作为 payload 下发(模拟 shuffle + 归并),worker 端按 key 分组调用 udf_reduce
  7. 三个场景分别验证:正常执行、worker 崩溃后正确性不受影响、backup task 把尾延迟从 3.01s 降到 1.24s。
  8. 最后用 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 映射;timeoutmapreduce.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 的样本数并不相等, 均值不可结合。

【代码做什么?】

  1. make_splits 生成 M=8 个 split,每个 500 条 <host, bytes> 记录;host 按 Zipf 分布抽取,模拟真实日志的数据倾斜
  2. run_job 是最小化的 MapReduce 引擎:每个 split 是一个 map 任务,本地按 key 聚成 local 字典(这一步是 map 的输出缓冲),然后是否调用 combinercombine_fn 是否为 None 决定。
  3. 每条被 emit 的记录都记账:recs_sent += 1bytes_sent += len(kk) + 8——这就是 shuffle 阶段真正过网络的数据量。
  4. 记录按 partition_of(MD5 稳定哈希)进入 R=4 个分区,每个分区各自做 reduce。
  5. 作业 1 用 sum:验证有/无 combiner 的结果完全相同,并打印流量压缩比。
  6. 作业 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 不维护任何可恢复的中间状态,只要求用户函数确定性,于是”重算”与”第一次算”结果相同,容错逻辑退化成一句”把它再做一遍”。简单,且正确性容易论证——这是分布式系统容错的经典范式。
  • 抽象边界的划法决定了系统的成败。 用户只写 mapreduce 两个纯函数;框架包办并行化 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 常见陷阱与注意事项

  1. 误以为”已完成的 map 任务不会重做”。
    • 为什么错:map 的输出写在 map worker 的本地磁盘上,那台机器一挂,这份数据就彻底不可访问了。若不重做,reduce 会少拉到一个 map 的全部数据,从而得到静默错误的(偏小的)聚合结果——这比崩溃危险得多。
    • 正确做法:worker 被判 failed 时,把它负责的所有 map 任务(无论 IDLE / IN_PROGRESS / COMPLETED)都重置为 idle;同时让所有仍在运行的 reduce 重新拉取数据(论文明确要求通知所有 reduce worker 改从新的 map worker 读取)。
  2. 把 combiner 当成”随便写都行”的优化。
    • 为什么错:combiner 会在任意次数任意分组的子集上被调用(每次溢写、每次归并都可能调),因此它必须满足结合律与交换律。用 meanmediancount 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 合并这两个量。
  3. 用 Python 内置的 hash() 做分区函数。
    • 为什么错:CPython 对 strhash() 默认启用随机化种子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 一致的)。
  4. reduce 里依赖 value 列表的顺序。
    • 为什么错list(v2) 的顺序由各 map 任务的完成顺序、shuffle 的到达顺序决定,不确定;而且重执行/备份任务会让顺序发生变化。依赖顺序的 reduce(例如”取第一个元素当代表”、”输出时保留输入顺序”)就变成了非确定性函数,从而丧失串行等价性保证。
    • 正确做法:把 reduce 写成一个对列表顺序不敏感的纯函数(求和、取最大、拼接但先排序等)。若确实需要顺序,就在 reduce 内显式排序,但要注意这会增加内存与时间开销。
  5. 以为 master 有高可用。
    • 为什么错:论文的实现里 master 是单点:它的数据结构(任务状态、中间文件位置表)只存在内存里,master 一挂,整个作业就中止,用户只能重跑。Hadoop 1.x 的 JobTracker 同样如此。
    • 正确做法:理解这一代系统的边界——master 靠”client 重试”而不是”自动接管”来容错。Hadoop 2.x/YARN 把它拆成全局 RM(checkpoint + 备用 RM)+ 每作业一个 AM,才真正改善了这一问题。
  6. 把 M 和 R 设得越大越好。
    • 为什么错:master 必须做 $O(M + R)$ 次调度决策并保存 $O(M \cdot R)$ 的状态;每个 map 任务的启动开销是固定的;每个 reduce 任务产生一个独立输出文件。$M$、$R$ 太大时,调度与元数据开销会超过并行化带来的收益,还会产生海量小文件。
    • 正确做法:按论文 §3.5 的经验值——让每个 map 任务处理 16–64MB 输入(保证局部性),让 $R$ 是预期机器数的小倍数
  7. 在 map 里产生副作用(写文件、发请求、累加全局计数器)。
    • 为什么错:map 任务可能被执行多次(故障重做、备份任务),副作用会被执行多次;而且备份任务的输出最终会被丢弃,副作用却已经发生了。论文 §4.5 明确要求用户自己保证副作用的原子性与幂等性
    • 正确做法:把副作用写成”先写临时文件、完成后原子 rename”;把计数工作交给框架提供的 counter 机制(master 会在聚合时消除重复执行带来的重复计数)。
  8. 认为 reduce 可以”看到全局信息”。
    • 为什么错:一个 reduce 任务只能看到属于它自己那个分区的 key;不同 reduce 任务之间没有任何通信,也不能共享全局变量。所以”求出出现次数最多的那个词”这种需要跨 reduce 比较的问题是一个 MR 作业做不到的(讲义中的练习正是问这个)。
    • 正确做法:串联第二个 MapReduce 作业——第一个作业输出 <word, count>,第二个作业的 map 把所有记录重新 key 成同一个 key(例如 1),只用一个 reduce;该 reduce 就能看到全部 <word, count> 并选出最大值。或者用自定义 partitioner 保证全局唯一 reduce。
  9. 忘记输入数据在作业期间不能变。
    • 为什么错:如果输入 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 任务了。这个优化显然更好。” 请指出这个想法错在哪里。

:错在它用错误的方式解决了一个已经被解决的问题,而且代价很高:

  1. 它破坏了局部性优化的收益。 map 输出是中间数据,体量往往与输入同量级(论文实测:3288TB 输入产生 758TB 中间数据)。把这些数据写进 HDFS 意味着一次完整的网络传输 + 3 副本复制(3 倍网络与磁盘写入)。而把中间数据写本地磁盘,之后只有真正需要它的那个 reduce 任务去拉取一次(1 次网络传输,且只传一份)。在 sort benchmark 中,正是因为输入和中间结果尽量走本地磁盘,输入速率才能达到 13–30 GB/s 而网络不至于被打爆。
  2. 它并不能消除重做。 即使中间数据在 HDFS 里,一旦某台 map worker 崩溃、它的 block 副本只剩 2 个,HDFS 会立刻启动副本重建——这同样是一次完整的网络传输,只是把成本从”重算”挪到了”复制”,而且复制量可能更大(重算只重算丢失的那几个任务,复制要维持 3 副本)。
  3. 重算本来就不贵。 论文实测杀掉 200/1746 个 worker 进程,总时间只增加 5%。原因正是 $M$ 远大于机器数:被重做的 map 任务可以立刻分散到大量空闲机器上并行重跑。
  4. 真正的设计哲学是”重算比恢复便宜”:本地写、远程读,加上”用重新执行代替状态恢复”,整体上是最优的工程权衡。如果想减少重算量,正确的方向是增大 $M$(把重算粒度变细)而不是把中间数据搬进 DFS。

Q3(推演题) 词频统计作业中,某个 map 任务的输入 split 是 "a b a a b",$R = 2$,分区函数 $P(w) = \text{len}(w) \bmod 2$(len("a")=1len("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 条。
  • (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 单独拆出来处理。