6.584-Lab1:MapReduce

前置知识/概念

Raft

是一个基于“Leader”的协议,能够保证分布式网路的一致性。

RPC(Remote Producer Call)

参考链接1
参考链接2
Go中RPC的简单实现

Golang中regexp正则表达式的用法

https://gukaifeng.cn/posts/golang-zheng-ze-biao-da-shi-regexp-de-ji-ben-yong-fa/index.html

Golang中自定义类型Sort

在提供的排序方法sort.Ints、sort.Floats、sort.Strings的底层都分别实现了三个函数Len()Less()Swap(),所以在我们实现自定义类型排序的时候要实现上述三个函数。

实现

参考链接

概览

在这里插入图片描述
如图,Lab1要实现两个部分分别是Map(映射)和Reduce(规约),Master负责调度,分配任务在代码中是Coordinator,而Worker则负责具体的map task和reduce task。
Map Task具体是统计文件中的单词生成对应的键值对[key, value],通过一个哈希函数Ihash将单词映射为key,存在中间文件,例如File1中有单词a、b对应键值对为[a, 1],[b, 1],分别存放在mr-out-1-ihash(a/b)%NReduce的中间文件中,其中NReduce是Reduce Task的个数,代码中为10。
Reduce Task具体是将中间文件mr-out-*-taskid中的键值对放入最终的文件mr-out-taskid中。

代码

rpc.go

Coordinator和Worker通过RPC进行通信,在文件rpc.go中需要设计相应的数据结构让其进行通信。分析可能的具体行为:Worker需要向Coordinator申请任务、Coordinator需要向Worker分配任务(map task & reduce task)、Worker向Coordinator回复任务的执行情况(成功、失败)、没有闲置任务分配时Coordinator告诉Worker等待(Wait)、所有任务完成告诉Worker结束(Shutdown)。
根据上述分析设计如下数据类型:

// 用不同数字表示不同信息的类别
type MsgType intconst (AskForTask    MsgType = iota // 表示worker向coordinator申请任务MapSucceed                   // 表示worker向coordinator传递Map Task完成MapFailed                    // 表示worker向coordinator传递Map Task失败ReduceSucceed                // 表示worker向coordinator传递Reduce Task完成ReduceFailed                 // 表示worker向coordinator传递Reduce Task失败MapAlloc                     // 表示coordinator向worker分配Map TaskReduceAlloc                  // 表示coordinator向worker分配Reduce TaskWait                         // 表示coordinator让worker休眠Shutdown                     // 表示coordinator让worker终止
)

从Worker视角出发,设计发送给Coordinator的信息结构体MsgSend,需要包含任务id以及任务执行情况:

type MsgSend struct {MsgType MsgTypeTaskId  int
}

从Coordinator视角出发,设计发送给Worker的信息的结构体MsgSend,需要包含任务的id即要处理的第几个文件用于中间文件的命名、任务类型、要处理的filename、NRduce:

type MsgReply struct {MsgType  MsgTypeNReduce  intTaskId   int    // 当worker发送MsgSend申请任务、coordinator回复任务的IDTaskName string //
}

worker.go

上面分析可知Worker的行为有:申请任务、执行任务、汇报任务执行情况。

// worker的任务就是不断请求任务、执行任务、报告执行状态
// main/mrworker.go calls this function.
func Worker(mapf func(string, string) []KeyValue,reducef func(string, []string) string) {// Your worker implementation here.for {// 不断请求replyMsg := CallForTask()switch replyMsg.MsgType {case MapAlloc: // coordinator分配了map taskerr := HandleMapTask(replyMsg, mapf)if err != nil { // Map Task任务完成_ = CallForReportStatus(MapFailed, replyMsg.TaskId)} else { // Map Task 任务失败_ = CallForReportStatus(MapSucceed, replyMsg.TaskId)}case ReduceAlloc:err := HandleReduceTask(replyMsg, reducef)if err != nil { // Map Task任务完成_ = CallForReportStatus(ReduceFailed, replyMsg.TaskId)} else { // Map Task 任务失败_ = CallForReportStatus(ReduceSucceed, replyMsg.TaskId)}case Wait:time.Sleep(time.Second * 10)case Shutdown:os.Exit(0)}time.Sleep(time.Second)}
}

其中申请任务、回复任务执行情况的函数均利用了rpc中的Call来实现向Coordinator进行通信。
在这里插入图片描述
在这里插入图片描述

其中执行Map Task 的执行函数HandleMapTask

func HandleMapTask(reply *MessageReply, mapf func(string, string) []KeyValue) error {file, err := os.Open(reply.TaskName)if err != nil {return err}defer file.Close()content, err := io.ReadAll(file)if err != nil {return err}kva := mapf(reply.TaskName, string(content))sort.Sort(ByKey(kva))tempFiles := make([]*os.File, reply.NReduce)encoders := make([]*json.Encoder, reply.NReduce)for _, kv := range kva {redId := ihash(kv.Key) % reply.NReduceif encoders[redId] == nil {tempFile, err := ioutil.TempFile("", fmt.Sprintf("mr-map-tmp-%d", redId))if err != nil {return err}defer tempFile.Close()tempFiles[redId] = tempFileencoders[redId] = json.NewEncoder(tempFile)}err := encoders[redId].Encode(&kv)if err != nil {return err}}for i, file := range tempFiles {if file != nil {fileName := file.Name()file.Close()newName := fmt.Sprintf("mr-out-%d-%d", reply.TaskID, i)if err := os.Rename(fileName, newName); err != nil {return err}}}return nil
}

在这里插入图片描述
由函数TemFile的描述可知,该函数在创建时会按照pattern+随机字符串作为临时文件名字,即是不同程序同时调用该函数会创建不同的临时文件所以不存在资源竞争是并发安全的。
先将Map Task映射的键值对存入临时文件中,等全部放入临时文件后再将临时文件重命名为需要的中间文件的名字,因为重命名操作时原子性的。
如果不采用临时文件直接存入目标中间文件的话,会出现存入中间文件之前中间文件中含有其他数据即脏数据,可能是上个Worker执行到一半因为某些原因而退出之前存入的,所以存入之前需要将中间文件清空一下,这样就比较浪费时间。

执行Reduce Task 的执行函数HandleReduceTask也是同样的问题,虽然从中间件读取的时候没有写入操作但写入最终文件时也同样需要像上面一样保证原子性,这里贴出没有使用临时文件的代码:

// 处理分配的 Reduce 任务,处理每个MapTask产生的mr-out-*-key_id
func HandleReduceTask(reply *MsgReply, reducef func(string, []string) string) error {key_id := reply.TaskId// todo:这里的key_id要 % NReduce吗 --答:传入的TaskId一定是小于NReduce的,规约的数量就是NRduce也即Reduce Task的数量files, err := ReadSpecificFile(key_id, "./")if err != nil {return err}// 从所有匹配的文件中读出Json格式的键值对[k1-value]/[k2-value],key哈希的值可以不一样但这里ihash(k1) % NReduce  = ihash(k2) % NReducek_vs := map[string][]string{}for _, file := range files {dec := json.NewDecoder(file)for { // 循环读JSON数据流var kv KeyValue // 将读出的JSON数据解码放入kvif err := dec.Decode(&kv); err != nil {break}k_vs[kv.Key] = append(k_vs[kv.Key], kv.Value)}file.Close()}keys := []string{} // 将map中的keys拿出排序,按照keys的字典序依次写入文件for k, _ := range k_vs {keys = append(keys, k)}sort.Strings(keys)oname := "mr-out-" + strconv.Itoa(reply.TaskId) // 将Reduce后的结果放入 mr-out-TaskIdofile, err := os.Create(oname)if err != nil {return err}defer ofile.Close()for _, key := range keys {output := reducef(key, k_vs[key])_, err := fmt.Fprintf(ofile, "%s %s\n", key, output) // 格式化写入文件if err != nil {return err}}CleanFileByReduceId(reply.TaskId, "./")return nil
}

利用临时文件保证原子性参考上面HandleMapTask函数的实现。其中需要注意的是最终写入文件的时候Key要按照字典升序排列,而map存储的Key不是有序的,所以把Key拿出来排序再放入文件。

coordinator.go

coordinator中要实现worker申请任务的函数AskForTask以及对worker报告任务执行情况后对相应任务状态的更新函数NoticeResult
任务的状态有闲置idle、成功finished、失败failed、超时、正在运行running。每次worker申请任务时都轮询一下所有任务,委派闲置、失败、超时的任务,而超时状态的判断则通过为每个任务打上运行开始的时间戳,若轮询到running的任务时判断当前时间与开始的时间戳比较若大于10s则可以判定为超时,可以再次委派给worker。

TaskInfo结构
type TaskStatue int// 定义类别
const (idle     TaskStatue = iota // 闲置finished                   // 完成running                    // 运行failed                     // 失败
)type MapTaskInfo struct {statue    TaskStatue // 任务状态TaskId    int        // 任务编号startTime int64      // 分配时间,当前时间-分配时间>10s表示超时
}
type ReduceTaskInfo struct {statue TaskStatue//TaskId reduce task的编号用数组下标表示startTime int64
}
type Coordinator struct {NReduce     int // reduce tasks可用的数量MapTasks    map[string]*MapTaskInfomu          sync.Mutex // 互斥锁ReduceTasks []*ReduceTaskInfo
}

初始化函数:

// 初始化函数
func (c *Coordinator) Init(files []string) {for idx, filename := range files {c.MapTasks[filename] = &MapTaskInfo{statue: idle, // 初始为闲置TaskId: idx,}}for idx := 0; idx < c.NReduce; idx++ {c.ReduceTasks[idx] = &ReduceTaskInfo{statue: idle,}}
}// create a Coordinator.
// main/mrcoordinator.go calls this function.
// nReduce is the number of reduce tasks to use.
func MakeCoordinator(files []string, nReduce int) *Coordinator {c := Coordinator{NReduce:     nReduce,MapTasks:    make(map[string]*MapTaskInfo),ReduceTasks: make([]*ReduceTaskInfo, nReduce),}c.Init(files)c.server()return &c
}
rpc相应函数AskForTask的实现

AskForTask:实现相对复杂,大致流程为:
1.每个任务初始化为闲置。
2.每当有worker申请任务的话就轮询所有任务。
3.若存在有idle、failed、超时的任务就可以分配。
4.若没有可分配的任务,判断完成任务的个数是否等于所有任务个数:
4.1 若相等,则表明所有任务完成,告知workershutdown
4.2 若不相等,则表明有任务还在进行且没有可分配的任务,告知workerWait

期间应保证共享资源的互斥,每个worker请求任务的同时要上锁。

func (c *Coordinator) AskForTask(req *MsgSend, reply *MsgReply) error {if req.MsgType != AskForTask { // 传入的不是“申请任务”类型的信息return NoMathMsgType}// 加锁,保证每个worker申请任务时互斥c.mu.Lock()defer c.mu.Unlock()// 选择一个失败or闲置or超时的任务分配给workerMapSuccessNum := 0 // Map task 完成个数for filename, maptaskinfo := range c.MapTasks {alloc := falseif maptaskinfo.statue == idle || maptaskinfo.statue == failed { // 该任务闲置或失败则可以分配alloc = true} else if maptaskinfo.statue == running { // 判断该任务是否超时,若超时则再分配if time.Now().Unix()-maptaskinfo.startTime > 10 {maptaskinfo.startTime = time.Now().Unix() // 再分配更新开始时间alloc = true}} else { // 该任务是已完成任务MapSuccessNum++}// 当前任务可以分配if alloc {reply.TaskId = maptaskinfo.TaskIdreply.TaskName = filenamereply.NReduce = c.NReducereply.MsgType = MapAllocmaptaskinfo.statue = runningmaptaskinfo.startTime = time.Now().Unix()return nil}}// 没有任务可以分配但所有任务没有完成if MapSuccessNum < len(c.MapTasks) {reply.MsgType = Waitreturn nil}// 运行到这里表明所有的Map任务都已经完成ReduceSuccessNum := 0for idex, reducetaskinfo := range c.ReduceTasks {alloc := falseif reducetaskinfo.statue == idle || reducetaskinfo.statue == failed {alloc = true} else if reducetaskinfo.statue == running {if time.Now().Unix()-reducetaskinfo.startTime > 10 {reducetaskinfo.startTime = time.Now().Unix()alloc = true}} else {ReduceSuccessNum++}if alloc {reply.TaskId = idexreply.NReduce = c.NReducereply.MsgType = ReduceAllocreducetaskinfo.statue = runningreducetaskinfo.startTime = time.Now().Unix()return nil}}if ReduceSuccessNum < len(c.ReduceTasks) {reply.MsgType = Waitreturn nil}// 运行到这里表明所有的任务都已完成reply.MsgType = Shutdownreturn nil
}
rpc相应函数NoticeResult的实现

只需要将worker传递过来的任务完成状态更新到coordinator的TaskInfo即可。

func (c *Coordinator) NoticeResult(req *MsgSend, reply *MsgReply) error {c.mu.Lock()defer c.mu.Unlock()if req.MsgType == MapSucceed {for _, taskinfo := range c.MapTasks {if taskinfo.TaskId == req.TaskId {taskinfo.statue = finished}}} else if req.MsgType == ReduceSucceed {c.ReduceTasks[req.TaskId].statue = finished} else if req.MsgType == MapFailed {for _, taskinfo := range c.MapTasks {if taskinfo.TaskId == req.TaskId {taskinfo.statue = failed}}} else if req.MsgType == ReduceFailed {c.ReduceTasks[req.TaskId].statue = failed}return nil
}
所有任务是否完成-Done函数

只需要轮询一遍任务数组,判断是否所有任务均完成。

// if the entire job has finished.
func (c *Coordinator) Done() bool {// Your code here.// 遍历所有任务,全部完成则返回true,否则返回falsefor _, taskinfo := range c.MapTasks {if taskinfo.statue != finished {return false}}for _, taskinfo := range c.ReduceTasks {if taskinfo.statue != finished {return false}}return true
}

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.rhkb.cn/news/473694.html

如若内容造成侵权/违法违规/事实不符,请联系长河编程网进行投诉反馈email:809451989@qq.com,一经查实,立即删除!

相关文章

脑机接口、嵌入式 AI 、工业级 MR、空间视频和下一代 XR 浏览器丨RTE2024 空间计算和新硬件专场回顾

这一轮硬件创新由 AI 引爆&#xff0c;或许最大受益者仍是 AI&#xff0c;因为只有硬件才能为 AI 直接获取最真实世界的数据。 在人工智能与硬件融合的新时代&#xff0c;实时互动技术正迎来前所未有的创新浪潮。从嵌入式系统到混合现实&#xff0c;从空间视频到脑机接口&…

Restful API接⼝简介及为什么要进⾏接⼝压测

一、RESTful API简介 在现代Web开发中&#xff0c;RESTful API已经成为一种标准的设计模式&#xff0c;用于构建和交互网络应用程序。本文将详细介绍RESTful API的基本概念、特点以及如何使用它来设计高效的API接口。 1. 基于协议 HTTP 或 HTTPS RESTful API通常使用HTTP&am…

面试经典 150 题:20、2、228、122

20. 有效的括号 参考代码 #include <stack>class Solution { public:bool isValid(string s) {if(s.size() < 2){ //特判&#xff1a;空字符串和一个字符的情况return false;}bool flag true;stack<char> st; //栈for(int i0; i<s.size(); i){if(s[i] ( |…

Python爬虫下载新闻,Flask展现新闻(2)

上篇讲了用Python从新闻网站上下载新闻&#xff0c;本篇讲用Flask展现新闻。关于Flask安装网上好多教程&#xff0c;不赘述。下面主要讲 HTML-Flask-数据 的关系。 简洁版 如图&#xff0c;页面简单&#xff0c;主要显示新闻标题。 分页&#xff0c;使用最简单的分页技术&…

基于Java和Vue实现的上门做饭系统上门做饭软件厨师上门app

市场前景 生活节奏加快&#xff1a;在当今快节奏的社会中&#xff0c;越来越多的人因工作忙碌、时间紧张而无法亲自下厨&#xff0c;上门做饭服务恰好满足了这部分人群的需求&#xff0c;为他们提供了便捷、高效的餐饮解决方案。个性化需求增加&#xff1a;随着人们生活水平的…

【配置后的基本使用】CMake基础知识

&#x1f308; 个人主页&#xff1a;十二月的猫-CSDN博客 &#x1f525; 系列专栏&#xff1a; &#x1f3c0;各种软件安装与配置_十二月的猫的博客-CSDN博客 &#x1f4aa;&#x1f3fb; 十二月的寒冬阻挡不了春天的脚步&#xff0c;十二点的黑夜遮蔽不住黎明的曙光 目录 1.…

Centos 7 安装wget

Centos 7 安装wget 最小化安装Centos 7 的话需要上传wget rpm包之后再路径下安装一下。rpm包下载地址&#xff08;http://mirrors.163.com/centos/7/os/x86_64/Packages/&#xff09; 1、使用X-ftp 或者WinSCP等可以连接上传的软件都可以首先连接服务器&#xff0c;这里我用的…

Linux最深刻理解页表于物理内存

目录 物理内存管理 页表设计 物理内存管理 如果磁盘上的内容加载到物理内存上&#xff0c;每次io都会按照4kb的方式进行加载(可能不同版本系统有些区别)。所以我们的物理内存上的内容也是4个字节进行管理的。 而每个页框都需要我们进行管理。所以自然物理内存就会对页框进行先…

几何合理的分片段感知的3D分子生成 FragGen - 评测

FragGen 来源于 2024 年 3 月 25 日 预印本的文章&#xff0c;文章题目是 Deep Geometry Handling and Fragment-wise Molecular 3D Graph Generation&#xff0c; 作者是 Odin Zhang&#xff0c;侯廷军&#xff0c;浙江大学药学院。FragGen 是一个基于分子片段的 3D 分子生成模…

PySpark——Python与大数据

一、Spark 与 PySpark Apache Spark 是用于大规模数据&#xff08; large-scala data &#xff09;处理的统一&#xff08; unified &#xff09;分析引擎。简单来说&#xff0c; Spark 是一款分布式的计算框架&#xff0c;用于调度成百上千的服务器集群&#xff0c;计算 TB 、…

基于Java Springboot编程语言在线学习平台

一、作品包含 源码数据库设计文档万字PPT全套环境和工具资源部署教程 二、项目技术 前端技术&#xff1a;Html、Css、Js、Vue、Element-ui 数据库&#xff1a;MySQL 后端技术&#xff1a;Java、Spring Boot、MyBatis 三、运行环境 开发工具&#xff1a;IDEA/eclipse 数据…

WebRTC视频 02 - 视频采集类 VideoCaptureModule

WebRTC视频 01 - 视频采集整体架构 WebRTC视频 02 - 视频采集类 VideoCaptureModule&#xff08;本文&#xff09; WebRTC视频 03 - 视频采集类 VideoCaptureDS 上篇 WebRTC视频 04 - 视频采集类 VideoCaptureDS 中篇 WebRTC视频 05 - 视频采集类 VideoCaptureDS 下篇 一、前言…

深度学习笔记14-卷积神经网络2

1.卷积神经网络的结构 卷积神经网络&#xff0c;是包含卷积运算且具有深度结构的前馈神经网络。在卷积神经网络中&#xff0c;包含卷积层、池化层和全连接层三种重要的结构。相比前馈神经网络&#xff0c;卷积层和池化层是新增的网络结构&#xff0c;在提取特征时&#xff0c;卷…

Python 正则表达式使用指南

Python 正则表达式使用指南 正则表达式&#xff08;Regular Expression, 简称 regex&#xff09;是处理字符串和文本的强大工具。它使用特定的语法定义一组规则&#xff0c;通过这些规则可以对文本进行匹配、查找、替换等操作。Python 提供了 re 模块&#xff0c;使得正则表达…

FPGA开发-逻辑分析仪的应用-数字频率计的设计

目录 逻辑分析仪的应用 数字频率计的设计 -基于原理图方法 主控电路设计 分频器设计 顶层电路设计 数字系统开发不但需要进行仿真分析&#xff0c;更重要的是需要进行实际测试。 逻辑分析仪的应用 测试方式&#xff1a;&#xff08;1&#xff09;传统的测试方式&#…

.NET 9.0 中 System.Text.Json 的全面使用指南

以下是一些 System.Text.Json 在 .NET 9.0 中的使用方式&#xff0c;包括序列化、反序列化、配置选项等&#xff0c;并附上输出结果。 基本序列化和反序列化 using System; using System.Text.Json; public class Program {public class Person{public string Name { get; se…

Linux 命令 | 每日一学,文本处理三剑客之awk命令实践

[ 知识是人生的灯塔&#xff0c;只有不断学习&#xff0c;才能照亮前行的道路 ] 0x00 前言简述 描述&#xff1a;前面作者已经介绍了文本处理三剑客中的 grep 与 sed 文本处理工具&#xff0c;今天将介绍其最后一个且非常强大的 awk 文本处理输出工具&#xff0c;它可以非常方便…

【第五课】Rust所有权系统(一)

目录 前言 所有权机制的核心 再谈变量绑定 主人变更-所有权转移 总结 前言 这节课我们来介绍下rust中最重要的一个点&#xff1a;所有权系统。这是网上经常说rust无gc的秘密所在。在开始之前&#xff0c;我们来想想JVM系语言&#xff0c;在做垃圾回收的过程&#xff0c;1.…

三周精通FastAPI:42 手动运行服务器 - Uvicorn Gunicorn with Uvicorn

官方文档&#xff1a;Server Workers - Gunicorn with Uvicorn - FastAPI 使用 fastapi 运行命令 可以直接使用fastapi run命令来启动FastAPI应用&#xff1a; fastapi run main.py如创建openapi.py文件&#xff1a; from fastapi import FastAPIapp FastAPI(openapi_url&…

任意文件下载漏洞

1.漏洞简介 任意文件下载漏洞是指攻击者能够通过操控请求参数&#xff0c;下载服务器上未经授权的文件。 攻击者可以利用该漏洞访问敏感文件&#xff0c;如配置文件、日志文件等&#xff0c;甚至可以下载包含恶意代码的文件。 这里再导入一个基础&#xff1a; 你要在网站下…