Scala语法(六) Akka与线程通信

2024-05-14 07:58
文章标签 线程 语法 通信 scala akka

本文主要是介绍Scala语法(六) Akka与线程通信,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

前言

在初期, Scala可以通过Akka来实现线程通信. 当然, 现在还支持使用Netty方式进行通信.

本章主要介绍使用Akka方式进行通信的写法.


正文

  • Master结点

import akka.actor.Actor
import akka.actor.ActorSystem
import com.typesafe.config.ConfigFactory
import akka.actor.Propsclass AkkaMaster extends Actor{// start 之前override def preStart() : Unit = {println("pre master invoke.")}// 用于接收消息override def receive:Receive = {case "connect" => {println("a client connected.")sender ! "reply"}case "hello" => {println("hello")}}
}
object AkkaMaster{def main(args: Array[String]): Unit = {// 使用创建ActorSystem来创建和监控下面的Actor对象. 单例的.val host = "127.0.0.1"val port = 8090// 准备配置val configStr = s"""|akka.actor.provider = "akka.remote.RemoteActorRefProvider"|akka.remote.netty.tcp.hostname = "$host"|akka.remote.netty.tcp.port = "$port"""".stripMarginval config = ConfigFactory.parseString(configStr)// 注意名称中间不要加空格val actorSysetm = ActorSystem("MasterSystem",config) // 创建Actorval master = actorSysetm.actorOf(Props(new AkkaMaster),"Master")master ! "hello"// 等待信号停止actorSysetm.awaitTermination()}
}// 顺利输出
//[INFO] [04/29/2019 16:43:20.512] [main] [Remoting] Starting remoting
//[INFO] [04/29/2019 16:43:20.770] [main] [Remoting] Remoting started; listening on addresses :[akka.tcp://MasterSystem@127.0.0.1:8090]
//[INFO] [04/29/2019 16:43:20.771] [main] [Remoting] Remoting now listens on addresses: [akka.tcp://MasterSystem@127.0.0.1:8090]
//pre master invoke.
//hello// 1. 名称中间不要加空格
//Exception in thread "main" java.lang.IllegalArgumentException: invalid ActorSystem name [Master System], must contain only word characters (i.e. [a-zA-Z0-9] plus non-leading '-' or '_')
//	at akka.actor.ActorSystemImpl.<init>(ActorSystem.scala:498)
//	at akka.actor.ActorSystem$.apply(ActorSystem.scala:142)
//	at akka.actor.ActorSystem$.apply(ActorSystem.scala:119)
//	at com.yanxml.quick_scala.multi.akka.AkkaMaster$.main(AkkaMaster.scala:33)
//	at com.yanxml.quick_scala.multi.akka.AkkaMaster.main(AkkaMaster.scala)
  • Worker结点
import akka.actor.Actor
import akka.actor.ActorSelection
import akka.actor.ActorSystem
import akka.actor.Props
import com.typesafe.config.ConfigFactoryclass AkkaWorker extends Actor{// 成员变量val master:ActorSelection = null// 建立链接override def preStart():Unit = {val master = context.actorSelection("akka.tcp://MasterSystem@127.0.0.1:8090/user/Master")println(master)master ! "connect"}override def receive:Receive = {case "reply" => {println("a reply from master")}}
}object AkkaWorker{def main(args: Array[String]): Unit = {// 使用创建ActorSystem来创建和监控下面的Actor对象. 单例的.val host = "127.0.0.1"val port = 8091// 准备配置val configStr = s"""|akka.actor.provider = "akka.remote.RemoteActorRefProvider"|akka.remote.netty.tcp.hostname = "$host"|akka.remote.netty.tcp.port = "$port"""".stripMarginval config = ConfigFactory.parseString(configStr)// 注意名称中间不要加空格val actorSysetm = ActorSystem("WorkerSystem",config) // 创建Actorval master = actorSysetm.actorOf(Props(new AkkaWorker),"Worker")master ! "hello"// 等待信号停止actorSysetm.awaitTermination()}
}

模拟RPC

在模拟RPC中主要有这样的流程.

其中主要包括两个结点: Worker结点&Master结点.

  • 运行流程:
    • Master结点先进行启动.
    • Worker结点后进行启动.
    • Worker结点Master结点发送注册消息.
    • Matser结点接收注册消息, 并进行记录. 并将主结点的地址返回给Worker结点(模拟Master是集群的情况).并记录,最后的通信时间作为心跳标志.
    • Worker结点接收主结点地址, 并形成通信链接. 开始通信. 并定时发送心跳消息.

改造上方的Demo代码. 其基本代码如下所示:

  • RemoteMessage

trait RemoteMessage  extends Serializable{}// Worker -> Master 用来封装Worker信息 
case class RegisterWorker(id:String,memory:Int,cores:Int)class WorkerInfo(val id:String, val memory:Int, val cores:Int){// 上一次心跳var heartbeatTime:String = _
}
  • Master

import akka.actor.Actor
import akka.actor.ActorSystem
import com.typesafe.config.ConfigFactory
import akka.actor.Props
import scala.collection.immutable.HashMapprivate [simulate] class AkkaMaster extends Actor{val idToWorker = new scala.collection.mutable.HashMap[String,WorkerInfo]()// start 之前override def preStart() : Unit = {println("pre master invoke.")}// 用于接收消息override def receive:Receive = {case "connect" => {println("a client connected.")sender ! "reply"}case "hello" => {println("hello")}// 传输样例类case RegisterWorker(id,memory,cores)=>{if(!idToWorker.contains(id)){idToWorker.put(id, new WorkerInfo(id,memory,cores))}sender ! "123"}}
}
object AkkaMaster{def main(args: Array[String]): Unit = {// 使用创建ActorSystem来创建和监控下面的Actor对象. 单例的.val host = "127.0.0.1"val port = 8090// 准备配置val configStr = s"""|akka.actor.provider = "akka.remote.RemoteActorRefProvider"|akka.remote.netty.tcp.hostname = "$host"|akka.remote.netty.tcp.port = "$port"""".stripMarginval config = ConfigFactory.parseString(configStr)// 注意名称中间不要加空格val actorSysetm = ActorSystem("MasterSystem",config) // 创建Actorval master = actorSysetm.actorOf(Props(new AkkaMaster),"Master")master ! "hello"// 等待信号停止actorSysetm.awaitTermination()}
}// 顺利输出
//[INFO] [04/29/2019 16:43:20.512] [main] [Remoting] Starting remoting
//[INFO] [04/29/2019 16:43:20.770] [main] [Remoting] Remoting started; listening on addresses :[akka.tcp://MasterSystem@127.0.0.1:8090]
//[INFO] [04/29/2019 16:43:20.771] [main] [Remoting] Remoting now listens on addresses: [akka.tcp://MasterSystem@127.0.0.1:8090]
//pre master invoke.
//hello// 1. 名称中间不要加空格
//Exception in thread "main" java.lang.IllegalArgumentException: invalid ActorSystem name [Master System], must contain only word characters (i.e. [a-zA-Z0-9] plus non-leading '-' or '_')
//	at akka.actor.ActorSystemImpl.<init>(ActorSystem.scala:498)
//	at akka.actor.ActorSystem$.apply(ActorSystem.scala:142)
//	at akka.actor.ActorSystem$.apply(ActorSystem.scala:119)
//	at com.yanxml.quick_scala.multi.akka.AkkaMaster$.main(AkkaMaster.scala:33)
//	at com.yanxml.quick_scala.multi.akka.AkkaMaster.main(AkkaMaster.scala)
  • Worker

import akka.actor.Actor
import akka.actor.ActorSelection
import akka.actor.ActorSystem
import akka.actor.Props
import com.typesafe.config.ConfigFactory
import com.yanxml.quick_scala.multi.akka.simulate.RegisterWorker
import java.util.UUIDprivate [simulate] class  AkkaWorker(val masterHost:String, val masterPort:String, val memory:Int, val cores:Int) extends Actor{// 成员变量val master:ActorSelection = null// 建立链接override def preStart():Unit = {// 和Master建立链接val master = context.actorSelection(s"akka.tcp://MasterSystem@$masterHost:$masterPort/user/Master")val workerId = UUID.randomUUID().toString()//println(master)// 向Master发送消息master ! RegisterWorker(workerId,memory,cores)}override def receive:Receive = {case "reply" => {println("a reply from master")}}
}object AkkaWorker{def main(args: Array[String]): Unit = {// 使用创建ActorSystem来创建和监控下面的Actor对象. 单例的.val host = "127.0.0.1"val port = 8091// 准备配置val configStr = s"""|akka.actor.provider = "akka.remote.RemoteActorRefProvider"|akka.remote.netty.tcp.hostname = "$host"|akka.remote.netty.tcp.port = "$port"""".stripMarginval config = ConfigFactory.parseString(configStr)// 注意名称中间不要加空格val actorSysetm = ActorSystem("WorkerSystem",config) // 创建Actorval master = actorSysetm.actorOf(Props(new AkkaWorker("127.0.0.1","8090",2,2)),"Worker")master ! "hello"// 等待信号停止actorSysetm.awaitTermination()}
}

注: 后续的通信逻辑就是丰富双方的receive()方法.

这篇关于Scala语法(六) Akka与线程通信的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/988198

相关文章

Java中实现线程的创建和启动的方法

《Java中实现线程的创建和启动的方法》在Java中,实现线程的创建和启动是两个不同但紧密相关的概念,理解为什么要启动线程(调用start()方法)而非直接调用run()方法,是掌握多线程编程的关键,... 目录1. 线程的生命周期2. start() vs run() 的本质区别3. 为什么必须通过 st

Linux实现线程同步的多种方式汇总

《Linux实现线程同步的多种方式汇总》本文详细介绍了Linux下线程同步的多种方法,包括互斥锁、自旋锁、信号量以及它们的使用示例,通过这些同步机制,可以解决线程安全问题,防止资源竞争导致的错误,示例... 目录什么是线程同步?一、互斥锁(单人洗手间规则)适用场景:特点:二、条件变量(咖啡厅取餐系统)工作流

Java中常见队列举例详解(非线程安全)

《Java中常见队列举例详解(非线程安全)》队列用于模拟队列这种数据结构,队列通常是指先进先出的容器,:本文主要介绍Java中常见队列(非线程安全)的相关资料,文中通过代码介绍的非常详细,需要的朋... 目录一.队列定义 二.常见接口 三.常见实现类3.1 ArrayDeque3.1.1 实现原理3.1.2

SpringBoot3中使用虚拟线程的完整步骤

《SpringBoot3中使用虚拟线程的完整步骤》在SpringBoot3中使用Java21+的虚拟线程(VirtualThreads)可以显著提升I/O密集型应用的并发能力,这篇文章为大家介绍了详细... 目录1. 环境准备2. 配置虚拟线程方式一:全局启用虚拟线程(Tomcat/Jetty)方式二:异步

如何解决Druid线程池Cause:java.sql.SQLRecoverableException:IO错误:Socket read timed out的问题

《如何解决Druid线程池Cause:java.sql.SQLRecoverableException:IO错误:Socketreadtimedout的问题》:本文主要介绍解决Druid线程... 目录异常信息触发场景找到版本发布更新的说明从版本更新信息可以看到该默认逻辑已经去除总结异常信息触发场景复

RabbitMQ工作模式中的RPC通信模式详解

《RabbitMQ工作模式中的RPC通信模式详解》在RabbitMQ中,RPC模式通过消息队列实现远程调用功能,这篇文章给大家介绍RabbitMQ工作模式之RPC通信模式,感兴趣的朋友一起看看吧... 目录RPC通信模式概述工作流程代码案例引入依赖常量类编写客户端代码编写服务端代码RPC通信模式概述在R

在Spring Boot中实现HTTPS加密通信及常见问题排查

《在SpringBoot中实现HTTPS加密通信及常见问题排查》HTTPS是HTTP的安全版本,通过SSL/TLS协议为通讯提供加密、身份验证和数据完整性保护,下面通过本文给大家介绍在SpringB... 目录一、HTTPS核心原理1.加密流程概述2.加密技术组合二、证书体系详解1、证书类型对比2. 证书获

Python模拟串口通信的示例详解

《Python模拟串口通信的示例详解》pySerial是Python中用于操作串口的第三方模块,它支持Windows、Linux、OSX、BSD等多个平台,下面我们就来看看Python如何使用pySe... 目录1.win 下载虚www.chinasem.cn拟串口2、确定串口号3、配置串口4、串口通信示例5

基于C#实现MQTT通信实战

《基于C#实现MQTT通信实战》MQTT消息队列遥测传输,在物联网领域应用的很广泛,它是基于Publish/Subscribe模式,具有简单易用,支持QoS,传输效率高的特点,下面我们就来看看C#实现... 目录1、连接主机2、订阅消息3、发布消息MQTT(Message Queueing Telemetr

mysql递归查询语法WITH RECURSIVE的使用

《mysql递归查询语法WITHRECURSIVE的使用》本文主要介绍了mysql递归查询语法WITHRECURSIVE的使用,WITHRECURSIVE用于执行递归查询,特别适合处理层级结构或递归... 目录基本语法结构:关键部分解析:递归查询的工作流程:示例:员工与经理的层级关系解释:示例:树形结构的数