
开始写正文。1. 为什么 6.824 的第一课是 MapReduce我在刷 MIT 6.824现在课程号已经改成 6.5840的时候第一个星期就被 Lab 1 锤了一顿。这门课第一课就叫 MapReduce作业也是让你自己实现一个简化版 MapReduce 框架。当时我不理解分布式系统这么大的主题为什么偏偏拿一个 2004 年的论文模型当开门课等我把 Lab 1 写完、跑完 crash test 之后才明白MapReduce 的抽象足够简单简单到你能在两周内手写一个能跑的版本但它的实现又足够典型把分布式系统里最恶心的两个问题——共享状态同步和节点失败处理——全部暴露出来了。你绕不过去只能正面去解决。这篇文章不是课程答案是我把自己实现 Lab 1 的完整思路摊开来讲RPC 协议怎么设计、Master 调度器的状态机怎么转、Worker 崩溃后为什么等 10 秒、以及最后怎么把 test-mr.sh 跑绿。我自己当时在本地反复调了两天踩过不少坑所以我会把那些实验指导文档里没写但真能救命的经验也混进去。适合正在肝这个 Lab、或者想从零理解 MapReduce 运行机制的人看。先说清楚 MapReduce 到底是什么。它本质上把数据处理拆成两个阶段Map 阶段把一个大输入拆成 M 个小片每个片交给一个 Worker 处理产出若干 key, value 中间结果Reduce 阶段把 key 相同的中间结果聚合到一起再输出最终结果。中间夹着一个 shuffle 过程负责让每个 Reduce 任务拿到属于它的那部分 key。你的 Lab 1 就是要实现一个 Mini 版一个 Master 进程负责任务调度一堆 Worker 进程通过 RPC 向 Master 要活干。关键的心智模型就一句话这不是一个程序这是一个分工系统。你写的 Master 不是帮 Worker 算数据而是管理谁在干什么、干了没有、干没干完。你写的 Worker 也不关心别人只关心我现在手里这个任务做完了没有、怎么把结果交回去。所有 bug 都来自对这个心智模型的误解。2. Lab 1 到底要求你做什么我一开始也理解偏了2.1 拿到手的代码骨架其实已经帮你圈好了边界6.824 的 Lab 1 会给你一个 Go 语言骨架里面有几个文件master.go、worker.go、rpc.go还有一个mrmaster.go/mrworker.go入口。你不需要管命令行参数怎么解析不需要关心测试脚本怎么调用进程。真正要你写的是三块Master 一侧分配任务、跟踪任务状态、处理超时重发、判断整个 Job 是否结束。Worker 一侧向 Master 请求任务、执行 Map 或 Reduce、把中间文件写到磁盘、汇报执行结果。中间落盘协议Map 产生的中间结果怎么命名、怎么按 key 分桶、Reduce 怎么找到自己需要的文件。我第一次写的时候有个误区以为要做一个类似 Hadoop 那样完整的资源调度系统。其实不用。这个 Lab 只要求你处理一个固定的 WordCount 作业但你要写的是框架不是应用。你可以把 map 函数和 reduce 函数看成插件框架负责调度和落盘。2.2 你要实现的最小闭环一个 Job 从启动到结束流程是这样Master 启动知道输入文件列表txtfiles和 Reduce 任务数量nReduce。Master 把输入文件切成M个 Map 任务M等于输入文件数量。Worker 启动后循环调用GetTask()Master 每次给它一个任务。如果是 Map 任务Worker 读对应输入文件调用用户提供的Map函数得到一组中间 key/value再按ihash(key) % nReduce写入磁盘上的中间文件。Master 等到所有 Map 任务都完成后才开始分配 Reduce 任务。如果是 Reduce 任务Worker 读取属于某个 bucket 的所有中间文件对 key 排序调用用户提供的Reduce函数把结果写入最终输出文件。所有 Reduce 完成后Master 返回完成测试脚本退出。注意第 5 条是硬约束Reduce 任务绝对不能在 Map 阶段还没结束时启动。因为 Reduce 要读全部 Map 产生的中间文件这个语义在 MapReduce 模型里是写死的。你不用自己去设计流式计算老老实实按阶段走就行。2.3 测试脚本里藏着的考点test-mr.sh会跑好几轮测试。我当初以为跑过 basic 测试就完事结果被crash test吊打。这份脚本大致验证这四件事正确性最终输出文件的单词统计结果是否和顺序执行的参考结果一致。并行性它统计了完成时间。如果你让 Worker 串行执行时间会超标。脚本里有一次是单 Worker 跑还有一次是 8 个 Worker 跑。容错性脚本会启动一堆 Worker然后随机杀掉正在运行的 Worker再补充新 Worker最后检查结果是否正确。这就是 crash test。幂等性重复跑多次输出稳定一致。我后来悟了这个 Lab 的评分点本质上是三个词——正确、并行、容错。前两个好理解容错才是拿高分的关键。3. 先把 RPC 协议定下来写并发系统前最重要的一步很多人的习惯是拿到骨架就开写 Worker 逻辑我强烈建议反过来先花半小时把 Master 和 Worker 之间的 RPC 交互协议定死。协议定清楚后面所有代码都只是在实现协议约定而已。3.1 用两个结构体画清交互边界Lab 1 至少需要两个 RPC 方向Worker 向 Master 请求任务Worker 向 Master 汇报任务结果。我用的结构体大概长这样// rpc.go type TaskRequest struct { WorkerID int // 可选的 worker 标识便于排查 } type TaskResponse struct { TaskType TaskType // MapTask / ReduceTask / Wait / Exit TaskID int // Map 任务对应输入文件编号Reduce 任务对应 bucket 编号 InputFile string // 仅 Map 任务需要 NReduce int // reduce 分桶数 NMap int // map 任务总数 } type ReportRequest struct { WorkerID int TaskType TaskType TaskID int Success bool } type ReportResponse struct { OK bool }这个协议有几个值得注意的点。第一任务描述里为什么必须带NMap和NReduce因为 Map 任务在写中间文件时要算ihash(key) % nReduce它必须知道桶的数量Reduce 任务在读取中间文件时必须知道总共有多少个 Map 任务它才知道要循环读mr-0-Y到mr-(NMap-1)-Y。这些参数与其让 Worker 猜不如直接在任务响应里带过去干净利落。第二TaskResponse里我加了一个Wait类型。当 Master 判断当前没有可分配的任务、但整个 Job 又没结束时比如 Map 阶段刚完成Reduce 还没就绪就返回一个稍等重试的响应Worker 收到后 sleep 一小会儿再来要任务。这个设计避免了 Worker 空转时疯狂打 RPC 打爆 Master。3.2 Master 状态机四个状态之间的流转我把每个任务不管 Map 还是 Reduce的状态建模成一个四态状态机Idle - InProgress - Completed中间加一个隐性的Expired恢复态。Idle任务还没发出去Master 可以分配。InProgress任务已分配给某个 Worker正在执行。Completed任务已汇报完成可以安全地从调度池里移除。如果InProgress超过 10 秒没有收到汇报Master 把状态强行改回Idle重新分配。代码上用两个 slice 去存状态一个管 Map 任务、一个管 Reduce 任务type Master struct { mu sync.Mutex mapTasks []TaskState reduceTasks []TaskState nMap int nReduce int done bool }启动时对所有输入文件初始化mapTasks[i].state IdlereduceTasks[i].state Idle但要等到allMapDone()为 true 才允许被分配。这个状态机看着简单写起来难的地方在于所有字段都可能被多个 RPC handler 并发读写你必须在分配任务、标记完成、轮询状态的所有路径上拿锁。具体怎么拿锁我放在下一节说。4. 调度器实现并发安全的细节都在这里4.1 任务分配策略active worker 轮流取活Master 分配任务的核心逻辑是一个getTask()函数每来一个 RPC 请求就执行一次。我用的策略是优先清空 Map 任务队列如果 Map 队列空了且全部完成就进入 Reduce 队列如果 Reduce 也全部完成就返回 Exit。func (m *Master) GetTask(req *TaskRequest, resp *TaskResponse) error { m.mu.Lock() defer m.mu.Unlock() // 第一阶段分配 map 任务 if !m.allMapDone() { for i : 0; i m.nMap; i { if m.mapTasks[i].State Idle { m.mapTasks[i].State InProgress m.mapTasks[i].StartTime time.Now() resp.TaskType MapTask resp.TaskID i resp.InputFile m.inputFiles[i] resp.NReduce m.nReduce resp.NMap m.nMap return nil } } resp.TaskType Wait return nil } // 第二阶段分配 reduce 任务 if !m.allReduceDone() { for i : 0; i m.nReduce; i { if m.reduceTasks[i].State Idle { // 类似地设置 InProgress、StartTime resp.TaskType ReduceTask resp.TaskID i resp.NReduce m.nReduce resp.NMap m.nMap return nil } } resp.TaskType Wait return nil } resp.TaskType Exit return nil }这段代码有两点容易翻车。第一如果你用了一个全局的nextTaskIndex来顺序发任务一旦某次 RPC 发完任务后 Worker 崩溃这个索引会乱掉。所以我每次都在请求到达时全量扫一遍状态表找第一个 Idle而不是维护一个游标。全量扫描在这个规模下性能没问题逻辑却简单很多。第二Wait响应必须配合合理的重试间隔。Worker 上面拿到 Wait 就time.Sleep(100 * time.Millisecond)再继续请求。间隔太小 Master 锁会被打爆间隔太大会拖慢整体完成时间100ms 是我试过比较舒服的值。4.2 你最可能踩的 race任务被重复分配这个 Lab 里 90% 的并发 bug 都出在同一件事任务被两个 Worker 同时执行。场景非常经典Worker A 拿到 Map 任务 3但执行到第八秒突然卡住Master 在第 10 秒发现超时把任务 3 重新标记为 Idle分配给 Worker B。这时候如果 Worker A 又缓过来了继续执行并写中间文件任务 3 就被执行了两次。Lab 指导文档里没有明说但测试脚本其实期望你接受这个情况。它允许重复执行但要求最终结果不因重复而错误。也就是说任务重复本身不是 bug重复后数据乱了才是 bug。具体怎么保证我在第 5 节展开。这里先说另一个坑Master 端发出去的任务在收到汇报前不能再发给别人——除了超时重发。超时重发本身要有单独的路径别和正常分配混在一起。我自己的实现里超时检查有两个时机。一个是在GetTask()的时候顺手扫一下把超时的 InProgress 任务强制改成 Idle另一个是独立 goroutine 周期扫描。两个方案都行但注意扫的时候必须持有m.mu而且一旦决定重发要让新任务从 Idle 状态重新走分配流程。4.3 等待判断allMapDone 和 allReduceDone 的写法allMapDone()不是返回是否所有 map 任务都分配完而是是否所有 map 任务都处于 Completed 状态。这两个概念经常被新手写混。如果只是分配完就开始 reduce你会发现 Reduce 任务读了半个中间文件就没数据了。我在这个函数上踩过一次很蠢的坑有一次我把条件写成了countInProgress 0结果 Map 任务全部进入 Wait 状态既不是 InProgress 也不是 Completed时也就误判为 done导致 Reduce 提前开始。正确写法是专门维护一个completedMapCount计数每收到一条成功的 Map 汇报就加一completedMapCount nMap才是真完成。因为 Master 的所有状态都在锁保护下这个计数是可靠的。Reduce 阶段同理维护completedReduceCount等于nReduce时把m.done置为 true。这个done标志会被Master.Done()接口读取测试脚本靠它判断进程是不是该结束了。这个环节有一个隐蔽要求Master 必须在所有任务完成后主动让Done()返回 true否则 mrmaster 进程不退出测试脚本会卡死或者直接判定失败。5. 10 秒超时与崩溃恢复这是 Lab 1 最核心的考点5.1 为什么偏偏是 10 秒实验指导里明确说Master 如果在 10 秒内没有收到某个任务的完成汇报就认为它挂了需要重新调度。10 秒不是为了等你慢慢算而是为了区分两种情况Worker 还在正常算但算得慢和 Worker 真的挂了。在测试脚本里crash test 会给每个 Worker 设置一个随机计时器到点就把进程杀掉。如果你把超时设得太短比如 1 秒那么正常的慢 Worker 也会被误杀然后反复重发任务整个 Job 永远跑不完如果设得太长比如 30 秒crash test 会在检测失败时拖很久。10 秒是课程权衡后的标准值我建议不要改。实现上超时重发不能等 Worker 来请求时才处理否则可能所有 Worker 都忙着手里的活没人来触发重发。我加了一个后台 goroutinefunc (m *Master) watchTimeout() { for { time.Sleep(500 * time.Millisecond) m.mu.Lock() now : time.Now() for _, t : range m.mapTasks { if t.State InProgress now.Sub(t.StartTime) 10*time.Second { t.State Idle } } for _, t : range m.reduceTasks { // 同理 } m.mu.Unlock() } }这样即使所有 Worker 都卡在崩溃状态Master 也能主动把任务救回来。注意watchTimeout用500ms的粒度做轮询10 秒超时配合 500ms 检查精度足够。5.2 Worker 崩溃后发生了什么Worker 崩溃分两种情况。第一种是执行 Map 任务时崩溃它可能已经写了一半中间文件也可能还没写任何文件。Master 超时后重新分配同一个 Map 任务给另一个 Worker新 Worker 会重新读取输入文件、重新计算、重新覆盖写中间文件。因为中间文件的命名是mr-X-YX 是 map 任务编号Y 是 reduce 桶编号所以覆盖写只会影响这一个 map 任务的输出其他 map 任务的文件不受影响。这是 MapReduce 能容忍重复执行的前提不同任务写不同文件相同任务重复写也写在同一个路径上。第二种是执行 Reduce 任务时崩溃。Reduce 任务的输入是所有 map 任务的中间文件它自己写的输出文件是mr-out-Y。如果这个文件写了一半下一轮重跑同一个 reduce 任务会把旧文件整个覆盖掉。这里有一个 key 点Reduce 写入必须具有全写全换的原子性至少不能把半截文件留给最终检查。我的做法是先写到一个临时文件比如mr-out-Y.tmp全部写完后用os.Rename原子改名成mr-out-Y。这样即使 Worker 写了半截Master 重新调度后也能得到一个干净的新文件测试脚本不会读到坏数据。还有个很反直觉的点超时重发后原先那个慢但没死的 Worker 如果最终完成了任务Master 依然会收到它的 RPC 汇报。你必须在ReportTask里对这种情况做幂等处理如果任务状态已经是 Completed就忽略这条汇报。否则日志会很难看而且可能引发状态错乱。5.3 中间文件用完不删是什么后果这一点实验指导几乎没强调但我在实际跑回归测试时吃了大亏。Reduce 任务每次要读取nMap个中间文件。如果 Map 任务被重复执行文件被覆盖没问题。但如果你的代码在 Map 阶段对文件名做了带随机后缀的写法比如mr-3-2-abc123那么旧文件和新文件会同时存在Reduce 会读到一个重复的 key统计结果直接翻倍。所以中间文件的命名必须严格可预期mr-mapTaskID-reduceBucket。不用带任何随机性。这个教训也说明crash test 为什么能检查出来重复执行的问题——它杀掉 Worker 后同一任务必然被重发如果你写的文件名不稳定最终结果就会错。5.4 自己模拟崩溃的调试方法我本地调试 crash 时从来不直接改测试脚本而是用一个取巧的办法在 Worker 的Map执行函数里人为加一个time.Sleep(11 * time.Second)然后把 Master 超时调到 3 秒。这样能快速模拟任务超时 - 重新分配 - 第二个 Worker 成功的路径。确认这条路走通后再把超时改回 10 秒。另一个有用的验证方法是跑一遍 8 Worker 的测试故意只留 1 个 Worker 存活其他全部kill -9。如果最终输出依然正确说明你的容错逻辑基本过关。6. 排序、分组与倒排打通 Lab 1 之后这些实训题都是一件事6.1 shuffle 阶段天然附带按 key 排序如果你用网上的实训题关键词搜过会看到大量MapReduce 排序—自定义排序分组排序倒排序索引之类的课程作业。这些题乍一看像是 Hadoop 生态里的独立技能但它背后的核心机制就是你在第 4 节实现的shuffle 过程中的 key 归约。Lab 1 里并没有强制要求按 key 排序但 Reduce 任务的正确姿势是读取所有属于某 bucket 的中间 key/value先按 key 排序然后遍历排序结果把相同 key 的值分组每组调用一次 Reduce。为什么排序因为 Map 阶段产出的中间结果是分散的不同 map 任务可能把相同 key 写到同一个 bucket但顺序是乱的。不排序你就没法高效地把相同 key 的 value 聚合到一起。WordCount 是最简单例子。Map 函数产出word, 1shuffle 把所有相同 word 的键值送到同一个 bucketReduce 遍历一组1,1,1加起来输出word, 3。你的 Reduce 函数能不能正确聚合完全取决于同一个 key 是否恰好落在同一个 bucket那就靠ihash(key) % nReduce保证。这个映射本身的构造很粗糙但它是可靠的同一个 key 哈希值永远相同取模结果永远相同。6.2 分组排序和自定义排序到底在排什么我在看过大量相关实训题后总结了一句规律凡是叫自定义排序的题核心就是让你在 Reduce 侧把 key 或 value 按你自己的比较规则排序后再输出。比如给数字排序Map 阶段按数值大小做 keyReduce 阶段对 key 做比较排序。分组排序则是在排序之后再分组先按第一个字段排序同组的按第二个字段排序Reduce 前要确保相同组的数据连续到达才能正确区分一组和下一组。这些题在 6.824 Lab 1 里不要求你专门实现因为课程让你把框架搭起来应用层的排序规则是用户自定义的Reduce函数内部的事。但你理解了 shuffle 的按桶切分 桶内排序机制后Hadoop 里的 Partitioner、Comparator 概念就一点也不神秘了。Partitioner 就是ihash(key) % nReduce的另一个名字Comparator 就是你在 Reduce 之前对 key 做的排序比较器。Lab 1 没有把这些东西拆成独立 API但原理完全相同。我建议你在完成 Lab 1 后自己写一个倒排序索引的练习Map 阶段输入若干文档产出word, 文档IDReduce 阶段输出每个词出现在哪些文档里。你会发现你根本不用再动框架代码只改 map、reduce 两个函数就行。这个练习能最大程度验证你对框架的理解——如果框架的底层是错的你的应用代码写得再对也白搭。6.3 招聘数据清洗这类实训的本质像招聘数据清洗这类综合案例本质是在 Map 阶段做转换过滤、Reduce 阶段做聚合统计。这和你做 WordCount 没有任何本质区别只是 Map 函数换成了解析一行 JSON、过滤掉非法字段、标准化日期格式Reduce 函数换成了求平均薪资、按城市分组。只要框架的调度、shuffle、落盘是稳的这些应用就是换皮。这也是为什么我一再强调6.824 Lab 1 值得认真写——它不是某个特定框架的教程而是所有分布式数据处理框架的地基。6.4 从 Lab 1 迁移到 Hadoop 的时候要注意什么很多人在做完 6.824 之后会去碰 Hadoop 系的实训题此时最容易犯的错误是把 6.824 里Worker 主动拉任务模型硬套到 Hadoop 的框架主动分配任务模型上。Hadoop 的 YARN 是 ResourceManager 主动找 NodeManager 分配 Container和 6.824 里 Worker 循环请求 Master 正好相反。但任务状态机、超时重发、分桶 shuffle 这些核心概念完全一致。你会发现在 6.824 里踩过的坑在 Hadoop 里还会再踩一遍只不过换了名字TaskAttempt、Speculative Execution、ShuffleHandler。理解了本质你就知道这些新名词都只是旧概念的产品化包装。7. 调试与验证我是怎么把 test-mr.sh 跑绿的7.1 先跑通最小路径别急着上并发我第一次写完直接跑go test -run 2D -race被一大屏报错砸晕。后来我学乖了先做单 Worker、单文件的验证开一个终端启动 Master另一个终端启动一个 Worker。跑完看mr-out-0的内容用sort mr-out-*和标准答案 diff。加日志看中间文件在 Map 函数写完文件后打一行日志确认每个 Map 任务写出的文件名和行数都对。确认单 Worker 能出结果后再上多 Worker。多 Worker 最容易出问题的不是正确性而是重复提交和轮询节奏。我建议在任何 RPC handler 的第一行加log.Printf把请求类型、任务 ID、当前状态打出来。日志多了会有点吵但排查并发问题时比单步调试高效得多。7.2 race 检测器才是你最该依赖的调试工具Go 自带的数据竞争检测其实非常实用go build -race之后跑测试一旦有并发读写冲突它会直接打印冲突的 goroutine 栈。这不是可选步骤是必须步骤。我几乎可以断定你不加-race就跑 crash test某些概率性 bug 会在毫无征兆的情况下冒出来而且极难复现加了-race才能让这些 bug 变成确定性报错。使用的时候注意race 只在运行时检测所以要做两件事一是把所有 worker 的数量拉满比如 8 个对应测试脚本的并发峰值二是故意制造超时重发。如果这两个场景下-race都干净基本可以认为没有数据竞争了。7.3 我们踩过最狠的三个坑三个坑值得单独说因为它们都不是语法错误而是语义错误。第一个坑是Worker 内部没有互斥。测试脚本要求每个 Worker 同一时刻只执行一个任务但如果你在 Worker 端用一个 goroutine 调度任务下一次 RPC 返回前上一个任务还没做完两个任务就会交叉执行把中间文件写坏。解决办法是 Worker 主循环用同步方式执行一次只处理一个任务如果确实要用并发必须加锁保证同一时刻只有一条执行路径。第二个坑是Master 提前退出。Master.Done()一定要等到所有 reduce 都完成才返回 true。我犯过的错是把done设定在reduce 任务全部分配完毕而不是全部完成结果测试脚本读输出文件时reduce 还在写文件不全。这个 bug 在单机小数据上可能因为时序侥幸通过但并发一高必然失败。第三个坑是文件句柄没有关闭。中间文件或输出文件写完后如果不Close()在 Linux 上短期内看不出问题但一旦 Worker 数量多、文件数目上千很可能由于句柄耗尽或缓冲未刷出导致丢失数据。我最后把所有落盘操作统一封装到一个函数里defer 关闭从根上解决。7.4 回归测试的正确姿势test-mr.sh跑起来其实有点慢crash test 要等待超时触发所以我建议先跑一个加速版把 Master 超时改成 2 秒跑完一遍确定逻辑没问题再改回 10 秒跑完整版。另外脚本结果要以退出码为准不要只盯着 Terminal 输出。很多时候你看到一堆数字以为跑过了其实脚本末尾报了 FAIL退出码非 0。最后说一句题外话。这个 Lab 做完之后我最大的收获不是学会了 Go 的锁和 RPC而是明白了分布式系统里失败是常态这件事。单机程序里函数要么返回要么崩溃逻辑是确定的分布式系统里一个任务可能执行了 5 分钟才完成但 Master 早以为你挂了把活派给了别人。你写的所有代码本质上都在回答一个问题当意外发生时系统还能不能保持正确带着这个问题去写每一行代码你才能真正把这个 Lab 吃透。