一、Akka介绍
1.1 什么是Akka
1.2 Akka核心架构

Akka的核心组件关系:
- Actor:核心并发单元
- ActorSystem:Actor的工厂和管理者
- ActorRef:Actor的代理/引用
- MailBox:消息队列(FIFO)
- Dispatcher Message:消息分发器(线程池)
- Supervisors:容错监督机制
二、Akka中Actor模型

2.1 Actor模型用于解决什么问题
经典案例:邮件系统

2.2 Actor模型及其说明
三、Actor模型工作机制说明
3.1 工作机制示意图

Actor间传递消息机制:
- 消息接收和处理
- 通过 sender() 方法可以得到发送消息的Actor的 ActorRef,通过这个 ActorRef,B Actor 也可以回复消息
3.2 Actor模型工作机制详解
3.3 Actor间传递消息机制
四、Actor模型快速入门
4.1 应用实例需求
运行效果:

4.2 Actor自我通讯机制原理图

流程说明:
- A Actor 通过 A ActorRef 发送消息
- Dispatcher Message 转发消息到 MailBox
- MailBox(实现了Runnable)是一个线程,一直运行并调用 Actor 的 receive 方法
- Actor 在 receive 方法中处理消息
4.3 代码实现
// SayHelloActorDemo.scala
class SayHelloActorDemo extends Actor {
override def receive: Receive = {
// 接受消息并处理,如果接收到exit,就退出
case "hello" => println("发送:hello\\t\\t回应: hello too:)")
case "ok" => println("发送:ok\\t\\t\\t回应: ok too:)")
case "exit" => {
println("接收到exit~指令,退出系统…. ")
context.stop(self) // 停止自己的actorRef
context.system.terminate() // 关闭ActorSystem
}
}
}
// SayHelloActor.scala (入口)
object SayHelloActor {
private val actoryFactory = ActorSystem("ActoryFactory")
private val sayHelloActorRef: ActorRef =
actoryFactory.actorOf(Props[SayHelloActorDemo], "sayHelloActor")
def main(args: Array[String]): Unit = {
// 给sayHelloActorRef发消息
sayHelloActorRef ! "hello"
sayHelloActorRef ! "ok"
sayHelloActorRef ! "exit"
}
}
4.4 小结和说明
当程序执行 aActorRef = actorFactory.actorOf(Props[AActor], "aActor"),会完成如下任务:
五、Actor模型应用实例 – Actor间通讯
5.1 应用实例需求
运行效果:

5.2 两个Actor的通讯机制原理图

小结和说明:
5.3 代码实现
// AActor.scala
class AActor(bActorRef: ActorRef) extends Actor {
var attack = 0
override def receive: Receive = {
case "start" => {
println("AActor(黄飞鸿) 开始游戏了")
bActorRef ! "我打"
}
case "我打" => {
attack += 1
println(s"AActor(黄飞鸿) 厉害 看我佛山无影脚~~~ 第${attack} 脚")
Thread.sleep(1000)
bActorRef ! "我打"
}
}
}
// BActor.scala
class BActor extends Actor {
var attack = 0
override def receive: Receive = {
case "我打" => {
attack += 1
println(s"BActor(乔峰): 挺猛的 看我降龙十八掌~ 第${attack} 掌")
Thread.sleep(1000)
sender() ! "我打" // 回复给发送者
}
}
}
// ActorsGame.scala (入口)
object ActorsGame {
def main(args: Array[String]): Unit = {
val actorFactory = ActorSystem("actorFactory")
val bActorRef = actorFactory.actorOf(Props[BActor], "bActor")
val aActorRef = actorFactory.actorOf(
Props(new AActor(bActorRef)), "aActor")
aActorRef ! "start"
}
}
5.4 如何理解 Actor 的 receive 方法被调用?
六、Akka网络编程
6.1 网络编程基本介绍
Akka 支持面向大并发后端服务程序,网络通信这块是服务端程序重要的一部分。
网络编程有两种:
6.2 网络编程基础知识
网线、网卡、无线网卡
计算机间要相互通讯,必须要求网线、网卡,或者是无线网卡。
协议(TCP/IP)
TCP/IP(Transmission Control Protocol/Internet Protocol)的简写,中文译名为传输控制协议/因特网互联协议,又叫网络通讯协议,这个协议是 Internet 最基本的协议、Internet 国际互联网络的基础,简单地说,就是由网络层的 IP 协议和传输层的 TCP 协议组成的。
OSI与TCP/IP参考模型
| 应用层(application) | 应用层(application):smtp, ftp, telnet, http |
| 表示层(presention) | |
| 会话层(session) | |
| 传输层(transport) | 传输层(transport):解释数据 |
| 网络层(ip) | 网络层(ip):定位ip地址和确定连接路径 |
| 数据链路层(link) | 链路层(link):与硬件驱动对话 |
| 物理层(physical) |
IP地址
每个 internet 上的主机和路由器都有一个 ip 地址,它包括网络号和主机号,ip 地址有 ipv4(32位)或者 ipv6(128位)。可以通过 ipconfig 来查看。

端口(Port)介绍
我们这里所指的端口不是指物理意义上的端口,而是特指 TCP/IP 协议中的端口,是逻辑意义上的端口。
如果把 IP 地址比作一间房子,端口就是出入这间房子的门。真正的房子只有几个门,但是一个 IP 地址的端口可以有 65535(即:256*256-1)个之多!端口是通过端口号来标记的。(端口号 0:Reserved)
端口分类
| 0 | 保留端口 | |
| 1-1024 | 固定端口(有名端口) | 被某些程序固定使用,一般程序员不使用 |
| 1025-65535 | 动态端口 | 程序员可以使用 |
常见固定端口:
- 22: SSH 远程登录协议
- 23: telnet 使用
- 21: ftp 使用
- 25: smtp 服务使用
- 80: iis 使用
- 7: echo 服务
端口使用注意
6.3 Socket编程模型
- 服务器端(Server)监听端口,等待客户端连接
- 客户端(Client)通过 TCP 链接连接到服务器
- 客户端通过 actorRef ! "hi!" 发送消息

七、Akka网络编程 – 小黄鸡客服
7.1 需求分析
运行效果:

7.2 程序网络拓扑图
- 服务器(Actor):127.0.0.1:9999
- 客户端(CustomerActor):127.0.0.1:9998
- 客户端(不同电脑):127.0.0.1:9997

7.3 代码实现
// MessageProtocol.scala (消息协议)
case class ClientMessage(mes: String) // 客户端 -> 服务端
case class ServerMessage(mes: String) // 服务端 -> 客户端
// YellowChickenServer.scala (服务端)
class YellowChickenServer extends Actor {
override def receive: Receive = {
case "start" => println("服务器在9999端口上监听了….")
case ClientMessage(mes) => {
println("客户咨询问题是:" + mes)
mes match {
case "大数据学费是多少" =>
sender() ! ServerMessage("15000RMB")
case "学校地址" =>
sender() ! ServerMessage("昌平区宏福大楼xxx路")
case "可以学哪些技术" =>
sender() ! ServerMessage("JavaEE 大数据 Python")
case _ =>
sender() ! ServerMessage("你说啥子~~")
}
}
}
}
// YellowChickenServerApp.scala (服务端入口)
object YellowChickenServerApp extends App {
val host = "127.0.0.1" // 服务端ip地址
val port = 9999
// 创建config对象,指定协议类型,监听的ip和端口
val config = ConfigFactory.parseString(
"akka.actor.provider=\\"akka.remote.RemoteActorRefProvider\\"\\n" +
"akka.remote.netty.tcp.hostname=127.0.0.1\\n" +
"akka.remote.netty.tcp.port=9999")
private val serverActorSystem = ActorSystem("Server", config)
private val yellowChickenServerActorRef: ActorRef =
serverActorSystem.actorOf(Props[YellowChickenServer], "YellowChickenServer")
}
// CustomerActor.scala (客户端)
class CustomerActor(serverHost: String, serverPort: Int) extends Actor {
var serverActorRef: ActorSelection = _
override def preStart(): Unit = {
serverActorRef = context.actorSelection(
"akka.tcp://Server@127.0.0.1:9999/user/YellowChickenServer")
println("this.serverActorRef=" + this.serverActorRef)
}
override def receive: Receive = {
case "start" => println("客户端启动了!!…")
case mes: String => {
println("开始咨询了")
serverActorRef ! ClientMessage(mes)
}
case ServerMessage(mes) => {
println(s"收到小黄鸡咨询老师(Server): " + mes)
}
}
}
// CustomerActorApp.scala (客户端入口)
object CustomerActorApp extends App {
val (host, port, serverHost, serverPort) =
("127.0.0.1", 9990, "127.0.0.1", 9999)
val config = ConfigFactory.parseString(
"akka.actor.provider=\\"akka.remote.RemoteActorRefProvider\\"\\n" +
"akka.remote.netty.tcp.hostname=127.0.0.1\\n" +
"akka.remote.netty.tcp.port=9990")
val clientActorSystem = ActorSystem("client", config)
val actorRef: ActorRef = clientActorSystem.actorOf(
Props(new CustomerActor(serverHost, serverPort)), "CustomerActor")
actorRef ! "start"
while (true) {
val mes = StdIn.readLine()
actorRef ! mes
}
}
八、Spark Master Worker 进程通讯项目
8.1 项目意义

8.2 项目需求分析
8.3 实现功能1 – Worker完成注册
功能要求: Worker 注册到 Master,Master 完成注册,并回复 Worker 注册成功
// MessageProtocol.scala
case class RegisterWorkerInfo(id: String, cpu: Int, ram: Int)
class WorkerInfo(val id: String, val cpu: Int, val ram: Int)
case object RegisteredWorkerInfo
// SparkMaster.scala
class SparkMaster extends Actor {
val workers = collection.mutable.HashMap[String, WorkerInfo]()
override def receive: Receive = {
case "start" => println("master服务,启动并开始监听端口….")
case RegisterWorkerInfo(workerId, cpu, ram) => {
if (!workers.contains(workerId)) {
println(workerId + " 注册ok…. ")
val workerInfo = new WorkerInfo(workerId, cpu, ram)
workers += ((workerId, workerInfo))
sender() ! RegisteredWorkerInfo
}
}
}
}
// SparkMaster.scala (入口)
object SparkMaster {
def main(args: Array[String]): Unit = {
val config = ConfigFactory.parseString(
"akka.actor.provider=\\"akka.remote.RemoteActorRefProvider\\"\\n" +
"akka.remote.netty.tcp.hostname=127.0.0.1\\n" +
"akka.remote.netty.tcp.port=10001")
val actorSystem = ActorSystem("sparkMaster", config)
val masterActorRef = actorSystem.actorOf(Props[SparkMaster], "master-01")
masterActorRef ! "start"
}
}
// SparkWorker.scala
class SparkWorker(masterUrl: String) extends Actor {
var masterProxy: ActorSelection = _
val workerId = UUID.randomUUID().toString
override def preStart(): Unit = {
masterProxy = context.actorSelection(masterUrl)
}
override def receive: Receive = {
case "start" => {
println(workerId + " 向master发出注册信息…")
masterProxy ! RegisterWorkerInfo(workerId, 1, 64 * 1024)
}
case RegisteredWorkerInfo => {
println(workerId + " 向master注册成功了…. ")
}
}
}
// SparkWorker.scala (入口)
object SparkWorker {
def main(args: Array[String]): Unit = {
val host = "127.0.0.1"
val port = 10002
val masterURL = "akka.tcp://sparkMaster@127.0.0.1:10001/user/master-01"
val workerName = "worker-01"
val config = ConfigFactory.parseString(
"akka.actor.provider=\\"akka.remote.RemoteActorRefProvider\\"\\n" +
"akka.remote.netty.tcp.hostname=127.0.0.1\\n" +
"akka.remote.netty.tcp.port=10002")
val actorSystem = ActorSystem("sparkWorker", config)
val workerActorRef = actorSystem.actorOf(
Props(new SparkWorker(masterURL)), workerName)
workerActorRef ! "start"
}
}
8.4 实现功能2 – Worker定时发送心跳
功能要求: Worker 定时发送心跳给 Master,Master 能够接收到,并更新 Worker 上一次心跳时间
// MessageProtocol.scala (新增)
case object SendHeartBeat // 触发发送心跳的消息
case class HeartBeat(id: String) // 心跳消息
// WorkerInfo新增字段
class WorkerInfo(val id: String, val cpu: Int, val ram: Int) {
var lastHeartBeatTime: Long = _ // 新增:最后一次心跳时间
}
// SparkWorker.scala (修改receive)
override def receive: Receive = {
case "start" => {
println(workerId + " 向master发出注册信息…")
masterProxy ! RegisterWorkerInfo(workerId, 1, 64 * 1024)
}
case RegisteredWorkerInfo => {
println(workerId + " 向master注册成功了…. ")
println(workerId + " 准备开始定时发送心跳消息给master… ")
import context.dispatcher
// 启动定时器,每隔3秒发送一次心跳
context.system.scheduler.schedule(
0 millis, 3000 millis, self, SendHeartBeat)
}
case SendHeartBeat => {
println(s"——- $workerId 发送心跳——-")
masterProxy ! HeartBeat(workerId)
}
}
// SparkMaster.scala (修改receive,增加心跳处理)
override def receive: Receive = {
case "start" => println("master服务,启动并开始监听端口….")
case RegisterWorkerInfo(workerId, cpu, ram) => { /* … */ }
case HeartBeat(workerId) => {
val workerInfo = workers(workerId)
workerInfo.lastHeartBeatTime = System.currentTimeMillis()
println(s"master: ${workerId} 更新了心跳时间…")
}
}
8.5 实现功能3 – Master启动定时任务,检测超时的Worker
功能要求: Master 启动定时任务,定时检测注册的 Worker 有哪些没有更新心跳,已经超时的 Worker,将其从 hashmap 中删除
// MessageProtocol.scala (新增)
case object StartTimeOutWorker
case object RemoveTimeOutWorker
// SparkMaster.scala (修改)
override def receive: Receive = {
case "start" => {
println("master服务,启动并开始监听端口….")
self ! StartTimeOutWorker // 启动定时检测
}
// … RegisterWorkerInfo 和 HeartBeat 处理 …
// 开启定时器,每隔一定时间检测是否有worker的心跳超时
case StartTimeOutWorker => {
import context.dispatcher
context.system.scheduler.schedule(
0 millis, 9000 millis, self, RemoveTimeOutWorker)
}
case RemoveTimeOutWorker => {
val workerInfos = workers.values
val currentTime = System.currentTimeMillis()
// 过滤心跳超时的worker(超过6秒未更新视为超时)
workerInfos
.filter(workerInfo => currentTime – workerInfo.lastHeartBeatTime > 6000)
.foreach(workerInfo => workers.remove(workerInfo.id))
println(s"—–还剩 ${workers.size} 存活的Worker—–")
}
}
8.6 实现功能4 – Master、Worker的启动参数运行时指定
功能要求: Master、Worker 的启动参数运行时指定,而不是固定写在程序中的。
// SparkMaster.scala (修改入口)
object SparkMaster {
def main(args: Array[String]): Unit = {
// 检验参数
if (args.length != 3) {
println("请输入参数: host port masterName")
sys.exit() // 退出程序
}
val host = args(0) // "127.0.0.1"
val port = args(1) // "10001"
val masterName = args(2) // "master-01"
// … 创建ActorSystem等 …
}
}
// SparkWorker.scala (修改入口)
object SparkWorker {
def main(args: Array[String]): Unit = {
if (args.length != 4) {
println("请输入参数: host port workerName masterURL")
sys.exit()
}
val host = args(0) // "127.0.0.1"
val port = args(1) // "10002"
val masterURL = args(2) // "akka.tcp://sparkMaster@127.0.0.1:10001/user/master-01"
val workerName = args(3) // "worker-01"
// … 创建ActorSystem等 …
}
}
运行项目(IDEA配置):
如果需要运行第二个 worker 服务,则需要修改参数,再运行。
九、总结
本章主要学习了:
Akka通过Actor模型实现了:
- 无锁并发:通过消息传递避免共享状态
- 轻量级:1GB内存可容纳百万级Actor
- 分布式:天然支持跨网络通讯
- 容错性:监督机制保证系统稳定




