项目概述
需求
目前大多数的分布式架构底层通信都是通过RPC实现的,RPC框架非常多,比如前我们学过的Hadoop项目的RPC通信框架,但是Hadoop在设计之初就是为了运行长达数小时的批量而设计的,在某些极端的情况下,任务提交的延迟很高,所以Hadoop的RPC显得有些笨重。
Spark 的RPC是通过Akka类库实现的,Akka用Scala语言开发,基于Actor并发模型实现,Akka具有高可靠、高性能、可扩展等特点,使用Akka可以轻松实现分布式RPC功能。
Akka简介
友情链接: Actors介绍: https://www.iteblog.com/archives/1154.html
Akka基于Actor模型,提供了一个用于构建可扩展的(Scalable)、弹性的(Resilient)、快速响应的(Responsive)应用程序的平台。
Actor模型:在计算机科学领域,Actor模型是一个并行计算(Concurrent Computation)模型,它把actor作为并行计算的基本元素来对待:为响应一个接收到的消息,一个actor能够自己做出一些决策,如创建更多的actor,或发送更多的消息,或者确定如何去响应接收到的下一个消息。

Actor是Akka中最核心的概念,它是一个封装了状态和行为的对象,Actor之间可以通过交换消息的方式进行通信,每个Actor都有自己的收件箱(Mailbox)。通过Actor能够简化锁及线程管理,可以非常容易地开发出正确地并发程序和并行系统,Actor具有如下特性:
(1)、提供了一种高级抽象,能够简化在并发(Concurrency)/并行(Parallelism)应用场景下的编程开发
(2)、提供了异步非阻塞的、高性能的事件驱动编程模型
(3)、超级轻量级事件处理(每GB堆内存几百万Actor)
项目实现
实战一:
利用Akka的actor编程模型,实现2个进程间的通信。
架构图

重要类介绍
ActorSystem**:**在Akka中,ActorSystem是一个重量级的结构,他需要分配多个线程,所以在实际应用中,ActorSystem通常是一个单例对象,我们可以使用这个ActorSystem创建很多Actor。
注意**:**
(1)、ActorSystem是一个进程中的老大,它负责创建和监督actor
(2)、ActorSystem是一个单例对象
(3)、actor负责通信
Actor
在Akka中,Actor负责通信,在Actor中有一些重要的生命周期方法。
(1)preStart()方法:该方法在Actor对象构造方法执行后执行,整个Actor生命周期中仅执行一次。
(2)receive()方法:该方法在Actor的preStart方法执行完成后执行,用于接收消息,会被反复执行。
具体代码
① Master****类
1package cn.itcast.rpc 2 3 import akka.actor.{Actor, ActorRef, ActorSystem, Props} 4 import com.typesafe.config.ConfigFactory 5 6 //todo:利用akka的actor模型实现2个进程间的通信-----Master端 7 class Master extends Actor{ 8 //构造代码块先被执行 9 println("master constructor invoked") 10 11 //prestart方法会在构造代码块执行后被调用,并且只被调用一次 12 override def preStart(): Unit = { 13 println("preStart method invoked") 14 } 15 16 //receive方法会在prestart方法执行后被调用,表示不断的接受消息 17 override def receive: Receive = { 18 case "connect" =>{ 19 println("a client connected") 20 21 //master发送注册成功信息给worker 22 sender ! "success" 23 } 24 } 25} 26 27 object Master{ 28 def main(args: Array[String]): Unit = { 29 //master的ip地址 30 val host=args(0) 31 //master的port端口 32 val port=args(1) 33 34 //准备配置文件信息 35 val configStr= 36 s""" 37 |akka.actor.provider = "akka.remote.RemoteActorRefProvider" 38 |akka.remote.netty.tcp.hostname = "$host" 39 |akka.remote.netty.tcp.port = "$port" 40 """.stripMargin 41 42 //配置config对象 利用ConfigFactory解析配置文件,获取配置信息 43 val config=ConfigFactory.parseString(configStr) 44 45 // 1、创建ActorSystem,它是整个进程中老大,它负责创建和监督actor,它是单例对象 46 val masterActorSystem = ActorSystem("masterActorSystem",config) 47 // 2、通过ActorSystem来创建master actor 48 val masterActor: ActorRef = masterActorSystem.actorOf(Props(new Master),"masterActor") 49 // 3、向master actor发送消息 50 //masterActor ! "connect" 51 } 52}
② Worker****类
1package cn.itcast.rpc 2 3 import akka.actor.{Actor, ActorRef, ActorSelection, ActorSystem, Props} 4 import com.typesafe.config.ConfigFactory 5 6 //todo:利用akka中的actor实现2个进程间的通信-----Worker端 7 class Worker extends Actor{ 8 println("Worker constructor invoked") 9 10 //prestart方法会在构造代码块之后被调用,并且只会被调用一次 11 override def preStart(): Unit = { 12 println("preStart method invoked") 13 14 //获取master actor的引用 15 //ActorContext全局变量,可以通过在已经存在的actor中,寻找目标actor 16 //调用对应actorSelection方法, 17 // 方法需要一个path路径:1、通信协议、2、master的IP地址、3、master的端口 4、创建master actor老大 5、actor层级 18 val master: ActorSelection = context.actorSelection("akka.tcp://masterActorSystem@172.16.43.63:8888/user/masterActor") 19 20 //向master发送消息 21 master ! "connect" 22 } 23 24 //receive方法会在prestart方法执行后被调用,不断的接受消息 25 override def receive: Receive = { 26 case "connect" =>{ 27 println("a client connected") 28 } 29 30 case "success" =>{ 31 println("注册成功") 32 } 33 } 34} 35 36 object Worker{ 37 def main(args: Array[String]): Unit = { 38 //定义worker的IP地址 39 val host=args(0) 40 //定义worker的端口 41 val port=args(1) 42 43 //准备配置文件 44 val configStr= 45 s""" 46 |akka.actor.provider = "akka.remote.RemoteActorRefProvider" 47 |akka.remote.netty.tcp.hostname = "$host" 48 |akka.remote.netty.tcp.port = "$port" 49 """.stripMargin 50 51 //通过configFactory来解析配置信息 52 val config=ConfigFactory.parseString(configStr) 53 54 // 1、创建ActorSystem,它是整个进程中的老大,它负责创建和监督actor 55 val workerActorSystem = ActorSystem("workerActorSystem",config) 56 // 2、通过actorSystem来创建 worker actor 57 val workerActor: ActorRef = workerActorSystem.actorOf(Props(new Worker),"workerActor") 58 59 //向worker actor发送消息 60 workerActor ! "connect" 61 } 62}
③ 运行
使用idea开发工具,配置参数时,多个参数之间用空格隔开


启动Master

启动Worker

实战二
使用Akka实现一个简易版的spark通信框架
架构图

具体代码
① Master****类
1package cn.itcast.spark 2 3 import akka.actor.{Actor, ActorRef, ActorSystem, Props} 4 import com.typesafe.config.ConfigFactory 5 import scala.collection.mutable 6 import scala.collection.mutable.ListBuffer 7 import scala.concurrent.duration._ 8 9 //todo:利用akka实现简易版的spark通信框架-----Master端 10 class Master extends Actor{ 11 //构造代码块先被执行 12 println("master constructor invoked") 13 14 //定义一个map集合,用于存放worker信息 15 private val workerMap = new mutable.HashMap[String,WorkerInfo]() 16 17 //定义一个list集合,用于存放WorkerInfo信息,方便后期按照worker上的资源进行排序 18 private val workerList = new ListBuffer[WorkerInfo] 19 20 //master定时检查的时间间隔 21 val CHECK_OUT_TIME_INTERVAL=15000 //15秒 22 23 //prestart方法会在构造代码块执行后被调用,并且只被调用一次 24 override def preStart(): Unit = { 25 println("preStart method invoked") 26 27 //master定时检查超时的worker 28 //需要手动导入隐式转换 29 import context.dispatcher 30 context.system.scheduler.schedule(0 millis,CHECK_OUT_TIME_INTERVAL millis,self,CheckOutTime) 31 } 32 33 //receive方法会在prestart方法执行后被调用,表示不断的接受消息 34 override def receive: Receive = { 35 //master接受worker的注册信息 36 case RegisterMessage(workerId,memory,cores) =>{ 37 //判断当前worker是否已经注册 38 if(!workerMap.contains(workerId)){ 39 //保存信息到map集合中 40 val workerInfo = new WorkerInfo(workerId,memory,cores) 41 workerMap.put(workerId,workerInfo) 42 43 //保存workerinfo到list集合中 44 workerList +=workerInfo 45 46 //master反馈注册成功给worker 47 sender ! RegisteredMessage(s"workerId:$workerId 注册成功") 48 } 49 } 50 51 //master接受worker的心跳信息 52 case SendHeartBeat(workerId)=>{ 53 //判断worker是否已经注册,master只接受已经注册过的worker的心跳信息 54 if(workerMap.contains(workerId)){ 55 //获取workerinfo信息 56 val workerInfo: WorkerInfo = workerMap(workerId) 57 58 //获取当前系统时间 59 val lastTime: Long = System.currentTimeMillis() 60 61 workerInfo.lastHeartBeatTime=lastTime 62 } 63 } 64 65 case CheckOutTime=>{ 66 //过滤出超时的worker 判断逻辑: 获取当前系统时间 - worker上一次心跳时间 >master定时检查的时间间隔 67 val outTimeWorkers: ListBuffer[WorkerInfo] = workerList.filter(x => System.currentTimeMillis() -x.lastHeartBeatTime > CHECK_OUT_TIME_INTERVAL) 68 //遍历超时的worker信息,然后移除掉超时的worker 69 for(workerInfo <- outTimeWorkers){ 70 //获取workerid 71 val workerId: String = workerInfo.workerId 72 //从map集合中移除掉超时的worker信息 73 workerMap.remove(workerId) 74 //从list集合中移除掉超时的workerInfo信息 75 workerList -= workerInfo 76 println("超时的workerId:" +workerId) 77 } 78 println("活着的worker总数:" + workerList.size) 79 80 //master按照worker内存大小进行降序排列 81 println(workerList.sortBy(x => x.memory).reverse.toList) 82 } 83 } 84} 85 86 object Master{ 87 def main(args: Array[String]): Unit = { 88 //master的ip地址 89 val host=args(0) 90 //master的port端口 91 val port=args(1) 92 93 //准备配置文件信息 94 val configStr= 95 s""" 96 |akka.actor.provider = "akka.remote.RemoteActorRefProvider" 97 |akka.remote.netty.tcp.hostname = "$host" 98 |akka.remote.netty.tcp.port = "$port" 99 """.stripMargin 100 101 //配置config对象 利用ConfigFactory解析配置文件,获取配置信息 102 val config=ConfigFactory.parseString(configStr) 103 104 // 1、创建ActorSystem,它是整个进程中老大,它负责创建和监督actor,它是单例对象 105 val masterActorSystem = ActorSystem("masterActorSystem",config) 106 // 2、通过ActorSystem来创建master actor 107 val masterActor: ActorRef = masterActorSystem.actorOf(Props(new Master),"masterActor") 108 // 3、向master actor发送消息 109 //masterActor ! "connect" 110 } 111}
② Worker****类
1package cn.itcast.spark 2 3 import java.util.UUID 4 import akka.actor.{Actor, ActorRef, ActorSelection, ActorSystem, Props} 5 import com.typesafe.config.ConfigFactory 6 import scala.concurrent.duration._ 7 8 //todo:利用akka实现简易版的spark通信框架-----Worker端 9 class Worker(val memory:Int,val cores:Int,val masterHost:String,val masterPort:String) extends Actor{ 10 println("Worker constructor invoked") 11 12 //定义workerId 13 private val workerId: String = UUID.randomUUID().toString 14 15 //定义发送心跳的时间间隔 16 val SEND_HEART_HEAT_INTERVAL=10000 //10秒 17 18 //定义全局变量 19 var master: ActorSelection=_ 20 21 //prestart方法会在构造代码块之后被调用,并且只会被调用一次 22 override def preStart(): Unit = { 23 println("preStart method invoked") 24 //获取master actor的引用 25 //ActorContext全局变量,可以通过在已经存在的actor中,寻找目标actor 26 //调用对应actorSelection方法, 27 // 方法需要一个path路径:1、通信协议、2、master的IP地址、3、master的端口 4、创建master actor老大 5、actor层级 28 master= context.actorSelection(s"akka.tcp://masterActorSystem@$masterHost:$masterPort/user/masterActor") 29 30 //向master发送注册信息,将信息封装在样例类中,主要包含:workerId,memory,cores 31 master ! RegisterMessage(workerId,memory,cores) 32 } 33 34 //receive方法会在prestart方法执行后被调用,不断的接受消息 35 override def receive: Receive = { 36 //worker接受master的反馈信息 37 case RegisteredMessage(message) =>{ 38 println(message) 39 40 //向master定期的发送心跳 41 //worker先自己给自己发送心跳 42 //需要手动导入隐式转换 43 import context.dispatcher 44 context.system.scheduler.schedule(0 millis,SEND_HEART_HEAT_INTERVAL millis,self,HeartBeat) 45 } 46 47 //worker接受心跳 48 case HeartBeat =>{ 49 //这个时候才是真正向master发送心跳 50 master ! SendHeartBeat(workerId) 51 } 52 } 53} 54 55 object Worker{ 56 def main(args: Array[String]): Unit = { 57 //定义worker的IP地址 58 val host=args(0) 59 //定义worker的端口 60 val port=args(1) 61 //定义worker的内存 62 val memory=args(2).toInt 63 //定义worker的核数 64 val cores=args(3).toInt 65 //定义master的ip地址 66 val masterHost=args(4) 67 //定义master的端口 68 val masterPort=args(5) 69 70 //准备配置文件 71 val configStr= 72 s""" 73 |akka.actor.provider = "akka.remote.RemoteActorRefProvider" 74 |akka.remote.netty.tcp.hostname = "$host" 75 |akka.remote.netty.tcp.port = "$port" 76 """.stripMargin 77 78 //通过configFactory来解析配置信息 79 val config=ConfigFactory.parseString(configStr) 80 // 1、创建ActorSystem,它是整个进程中的老大,它负责创建和监督actor 81 val workerActorSystem = ActorSystem("workerActorSystem",config) 82 // 2、通过actorSystem来创建 worker actor 83 val workerActor: ActorRef = workerActorSystem.actorOf(Props(new Worker(memory,cores,masterHost,masterPort)),"workerActor") 84 85 //向worker actor发送消息 86 workerActor ! "connect" 87 } 88}
③ WorkerInfo****类
1package cn.itcast.spark 2 3 //封装worker信息 4 class WorkerInfo(val workerId:String,val memory:Int,val cores:Int) { 5 //定义一个变量用于存放worker上一次心跳时间 6 var lastHeartBeatTime:Long=_ 7 8 override def toString: String = { 9 s"workerId:$workerId , memory:$memory , cores:$cores" 10 } 11}
④ 样例类
1package cn.itcast.spark 2 3 trait RemoteMessage extends Serializable{} 4 5 //worker向master发送注册信息,由于不在同一进程中,需要实现序列化 6 case class RegisterMessage(val workerId:String,val memory:Int,val cores:Int) extends RemoteMessage 7 8 //master反馈注册成功信息给worker,由于不在同一进程中,也需要实现序列化 9 case class RegisteredMessage(message:String) extends RemoteMessage 10 11 //worker向worker发送心跳 由于在同一进程中,不需要实现序列化 12 case object HeartBeat 13 14 //worker向master发送心跳,由于不在同一进程中,需要实现序列化 15 case class SendHeartBeat(val workerId:String) extends RemoteMessage 16 17 //master自己向自己发送消息,由于在同一进程中,不需要实现序列化 18 case object CheckOutTime
⑤ 运行
配置参数时,多个参数之间用空格隔开


首先启动Master_Spark
启动work_spark-01

启动work_spark-02,然后关闭

观察Master_Spark 输出
