lab / 2022.09.22

6.824 Lab1 MapReduce

本实验要求我们构建一个分布式的 MapReduce 系统,并实现 paper 中提到的文章字数统计算法。

本实验要求我们构建一个分布式的 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 找出一个 idletimeout 的任务并分配;如果没有这样的任务,则需通知 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 需不断请求任务,然后根据回复内容执行对应的操作:

  1. 如果所有任务已结束,reply 告知 Over,关闭线程;

  2. 如果没有任务能分配,则调用 time.Sleep() 等待一段时间;

  3. 如果收到一个 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() 在第二次及以后打开的时候不会在末尾添加,而是直接覆盖。

  4. 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}