本实验要求我们构建一个分布式的 MapReduce 系统,并实现 paper 中提到的文章字数统计算法。
参考资料
实现
由于 paper 中已经给了算法流程图,只需要严格遵循该图即可。
Coordinator
其实就是 paper 中的 Master
因为不涉及 MapReduce 具体操作,总体而言比较简单。
结构体
Coordinator 负责任务调度以及维护任务状态,正如 paper 中提到的那样。同时,为了记录一个 Worker 是否超时未回应,需维护每个任务开始的时间戳。
于是结构体可以如下定义:
1type TaskState int
2
3const (
4 IDLE TaskState = iota
5 IN_PROGRESS
6 COMPLETED
7)
8
9type Coordinator struct {
10 files []string
11
12 MapTask []TaskState
13 ReduceTask []TaskState
14
15 MapStartTimeStamp []time.Time
16 ReduceStartTimeStamp []time.Time
17
18 M int // total map tasks
19 R int // total reduce tasks
20 Mcnt int // completed map tasks
21 Rcnt int // completed reduce tasks
22
23 State TaskType // current Execution State
24 lock sync.RWMutex
25}
任务调度
根据规则 “reduces can’t start until the last map has finished”,任务调度应当基于当前 MapReduce 执行阶段,即在所有 map 任务被标记为 completed 后才能进行 reduce 任务的调度,否则 reduce Worker 将未能读取部分数据。
各个阶段中任务调度的思想都是一样的,收到一个 Arrange RPC 时,Coordinator 找出一个 idle 或 timeout 的任务并分配;如果没有这样的任务,则需通知 Worker 等待。
1func (m *Coordinator) Arrange(message *Args, reply *Reply) error {
2 m.lock.Lock()
3 defer m.lock.Unlock()
4
5 if m.Rcnt == m.R { // all tasks completed
6 reply.Over = true
7 return nil
8 }
9 if m.State == MAP {
10 for i := 0; i < m.M; i++ {
11 // the task_i is as-yet-unstarted or time-out
12 if m.MapTask[i] == IDLE || (m.MapTask[i] == IN_PROGRESS && time.Since(m.MapStartTimeStamp[i]) > 10*time.Second) {
13 reply.Task = MAP
14 reply.Wait = false
15 reply.Filename = m.files[i]
16 reply.R = m.R
17 reply.MapTaskNumber = i
18 reply.TimeStamp = time.Now()
19
20 m.MapStartTimeStamp[i] = reply.TimeStamp
21 m.MapTask[i] = IN_PROGRESS
22 return nil
23 }
24 }} else if m.State == REDUCE {
25 for i := 0; i < m.R; i++ {
26 // the task_i is as-yet-unstarted or time-out
27 if m.ReduceTask[i] == IDLE || (m.ReduceTask[i] == IN_PROGRESS && time.Since(m.ReduceStartTimeStamp[i]) > 10*time.Second) {
28 reply.Task = REDUCE
29 reply.M = m.M
30 reply.ReduceTaskNumber = i
31 reply.TimeStamp = time.Now()
32 m.ReduceStartTimeStamp[i] = reply.TimeStamp
33 m.ReduceTask[i] = IN_PROGRESS
34 m.lock.Unlock()
35 return nil
36 }
37 }
38 }
39
40 // no more as-yet-unstarted tasks
41 reply.Wait = true
42 return nil
43}
完成通知
同样的,需要忽略那些过期回复。考虑这种场景,Coordinator 发给 WorkerA 的任务超时未完成,然后将该任务调度给 WorkerB,但从 WorkerA 的视角来看,它已经完成了该任务,并发送了完成通知,只不过该通知因为网络拥塞或其他原因过了很久才被 Coordinator 收到,Coordinator 如何判断这个通知是不是当前正在执行该任务的 Worker 发送的呢?一个思路是可以维护每个任务当前被哪个 Worker 执行,在这种情况下,任务的当前执行者为 B,那么收到 A 的通知后理所当然会将其忽略。
但还不足够,设想一下,如果 B 也因为同样原因超时,任务再一次调度给了 A 呢?虽然任务当前执行者为 A,但 A 并未完成,Coordinator 收到的是非常古老的那条消息,此时将任务设为 completed,显然会出现问题——A 并没有完成任务,尽管它"完成"过一次。
考虑到 Coordinator 结构体里维护了每个任务的开始时间戳,不妨令任务调度与完成通知中都捎带本次任务的 StartTime,这样就可以进行检查,如果一致则正确接收;反之视为过期消息。这样做减少了额外信息维护,也提高了安全性。
这其实就相当于分配了一个递增的版本号了。
1// an RPC handler to tell the Coordinator that the worker finishes the task
2func (m *Coordinator) Finished(message *Args, reply *Reply) error {
3 m.lock.Lock()
4 defer m.lock.Unlock()
5
6 if message.Finished == MAP {
7 if m.MapStartTimeStamp[message.MapTaskNumber].Sub(message.TimeStamp) != 0 {
8 reply.Wait = true
9 return nil
10 }
11 m.MapTask[message.MapTaskNumber] = COMPLETED
12 m.Mcnt++
13 if m.Mcnt == m.M {
14 m.State = REDUCE
15 }
16 } else {
17 if m.ReduceStartTimeStamp[message.ReduceTaskNumber].Sub(message.TimeStamp) != 0 {
18 reply.Wait = true
19 return nil
20 }
21 m.ReduceTask[message.ReduceTaskNumber] = COMPLETED
22 m.Rcnt++
23 }
24 return nil
25}
Worker
Worker 就涉及具体的 MapReduce 操作了,好在课程给的代码提供了 map/reduce 函数,我们只需要关注对输入输出的处理即可。在本系统中,任务的调度采用 Worker pull 而非 Coordinator push 的策略。Worker 需不断请求任务,然后根据回复内容执行对应的操作:
如果所有任务已结束,reply 告知 Over,关闭线程;
如果没有任务能分配,则调用
time.Sleep()等待一段时间;如果收到一个 map 任务,此时 reply 会捎带所要操作的文件名,然后对文件进行 map 操作(这里需要去 mrsequential.go 里参考一下代码)。由于它这里需要考虑线程崩溃,先把结果写到临时文件,全部写完后(说明没有发生 crash)再输出到
/mr-tmp中,那么 mrsequential.go 里的代码就不能完全照搬了。我的做法是,开一个长度为 nReduce,类型为 []KeyValue 的切片
temp,其中temp[i]存放输出到mr-X-i中的所有 kv 对,先把所有 kv 对按照 key 的hash 值写到对应的temp[hash(key) % nReduce]里,再逐个写到临时文件中。示例提供的 json.Encoder() 在第二次及以后打开的时候不会在末尾添加,而是直接覆盖。
reduce 任务大体类似,就是读取
mr-X-Y文件。因为 Y 是固定的,所以 reply 要捎带 nMap。也是要对 mrsequential.go 里的代码进行一些修改。这里遇到一个坑点,和 paper 中所描述的产生了冲突。就是输出文件如果存在就不进行
os.Rename(),否则在 crash test 中会出现mr-X-?都有内容而mr-out-?没有内容的情况。合理猜测是读取 mr-X-Y 的时候出了点什么问题。
执行
1func Worker(mapf func(string, string) []KeyValue,
2 reducef func(string, []string) string) {
3
4 // One way to get started is to modify mr/worker.go's Worker()
5 // to send an RPC to the coordinator asking for a task.
6 for {
7 reply := AskForTask()
8 if reply.Over {
9 return
10 }
11
12 if reply.Wait {
13 time.Sleep(WorkerWaitTime)
14 } else if reply.Task == MAP {
15 // map phase
16 mapTaskNumber := reply.MapTaskNumber
17
18 // read file
19 file, err := os.Open(reply.Filename)
20 if err != nil {
21 log.Fatalf("cannot open %v", reply.Filename)
22 }
23 content, err := ioutil.ReadAll(file)
24 if err != nil {
25 log.Fatalf("cannot read %v", reply.Filename)
26 }
27 file.Close()
28 kvs := mapf(reply.Filename, string(content))
29
30 // write kv into file buckets
31 // key -> filename: "ihash(key)"
32 temp := make([][]KeyValue, reply.R)
33 for i := range temp {
34 temp[i] = make([]KeyValue, 0)
35 }
36
37 for _, kv := range kvs {
38 hash := ihash(kv.Key) % reply.R
39 temp[hash] = append(temp[hash], kv)
40 }
41
42 for i := 0; i < len(temp); i++ {
43 ofile, _ := ioutil.TempFile("./mr/mapfile", fmt.Sprintf("%d", i))
44 enc := json.NewEncoder(ofile)
45 for _, kv := range temp[i] {
46 enc.Encode(&kv)
47 }
48
49 old_path := ofile.Name()
50 new_path := fmt.Sprintf("../main/mr-tmp/mr-%d-%d", mapTaskNumber, i)
51
52 os.Rename(old_path, new_path)
53 ofile.Close()
54 }
55
56 // tell the master that the map job is done
57 CallFinish(MAP, reply.TimeStamp, mapTaskNumber, 0)
58 } else {
59 // reduce phase
60 reduceTaskNumber := reply.ReduceTaskNumber
61
62 // read file
63 ofile, _ := ioutil.TempFile("./mr/reducefile", fmt.Sprintf("%d", reduceTaskNumber))
64 var kva []KeyValue
65 for i := 0; i < reply.M; i++ {
66 iFilename := fmt.Sprintf("../main/mr-tmp/mr-%d-%d", i, reduceTaskNumber)
67 iFile, err := os.Open(iFilename)
68 if err == nil {
69 dec := json.NewDecoder(iFile)
70 for {
71 var kv KeyValue
72 if err := dec.Decode(&kv); err != nil {
73 break
74 }
75 kva = append(kva, kv)
76 }
77 } else {
78 log.Fatal(err)
79 }
80 }
81 sort.Sort(ByKey(kva))
82 i := 0
83 for i < len(kva) {
84 j := i + 1
85 for j < len(kva) && kva[j].Key == kva[i].Key {
86 j++
87 }
88 values := []string{}
89 for k := i; k < j; k++ {
90 values = append(values, kva[k].Value)
91 }
92 output := reducef(kva[i].Key, values)
93
94 // this is the correct format for each line of Reduce output.
95 fmt.Fprintf(ofile, "%v %v\n", kva[i].Key, output)
96
97 i = j
98 }
99
100 old_path := ofile.Name()
101 new_path := fmt.Sprintf("../main/mr-tmp/mr-out-%d", reduceTaskNumber)
102
103 // no reduce task finished yet before
104 if _, err := os.Open(new_path); err != nil {
105 os.Rename(old_path, new_path)
106 }
107 ofile.Close()
108
109 // tell the master that the reduce job is done
110 CallFinish(REDUCE, reply.TimeStamp, 0, reduceTaskNumber)
111 }
112 }
113}
RPC
综上所述,RPC 的结构体定义就呼之欲出了。
1type Args struct {
2 Finished TaskType
3
4 TimeStamp time.Time
5
6 MapTaskNumber int
7 ReduceTaskNumber int
8}
9
10type Reply struct {
11 Task TaskType
12 Wait bool // true for wait
13 Over bool // true for done
14
15 Filename string
16 M int
17 R int
18
19 MapTaskNumber int
20 ReduceTaskNumber int
21
22 TimeStamp time.Time
23}