目录 · 第 8 / 16 章
ActorPart III · 手写运行时

08四大优势的实测验证

无锁并发精确、let-it-crash 自愈且状态不丢、有界背压拖不垮、断线续传零外部依赖。

demo/session.go:onSubscribe

上一章我们搭出了三层监督树:Manager → Session → Worker,SSE 边界打通,断线续传的接口也留好了。空谈无益,眼见为实——这一章不再讲道理,只做一件事:把第 2 章立下的四条金律逐条拿到真实运行的服务上验证,再往下多问一句「为什么会这样」,把结果背后的机制也交代清楚。所有输出都来自真实命令(go test -race ./actor/...,或 PORT=9911 go run -race ./cmd 起服务后打 curl),原样贴出,不做美化。

最省事的验证方式是先不开服务器,直接把单测跑一遍。四个用例分别对应下面要讲的四条优势,而且是在 -race 探测器全程盯梢的情况下全绿——这本身就是「无锁」与「自愈」两条优势最干净的编译期预演:

Terminal window
=== RUN TestNoLockConcurrency
--- PASS: TestNoLockConcurrency (0.00s)
=== RUN TestSupervisionRestartsWithCleanState
--- PASS: TestSupervisionRestartsWithCleanState (0.05s)
=== RUN TestSupervisionStopsAfterMaxRestarts
--- PASS: TestSupervisionStopsAfterMaxRestarts (0.11s)
=== RUN TestBackpressureBusySignal
--- PASS: TestBackpressureBusySignal (0.02s)
PASS
ok actor-notes/actor 1.835s

接下来起服务打 curl,把这四条从单测的「隔离环境」搬到真实的 HTTP + SSE 链路上,逐条验证,并搞清楚每一条背后到底是什么在起作用。

🔑 设计钥匙

下面四条优势不是我们「实现」出来的四个独立特性,而是三个设计决策的自然结果:mailbox 把并发关进队列,于是无锁;私有状态活在处理它的 Actor 里,于是崩溃可以带着状态重启;emit 是阻塞调用,于是背压自动成立、续传自动退化成一次内存重放。金律先立住,优势才会自己长出来。

无锁并发,结果精确

1000 个 goroutine 并发给同一个会话发 Bump,会话内部是裸的 s.counter++——没有 mutex,没有 atomic。如果 Actor 真的只是「又一层封装」,这里就该和普通共享内存模型一样,-race 报警、计数还对不上。实测:

Terminal window
curl -s -X POST "http://localhost:9911/bump?session=s1b&n=1000"
# {"sent":1000,"session":"s1b"}
curl -s "http://localhost:9911/count?session=s1b"
# {"count":1000,"session":"s1b"}

精确 1000,-race 全程无报警,单测 TestNoLockConcurrency 佐证了同一件事。原因不是「运气好没撞上」,而是从类型系统层面就没给「撞车」留门:PID[T] 只暴露 Tell / TellBlocking 两个方法,内部持有的是不导出的 proc *process[T]——调用方拿到的从来不是会话对象本身,而是一张「地址卡」,物理上够不到 s.counter 这个字段。唯一能碰到它的,是会话自己的那一个 Actor goroutine。

再往下一层:每个 Actor 严格是「一个 goroutine + 一个有界 mailbox channel」。Receive 每次只取一条消息、跑到底、再取下一条——run-to-completion。所以 1000 个并发的 TellBlocking 调用,竞争的只是「往 channel 里塞消息」这一步(Go channel 自身的内部同步机制保证了这一步的安全性),而真正执行 s.counter++ 的,永远只有会话那唯一一个消费者。并发被整体挪到了 mailbox 这一层,++ 本身从未离开单线程,因此也谈不上锁不锁。

对照一下就更清楚:如果把 SessionActor 换成一个裸的 struct{ Counter int },让 1000 个 goroutine 直接 obj.Counter++,这是教科书式的读—改—写数据竞争——-race 会立刻报警,最终计数也大概率小于 1000(丢更新)。Actor 模型不是「更小心地加锁」,而是从签名上删掉了「绕过 mailbox 直接改状态」这条路,想犯这个错都没有入口。

// Bump 让会话的私有计数器加一——用于演示「单线程 Actor 无需锁也无数据竞争」。
type Bump struct{}

let-it-crash 自愈,状态不丢

让 worker 在第 5 个 token 处故意 panic,随后订阅流。传统「线程绑状态」的模型里,崩溃约等于进度全丢、这次生成宣告失败。

Terminal window
curl -s -X POST "http://localhost:9911/chat?session=crashdemo&crash_at=5"
curl -sN "http://localhost:9911/stream?session=crashdemo"
# ...(中间省略)...
# id: 19
# event: token
# data: tok-19
#
# id: 20
# event: done
# data: generated 20 tokens

尽管在 token 5 崩了,客户端依然完整收到全部 20 个 token + done。服务端日志摊开了自愈的每一步:

[worker crashdemo] CRASHED: simulated worker failure at token 5 -> supervisor will restart
[worker crashdemo] started, reporting ready
[session crashdemo] worker RESTARTED (#2), resume from seq=5 (state survived crash)
[session crashdemo] generation finished, 20 tokens

崩溃 → 监督者重启 → 新 worker 上报 ready → 会话从自己的私有状态里算出续传点 seq=5 → 无缝接着吐完,全程没有一行 try/catch,也没有一次外部存储读写。拆开看,关键在 onWorkerReady:它先算 fromSeq := len(s.tokens),如果 workerStart > 1(说明这是崩溃后的重启实例,不是首次启动),就把 crashAt 强制置为 -1,再把 FromSeq 一并下发给新 worker——续传点不是另外记的一个「断点」,它就是 tokens 切片当下的长度,状态本身即检查点。

这背后是「状态与执行体分离」:worker 是无状态、可随时报废的生成器,真正值钱的进度全部活在会话这一侧的私有字段里,worker 崩溃摧毁的只是它自己的调用栈,从不牵连会话状态半分。这也是 let-it-crash 与传统 try/catch 的分野——try/catch 要求业务代码在每个可能出错的地方显式预判并处理,恢复逻辑和业务逻辑焊在一起;而 onToken/onFinished 里干干净净、一行错误处理都没有,恢复能力来自监督树这一结构性声明(newWorkerProps(...).WithSupervisor(actor.RestartStrategy(3))),对任何一次崩溃都一体适用。

当然自愈不是无限续命:demo 里 worker 挂在 RestartStrategy(3) 下,最多重启 3 次;actor 包自己的单测 TestSupervisionStopsAfterMaxRestarts 用另一组更紧的阈值(RestartStrategy(2))验证了硬币的另一面——连续崩溃超过上限后,Actor 永久停止,绝不会陷入「崩溃—重启」的死循环,这正是第 5 章 Flink restart-strategy 的同构。TestSupervisionRestartsWithCleanState 则确认了重启后拿到的是干净初始状态(脏值 42 被丢弃、Started 恰好被调用 2 次)——干净,是因为重启的是「执行体」,不是「状态」。

有界背压,慢客户端拖不垮

slow=400 模拟一个每 400ms 才读一次的慢客户端,再在生成过程中打一次监督树快照:

Terminal window
curl -s -X POST "http://localhost:9911/chat?session=bpdemo"
curl -sN "http://localhost:9911/stream?session=bpdemo&slow=400" & # 慢客户端,后台跑
curl -s "http://localhost:9911/debug/actors"

快照里 session-bpdemo 这一行:

{"pid":"/manager/session-bpdemo","state":"running","handled":14,"mailbox_len":9,"mailbox_cap":64}

mailbox_len:9——消息在 mailbox 里排队堆积,却被 cap:64 稳稳兜住,既没有无限膨胀,也没有 OOM。这条链路值得拆开看一遍:订阅者的 SSE 缓冲区 sub.events 容量被刻意调到 sseBuffer = 4,慢客户端每 400ms 才读一次,缓冲区很快见底;会话在 emit 里对 sub.events <- ev 做的是阻塞发送,一旦缓冲区满,emit 就卡住;而 emit 是会话 Receive 内部的调用,它一卡,session 这一个 Actor goroutine 就没法回去继续从自己的 mailbox 取下一条 tokenProduced——于是压力从「客户端读得慢」一路顶回「会话 mailbox 堆积」,肉眼可见地体现在 mailbox_len 上。

想直观看这条链路怎么一步步顶回去,下面这个动画可以亲手拨一拨:

💡 小贴士

Actor 的投递原语其实给了两种诚实的选择:Tell 非阻塞,mailbox 满了就直接返回 false、消息记一笔死信,把「降速还是丢弃」的裁量权交还调用方;TellBlocking 相反,宁可让调用方等,也不丢消息——handleBump 正是用它保证 1000 次 Bump 一条不少。TestBackpressureBusySignal 验证的就是前一种:非阻塞 Tell 在 mailbox 满时如实返回忙信号,而不是悄悄吞掉消息或无限扩容——背压从「猜」变成了「量出来的数字」。

// emit 把事件推给当前订阅者。阻塞发送 = 背压:客户端慢则此处等待,
// 进而 mailbox 堆积、worker 收到忙信号降速。客户端断开则通过 done 立即解除阻塞。
func (s *SessionActor) emit(ev SSEEvent) {
sub := s.sub
if sub == nil {
return
}
select {
case sub.events <- ev:
case <-sub.done:
s.sub = nil // 订阅者已断开
}
}
// handleStream 是 SSE 端点。它把自己注册为会话的订阅者,然后把会话推来的事件写给客户端。
// Last-Event-ID(或 ?from= 查询参数)决定从哪个序号开始重放——这是断线续传的入口。
func (s *Server) handleStream(w http.ResponseWriter, r *http.Request) {
id := sessionID(r)
sess, err := s.session(id)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
flusher, ok := w.(http.Flusher)
if !ok {
http.Error(w, "streaming unsupported", http.StatusInternalServerError)
return
}
fromSeq := lastEventID(r)
sub := &subscriber{
events: make(chan SSEEvent, sseBuffer),
done: make(chan struct{}),
fromSeq: fromSeq,
}
defer close(sub.done)
sess.Tell(Subscribe{Sub: sub})
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
w.Header().Set("Connection", "keep-alive")
flusher.Flush()
// 可选:慢客户端模拟。?slow=<ms> 让每次写出前 sleep,制造下游背压。
slow := time.Duration(intQuery(r, "slow", 0)) * time.Millisecond
ctx := r.Context()
for {
select {
case <-ctx.Done():
return // 客户端断开
case ev := <-sub.events:
if slow > 0 {
time.Sleep(slow)
}
fmt.Fprintf(w, "id: %d\nevent: %s\ndata: %s\n\n", ev.Seq, ev.Name, ev.Data)
flusher.Flush()
if ev.Name == "done" {
return
}
}
}
// … 省略 1 行;完整声明 L102–150,点击上方「浏览完整文件」

断线续传,零外部依赖

Last-Event-ID: 4 重连,语义是「我已经收到 4 号了,从 5 号开始给我」:

Terminal window
curl -s -X POST "http://localhost:9911/chat?session=resumedemo"
curl -sN -H "Last-Event-ID: 4" "http://localhost:9911/stream?session=resumedemo"
# id: 5
# data: tok-5
# id: 6
# data: tok-6

重连流精确从 id: 5 起步,0–4 不再重复推送。实现这份续传,我们没动 Redis、没用 MQ、没碰 sticky session——会话 Actor 自己的 tokens 切片就是那份检查点,onSubscribe 径直从 fromSeq 重放:HTTP 层的 lastEventID 先把 Last-Event-ID 头(+1)或 ?from= 参数解析成一个整数,handleStream 把它塞进 subscriber.fromSeq,再用 Tell(Subscribe{...}) 交给会话。

容易被忽略的一点是:onSubscribe 和产出新 token 的 onToken 跑在同一个会话 Actor 的同一条 Receive 里,严格串行。这意味着「把 0 到 4 号重放给新订阅者」和「继续产出第 5、第 6 号」之间不存在竞争窗口——不会出现旧连接刚断、新订阅还没挂上时溜走一个 token 的经典 bug,也不会出现同一个 token 被重放又被实时推送两次。第 1 章说续传要靠「存储组件 + 会话粘滞 + 分布式任务管理」三者配合,在 Actor 模型里退化成了一个 Actor 的一次内存重放,而且这次重放天然线性一致。

func (s *SessionActor) onSubscribe(ctx *actor.Context[SessionMsg], sub *subscriber) {
s.sub = sub
// 断线续传 / 重放:直接从私有状态里把 fromSeq 之后的 token 补发给新连接。
// 这一步完全不需要外部存储——会话状态本身就是「检查点」。
for seq := sub.fromSeq; seq < len(s.tokens); seq++ {
s.emit(SSEEvent{Seq: seq, Name: "token", Data: s.tokens[seq]})
}
if s.done {
s.emit(SSEEvent{Seq: len(s.tokens), Name: "done", Data: "already finished"})
}
ctx.Logger().Info("subscriber attached",
"session", s.id, "from_seq", sub.fromSeq, "replayed", max(0, len(s.tokens)-sub.fromSeq))
}
// lastEventID 读取 SSE 断线续传起点:优先 Last-Event-ID 头,其次 ?from= 参数。
func lastEventID(r *http.Request) int {
if h := r.Header.Get("Last-Event-ID"); h != "" {
if n, err := strconv.Atoi(h); err == nil {
return n + 1 // 从「已收到的最后一个」的下一个开始
}
}
return intQuery(r, "from", 0)
}

意外收获:一张免费的系统全景

跑完上面几个场景,/debug/actors 给出整棵监督树的快照,节选长这样:

[
{"pid":"/manager","state":"running","handled":9,"mailbox_len":0,"mailbox_cap":64},
{"pid":"/manager/session-s1b","state":"running","handled":1001,"mailbox_len":0},
{"pid":"/manager/session-crashdemo/worker","state":"running","restarts":1,"handled":23},
{"pid":"/manager/session-bpdemo","state":"running","handled":14,"mailbox_len":9}
]

一眼便能读出:session-s1b 处理了 1001 条消息(1000 次 Bump + 1 次 Count)、crashdemo/worker 重启过 1 次(restarts:1)、bpdemo 正被背压卡着(mailbox_len:9)。这份快照不是专门为本章现搭的旁路——它就是 actor.Actor.Snapshot() 遍历系统里每一个 process 直接读出的 ActorInfo(pid/parent/state/restarts/handled/mailbox_len/mailbox_cap),源码注释管它叫「具有业务语义的 pprof 火焰图」:普通锁 + 共享内存模型里,你能免费拿到 goroutine 数、能拿到锁竞争,却几乎不可能免费拿到「哪个会话」「重启过几次」「堆积了多少」这类业务量级的答案。

小结

优势验证手段关键证据
无锁并发1000 并发 Bump + -race计数精确 1000,零 data race
let-it-crash 自愈crash_at=5 + 流式订阅崩溃后仍收全 20 token,日志显示续传 seq=5
有界背压slow=400 慢客户端 + 快照mailbox_len:9 / cap:64,无 OOM
断线续传Last-Event-ID: 4 重连精确从 id:5 续传,零外部存储
  • 四条金律不是四个功能,而是三处设计决策的必然推论:mailbox 把并发关进队列 → 无锁;私有状态与执行体分离 → 自愈不丢状态;emit 阻塞 → 背压自动成立;状态本身就是检查点 → 续传退化成一次内存重放。
  • 每一条「结果」背后都能拆开看到「机制」——这正是本章比第 2 章多做的一步:不只给出证据,还给出证据为什么必然如此。
  • 教学玩具已经验证了「对」,但还没验证「稳」——下一章把它加固到生产可用:消灭 goroutine 泄漏、死信队列、优雅停止、退避 + 失败率监督、快照持久化、可观测性。
源码

正在读取完整文件…