一批输入有十八条记录。某个工作者写出半个结果文件后退出,调度器把同一份输入交给另一个进程。第二个进程成功了,最终结果能否保证每条记录只出现一次?如果第一个进程没有退出,只是回复得晚,它又该如何处理自己的结果?

上一篇的 RPC 重试与请求去重讨论了一次调用可能执行多次。批处理框架面对相同的不确定性,却有一个有利条件:许多任务可以从固定输入重新计算。代价是框架必须规定哪一次执行的输出能进入下游,不能把工作者留下的所有文件直接拼起来。

MapReduce 把这件事与任务切分、数据分组、调度放进一个计算框架。本文沿用 2004 年论文的模型讨论保证,再用独立 Go 实验检查进程被杀和旧 attempt 晚到的发布边界。实验有真实子进程和文件操作,不连接分布式存储,不模拟成生产集群。

先确定计算结果是什么

输入记录为 (ID, key),ID 从 0 到 17,key 表示 ID 的奇偶性。目标是按 key 收集 ID,并在每组内排序。正确结果只有两组:偶数组包含 0,2,...,16,奇数组包含 1,3,...,17。这个任务很小,但能检查重试后是否多算、漏算或分错组。

将输入按 ID mod 3 分给三个 map task。每个 task 处理一份互不重叠的切片,输出原记录;这里使用恒等 map,刻意保留每条记录的身份。两个 reduce task 各取一种 key,收集三个 map task 对应的数据并排序。对于相同的逻辑输入,执行多少次都应产生同样的逻辑输出。

1
2
3
输入切片 0: 0,3,6,9,12,15  ── map 0 ──┐
输入切片 1: 1,4,7,10,13,16 ── map 1 ──┼── 按 key 分组 ── reduce 0: 偶数 ID
输入切片 2: 2,5,8,11,14,17 ── map 2 ──┘ └─ reduce 1: 奇数 ID

在一般的 MapReduce 程序中,map 可以生成零条、一条或多条中间记录;reduce 对同一 key 的所有 value 计算结果。框架负责将相同 key 分发给相同 reducer,这个数据搬运与重组阶段称为 shuffle。示例的奇偶分区规则易于手算,实际系统通常用确定性的分区函数分散 key。

一个 key 即使拥有大量记录,也不能随意拆给多个互不协调的 reducer,然后声称得到了完整聚合。必须先确定聚合能否分解、局部结果如何合并。例如求平均值应传递总和与数量,仅平均各分区的平均值,在分区大小不等时会出错。能够恢复任务,并不意味着任意应用逻辑都会得到正确答案。

已有的 Hadoop 编程模型篇Hadoop Shuffle 篇可作为实现术语的扩展阅读。这里的证明和实验不依赖某个 Hadoop 版本,也不把 2004 年 Google 实现、Hadoop 和本地教学代码视作同一个系统。

MIT 6.5840 Spring 2026以 MapReduce 阅读与实验起步,随后介绍 RPC、GFS 和复制;系列为便于编程而把 RPC 前置。Stanford CS244B Spring 2024包含 Spark 与 Ray,后续可比较这些计算模型,不据此声称它单独开设了 MapReduce 单元。

任务与执行尝试的身份

任务描述要回答“算哪一份输入”。执行尝试还要回答“这是该任务第几次运行”。逻辑任务 ID 在重试时不变,attempt ID 每次重派都变化。只记录工作者 PID 不够:工作者可能连续处理多个任务,进程重启之后也不能据此确定旧消息属于哪轮运行。

给 map 0 的第一次执行记为 (map,0,1),重试记为 (map,0,2)。两次执行使用不同临时路径,不能同时向 map-0.json 追加内容。临时文件包含的只是某次尝试的产物,文件存在不表示任务已经被调度器接受。

调度器至少保存任务当前 attempt、运行或完成状态、已接受输出的位置。为便于讨论,可以把状态转移写成:

1
2
3
4
Pending ── 分派 attempt 1 ──> Running(1)
Running(1) ── 放弃旧尝试并重派 ──> Running(2)
Running(2) ── 接受 attempt 2 输出 ──> Done(output 2)
Done(output 2) ── 收到 attempt 1 完成 ──> 保持原状态

超时只能促使调度器放弃等待,不能证明旧进程已经停止。旧进程可能仍在计算、写文件或发送完成消息。因此,重派后两次执行可以并存,正确性必须允许这一重叠。要求每个任务在物理上只运行一次,会把恢复工作变成无法实现的前提。

这一状态机还没有解决调度器自身故障。若所有状态只在内存中,调度器退出就丢掉了哪些输出被接受的信息。持久化任务表、恢复日志、重复作业提交等问题需要另行设计;工作者重试能力不能替代这些机制。

2004 年实现如何恢复与提交

Dean 与 Ghemawat 的原始论文在 §3.3 分别处理工作者与 master 故障。工作者失效后,未完成任务需要重派;已完成 map 也可能重算,因为其中间输出保存在工作者本地。已经完成的 reduce 输出位于全局文件系统,不因计算工作者失效而同样丢失。该实现的 master 失效会中止作业,不应把论文提到的检查点可能性写成已实现的自动恢复。

§3.3 对提交的处理并不统一。map 完成时,工作者报告一组临时文件,master 只接受一个已完成结果;reduce 完成时,工作者把临时文件原子重命名为最终文件。多个 reduce 尝试可能都执行 rename,底层原子操作保证最终文件来自一次完整执行。MIT 2026 讲义也区分这两条路径。

这个差异来自输出的消费者。map 输出是框架内部数据,master 可以决定向 reducer 公布哪组文件。reduce 输出要以约定名称交给外部使用,文件系统负责一次文件名替换的原子可见性。不能把两条路径简化为“所有工作者都写同一文件,最后完成的随便覆盖”。临时隔离和原子发布正是避免部分写入混合的条件。

设 attempt A 写完两行,attempt B 写完五行。如果它们逐行追加同一个最终文件,读者可能得到七行混合内容。如果各写私有文件,再原子替换最终名称,某次打开该名称得到的应是一个完整版本。原子替换不要求两次执行之间没有重叠,它要求一个名称的可见状态不暴露半次替换。

POSIX.1-2024 的 rename 规范描述了替换名称时旧文件或新文件的可见性。这是文件命名操作的保证,不代表若干个 reducer 文件同时成为一个事务;也不能直接推导文件在断电后持久保存。Go 的 os.Rename 文档还明确提醒平台差异:非 Unix 平台即使同目录也不保证原子。下文只在 Unix 本机同一临时目录使用这条路径。

不重不漏的条件与证明

对本篇任务,可以把“不重不漏”写成一个具体不变量:每个原始 ID 在最终两份输出的并集中恰好出现一次,且所在分区与它的 key 匹配。只校验最终数量为十八不够,漏掉 ID 3 又多出 ID 5,也会得到十八条。

证明从输入切片开始。切片互不重叠,且覆盖全部输入,所以每条记录只属于一个逻辑 map task。每个 map task 只向下游公布一个已接受版本,重试产生的其他文件不在消费清单中。因此,每个原始 ID 在被接受的中间数据里出现一次。

分区函数再把每个 key 映射到唯一 reducer。每个 reducer 读取全部已接受 map 输出中属于自己分区的记录;本例的 reduce 只排序而不增加、删除 ID。于是每条中间记录恰好进入一个 reduce 结果。每个 reduce 结果通过完整文件发布,最终枚举固定的两份逻辑输出即可得到目标集合。

证明中的每一步都有对应失败方式:输入切片重叠会造成重复;漏读某个 map 会丢数据;从目录里通配读取所有 attempt 文件会重复;分区函数不一致会分错组;直接发布半文件会丢失完整性。这些错误不可能仅靠“多重试几次”消除。

确定性还决定重新计算是否可以替代丢失的输出。论文 §3.3 对确定性 map/reduce 给出与无故障顺序执行相同的结果语义。对非确定性计算,保证更弱:不同 reduce 结果可能对应不同顺序运行。这不等价于所有输出共同来自一次统一运行。

一个反例是 map 每次执行生成一个批次标记,并分别向两个 reduce 分区各输出一条标记。attempt A 生成 batch=A,attempt B 生成 batch=B。一个 reducer 已经读到 A,另一个因故障恢复读到 B,最终就可能得到 (A,B)。任意一次单独运行只会生成 (A,A)(B,B)。每个 reducer 都读到了完整文件,跨分区的一致关系仍被破坏。

该反例是根据论文弱语义构造的手算情景,不是下文代码的观察结果。教学代码先完成所有 map 的选择再进入 reduce,不支持这类早期读取与丢失重算交错;实验通过也不能否定论文描述的弱语义。

外部副作用与活性边界

任务向外部数据库执行 counter += 1,再向调度器报告成功。若累加成功但回复丢失,重试又会累加一次。最终输出文件只选一份,并不能回滚数据库里多出的一次修改。外部服务不知道框架选择了哪个 attempt,也不会因为另一个 attempt 被拒绝就自动撤销副作用。

论文 §4.5 把额外副作用的原子性和幂等性责任交给应用,并明确没有为单个任务产生的多个输出文件提供原子两阶段提交。工程上可以让外部操作携带稳定的业务幂等键,也可以先产生纯数据,由另一个具备事务边界的环节消费;但必须同时处理去重记录的持久化与保存期限。不能仅把 attempt ID 当作业务幂等键,因为每次重试的 attempt ID 恰好不同。

安全性说的是不接受错误组合,活性说的是作业最终完成。后者至少需要某次尝试能够终止、工作者与存储持续可用、调度器继续分派。若某条输入让所有工作者确定性崩溃,换再多机器也不会解决;跳过该记录则改变了结果契约,不能继续声称原输入不重不漏。

慢任务还会拖长作业完成时间。启动备份尝试可能让较快副本先完成,但它增加计算和存储流量;若数据倾斜让同一 key 天然很大,重复执行同样的工作也不一定有效。把任务数细分到多于工作者数,有助于动态分担负载,却会增加调度状态和中间文件数量。这里没有性能测量,不能据此给出固定的扩容收益。

本机实验:半文件、被杀进程与晚到结果

代码位于仓库 examples/distributed-systems/mr02/main.go。父进程作为唯一调度器,通过标准输入向自身启动的 worker 子进程发送 JSON 请求。每个请求包含阶段、分区、输入路径、私有输出路径和可控暂停点。三个 map task 与两个 reduce task 顺序调度,晚到情景中同一 task 的两个子进程重叠运行。

这里刻意不实现网络 RPC、并行工作者池或完整 shuffle 服务。reduce 读取三份完整 map 文件并过滤自己的分区,适用于十八条记录的检查;实际大数据框架应按分区组织中间文件,避免每个 reducer 扫描所有中间数据。简化后仍保留了真实进程退出、重派、文件隔离与发布检查。

父进程只接受当前 attempt 的结果,并且拒绝已经完成任务的后续提交:

1
2
3
4
5
6
7
8
9
func publish(t *task, attempt int, path string) bool {
if t.done || attempt != t.attempt {
return false
}
_ = readRecords(path)
must(os.Rename(path, t.output))
t.done = true
return true
}

完整 JSON 解析发生在 rename 之前,拒绝半份文件;它不验证任意恶意工作者是否算对了结果。代码假设工作者计算正确,只注入退出和延迟故障。attempt 检查与 rename 都在父进程的同一串行路径上,没有并发提交处理器;如果将来并行化,必须保持检查与发布之间不被另一次重派或提交穿插。

这个提交者安排与原论文的 reduce worker 自行 rename 不同。由父进程统一发布,可以直接演示旧 attempt 被拒绝,但不能宣称复现了 Google 2004 实现。父进程若在 rename 后、done=true 前崩溃,内存状态与文件之间会有恢复问题;当前程序不注入或修复父进程故障。

在累计实验目录运行:

1
2
cd examples/distributed-systems
GOCACHE=/private/tmp/ds-go-cache go run -race ./mr02

默认执行三种情景,每种都新建独立临时目录,结束时删除该目录中的实验数据。程序只杀自己启动的工作者,不接受外部 PID。单次子进程有十五秒保护期限,用于防止实验挂起;保护期限不是协议中的故障检测证明。

正常情景作为对照。半写情景要求 map 0 的 attempt 1 先写出不完整 JSON,再通过标准输出报告 READY。父进程检查文件确实不是完整 JSON,调用 Process.Kill(),通过 Wait() 确认信号终止,再启动 attempt 2。握手确保故障发生在指定位置,不依赖猜测几毫秒的休眠。

晚到情景使用另一个独立作业。attempt 1 写完私有文件后暂停,父进程将 task 的当前 attempt 更新为 2。attempt 2 完成并发布后,父进程释放 attempt 1,使其正常退出并提交旧结果。发布函数必须返回 false,最终文件必须保持不变。被杀进程不会复活,因而这个情景不能与前一个混成同一次执行。

一次本地运行得到以下输出,其中 PID 随运行变化:

1
2
3
4
5
PASS scenario=normal maps=3 reduces=2 records=18 exact=true
KILLED pid=3503 stage=map task=0 attempt=1 partial=true
PASS scenario=kill maps=3 reduces=2 records=18 exact=true
STALE stage=map task=0 attempt=1 current=2 rejected=true
PASS scenario=late maps=3 reduces=2 records=18 exact=true

exact=true 来自逐条比较排序后的 (ID,key) 与原始输入,并检查每个输出分区;不是只对比记录数量。三个情景只覆盖选择的故障位置,不是对所有并发执行的穷尽证明。race 运行通过也只说明本次覆盖的执行没有被检测出数据竞争。

实验还对临时副本做了一次变异检查:去掉 publish 中的完成状态与 attempt 守卫,运行晚到情景,程序因 stale attempt accepted 失败。这个检查说明断言能识别旧结果被接受;它不验证进程崩溃后的断电持久性。正式代码保持守卫不变。

本地环境为 Go 1.27.0、macOS 27.0、arm64。go vet ./mr02、帮助命令和非法场景退出检查也已执行。源码声明的 Go 1.23 最低版本未单独实跑;没有 Windows、网络分区、磁盘故障或分布式文件系统验证。逐项命令与边界保存在 实验记录,资料定位保存在 论断与证据

两道练习

根据输出集合定位重复

设一个保留所有 attempt 文件的实现变体中,map 0 的 attempt 1 与 attempt 2 都正确地产生六条记录。调度器只接受 attempt 2,但 reducer 用 map-*-attempt-*.json 通配读取所有文件。最终结果会怎样?在不修改 map 计算函数的条件下,修复位置应在哪里?

推导:map 0 的六个 ID 会出现两次。对于本例,ID 0、3、6、9、12、15 各重复一次,总条数变成二十四。调度器正确拒绝旧完成消息并不够,下游还必须只消费提交清单里明确选定的输出。修复应让 reducer 获取三个逻辑 map 的已接受路径,不能依靠扫描目录猜测任务状态。

区分丢回复与丢持久化

设工作者向外部服务发送一次带业务键 job-7/map-0 的累加,服务成功处理后回复丢失。重试使用相同业务键。服务端在内存中记住该键,但随后重启,之后再次收到重试。结果是否只会累加一次?怎样修改才能建立更强保证?

推导:重启销毁去重状态后,服务端可能再次累加。稳定的键只是标识,保证还依赖去重记录与业务效果的原子持久化,以及记录覆盖全部可能重试的保存期限。如果外部服务无法提供这种边界,框架必须承认结果不确定,或把待处理事件交给能够履行这一契约的后续环节。

任务恢复依赖一个尚未展开的前提:原始输入可再次读取,已发布结果有可靠的存储位置。下一篇 GFS 将分别检查元数据、数据副本、租约和记录追加,说明计算框架从存储层取得了什么保证。

参考资料