深入 Play E05:Pekko Actor 的邮箱、状态与生命周期
一次 30 ms 的 ask 超时后,HTTP 客户端收到 504;300 ms 的 Actor 工作仍然完成,迟到的回复进入 dead letters。另一组实验向容量为 4 的邮箱发送 20 条业务消息,只有 4 条执行,16 条被丢弃。把状态放进 Actor,可以限制谁修改它,却不会自动获得取消、背压或持久化。
在 Play 中,控制器交付响应、ask 收到回复、Actor 结束工作与应用完成停机是不同的事件。诊断记录需要用同一个消息 ID 将它们关联起来。
版本与独立实验
工程 examples/play-electives/e05-typed-actors 固定 Play 3.0.6、Pekko Actor Typed 1.0.3、Scala 2.13.15、sbt 1.10.7 与 JDK 21。Play 源码提交为 2e56aff7d4e7a74af61e4bd39ec9e3ed7f300cd6,Pekko 为 4f77c8108aaf548a65531d2c8807da13dbba8146。实际 stage 中的 JAR 名称与摘要由 stage-manifest.json 保存,避免把文档示例版本当成运行版本。
Pekko 网站的 1.0 文档页面仍可能展示 1.0.1 依赖示例。本实验显式依赖 1.0.3,并检查实际产物,不照搬页面中的版本号。Java API 使用 org.apache.pekko.actor.typed;与 Play 现有 ActorSystem 交互时,通过 javadsl.Adapter 适配 Classic API。Pekko Typed 与 Classic 共存文档
2 项 JUnit、production stage 和真实 HTTP 验证均已通过。隔离运行 run-01 保存 22 次 HTTP 观察,其中包含等待终态的轮询;核心场景是正常 ask、超时后处理、邮箱溢出、显式重启和停止后发消息。原始记录位于 examples/play-electives/evidence/e05/isolated/,页面生成与浏览器验收另行核对。
服务绑定回环地址,没有数据库或远程集群。它使用合成密钥、关闭鉴权,并通过 GET 触发故障场景;这些端点只适合有限实验。外部诊断账本也只保存在当前进程内,不承担业务持久化。
使用 Play 管理的 ActorSystem
ActorController 注入 org.apache.pekko.actor.ActorSystem,调用 Adapter.spawn 创建两个 Typed Actor。正常计数器设置显式 restart supervision,另一个计数器设置四槽邮箱。两者使用 actor-lab-dispatcher,固定两个平台线程、throughput=1,将实验中的故意阻塞与 HTTP 默认 dispatcher 分开。
flowchart LR
H[Play Controller] --> A[Adapter.spawn]
P[Play ActorSystem] --> A
A --> N[正常计数器 显式 restart]
A --> B[有界计数器 邮箱容量 4]
N --> D[专用 dispatcher 两线程]
B --> D
B -->|邮箱满| L[DeadLetter 观察者]
N -->|迟到回复| L
P --> S[应用停止与 Actor 生命周期]
固定源码的 Adapter.spawn 最终把 Typed 行为适配成 Classic actor,由传入的系统创建。生产控制器没有调用 ActorSystem.create,因此不会额外生成一套调度线程和终止责任。单元测试和独立示例确实创建自己的系统,但都在 finally 中调用 terminate 并等待终止。固定 Adapter 源码
Actor 不是专属线程。运行记录中,同一个正常计数器的 baseline 消息在线程尾号 6 上执行,late 消息在线程尾号 5 上执行;状态仍通过该 Actor 的消息处理顺序维护。需要绑定的是状态访问协议,而不是某个 Java 线程 ID。
Counter.Command 包含 Work、Read、Crash 与实验用 Gate。业务 Work 携带 ID、有限延迟和回复地址;Reply 携带当前计数与行为代次。普通生产协议可以保留不可变消息,但应去掉 Gate 暴露的 latch 与 Future,它们只是制造可重复邮箱积压的测试设施。
ask 超时只结束等待
控制器使用 AskPattern.ask,由消息工厂取得临时 replyTo 后构造 Work。固定 AskPattern 源码解释了这层适配:隐藏的 ActorRef 与 Future 关联,接收方通过该地址回复。它不是把业务方法直接调用成同步 RPC。固定 AskPattern 源码
实验先执行 baseline,得到 count=1。随后请求 late,消息处理故意休眠 300 ms,ask 预算为 30 ms。HTTP 客户端在 49.819 ms 收到 504;控制器将 ask 的异常映射为该状态,HTTP 总时延包含调度与响应传输,不应声称它精确等于 30 ms。
sequenceDiagram
participant C as HTTP 客户端
participant H as Controller 与 ask
participant A as Counter Actor
participant L as 外部实验账本
participant D as DeadLetter 观察者
C->>H: late,ask 预算 30 ms
H->>A: Work,延迟 300 ms,replyTo
H-->>C: Future 超时,HTTP 504
C->>H: 查询处理记录
H-->>C: late 尚未完成
A->>L: 记录 late,count 2
A->>D: 回复目标已失效,记录 late reply
客户端收到 504 后立即查询 /stats,账本尚无 late。有界轮询随后观察到同一 ID 的处理记录与迟到回复的 dead-letter 记录。因此这次实验能区分“等待先结束”和“业务之后完成”,而不是仅凭超时异常推测后台状态。
自动重试会引入新的消息。如果两条消息都执行计数、写库或发送通知,超时重试可能重复产生效果。需要取消时,应把取消设计成协议的一部分,例如使用工作 ID、截止时间与可中断阶段;需要重试时,应在业务存储或协议层定义幂等。ask 本身没有替这些策略作决定。Pekko 请求与回复模式
在 Actor 内部发起异步请求,还需要避免 Future 回调直接修改私有状态。回调可能运行在另一个执行上下文,绕过邮箱串行处理。Typed 的 ActorContext.ask 或 pipeToSelf 可以把异步结果转成自己的协议消息,再由行为处理;本文控制器位于 Actor 外部,因此使用 AskPattern。
四槽邮箱与二十条消息
MailboxSelector.bounded(4) 的固定源码把满邮箱行为定义为丢弃新消息并转交 dead letters。它没有让 tell 等待空位,也没有向调用方返回业务拒绝值。若调用方需要 HTTP 503 或可恢复的排队反馈,还要增加显式准入或确认协议。固定 MailboxSelector 源码
为了消除发送速度与消费速度的竞争,实验先让一条 Gate 消息进入处理。它通过 Future 通知已进入,再等待最多两秒的 latch。控制器收到通知后连续发送 burst_0 至 burst_19,最后在 finally 中释放 gate。正在处理的 gate 不占四个排队槽位。
| 分类 | 实际消息 ID | 数量 |
|---|---|---|
| 处理成功 | burst_0 至 burst_3 |
4 |
| 满邮箱转 dead letters | burst_4 至 burst_19 |
16 |
| 重复触发同一场景 | 第二次 /burst 返回 409 |
0 条新增业务消息 |
验证脚本检查两个 ID 集合互不相交,且并集覆盖全部 20 个 ID。只看到四条成功日志,不足以说明另外十六条发生了什么;完整 ID 对账才把“未执行”落实为本地观察到的邮箱丢弃。
DeadLetter 观察者订阅当前 ActorSystem 的 event stream,只记录本实验 Work 与 Reply 类型。应用启动后先等待订阅完成,再发实验消息,避免把订阅建立之前的空窗误读为零丢弃。这个观察者证明的是此次本地运行,不把 dead letters 当作持久审计队列或可保证投递的补偿通道。Pekko 邮箱文档
邮箱上限也不是端到端背压。消息已经可能经过 HTTP 接入、反序列化或上游缓冲,进入满邮箱后才被丢弃。第 18 篇讨论的流需求协议、应用准入和 Actor 的 bounded mailbox 约束处于不同位置,不能因为它们都限制积压就互相替代。
restart 重新创建状态,不撤销已经发生的效果
正常计数器显式使用 Behaviors.supervise(Counter.create(ledger)).onFailure(IllegalStateException.class, SupervisorStrategy.restart())。Counter.create 内部用 Behaviors.setup 创建新计数状态,每次建立行为时将外部实验代次加一。故障消息抛出 synthetic-restart,该堆栈是实验预期输入,不是测试失败。
| 时点 | Actor 内部 count | 行为 generation | 外部实验账本 |
|---|---|---|---|
| baseline 与 late 完成后 | 2 | 1 | 两条已完成记录 |
| 显式 restart 后第一次 Read | 0 | 2 | 原来两条记录仍在 |
| 新消息 after_restart 完成 | 1 | 2 | 新增第三条记录 |
固定 supervision 源码要求行为实例在重启时能够重新创建,使用 Behaviors.setup 正是为了把状态构造放进这一步。实验中的 count 位于 setup 内,因此重启后归零;外部账本由控制器持有,因此不随这个 Actor 的行为重建而清空。固定 supervision 实现
这个账本只展示已发生效果与 Actor 内存状态的区别,进程退出时同样会丢失。真实数据库写入或外部通知也不会因为 Actor restart 自动回滚;那需要各自的事务与补偿设计。本文没有实现事件溯源、快照或恢复,不把非持久 Actor 的 restart 描述成数据恢复。
把可变状态分配在 setup 外,再让重启后的行为继续引用它,会得到不同结果。重启边界由对象所有权和构造位置共同决定。Actor 类型系统约束可接收的消息种类,不会自动识别哪个对象应该重新初始化。
可编译的状态协议
下面的完整类将状态编码在下一份 Behavior 中。示例在独立 ActorSystem 内运行,读取到 count=1 后终止系统。实际 Play 控制器应沿用注入的系统,而不是照搬示例中的系统创建方式。
1 | |
该文件已经使用 stage 的实际依赖编译并运行,输出包含 count=1 和系统终止记录。ActorTest 的两项测试另外验证了连续消息的状态更新,以及 restart 后 count 归零、generation 增加、外部账本仍存在。
示例只覆盖本地顺序提交。多个发送者并发提交消息、远程重试或节点故障需要各自的顺序与投递约定,不能从本地 FIFO 观察推广成全局有序或恰好一次处理。
停止 Actor 与停止应用
/stop 先通过 Classic 系统停止有界 Actor,等待该行为的 PostStop 信号,再发送 after_stop。脚本确认该 ID 出现在 dead letters,并且正常计数器仍未停止。停止一个 Actor 不会关闭整个 Play 应用,也不会让已经保存的 ActorRef 重新指向一个可用行为。
没有被引用不代表 Actor 会自动停止。应用需要明确所有权;系统关闭时会停止它管理的 Actor,局部资源则可由相应生命周期信号释放。Pekko Actor 生命周期
控制器注册 ApplicationLifecycle 关闭钩子,停止自己创建的两个 Actor 与 dead-letter 观察者,并等待三份实际停止信号。它没有调用共享 ActorSystem 的 terminate;系统本身的终止由 Play 管理。观察者在 postStop 中取消事件订阅。
实际 production 进程收到 SIGTERM 后以 143 退出,没有强杀。cleanup.json 记录两个 Actor 与观察者均已停止,业务记录最终为正常计数器 3 条、有界计数器 4 条;脚本另验证监听端口拒绝连接。只看到 JVM 退出码,无法代替这些应用终态。
复现与验证范围
下载本篇源码与实验记录,并核对SHA256SUMS。解压后的 play-electives/e05-typed-actors 是独立应用目录;接入仓库后重新执行2项JUnit、stage与生产网络,22个应用及生成class与生产jar字节一致,ask迟到、4/16邮箱分流、重启与停止断言再次通过,见 evidence/e05/shared-build.json 与 shared-http。
在工程目录设置公开的本地 JDK 与 launcher 路径,使用新的输出目录运行脚本。验证器不会启动数据库或其他服务,也不会覆盖已有证据。
1 | |
late.json 保存 HTTP 超时与后续处理,mailbox.json 保存已处理和已丢弃 ID 集合,summary.json 保存重启前后状态与停机结果,observations.json 保留单次 HTTP。编译成功只能证明 API 与类型匹配,不能证明这些运行时边界。
| 场景 | 本次状态 | 证据与限制 |
|---|---|---|
| ask 超时后继续处理 | VERIFIED_LOCAL | 同 ID 先无完成记录,随后完成并出现迟到回复 |
| 四槽邮箱积压 | VERIFIED_LOCAL | 20 条消息分为 4 处理、16 dead letters,集合完整对账 |
| 显式监督重启 | VERIFIED_LOCAL | count 2→0、generation 1→2,外部账本保留 |
| Actor 与应用停止 | VERIFIED_LOCAL | PostStop、observer 停止、进程退出与端口关闭 |
| 跨节点、持久恢复、业务取消与重试幂等 | NOT_RUN | 没有集群或持久化保证 |
练习
将故障监督策略由 restart 改成 stop,保留相同消息与 HTTP 验证。观察故障之后 Read 的结果、PostStop 时点和 dead letters;不要把“ask 超时”直接解释成“接收方很慢”。这个变体尚未运行。
在 Work 中增加明确截止时间,处理前判断是否已经过期,并用业务回复表示拒绝。让消息在 gate 后排队超过截止时间,检查过期消息是否仍修改 count。该方案只能拒绝尚未开始的工作;已经执行的副作用是否可取消,需要另外的协议。这个变体也需要新的证据。
系列导航
- 系列入口:最小应用
- 前一篇:E04,不同 HTTP 后端的行为边界(待发布)。
- 后一篇:E06,OpenTelemetry 与跨执行上下文观测(待发布)。
