Pekko Actor 管理器

因为都是用面向对象的设计实现没接触过 Actor 模式, 所以对 Actor 模式理解的不够深刻, 这里也是围绕 Pekko Actor 模式进行总结

Actor 模式总是和大量 Java 组件设计模式耦合在一起, 所以这篇篇章主要是总结常用开发概念和设计思路

首先要知道的是 Actor 最佳实现是 Erlang 语言, 很多概念都是可以从其中获取出来, 这里说下常见 Actor 维护模式(监督体系)

  • OneForOneStrategy(one_for_one): 一对一策略, 只有 发生故障的单个子 Actor 会执行恢复动作

    • 设计思想: 故障隔离, 单个实例失败不影响整体服务可用性
    • 适用场景: 子 Actor 之间完全独立、无状态依赖, 互相之间隔离所有状态
    • 每个用户连接对应一个 Session Actor, 但是需要自己定制监控规则, 配置最大重试次数、重试时间窗口、按异常类型匹配不同处理动作
  • OneForAllStrategy(one_for_all): 一对多策略, 只要有 任意一个子 Actor 故障, 所有同级子 Actor 都会执行完全相同的恢复动作

    • 设计思想: 强依赖一致性,局部故障会导致整个工作组重置
    • 适用场景: 子 Actor 之间有强状态依赖, 一个失效则整体不可用
    • 游戏场景(房间)都依赖于某个房间管理 Actor, 如果场景管理 Actor 故障, 则场景都会受到影响

Actor 监督体系就如上面所说, 日常使用有几种方向

  • one_for_one:

    • 每个用户连接对应一个 Session Actor, 用于处理用户远程会话逻辑
    • 每个独立任务对应一个 Worker Actor, 用于处理具体业务逻辑
    • 游戏玩家单独的线上实体对象, 每个玩家会话都是独立的 actor 对象
  • one_for_all:

    • 需要广播消息给所有子 Actor
    • 游戏地图场景的地区广播
    • 聊天室的房间维护

Pekko(这里采用 typed 版本) 当中自定义监督模式如下

1
2
3
4
5
6
7
8
9
10
11
12
// 10秒内最多重试3次的 Actor 管理器策略
RestartSupervisorStrategy strategy = SupervisorStrategy.restart()
.withLoggingEnabled(true) // 启用日志记录
.withStopChildren(true) // 重启时一并停止所有子 Actor, 默认为 true(可以不设置)
.withLimit(3, Duration.ofSeconds(10)); // 异常最多重启3次, 10秒内超限则停止

Behavior<Object> behavior = Behaviors.supervise(Behaviors.empty()) // Behaviors.empty() 替换成具体 Actor
.onFailure(RuntimeException.class, strategy) // 运行时异常: 重启策略, 10秒内最多重试3次, 超限则停止
.onFailure(IllegalArgumentException.class, SupervisorStrategy.stop()) // 参数异常: 直接停止, 不重试
;
// 需要拦截不同异常可以自定义设置 onFailure 回调处理
// 不需要自动 Behaviors.setup, 直接移交 ActorSystem 即可

官方测试样例都大量简化监督行为初始化参数, 所以看起来 actorOf 创建 Actor 就可以了, 但是涉及到调优就需要手动配置这些参数

监督策略

在 Pekko 内部 OneForOneStrategy 和 OneForAllStrategy 其实初始化是类似的, 只是构建行为方面有差别

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
/**
* 房间命令接口
*/
public sealed interface IRoomCommand {

// 加入房间
record JoinRoom(long uid) implements IRoomCommand {
}

// 广播消息
record Broadcast(String message) implements IRoomCommand {
}

// 离开房间
record LeaveRoom(long uid) implements IRoomCommand {
}
}


/**
* 房间管理器
*/
public static class RoomMonitor extends AbstractBehavior<IRoomCommand> {


/**
* 构造函数
*/
private RoomMonitor(ActorContext<IRoomCommand> context) {
super(context);
}

/**
* 静态实例化
*/
public static Behavior<IRoomCommand> create() {
// 10秒内最多重试3次的 Actor 管理器策略
RestartSupervisorStrategy strategy = SupervisorStrategy.restart()
.withLoggingEnabled(true) // 启用日志记录
.withStopChildren(true) // 重启时一并停止所有子 Actor, 默认为 true(可以不设置)
.withLimit(3, Duration.ofSeconds(10)); // 异常最多重启3次, 10秒内超限则停止

// 生成对应监督策略 Actor 管理器, 这里就是具体 OneForAllStrategy 初始化
return Behaviors.supervise(Behaviors.setup(RoomMonitor::new))
.onFailure(RuntimeException.class, strategy) // 运行时异常: 重启策略, 10秒内最多重试3次, 超限则停止
.onFailure(IllegalArgumentException.class, SupervisorStrategy.stop()) // 参数异常: 直接停止, 不重试
;
}


/**
* 用户列表
*/
final Map<Long, ActorRef<IRoomCommand>> users = new ConcurrentHashMap<>();

@Override
public Receive<IRoomCommand> createReceive() {
return newReceiveBuilder()

// 系统处理
.onSignal(Terminated.class, this::onTerminated)
.onSignal(PreRestart.class, this::onPreRestart)

// 业务处理
.onMessage(IRoomCommand.JoinRoom.class, this::onJoinRoom)
.onMessage(IRoomCommand.Broadcast.class, this::onBroadcast)
.onMessage(IRoomCommand.LeaveRoom.class, this::onLeaveRoom)
.build();
}

/**
* 离开房间
*/
private Behavior<IRoomCommand> onLeaveRoom(IRoomCommand.LeaveRoom leaveRoom) {
ActorRef<IRoomCommand> user = users.remove(leaveRoom.uid());
if (user != null) getContext().stop(user);
return this;
}

/**
* 广播消息
*/
private Behavior<IRoomCommand> onBroadcast(IRoomCommand.Broadcast broadcast) {
users.values().forEach(user -> user.tell(broadcast));
return this;
}

/**
* 监听 PreRestart 信号
* restart 重建实例后 sessions 为全新空映射, 旧会话全部作废
*/
private Behavior<IRoomCommand> onPreRestart(PreRestart preRestart) {
users.values().forEach(ref -> getContext().stop(ref));
users.clear(); // 清空用户列表等待重新复用
return this;
}

/**
* 监听子 Actor 停止信号
*/
private Behavior<IRoomCommand> onTerminated(Terminated terminated) {
ActorRef<?> ref = terminated.getRef();
var iterator = users.entrySet().iterator();
while (iterator.hasNext()) {
var entry = iterator.next();
if (entry.getValue().equals(ref)) {
iterator.remove();

// 打印内容
getContext().getLog().info("Removing session with uid {}", entry.getKey());
return this;
}
}
return this;
}

/**
* 加入房间
*/
private Behavior<IRoomCommand> onJoinRoom(IRoomCommand.JoinRoom command) {
if (users.containsKey(command.uid)) return this; // 已经加入会话


// 包装监督策略: 1分钟内最多重启3次, 超过则停止
// 一般子 Actor 推荐自定义再初始化一次
Behavior<IRoomCommand> behavior = Behaviors.supervise(RoomSession.create(command.uid()))
.onFailure(SupervisorStrategy
.restart()
.withLimit(3, Duration.ofMinutes(1))
.withLoggingEnabled(true));

// 生成子Actor 纳入当前监督者管理
ActorRef<IRoomCommand> ref = getContext().spawn(behavior, "session-" + command.uid);
getContext().watch(ref); // 建立生命周期监视

// 写入管理表
users.put(command.uid, ref);
return this;
}
}


/**
* 房间会话实体, 也就是最底层的 one 对象
*/
public static class RoomSession extends AbstractBehavior<IRoomCommand> {

final long uid;

private RoomSession(ActorContext<IRoomCommand> context, long uid) {
super(context);
this.uid = uid;
}

public static Behavior<IRoomCommand> create(long uid) {
return Behaviors.setup(context -> new RoomSession(context, uid));
}

@Override
public Receive<IRoomCommand> createReceive() {
return newReceiveBuilder()
.onSignal(PostStop.class, (ignore) -> {
// 退出信号
getContext().getLog().info("RoomSession stopped, uid: {}", uid);
return this;
})
.build();
}
}

核心关键点是 RoomMonitor.create 和 RoomMonitor.onJoinRoom 初始化监督管理器和子 Actor 构建

这也是传统生产工程化的使用方法, 创建监督者和子 Actor 构建监督管理关系, 强依赖一致性保证出现故障直接整体重置使用

getContext().watch(ref) 用于监控子 Actor 死亡信号, 如果子 Actor 正常停止或异常终止时, 父 Actor 会收到 Terminated 信号

这套实现可以直接面向游戏服务端和实时通讯系统, 后续就是可以补充好相关业务细节

任务池调用

上面提出的聊天室实例一般是带有状态调度的情况, 除了带有状态的 Actor 还有用于一次性的 无状态 Actor, 这部分应用大致如下:

  • 业务运算, 只需要传递数据然后运行

  • 数据落地, 将数据异步保存到指定数据源

  • AI 行为决策, 游戏 NPC 的逻辑运算

  • 路径计算, 获取地图的自动寻路路线

  • …

  • 其他的计算业务就如下所示:

  • 协议编解码、数据转换、规则引擎执行

  • 游戏伤害计算、属性校验、匹配算法

常规的状态相关业务使用维度比较

维度 有状态 Actor(房间/会话模型) 无状态 Actor 任务池
核心职责 管理状态、生命周期、业务调度 纯执行、无残留状态、用完即复用
生命周期 长生命周期,和业务实体绑定 短生命周期,池化复用,动态伸缩
故障影响 重启会丢失状态,影响关联实体 重启无副作用,仅影响当前单条任务
监督策略 按需组合 OneForOne / 一对全效果 天然 OneForOne,单个故障不影响池内其他 Worker
典型场景 房间管理、玩家会话、状态机 数值计算、数据落地、协议转换、异步通知

Pekko 针对这种情况专门设计单独 Actor 池模式, 不需要自己去设计负载均衡、生命周期管理、故障自动补充能力的 Actor 池

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
/**
* 状态命令
*/
public sealed interface StorageCommand {
record SaveSnapshot(StorageSnapshot snapshot) implements StorageCommand {
}
}

/**
* 状态 Actor
*/
public static class StorageActor extends AbstractBehavior<StorageCommand> {


/**
* 构造函数
*/
private StorageActor(ActorContext<StorageCommand> context) {
super(context);
}

/**
* 直接注入 ActorSystem 而不需要手动 Behavior 挂载
*/
public static ActorRef<StorageCommand> create(ActorContext<StorageCommand> context, String name, int size) {
// Routers 就是内置 Actor 池功能, withRoundRobinRouting 就是声明采用负载均衡
return context.spawn(
Routers.pool(size, Behaviors.setup(StorageActor::new))
.withRoundRobinRouting(), name
);
}


@Override
public Receive<StorageCommand> createReceive() {
return newReceiveBuilder()
.onMessage(StorageCommand.SaveSnapshot.class, this::onSaveSnapshot)
.build();
}

private Behavior<StorageCommand> onSaveSnapshot(StorageCommand.SaveSnapshot saveSnapshot) {
// 处理数据落地业务
return this;
}
}

这种没有状态的运算池能够充分调度服务器的计算能力, 可以用于分拆大量独立任务来做并行计算, 充分释放服务器硬件的潜力

pekko 任务池默认自带 OneForOne 监督模式(也可以自定义), 其中路由策略有以下方式

路由策略 写法 适用场景
轮询路由 withRoundRobinRouting() CPU 密集型均匀计算,比如伤害计算、属性校验、协议编解码
随机路由 withRandomRouting() IO 密集型任务,比如数据落地、第三方接口调用、缓存更新
一致性哈希路由 withConsistentHashingRouting(...) 按 key 分片调度,比如同玩家数据固定路由到同一个 Worker,减少数据库连接竞争
最小邮箱路由 withSmallestMailboxRouting() 任务耗时不均的场景,优先分配给当前最空闲的 Worker

路由策略部分按照需要选择指定功能实现即可

游戏服务端经典分层架构: 有状态 Actor 负责状态管理与调度 → 无状态任务池负责纯算力执行 → 结果回流串行更新状态

以场景/房间战斗场景为例, 完整流程如下:

  1. 场景/房间 Actor(有状态)收到玩家攻击消息

  2. 场景/房间 Actor 把伤害计算任务异步发给计算池, 自身不阻塞, 继续处理其他消息

  3. 计算池 Worker 完成纯数值计算, 把结果发回给房间 Actor

  4. 场景/房间 Actor 收到结果, 串行更新房间/玩家状态, 天然保证状态一致性

  5. 场景/房间 Actor 把最新状态广播给房间内所有玩家

  6. 同时把玩家数据变更 异步发给存储池异步落库

需要注意, pekko-pool 默认读取 dispatcher 线程数(默认等于 CPU 核心数), 阻塞 IO 会迅速耗尽线程从而拖垮整个 Actor 系统

所以推荐引入 pekko-pool 的时候, 单独绑定任务池的 dispatcher(隔离主业务线程池)

1
2
3
4
// 代码中为任务池指定独立 dispatcher( pekko.actor.io-dispatcher )
Routers.pool(size, Behaviors.setup(StorageActor::new))
.withRoundRobinRouting()
.withDispatcher("pekko.actor.io-dispatcher")

配套 application.conf 设置如下, 设置出单独的任务池

1
2
3
4
5
6
7
8
pekko.actor.io-dispatcher {
type = Dispatcher
executor = "thread-pool-executor"
thread-pool-executor {
fixed-pool-size = 16
}
throughput = 100
}

任务池支持两种标准调用方式, 对应不同业务需求:

  • tell 模式(发后即忘): poolRef.tell(command) 无返回值, 适合数据落地、日志上报、异步通知

  • ask 模式(异步响应): poolRef.ask(...) 返回 CompletionStage, 适合需要结果的计算任务, 比如伤害计算、路径查询

1
2
3
4
5
// ask 模式示例: 异步获取伤害计算结果
CompletionStage<DamageCalculateResp> result = computePool.ask(
replyTo -> new DamageCalculateReq(1001, 2002, 500, replyTo),
Duration.ofSeconds(3)
);

注意: 不推荐手动去实现任务池功能, 最好直接复用 pekko 集成好的工程化组件

当单节点算力不足时, Pekko Cluster 支持 集群级分布式任务池, 任务可以透明路由到集群任意节点的 Worker 上执行:

1
2
3
// 集群分布式组池,自动发现集群所有节点的 Worker
Routers.group(ComputeWorker.class.getSimpleName())
.withRoundRobinRouting()

具体可以参照官网文档来设计处理

连接授权

很多时候依托 Java 生态圈集成外部网络连接服务(TCP/UDP/WebSocket), 但集成时会出现很多开发的细节问题

这里假设外部集成 WebSocket 做数据转发, 主要核心是会话内部的处理问题, 所以先编写伪代码用于演示

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
/**
* WebSocket 运行时
* 这是一套伪代码运行时
*/
public static class WebSocketRunner<T> {

final ActorSystem<Void> system = ActorSystem.create(Behaviors.empty(), "WebSocketRunner");

protected ActorRef<T> actor = null;

/**
* 会话连接
*/
public void onConnected(String sessionId) {
// 假设已经创建好 WebSocketManager 管理器实现了 one_for_all
// 后面将 Behaviors.empty() 替换成具体 WebSocketManager 实现
this.actor = system.systemActorOf(Behaviors.empty(), "session", Props.empty());
}

/**
* 会话断开
*/
public void onDisconnected(String sessionId, Throwable throwable) {
if (Objects.nonNull(actor)) {
// 通知退出 PostStop 停止信号
//this.actor.tell({自定义断开 Command});
}
}

/**
* 消息传递
*/
public void onMessage(String sessionId, String message) {
if (Objects.nonNull(actor)) {
//actor.tell(sessionId,message);
}
}
}

接下来就是具体会话管理器实现功能, 这里我尽量简单点防止篇幅过长:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
/**
* WebSocket 命令
*/
public sealed interface IWebSocketCommand {
record Connected(String sessionId) implements IWebSocketCommand {
}

record Disconnected(String sessionId, Throwable throwable) implements IWebSocketCommand {
}

record Message(String sessionId, String message) implements IWebSocketCommand {
}
}


/**
* WebSocket 管理器
*/
public static class WebSocketManager extends AbstractBehavior<IWebSocketCommand> {

public WebSocketManager(ActorContext<IWebSocketCommand> context) {
super(context);
}

/**
* 静态实例化
*/
public static Behavior<IWebSocketCommand> create() {
// 10秒内最多重试3次的 Actor 管理器策略
RestartSupervisorStrategy strategy = SupervisorStrategy.restart()
.withLoggingEnabled(true) // 启用日志记录
.withStopChildren(true) // 重启时一并停止所有子 Actor, 默认为 true(可以不设置)
.withLimit(3, Duration.ofSeconds(10)); // 异常最多重启3次, 10秒内超限则停止

// 生成对应监督策略 Actor 管理器, 这里就是具体 OneForAllStrategy 初始化
return Behaviors.supervise(Behaviors.setup(WebSocketManager::new))
.onFailure(RuntimeException.class, strategy) // 运行时异常: 重启策略, 10秒内最多重试3次, 超限则停止
.onFailure(IllegalArgumentException.class, SupervisorStrategy.stop()) // 参数异常: 直接停止, 不重试
;
}

/**
* 创建接收器
*/
@Override
public Receive<IWebSocketCommand> createReceive() {
return newReceiveBuilder()
.onMessage(IWebSocketCommand.Connected.class, this::onConnected)
.onMessage(IWebSocketCommand.Disconnected.class, this::onDisconnected)
.build();
}

final Map<String, ActorRef<IWebSocketCommand>> sessions = new ConcurrentHashMap<>();

/**
* 断开连接
*/
private Behavior<IWebSocketCommand> onDisconnected(IWebSocketCommand.Disconnected disconnected) {
ActorRef<IWebSocketCommand> session = sessions.remove(disconnected.sessionId);
if (Objects.nonNull(session)) {
getContext().stop(session);
}

return this;
}

/**
* 处理连接
*/
private Behavior<IWebSocketCommand> onConnected(IWebSocketCommand.Connected command) {
if (sessions.containsKey(command.sessionId())) return this; // 已经加入会话


// 包装监督策略: 1分钟内最多重启3次, 超过则停止
// 一般子 Actor 推荐自定义再初始化一次
Behavior<IWebSocketCommand> behavior = Behaviors.supervise(WebSocketSession.create(command.sessionId()))
.onFailure(SupervisorStrategy
.restart()
.withLimit(3, Duration.ofMinutes(1))
.withLoggingEnabled(true));

// 生成子Actor 纳入当前监督者管理
ActorRef<IWebSocketCommand> ref = getContext().spawn(behavior, "session-" + command.sessionId());
getContext().watch(ref); // 建立生命周期监视

// 写入管理表
sessions.put(command.sessionId(), ref);

return this;
}
}


/**
* 具体会话对象
*/
public static class WebSocketSession extends AbstractBehavior<IWebSocketCommand> {

protected final String sessionId;

public WebSocketSession(ActorContext<IWebSocketCommand> context, String sessionId) {
super(context);
this.sessionId = sessionId;
}

public static Behavior<IWebSocketCommand> create(String sessionId) {
return Behaviors.setup(context -> new WebSocketSession(context, sessionId));
}

@Override
public Receive<IWebSocketCommand> createReceive() {
return newReceiveBuilder().build();
}
}

调用链路是 WebSocketRunner(消息层) -> WebSocketManager(管理器) -> WebSocketSession(具体会话)

那么问题就来了, 业务是不开放暴露的, 必须通过实体验证才允许进行业务通讯处理, 而在此之前网络连接是一直占用的状态

你这时候应该理解到问题所在了, 如果按照上面直接跑运行, 客户端的网络连接资源一直未释放状态

端口访问可以容纳的会话是有限的, 直接占用而不做任何业务处理只会干扰用户正常游玩逻辑

而除了这部分还有授权验证问题, 这里就是偏业务逻辑当中的设计, 一般来说用户主动推送的首个消息包必然是授权验证包, 消息结果如下

lines
1
2
3
4
5
6
7
8
9
// 这里用 JSON 格式做演示
{
// 服务器ID, 没有可传 0
"sid": 0,
// 用户ID
"uid": 10001,
// 第三方授权 token
"token": "abc123"
}

这里用 JSON 格式做演示, 实际上项目大部分情况用的是 Protobuf 等二进制序列化做传递

授权验证工程化都会采用 Web 中间层去验证, 比如 OAuth2.0 和自己的 Web 服务验证完成就在数据库创建用户实体并分配自己 token

token 验证最方便的方式是使用两方约定俗称密钥做数据签名编码解码, 比如使用 HMAC-SHA256 签名验证, 服务端只要解密验证即可

但是某些情况是需要内部 Actor 二次转发到第三方 Web 接口, 从而节省中间自己搭建 Web 接口服务的时间

大部分项目起步初期都不会搭建专门自己的 Web 接口服务, 可能回去自己做 Web 二次验证, 所以 Actor 内部会用到 HTTP 转发

不过好处是 HTTP 转化请求本身属于无状态, 所以直接沿用之前的 PekkoPool 任务池来处理即可

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
/**
* 授权命令
*/
public sealed interface IAuthorizationCommand {

/**
* 验证授权
*/
record Verify(
long uid,
String token,
ActorRef<IAuthorizationCommand> replyTo
) implements IAuthorizationCommand {

}

/**
* 验证结果
*/
record Result(
boolean success,
long uid,
int state,
String body
) implements IAuthorizationCommand {

}
}

/**
* 权限验证池
*/
public static class AuthorizationWorker extends AbstractBehavior<IAuthorizationCommand> {

/**
* 构造函数
*/
private AuthorizationWorker(ActorContext<IAuthorizationCommand> context) {
super(context);
this.client = HttpClient.newHttpClient();
}

/**
* 创建授权验证池
*/
public static <T> ActorRef<IAuthorizationCommand> create(ActorContext<T> ctx, int size, String name) {
return ctx.spawn(
Routers.pool(size, Behaviors.setup(AuthorizationWorker::new)).withRoundRobinRouting(),
name
);
}

@Override
public Receive<IAuthorizationCommand> createReceive() {
return newReceiveBuilder()
.onMessage(IAuthorizationCommand.Verify.class, this::onVerify)
.build();
}


/**
* Http 异步客户端
*/
private final HttpClient client;

/**
* 验证授权
*/
private Behavior<IAuthorizationCommand> onVerify(IAuthorizationCommand.Verify verify) {

// 异步 HTTP 调用
HttpRequest request = HttpRequest.newBuilder()
.uri(URI.create("https://example.com"))
.timeout(Duration.ofSeconds(3))
.header("Authorization", "Bearer " + verify.token)
.header("Content-Type", "application/json")
.POST(HttpRequest.BodyPublishers.ofString(
"{\"token\":\"" + verify.token() + "\"}"
))
.build();


// 异步推送
// 验证数据验证之后返回给 Actor
client.sendAsync(request, HttpResponse.BodyHandlers.ofString())
.thenApply(response -> {
// 解析验证结果
return new IAuthorizationCommand.Result(
response.statusCode() == 200,
verify.uid,
response.statusCode(),
response.body()
);
}).exceptionally(ex -> new IAuthorizationCommand.Result(
false,
verify.uid,
500,
ex.getMessage()
)).thenAccept(verify.replyTo::tell);// 完成数据验证后回写给来源 ActorRef

return Behaviors.same();
}
}

这里就是实现编写 Actor 风格的远程授权验证器样例, 之后就是核心的 Session 管理器的超时验证处理, 避免无意义连接占用服务器资源

需要重写 WebSocketSession 这部分会话处理, 简略代码如下

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
/**
* 具体会话对象
*/
public static class WebSocketSession extends AbstractBehavior<IWebSocketCommand> {

protected final String sessionId;

public WebSocketSession(ActorContext<IWebSocketCommand> context, String sessionId) {
super(context);
this.sessionId = sessionId;
}

public static Behavior<IWebSocketCommand> create(String sessionId) {
return Behaviors.setup(context -> {
WebSocketSession that = new WebSocketSession(context, sessionId);
// 这里是个小技巧, 让 createReceive 指定运行消息处理器
return that.unauthenticated();
});
}

@Override
public Receive<IWebSocketCommand> createReceive() {
return newReceiveBuilder().build();
}


/**
* 已授权的消息处理器
*/
private Behavior<IWebSocketCommand> authenticated() {
// 只放置已授权的业务
return newReceiveBuilder()
.build();
}

/**
* 未授权的消息处理器
*/
private Behavior<IWebSocketCommand> unauthenticated() {
// 只放置未授权的业务
return newReceiveBuilder()
.build();
}

}

这里用到 双状态机 严格隔离权限(自定义 authenticated 和 unauthenticated 状态切换)

不过后续的 Behavior 不建议直接返回 this:

1
2
3
4
5
6
// 示例: 普通业务消息处理
private Behavior<IWebSocketCommand> onMessage(IWebSocketCommand.Message msg) {
// 下面两种写法效果基本一致, 但是更加推荐 Behaviors.same() 返回值
return Behaviors.same(); // 标准写法
return this; // 不推荐
}

后续这里授权只需要在 unauthenticated 消息拦截追加授权验证逻辑即可, 完成之后切换到 authenticated 状态:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
/**
* 认证超时定时器
*/
private Cancelable timeoutTimer;


// ========================================================================
// 未认证态 —— 只处理认证消息, 超时强制释放连接
// ========================================================================
private Behavior<IWebSocketCommand> unauthenticated() {
// 启动认证超时定时器: 5秒未完成认证直接断开
timeoutTimer = getContext().scheduleOnce(
Duration.ofSeconds(5),
getContext().getSelf(),
new IWebSocketCommand.AuthTimeout()
);

return newReceiveBuilder()
// 处理客户端消息
.onMessage(IWebSocketCommand.Message.class, this::onUnauthenticatedMessage)
// 认证超时
.onMessage(IWebSocketCommand.AuthTimeout.class, this::onAuthTimeout)
// 异步验证结果回调
.onMessage(IWebSocketCommand.AuthResult.class, this::onAuthResult)
.build();
}


/**
* 认证超时: 主动关闭连接, 销毁会话并释放端口资源
*/
private Behavior<IWebSocketCommand> onAuthTimeout(IWebSocketCommand.AuthTimeout timeout) {
getContext().getLog().warn("会话认证超时, 释放连接: {}", sessionId);
//writeRef.tell(new IWriteCommand.CloseConnection(408, "认证超时")); // 推送 socket/websocket 数据
return Behaviors.stopped(); // 直接停止当前 Actor
}

/**
* 处理认证结果: 在 Actor 主线程串行执行, 保证线程安全
*/
private Behavior<IWebSocketCommand> onAuthResult(IWebSocketCommand.AuthResult result) {
// 取消超时定时器
timeoutTimer.cancel();

if (result.success()) {
this.uid = result.uid();
getContext().getLog().info("会话认证通过: uid={}, session={}", uid, sessionId);
// 通知客户端认证成功
writeRef.tell(new IWriteCommand.WriteMessage("{\"code\":0,\"msg\":\"auth ok\"}"));
// 核心: 切换到 [已认证业务态], 这就是典型的状态机切换状态行为
return authenticated();
} else {
// 认证失败: 直接断开连接, 释放资源
writeRef.tell(new IWriteCommand.CloseConnection(401, "认证失败: " + result.body()));
return Behaviors.stopped();
}
}

这里授权验证转发逻辑可以自己动手实现, 其实也没什么难度, 只需要根据具体需求编写相应的逻辑即可

消息解析处理

虽然写了授权验证逻辑, 但是核心的消息解析逻辑设计还没有说明, 现在就是补充这部分

一般来说客户端提交的数据从性能出发都是二进制数据并且在头部包含数据长度, 简单来说就是这类数据结构

1
2
3
4
5
数据包A: [int32:二进制数据长度] [int32:数据包ID] [bytes:二进制数据]
└── 帧头(4字节) └─ 自定义消息ID └── 长度值对应的完整业务包(通常为 Protobuf/自定义二进制)

数据包B: [int32:二进制数据长度] [int32:数据包ID] [bytes:二进制数据]
以此类推....

外部不管是 TCP/UDP/WebSocket 都不是 Actor 需要关心的, Actor 只需要留意传递过来 sessionId + msgId + byte[] 对象就行

而其中依据网络服务端架构来说, 解析层应该放置于 SessionActor 层面处理, 解析成功直接将解包的数据传递给业务逻辑层的 Actor

而 MsgId 则是另外需要关心的情况, 因为单纯传递二进制数据并不知道用什么解析类处理

1
2
3
4
5
6
7
8
9
10
11
12
// 心跳包消息
public class HeartbeatMessage {

}
// 授权消息
public class AuthorizationMessage {

}

byte[] message = ...; // 假设这是接收到的二进制数据

// 除了遍历所有消息确定哪个消息类解析, 没办法高效解析处理消息内容

所以也就需要在客户端消息提交的时候附带上消息 ID, 服务端根据消息 ID 调用对应的解析类进行解析, 采用 Protobuf 作为协议载体来划分以下功能

  • common: 共享所有数据结构

  • 其他具体的协议文件

这样分配出来的文件如下

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
// common.proto 定义文件
// ------------------------------------------------------------------------------
// 帧格式: 4 字节大端长度前缀 + 序列化帧体. 帧体是 oneof, 单一解析器即可分发全部客户端/服务端消息
syntax = "proto3";

// 设置 Java 相关构建信息, 将其构建的协议生成文件放置在指定 Java 包内
option java_multiple_files = true;
option java_package = "me.meteorcat.nova.protocol"; // 指定包
option optimize_for = SPEED; // 优先序列化速度
option java_generic_services = false; // 关闭无用RPC代码生成
option java_string_check_utf8 = true; // 强制UTF8校验


// 心跳包
message Heartbeat{
int64 ping = 1; // 客户端发起时间: 毫秒
int64 pong = 2; // 服务端响应时间: 毫秒
}


// session.proto 定义文件
// ------------------------------------------------------------------------------
import "common.proto"; // 引入公用定义

// 如果提示 warning: Import common.proto is unused. 错误信息的话将共用定义文件注释掉
// 然后采用手动重写声明如下信息
// syntax = "proto3";
// option java_multiple_files = true;
// option java_package = "me.meteorcat.nova.protocol";
// option optimize_for = SPEED; // 优先序列化速度
// option java_generic_services = false; // 关闭无用RPC代码生成
// option java_string_check_utf8 = true; // 强制UTF8校验
// 这里的 C2S 代表 client to server 客户端到服务器, S2C 代表 server to client 服务器到客户端


// 授权请求
message C2SLogin {
uint32 uid = 1;
string token = 2;
}

// 授权响应
message S2CLogin {
string sessionId = 1;
uint32 status = 2;
string command = 3;
}

然后就是具体的 SessionActor 数据解析功能处理, 以上面的 onUnauthenticatedMessage 方法来说说明

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
private static final int MIN_FRAME_LENGTH = 8;       // 最小帧长: 4字节长度+4字节msgId
private static final int MAX_FRAME_LENGTH = 1024 * 1024; // 最大单包1MB
private static final int MSG_ID_C2S_LOGIN = 10001; // 未认证唯一允许的msgId

/**
* 消息解包并启动验证
*/
private Behavior<SessionCommand> onUnauthenticatedMessage(SessionCommand.Message command) {
final String connectionId = command.connection().id();

// 1. 显式大端序
final ByteBuffer buffer = ByteBuffer
.wrap(command.message())
.order(ByteOrder.BIG_ENDIAN);

// 空包直接丢弃
if (!buffer.hasRemaining()) {
return Behaviors.same();
}


// 2. 最小帧长校验
if (buffer.remaining() < MIN_FRAME_LENGTH) {
logger.warnf("session actor %s frame parse error", connectionId);
return Behaviors.stopped();
}

// 3. 读取帧体总长度(=msgId 4字节 + 业务载荷长度)
int bodyLength = buffer.getInt();
if (bodyLength < 4 || bodyLength > MAX_FRAME_LENGTH) {
logger.warnf("session actor %s frame invalid length: length = %d", connectionId, bodyLength);
return Behaviors.stopped();
}
if (buffer.remaining() != bodyLength) {
logger.warnf("session actor %s frame length mismatch: length = %d, remaining = %d", connectionId, bodyLength, buffer.remaining());
return Behaviors.stopped();
}

// 4. 读取msgId, 未认证态白名单拦截
int msgId = buffer.getInt();
if (msgId != MSG_ID_C2S_LOGIN) {
logger.warnf("session actor %s received illegal msgId: %d", connectionId, msgId);
return Behaviors.stopped();
}

// 5. 提取纯业务载荷(跳过msgId的4字节)
int payloadLength = bodyLength - 4;
byte[] payload = new byte[payloadLength];
buffer.get(payload);

// 用指定构建器来转化数据实体
// 异常直接终端 Actor 会话
try {
C2SLogin protobuf = C2SLogin.parseFrom(payload);
// todo: 转发到验证 Actor

// 通知验证器处理消息
final var ref = verifier.actorRef();
ref.tell(new VerifierCommand.VerifyRequest(1001, "meteorcat", getContext().getSelf()));
return Behaviors.same();
} catch (InvalidProtocolBufferException e) {
logger.errorf(e, "session actor %s protobuf parse error", command.connection().id());
return Behaviors.stopped();
}
}

这里只是粗浅编写这部分逻辑, 实际上工程项目设计需要更加细致, 正式生产的时候需要对数据流程写入和切分(这部分后续再将)

后续需要对数据流进行切流, 要把粘包拆包逻辑挪到连接级缓冲区, 不要放在 Session Actor 里做流式拼接

这里就是简单的处理逻辑, 实际上需要对 Protobuf 再包一层来做通用映射消息处理组件

广播与订阅

之前的 RoomMonitor 看过简单的广播逻辑业务

1
2
3
4
5
6
7
/**
* 广播消息
*/
private Behavior<IRoomCommand> onBroadcast(IRoomCommand.Broadcast broadcast) {
users.values().forEach(user -> user.tell(broadcast));
return this;
}

虽然这样也能直接广播消息, 但是业务全部挤压在核心管理器当中(职责越界: 消息处理应该交给自己旗下 Actor 自己去发送)

而 Pekko 当中提供事件总线 EventStream 就是专门来处理 发布-订阅模式 的处理方式

EventStream 完全解耦发布方和接收方, 专门针对 one_for_all 的广播处理, 是广泛用于跨模块广播的标准方案

EventStream 不同版本实现可能不一样, 请根据具体版本文档进行实现, 目前订阅发布机制如下处理

操作 写法 说明
发布事件 eventStream.tell(new EventStream.Publish<>(event)) 发送 Publish 命令到事件总线 Actor
订阅事件 eventStream.tell(new EventStream.Subscribe<>(subscriberRef)) 发送 Subscribe 命令, 按事件类型订阅
取消订阅 eventStream.tell(new EventStream.Unsubscribe<>(subscriberRef)) 发送 Unsubscribe 命令

而 EventStream 使用也是十分简单, 基本上直接用在父级管理器当中调用消息处理

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
/**
* 房间管理器 - 消息发布端
*/
public static class RoomMonitor extends AbstractBehavior<IRoomCommand> {

// 其他略

/**
* 广播消息
*/
private Behavior<IRoomCommand> onBroadcast(IRoomCommand.Broadcast broadcast) {
// 直接发布消息 , 让所有订阅者接收(也就是子 Actor 接收)
getContext().getSystem().eventStream().tell(new EventStream.Publish<>(broadcast));
return Behaviors.same();
}
}


/**
* 房间会话实体, 也就是最底层的 one 对象 - 消息订阅端
*/
public static class RoomSession extends AbstractBehavior<IRoomCommand> {

// 其他略

/**
* 构造函数
*/
private RoomSession(ActorContext<IRoomCommand> context, long uid) {
super(context);
this.uid = uid;

// 订阅句柄
final ActorRef<EventStream.Command> event = context.getSystem().eventStream();

// 订阅父级广播消息
event.tell(new EventStream.Subscribe<>(
IRoomCommand.Broadcast.class, getContext().getSelf().narrow()
));
}


@Override
public Receive<IRoomCommand> createReceive() {
return newReceiveBuilder()
.onSignal(PostStop.class, (ignore) -> {
// 退出信号
getContext().getLog().info("RoomSession stopped, uid: {}", uid);

// 订阅句柄, 退订广播消息
// 取消订阅一个 Actor 时, 默认移除该 Actor 在事件总线上的全部订阅
final ActorRef<EventStream.Command> event = getContext().getSystem().eventStream();
event.tell(new EventStream.Unsubscribe<>(getContext().getSelf()));
return this;
})
.onMessage(IRoomCommand.Broadcast.class, this::onBroadcast) // 获取订阅的消息
.build();
}
}

EventStream 是全局单例, 也就是所有旗下 Actor 共享同一个事件总线, 如果推送比较频繁可能还是有点性能问题

这里就可以实现最基础的 PekkoActor 系统, 后续就是业务方面的逻辑问题, 后续有时间再继续完善