JavaPekkoPekko Actor 管理器
MeteorCat
因为都是用面向对象的设计实现没接触过 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
| RestartSupervisorStrategy strategy = SupervisorStrategy.restart() .withLoggingEnabled(true) .withStopChildren(true) .withLimit(3, Duration.ofSeconds(10));
Behavior<Object> behavior = Behaviors.supervise(Behaviors.empty()) .onFailure(RuntimeException.class, strategy) .onFailure(IllegalArgumentException.class, SupervisorStrategy.stop()) ;
|
官方测试样例都大量简化监督行为初始化参数, 所以看起来 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() { RestartSupervisorStrategy strategy = SupervisorStrategy.restart() .withLoggingEnabled(true) .withStopChildren(true) .withLimit(3, Duration.ofSeconds(10));
return Behaviors.supervise(Behaviors.setup(RoomMonitor::new)) .onFailure(RuntimeException.class, strategy) .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; }
private Behavior<IRoomCommand> onPreRestart(PreRestart preRestart) { users.values().forEach(ref -> getContext().stop(ref)); users.clear(); return this; }
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;
Behavior<IRoomCommand> behavior = Behaviors.supervise(RoomSession.create(command.uid())) .onFailure(SupervisorStrategy .restart() .withLimit(3, Duration.ofMinutes(1)) .withLoggingEnabled(true));
ActorRef<IRoomCommand> ref = getContext().spawn(behavior, "session-" + command.uid); getContext().watch(ref);
users.put(command.uid, ref); return this; } }
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 { } }
public static class StorageActor extends AbstractBehavior<StorageCommand> {
private StorageActor(ActorContext<StorageCommand> context) { super(context); }
public static ActorRef<StorageCommand> create(ActorContext<StorageCommand> context, String name, int size) { 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 负责状态管理与调度 → 无状态任务池负责纯算力执行 → 结果回流串行更新状态
以场景/房间战斗场景为例, 完整流程如下:
-
场景/房间 Actor(有状态)收到玩家攻击消息
-
场景/房间 Actor 把伤害计算任务异步发给计算池, 自身不阻塞, 继续处理其他消息
-
计算池 Worker 完成纯数值计算, 把结果发回给房间 Actor
-
场景/房间 Actor 收到结果, 串行更新房间/玩家状态, 天然保证状态一致性
-
场景/房间 Actor 把最新状态广播给房间内所有玩家
-
同时把玩家数据变更 异步发给存储池异步落库
需要注意, 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
|
public static class WebSocketRunner<T> {
final ActorSystem<Void> system = ActorSystem.create(Behaviors.empty(), "WebSocketRunner");
protected ActorRef<T> actor = null;
public void onConnected(String sessionId) { this.actor = system.systemActorOf(Behaviors.empty(), "session", Props.empty()); }
public void onDisconnected(String sessionId, Throwable throwable) { if (Objects.nonNull(actor)) { } }
public void onMessage(String sessionId, String message) { if (Objects.nonNull(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 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
|
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 { } }
public static class WebSocketManager extends AbstractBehavior<IWebSocketCommand> {
public WebSocketManager(ActorContext<IWebSocketCommand> context) { super(context); }
public static Behavior<IWebSocketCommand> create() { RestartSupervisorStrategy strategy = SupervisorStrategy.restart() .withLoggingEnabled(true) .withStopChildren(true) .withLimit(3, Duration.ofSeconds(10));
return Behaviors.supervise(Behaviors.setup(WebSocketManager::new)) .onFailure(RuntimeException.class, strategy) .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;
Behavior<IWebSocketCommand> behavior = Behaviors.supervise(WebSocketSession.create(command.sessionId())) .onFailure(SupervisorStrategy .restart() .withLimit(3, Duration.ofMinutes(1)) .withLoggingEnabled(true));
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(具体会话)
那么问题就来了, 业务是不开放暴露的, 必须通过实体验证才允许进行业务通讯处理, 而在此之前网络连接是一直占用的状态
你这时候应该理解到问题所在了, 如果按照上面直接跑运行, 客户端的网络连接资源一直未释放状态
端口访问可以容纳的会话是有限的, 直接占用而不做任何业务处理只会干扰用户正常游玩逻辑
而除了这部分还有授权验证问题, 这里就是偏业务逻辑当中的设计, 一般来说用户主动推送的首个消息包必然是授权验证包, 消息结果如下
lines1 2 3 4 5 6 7 8 9
| { "sid": 0, "uid": 10001, "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(); }
private final HttpClient client;
private Behavior<IAuthorizationCommand> onVerify(IAuthorizationCommand.Verify verify) {
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();
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);
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); 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) { 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() { 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); return Behaviors.stopped(); }
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
|
syntax = "proto3";
option java_multiple_files = true; option java_package = "me.meteorcat.nova.protocol"; option optimize_for = SPEED; option java_generic_services = false; option java_string_check_utf8 = true;
message Heartbeat{ int64 ping = 1; int64 pong = 2; }
import "common.proto";
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; private static final int MAX_FRAME_LENGTH = 1024 * 1024; private static final int MSG_ID_C2S_LOGIN = 10001;
private Behavior<SessionCommand> onUnauthenticatedMessage(SessionCommand.Message command) { final String connectionId = command.connection().id();
final ByteBuffer buffer = ByteBuffer .wrap(command.message()) .order(ByteOrder.BIG_ENDIAN);
if (!buffer.hasRemaining()) { return Behaviors.same(); }
if (buffer.remaining() < MIN_FRAME_LENGTH) { logger.warnf("session actor %s frame parse error", connectionId); return Behaviors.stopped(); }
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(); }
int msgId = buffer.getInt(); if (msgId != MSG_ID_C2S_LOGIN) { logger.warnf("session actor %s received illegal msgId: %d", connectionId, msgId); return Behaviors.stopped(); }
int payloadLength = bodyLength - 4; byte[] payload = new byte[payloadLength]; buffer.get(payload);
try { C2SLogin protobuf = C2SLogin.parseFrom(payload);
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) { getContext().getSystem().eventStream().tell(new EventStream.Publish<>(broadcast)); return Behaviors.same(); } }
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);
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 系统, 后续就是业务方面的逻辑问题, 后续有时间再继续完善