• MapReduce极简实现


    0 概述

    MapReduce是一种广泛运用的分布式-大数据计算编程模型,最初由Google发表,其开源实现为Hadoop。

    MapReduce 的编程模型非常简单,正如名字一样,用户仅仅需要实现一个 Map 函数,一个 Reduce 函数。

    • Map 函数,即映射函数:它会接受一个 key-value 对,然后把这个 key-value 对转换成 0 到多个新的 key-value 对并输出出去。

      map (k1, v1) -> list (k2, v2)
    • Reduce 函数,即化简函数:它接受一个 Key,以及这个 Key 下的一组 Value,然后化简成一组新的值 Value 输出出去。

      reduce (k2, list(v2)) -> list(v3)

    可以解决的任务例子:

    • 分布式 grep;
    • 统计 URL 的访问频次;
    • 反转网页 - 链接图;
    • 分域名的词向量;
    • 生成倒排索引;
    • 分布式排序。

    1 MapReduce结构

    一图胜千言:

    截屏2022-06-29 21.47.52

    2 总体设计

    以完成6.8242021Spring的lab1为目标。

    可以通过以下git命令:clone代码:

    git clone git://g.csail.mit.edu/6.824-golabs-2021 6.824

    master采用lazy分配任务方法,由worker主动去触发任务分配、任务结束等操作。master分配不同的块给不同的worker执行。

    因此worker需要实现获取任务,任务结束等RPC,代码如下:

    1. type GetTaskArgs struct {
    2. }
    3. type GetTaskReply struct {
    4. Type TaskType
    5. Filenames []string
    6. Task_no int
    7. NReduce int
    8. Err Errno
    9. }
    10. type FinishTaskArgs struct {
    11. Type TaskType
    12. Task_no int
    13. }
    14. type FinishTaskReply struct {
    15. Err Errno
    16. }

    3 worker设计

    worker的工作就是不断获取任务,若任务完成则提交之。

    其主要代码为:

    1. func Worker(mapf func(string, string) []KeyValue,
    2. reducef func(string, []string) string) {
    3. // Your worker implementation here.
    4. for {
    5. args := GetTaskArgs{}
    6. reply := GetTaskReply{}
    7. log.Printf("get task request: %v\n", args)
    8. ok := CallGetTask(&args, &reply)
    9. log.Printf("recv get task reply: %v\n", reply)
    10. if !ok || reply.Type == STOP {
    11. break
    12. }
    13. // handle map fynction
    14. switch reply.Type {
    15. case MAP:
    16. if len(reply.Filenames) < 1 {
    17. log.Fatalf("don't have filename")
    18. }
    19. DoMAP(reply.Filenames[0], reply.Task_no, reply.NReduce, mapf)
    20. // map complete, send msg to master
    21. finish_args := FinishTaskArgs{
    22. Type: MAP,
    23. Task_no: reply.Task_no,
    24. }
    25. finish_reply := FinishTaskReply{}
    26. log.Printf("finish request: %v\n", finish_args)
    27. CallFinishTask(&finish_args, &finish_reply)
    28. log.Printf("recv finish reply: %v\n", finish_reply)
    29. // time.Sleep(time.Second)
    30. case REDUCE:
    31. if len(reply.Filenames) < 1 {
    32. log.Fatalf("don't have filenames")
    33. }
    34. DoReduce(reply.Filenames, reply.Task_no, reducef)
    35. // reduce complete, send msg to master
    36. finish_args := FinishTaskArgs{
    37. Type: REDUCE,
    38. Task_no: reply.Task_no,
    39. }
    40. finish_reply := FinishTaskReply{}
    41. log.Printf("finish request: %v\n", finish_args)
    42. CallFinishTask(&finish_args, &finish_reply)
    43. log.Printf("recv finish reply: %v\n", finish_reply)
    44. // time.Sleep(time.Second)
    45. case WAIT:
    46. log.Printf("wait task\n")
    47. time.Sleep(time.Second)
    48. default:
    49. time.Sleep(time.Second)
    50. }
    51. }
    52. }

    其中分MAP、REDUCE、WAIT和STOP四个状态:

    • MAP:进行MAP操作
    • REDUCE:进行REDECE操作
    • WAIT:等待其他worker完成任务(比如等待在总体MAP任务的收尾上,以及没有更多的MAP任务可以分配了)
    • STOP:worker停止、退出

    其中最重要的为map和reduce任务的执行。

    map任务的执行实现代码如下:(对应上图中的2、3、4步)

    1. func DoMAP(filename string, task_no int, nReduce int, mapf func(string, string) []KeyValue) {
    2. file, err := os.Open(filename)
    3. if err != nil {
    4. log.Fatalf("cannot open %v", filename)
    5. }
    6. content, err := ioutil.ReadAll(file)
    7. if err != nil {
    8. log.Fatalf("cannot read %v", filename)
    9. }
    10. file.Close()
    11. kva := mapf(filename, string(content))
    12. sort.Sort(ByKey(kva))
    13. log.Println("encode to json")
    14. files := make([]*os.File, nReduce)
    15. encoders := make([]*json.Encoder, nReduce)
    16. for i := 0; i < nReduce; i++ {
    17. ofile, err := ioutil.TempFile("", "mr-tmp*")
    18. if err != nil {
    19. log.Fatalf("cannot create temp file")
    20. }
    21. defer ofile.Close()
    22. encoder := json.NewEncoder(ofile)
    23. encoders[i] = encoder
    24. files[i] = ofile
    25. }
    26. var index int
    27. for _, kv := range kva {
    28. index = ihash(kv.Key) % nReduce
    29. err = encoders[index].Encode(&kv)
    30. if err != nil {
    31. log.Fatalf("cannot encode %v", kv)
    32. }
    33. }
    34. // atomically rename
    35. for i := 0; i < nReduce; i++ {
    36. filename_tmp := fmt.Sprintf("mr-%d-%d", task_no, i)
    37. err := os.Rename(files[i].Name(), filename_tmp)
    38. if err != nil {
    39. log.Fatalf("cannot rename %v to %v", files[i].Name(), filename_tmp)
    40. }
    41. }
    42. }

    比较有意思的是map需要通过一个hash函数将相同的条目分布在同一输出文件中:

    1. func ihash(key string) int {
    2. h := fnv.New32a()
    3. h.Write([]byte(key))
    4. return int(h.Sum32() & 0x7fffffff)
    5. }
    6. var index int
    7. for _, kv := range kva {
    8. index = ihash(kv.Key) % nReduce
    9. err = encoders[index].Encode(&kv)
    10. if err != nil {
    11. log.Fatalf("cannot encode %v", kv)
    12. }
    13. }

    reduce任务的执行实现代码如下:(对应上图中的5、6步)

    1. func DoReduce(filenames []string, task_no int, reducef func(string, []string) string) {
    2. // read data from mid-file
    3. kva := make([]KeyValue, 0)
    4. for _, filename := range filenames {
    5. file, err := os.Open(filename)
    6. if err != nil {
    7. log.Fatalf("cannot open %v", filename)
    8. }
    9. defer file.Close()
    10. dec := json.NewDecoder(file)
    11. for {
    12. var kv KeyValue
    13. if err := dec.Decode(&kv); err != nil {
    14. break
    15. }
    16. kva = append(kva, kv)
    17. }
    18. }
    19. sort.Sort(ByKey(kva))
    20. // call Reduce on each distinct key in kva[],
    21. // and print the result to mr-out-0.
    22. ofile, err := ioutil.TempFile("", "mr-out-tmp*")
    23. if err != nil {
    24. log.Fatalf("cannot create temp file")
    25. }
    26. defer ofile.Close()
    27. i := 0
    28. for i < len(kva) {
    29. j := i + 1
    30. for j < len(kva) && kva[j].Key == kva[i].Key {
    31. j++
    32. }
    33. values := []string{}
    34. for k := i; k < j; k++ {
    35. values = append(values, kva[k].Value)
    36. }
    37. output := reducef(kva[i].Key, values)
    38. // this is the correct format for each line of Reduce output.
    39. fmt.Fprintf(ofile, "%v %v\n", kva[i].Key, output)
    40. i = j
    41. }
    42. output_filename := fmt.Sprintf("mr-out-%d", task_no)
    43. err = os.Rename(ofile.Name(), output_filename)
    44. if err != nil {
    45. log.Fatalf("cannot rename %v to %v", ofile.Name(), output_filename)
    46. }
    47. }

    按道理应该是要在GFS上读写文件的,条件不允许,就直接采用UNIX的文件系统了。

    4 master设计

    master的设计还是比较简单的,只包含很少的信息:

    1. type Coordinator struct {
    2. tasks []Task
    3. nReduce int
    4. nMap int
    5. status CoordinatorStatus
    6. mu sync.Mutex
    7. }

    对所需要进行的任务信息进行定义,如下:

    1. type TaskStatus int
    2. const (
    3. idle TaskStatus = iota
    4. in_progress
    5. completed
    6. )
    7. type Task struct {
    8. tno int
    9. filenames []string
    10. status TaskStatus
    11. startTime time.Time
    12. }

    其主要就是接受worker的两个RPC请求。

    获取任务的RPC handler实现如下:

    • 对于长时间(10s)未完成的任务,重新制定一个worker执行此任务。
    1. func (c *Coordinator) GetTask(args *GetTaskArgs, reply *GetTaskReply) error {
    2. c.mu.Lock()
    3. defer c.mu.Unlock()
    4. finish_flag := c.IsAllFinish()
    5. if finish_flag {
    6. c.NextPhase()
    7. }
    8. for i := 0; i < len(c.tasks); i++ {
    9. if c.tasks[i].status == idle {
    10. log.Printf("send task %d to worker\n", i)
    11. reply.Err = SuccessCode
    12. reply.Task_no = i
    13. reply.Filenames = c.tasks[i].filenames
    14. if c.status == MAP_PHASE {
    15. reply.Type = MAP
    16. reply.NReduce = c.nReduce
    17. } else if c.status == REDUCE_PHASE {
    18. reply.NReduce = 0
    19. reply.Type = REDUCE
    20. } else {
    21. log.Fatal("unexpected status")
    22. }
    23. c.tasks[i].startTime = time.Now()
    24. c.tasks[i].status = in_progress
    25. return nil
    26. } else if c.tasks[i].status == in_progress {
    27. curr := time.Now()
    28. if curr.Sub(c.tasks[i].startTime) > time.Second*10 {
    29. log.Printf("resend task %d to worker\n", i)
    30. reply.Err = SuccessCode
    31. reply.Task_no = i
    32. reply.Filenames = c.tasks[i].filenames
    33. if c.status == MAP_PHASE {
    34. reply.Type = MAP
    35. reply.NReduce = c.nReduce
    36. } else if c.status == REDUCE_PHASE {
    37. reply.NReduce = 0
    38. reply.Type = REDUCE
    39. } else {
    40. log.Fatal("unexpected status")
    41. }
    42. c.tasks[i].startTime = time.Now()
    43. return nil
    44. }
    45. }
    46. }
    47. reply.Err = SuccessCode
    48. reply.Type = WAIT
    49. return nil
    50. }

    完成任务的RPC handler实现如下:

    1. func (c *Coordinator) FinishTask(args *FinishTaskArgs, reply *FinishTaskReply) error {
    2. c.mu.Lock()
    3. defer c.mu.Unlock()
    4. if args.Task_no >= len(c.tasks) || args.Task_no < 0 {
    5. reply.Err = ParaErrCode
    6. return nil
    7. }
    8. c.tasks[args.Task_no].status = completed
    9. if c.IsAllFinish() {
    10. c.NextPhase()
    11. }
    12. return nil
    13. }

    检查全部任务是否完成,完成就进入下一个阶段:

    1. func (c *Coordinator) IsAllFinish() bool {
    2. for i := len(c.tasks) - 1; i >= 0; i-- {
    3. if c.tasks[i].status != completed {
    4. return false
    5. }
    6. }
    7. return true
    8. }
    9. func (c *Coordinator) NextPhase() {
    10. if c.status == MAP_PHASE {
    11. log.Println("change to REDUCE_PHASE")
    12. c.MakeReduceTasks()
    13. c.status = REDUCE_PHASE
    14. } else if c.status == REDUCE_PHASE {
    15. log.Println("change to FINISH_PHASE")
    16. c.status = FINISH_PHASE
    17. } else {
    18. log.Println("unexpected status change!")
    19. }
    20. }

    客户端查看MapReduce任务是否完成:

    1. func (c *Coordinator) Done() bool {
    2. c.mu.Lock()
    3. defer c.mu.Unlock()
    4. if c.status == FINISH_PHASE {
    5. return true
    6. }
    7. return false
    8. }

    5 客户端如何使用呢?

    写两个函数(Map和Reduce)就行啦:

    1. //
    2. // The map function is called once for each file of input. The first
    3. // argument is the name of the input file, and the second is the
    4. // file's complete contents. You should ignore the input file name,
    5. // and look only at the contents argument. The return value is a slice
    6. // of key/value pairs.
    7. //
    8. func Map(filename string, contents string) []mr.KeyValue {
    9. // function to detect word separators.
    10. ff := func(r rune) bool { return !unicode.IsLetter(r) }
    11. // split contents into an array of words.
    12. words := strings.FieldsFunc(contents, ff)
    13. kva := []mr.KeyValue{}
    14. for _, w := range words {
    15. kv := mr.KeyValue{w, "1"}
    16. kva = append(kva, kv)
    17. }
    18. return kva
    19. }
    20. //
    21. // The reduce function is called once for each key generated by the
    22. // map tasks, with a list of all the values created for that key by
    23. // any map task.
    24. //
    25. func Reduce(key string, values []string) string {
    26. // return the number of occurrences of this word.
    27. return strconv.Itoa(len(values))
    28. }

    6 附录

    详细代码可以参考:

    仓库

    commit

  • 相关阅读:
    vue 封装菜单组件 来回跳转使菜单高亮
    音视频技术开发周刊 | 257
    若依分离版——使用Knife4j 自动生成接口文档
    信奥中的数学基础:多边形内角和 编程常用英语词汇
    基于Matlab的汽车安全应用轨道融合仿真(附源码)
    c++ 中头文件
    6.Vgg16--CNN经典网络模型详解(pytorch实现)
    vue3+ts withDefaults的使用
    Leetcode225.用队列实现栈
    P4068 [SDOI2016]数字配对
  • 原文地址:https://blog.csdn.net/hrbust_cxl/article/details/125532595