07Demo 实战:用 Actor 建模 AI 流式会话
三层监督树:Manager 按 sessionID 路由、Session 持有 keyed state、Worker 无状态易崩;SSE 边界与断线续传。
demo/session.go:SessionActordemo/worker.go:GenWorker运行时的骨架已经就位,这一章把它接上一个更贴近真实场景的系统:一个支持 SSE 流式输出、断线续传、故障自愈的对话服务。整套 demo 只用了三种 Actor——Manager、Session、Worker——却把第 1 章提出的三重困境(流式建模、状态管理、错误与超时)逐一化解。这是「手写运行时」到「验证它真的好用」之间关键的一跳:接下来每一节都会对照 demo/ 目录下的真实代码,而不是停留在抽象讨论。
三层监督树:Manager → Session → Worker
整个 demo 是一棵三层监督树:/manager 按 sessionID 路由并按需创建会话,是所有会话的监督者;/manager/session-<id> 是一个有状态会话,持有该会话的全部进度;挂在其下的 .../worker 才真正吐 token。三者分工很鲜明,而且这个分工不是随意的——它对应着监督树里三种不同的职责:路由、状态、执行。
- ManagerActor 只做一件事——把请求分发到正确的 SessionActor,必要时创建它:
// ManagerActor 持有 sessionID -> 会话 PID 的私有映射。单线程访问,无需锁。type ManagerActor struct { sessions map[string]*actor.PID[SessionMsg]}它的私有状态只有一张 sessionID -> PID 的映射表。GetOrCreate 消息命中就直接把已有 PID 回传,不命中就 SpawnChild 出一个新会话并纳入监督(最多重启 5 次)。这正是第 5 章 Flink keyBy 分区器的 Actor 版本:同一个 sessionID 永远落在同一个 Session 实例上,状态因 key 天然隔离,不需要一致性哈希,也不需要外部路由表——sessions 这张 map 本身就是路由表,而且只有 Manager 自己会碰它。
- SessionActor 持有会话的全部状态:
prompt、已产出的 token、进度、当前订阅者。看它的私有字段:
// SessionActor 是一个有状态会话。它实现 actor.Receiver[SessionMsg] 与 actor.Lifecycle。//// 所有字段都是私有状态,只被本 Actor 的 Receive 单线程访问——所以【没有一把锁】。type SessionActor struct { id string
// 生成规格与进度(会话的核心状态,worker 崩溃也不会丢)。 prompt string total int crashAt int tokens []string // 已产出的 token,index 即 seq;断线重放/崩溃续传的数据来源 started bool done bool workerStart int // worker 启动次数:>1 即发生过崩溃重启
worker *actor.PID[WorkerMsg] sub *subscriber // 当前连接的订阅者,nil 表示无人连接
counter int // 演示无锁并发的私有计数器
recovered bool // 是否已尝试从快照恢复(懒加载,首条消息时触发)}tokens []string 是这一章最关键的一行——它的下标即 seq,断线重放、崩溃续传全靠对它做切片;workerStart 记录 worker 一共启动过几次,大于 1 就意味着刚发生过一次崩溃重启。SessionActor 自己不生成任何 token,只负责记账、转发事件、处理续传与重放——这种「甘当传达室」的克制,正是它能在 worker 崩溃时全身而退的前提。
- GenWorker 才是真正吐 token 的执行体,而且刻意设计成无状态、可随时重建:
// GenWorker 模拟一个 LLM 生成 worker。它是无状态可重建的:// 崩溃后监督者用工厂函数造一个全新实例,进度由父会话 Actor 持有并续传。type GenWorker struct { logger *slog.Logger sessionID string session *actor.PID[SessionMsg] // 回报进度的目标(父会话)
total int crashAt int nextSeq int}它的字段只剩三个进度量——total、crashAt、nextSeq——没有一个值得跨崩溃保留。它会在 nextSeq 撞上指定的 crashAt 时故意 panic,专门用来演示监督树的自愈。
「有状态的父 + 无状态易崩的子」,是整章的设计精髓——也是下面这张图想说的事:
flowchart TD
M["ManagerActor<br/>按 sessionID 路由"] -->|SpawnChild| S["SessionActor<br/>keyed state · tokens/进度"]
S -->|"SpawnChild(建立监督)"| W["GenWorker<br/>无状态 · 易崩"]
W -->|"Started() → workerReady"| S
S -->|"beginGen(FromSeq)"| W
W -->|"tick → tokenProduced(seq)"| S
S -->|"emit"| C(["SSE 订阅者"])
W -.->|"panic@crashAt"| Crash{{"let-it-crash<br/>supervisor 重启"}}
Crash -.->|"新实例 Started()"| W
三层监督树与一次生成的消息流
这张图里只有一个节点会崩——Worker。Manager 与 Session 之间、Session 与 Worker 之间的箭头全是稳定的监督/转发关系,崩溃被死死摁在树的最底层,不会向上传染:Manager 不知道也不需要知道某个 Session 底下的 Worker 崩过几次。
Worker 的自驱动:给自己发 tick,而不是 sleep
一个 Actor 的 Receive 绝不该长时间阻塞——如果 worker 用一个 for 循环加 time.Sleep 把 20 个 token 一次吐完,这段时间它 mailbox 里其它消息(停止指令、查询)全被饿死。Actor 处理周期性任务的惯用法是反过来:给自己发一条定时消息。
// tick 是 worker 发给自己的自驱动信号:产出下一个 token。// worker 不用 time.Sleep 阻塞 mailbox,而是「给自己发定时消息」——// 这样在吐流间隙依然能响应其它控制消息,是 Actor 里处理周期性任务的惯用法。type tick struct{}type tick struct{}
func (w *GenWorker) scheduleTick(ctx *actor.Context[WorkerMsg]) { ctx.Schedule(tokenInterval, tick{})}beginGen 把 Total、FromSeq、CrashAt 一次性交给 worker 后,Receive 里的分支就只剩两条:beginGen 触发第一次 scheduleTick,tick 触发 produce。produce 每次只做一件事——判断是否已到 total、判断是否命中 crashAt、吐一个 token、再排下一个 tick。整个循环没有一次显式等待,时间推进全部交给调度器驱动。每产出一个 token,就安排下一个 tick;两次 tick 之间 mailbox 是空闲的,能随时响应其它消息——这是 Actor 版的协作式调度。
更进一步,ctx.Schedule 不是「起一个 goroutine 睡一觉」,而是交给运行时的集中调度器:它用单个 goroutine 承载所有 Actor 的定时任务,并把定时器与 worker 的生命周期绑定——崩溃重启或停止时自动取消,旧实例的 tick 绝不会打到新实例头上。第 9 章会把这个调度器讲透。
tokenInterval 是 60 毫秒,DefaultTokenCount 是 20,一次完整生成大约要跑一秒多——慢到客户端在 SSE 流上能看清 token 一个个到达,也慢到把 crash_at 设在这段窗口正中间时,崩的是一次真实发生在生成半途的事故,而不是纸上谈兵。
💡 小贴士
「给自己发消息」这个模式不止用在 tick 上——心跳、重试、超时,在这套运行时里全是同一招:把「过一段时间该做的事」表达成一条消息加一次
Schedule,而不是一个阻塞的 goroutine。谁在处理这条消息,谁的 mailbox 就有能力在两次触发之间响应别的事——这是 Actor 模型把「时间」也纳入消息驱动的方式。
一次生成的完整消息流:从 StartGen 到崩溃续传
HTTP 层一次 POST /chat 变成一条 StartGen 消息发给 SessionActor,之后完全是消息驱动:
- Session 收到
StartGen,SpawnChild出一个 worker,建立监督关系(最多重启 3 次)。 - worker 的
Started()钩子上报workerReady。 - Session 用
len(tokens)算出fromSeq,回一条beginGen(FromSeq)。 - worker 进入 tick 循环:每个 tick 产出一个 token,
TellBlocking一条tokenProduced(seq)给 Session——
func (s *SessionActor) onToken(ctx *actor.Context[SessionMsg], m tokenProduced) { if m.Seq != len(s.tokens) { return // 只接受按序到达的 token,防御乱序/重复 } s.tokens = append(s.tokens, m.Text) s.persist(ctx) // 每产出一个 token 落一次快照,进程崩溃也能从此处续传 // 向当前订阅者转发。emit 是阻塞式的——这正是背压的源头。 s.emit(SSEEvent{Seq: m.Seq, Name: "token", Data: m.Text})}只接受按序到达的 seq,顺手落一次快照,再把事件转发给当前订阅者。
5. 如果 nextSeq 撞上 crashAt,worker 直接 panic——运行时 recover 住,监督者按策略重启。
6. 新 worker 实例的 Started() 再次上报 workerReady,但这次 workerStart 大于 1:
func (s *SessionActor) onWorkerReady(ctx *actor.Context[SessionMsg]) { s.workerStart++ fromSeq := len(s.tokens) crashAt := s.crashAt if s.workerStart > 1 { crashAt = -1 // 这是崩溃后的重启实例:不再故意崩溃,从已产出处续传 ctx.Logger().Info("worker restarted, resume from seq (state survived crash)", "session", s.id, "restart", s.workerStart, "from_seq", fromSeq) } s.worker.Tell(beginGen{Prompt: s.prompt, Total: s.total, FromSeq: fromSeq, CrashAt: crashAt})}Session 把 crashAt 置为 -1、fromSeq 设成已产出的 token 数,续传信号原样回传——这是整个 demo 里唯一一处「判断是否发生过崩溃」的代码,而它写在 Session 里,不在 worker 里,因为只有 Session 记得。
7. worker 从续传点接着吐,直到 genFinished,Session 把 done 事件推给客户端。
下面这个交互演示把上面七步具象化:worker 在第 5 个 token 处 panic,监督者重启出一个全新实例,新实例从 seq=5 续传,客户端最终无感收满 20 个 token——留意画面里 tokens 计数在崩溃前后始终没有归零。
整个自愈过程没有一行 try-catch,也没有一次 Redis 读写。
状态与执行分离:崩溃为什么不丢进度
worker 重启时是彻头彻尾的「全新实例」——nextSeq、total 全部归零,因为监督者用工厂函数重新 new 了一个 GenWorker。它能若无其事地续传,靠的不是自己记性好,而是压根不用记:进度这件事从一开始就没托付给它。
设想反过来的做法:把 nextSeq 和 token 历史直接塞进 worker 自己身上。它一 panic,这些数据就跟着一起没了——监督者可以重启这个 goroutine,但除了「从 token 0 重新开始」,没有别的东西可以交给新实例。客户端要么看到重复错乱的 token,要么看到一次生成悄悄从头重来,HTTP 层再多重试逻辑也补不回来。两种设计里,监督策略、重启机制一点没变,变的只是【状态放在哪里】——而这一个放置决定,就是「续传」和「重来」之间的全部差距。
🔑 设计钥匙
把「易崩的执行」与「要保住的状态」拆成两个 Actor,是这一章真正的设计钥匙。worker 无状态、可随手抛弃,崩了就地重启;进度、token 历史这些必须保住的东西,全部托管给更稳固的父——SessionActor。父不崩,状态就不丢;子崩了,父凭自己的私有状态算出续传点喂回去。这不是「事后加了一层容错」,而是从一开始就把状态和执行拆成了两个生命周期不同的对象。
HTTP/SSE 边界:Actor 系统如何对接外部世界
Actor 系统内部全是消息,但外部世界是 HTTP。server.go 就是这层翻译,一共挂了七个端点,每一个都在把「HTTP 请求-响应」和「Actor 消息传递」这两套语义对齐:
POST /chat把一次生成请求翻译成StartGen,顺带把crash_at查询参数透传进去;GET /stream是 SSE 端点,把自己注册成会话的订阅者,Last-Event-ID(或?from=)决定从哪个 seq 开始重放——这正是断线续传的入口。订阅者的事件 channel 缓冲故意调得很小(sseBuffer = 4),是刻意留出的背压观测点:客户端读得慢,这个 channel 就会塞满,SessionActor.emit随之阻塞,会话 mailbox 开始堆积,worker 的TellBlocking跟着降速——一条从 HTTP 响应体一路顶到生成速度的背压链;POST /bump演示无锁并发:1000 个 goroutine 并发给同一个会话发消息,SessionActor 单线程串行处理,无需一把锁;GET /count用actor.Ask把「查询私有状态」变成一次 Request-Reply,调用方拿到的是强类型的整数,不是裸 channel;GET /debug/actors、GET /debug/deadletters、GET /metrics把监督树本身变成可观测对象——业务语义的快照、死信记录、Prometheus 指标,第 9 章会把死信队列与可观测性这一层继续加固。
这几个端点合起来,把第 1 章的三重困境逐一化解:SSE + 重放解决流式建模,会话私有状态解决状态管理,监督树 + let-it-crash 解决错误与超时建模。
小结
- 三层监督树把「路由」「状态」「执行」拆成三种职责单一的 Actor:Manager 按 sessionID 路由,Session 持有 keyed state,Worker 无状态、可随时重启。
- worker 用「给自己发 tick」代替 sleep 循环,配合运行时的集中调度器,既不阻塞 mailbox,也不产生 goroutine 泄漏。
- 状态与执行分离,是 let-it-crash 真正好用的前提:子崩溃不丢进度,因为进度从没交给会崩溃的那一方。
- 下一章,我们把无锁并发、自愈、背压、断线续传这四大优势逐一跑起来,用数据说话。