06手写一个极简 Actor 运行时
密封接口消息、有界 channel mailbox、串行 run 循环、recover 兜住 panic、spawn 与父子监督树、Ask。
actor/actor.go:processactor/actor.go:deliveractor/actor.go:Ask理论已经铺够,是时候动手了。本章拆开 actor/ 目录里那个几百行的运行时——它足够小,可以一口气读完;又足够真,泛型、监督树、背压、可观测性一样不缺。全书骨架可以压成一句话:一个 Actor = 一个 goroutine(run 循环)+ 一个有界 channel(mailbox)。剩下的一切,都是围绕这句话展开的类型安全与容错封装。带着这句话逐节读下去:地址与消息如何做到编译期类型安全、mailbox 如何把背压变成两个函数签名、run 循环里一次 select 如何同时扛起业务与容错、崩溃后 recover 怎样把决定权交还给监督策略、spawn 怎样把父子关系落成一棵可级联停止的树,以及 Ask 如何在纯异步的 Tell 之上长出请求-响应。
地址与消息:泛型优先,any 只留两个边界
这套运行时最重要的设计决策是拒绝 any 满天飞——这直接呼应 Go 1.18+ 社区共识与 Akka Typed 的演进方向:早期 Akka 的 ActorRef 收 Any,发错消息只能运行时炸;Akka Typed 引入 ActorRef[T] 后,发错类型编译期就报错。PID[T] 就是这条思路在 Go 里的落地——一个 Actor 的强类型地址,只能向它投递类型为 T 的消息:
type PID[T any] struct { id string proc *process[T]}你永远拿不到 Actor 对象本身,只能拿到它的 PID,这是「状态私有」的物理保证:没有对象引用,就没法直接读写别人的字段。泛型参数 T 让「地址」自带类型信息——SessionPID = PID[SessionMsg] 在编译期就无法接收 WorkerMsg,发错类型直接报错,不必等到运行时才炸。id 字段是位置无关的路径(如 /user/session-abc),分布式实现里同一个 ID 可以指向另一台机器上的 Actor,调用方代码不变——这就是「位置透明」。
但真实 Actor 通常要处理一族消息(开始、查询、上报……),不能退回 any。答案是密封接口(sealed interface)加私有标记方法:只有实现了私有方法 isSessionMsg() 的类型才能发给 PID[SessionMsg]——这是 Go 表达「消息联合体」的惯用法,私有方法保证别的包无法伪造成员:
type SessionMsg interface{ isSessionMsg() }
func (StartGen) isSessionMsg() {}func (Subscribe) isSessionMsg() {}func (Bump) isSessionMsg() {}func (GetCount) isSessionMsg() {}于是 PID[SessionMsg].Tell(WorkerMsg{...}) 在编译期就报错;Receive 里用 type switch 分派,每个分支类型都是明确的:
func (s *SessionActor) Receive(ctx *actor.Context[SessionMsg], msg SessionMsg) { switch m := msg.(type) { case StartGen: s.onStartGen(ctx, m) case Bump: s.counter++ // 无锁自增:单线程语义保证绝不发生数据竞争 case GetCount: m.Reply <- s.counter }}🔑 设计钥匙
any在这个运行时里只留两个正当边界:recover()捕获的 panic 值——它的类型在语言层面天然未知;以及 Actor 系统内部的异构注册表——需要在同一张 map 里存放消息类型各异的 Actor(PID[SessionMsg]、PID[WorkerMsg]……)。业务代码永远只看到强类型的PID[T]和密封接口消息,any被死死关在运行时内部。
第二个边界体现在 procHandle——一个只暴露「管理 / 监督 / 可观测」行为、完全不涉及消息类型 T 的小接口,让不同的 process[T] 能装进同一张 map[string]procHandle:
type procHandle interface { id() string parentID() string stateName() string restartCount() int64 handledCount() int64 mailboxStats() (queueLen, capacity int) stopTree() escalate(reason any) requestRestart(reason any) restartChildrenExcept(exceptID string, reason any)}这是 Go 处理泛型异构集合的标准招式:对外泛型,对内接口擦除。业务侧写代码时,T 从来不会退化;只有运行时内部的管理/监督/可观测逻辑,才需要经由这层擦除去碰触「任意类型的 Actor」。
Mailbox:有界 channel 就是 FIFO + 背压
Mailbox 是最能体现「share memory by communicating」的地方——直接用带缓冲 channel,白得 FIFO 顺序与并发安全,而「有界」本身就是背压之源:
post 是非阻塞投递:先判一次 done 是否已关闭,避免向已停止的 Actor 投递;再对 user channel 做一次非阻塞 select,满了就直接返回 false——这是「忙信号」,交由上游决定降速、丢弃还是重试。与它对称的是 postBlocking:
func (m *mailbox[T]) postBlocking(e envelope[T]) bool { select { case m.user <- e: return true case <-m.done: return false }}post(满则丢)与 postBlocking(满则等)这一对语义,正是前一章 credit-based 背压里「丢弃/降级」与「等待 credit」两种策略在 API 层面的落地——调用方在编译期就要显式选择自己要哪种容错姿态,不存在「默默阻塞」或「默默丢弃」的灰色地带。控制信令(停止、上报、兄弟重启)则完全绕开这条业务队列,走独立的 ctrl channel——这保证了「停止」这种管理指令的优先级天然高于业务消息,不会被一屋子排队的业务消息饿死。真正投递消息时还会裹一层 envelope[T],带上 traceID,让一条消息从进入系统到被处理的整条因果链可追踪。
run 循环:串行就是无锁之源
process[T] 把一个 Actor 实例的全部运行时状态收拢在一处——当前实例、复用的 Context、mailbox、ctrl、子 Actor 表、失败历史,一样不少:
// process[T] == 一个 goroutine(run 循环)+ 一个有界 mailbox。这是全部魔法所在。type process[T any] struct { system *Actor pid *PID[T] parent procHandle props *Props[T] mailbox *mailbox[T] ctrl chan ctrlMsg
actor Receiver[T] // 当前实例;重启时被换成全新实例 rctx *Context[T] // 复用的 Context,每次 Receive 前刷新 msg / traceID
state atomic.Int32 restarts atomic.Int64 handled atomic.Int64
failures []time.Time // 最近失败时间戳,供 failure-rate 策略判定
childMu sync.Mutex children map[string]procHandle
schedMu sync.Mutex scheduled []CancelFunc // 本 Actor 安排的定时任务,停止时统一取消,防止 stop 后仍触发}run 循环是整个运行时唯一「跑起来」的地方:
// run 是 Actor 的处理循环:一个 goroutine 串行地取消息、调用 Receive。// 「串行」是关键——它保证 Receive 永不被并发调用,Actor 内部因此天然无锁。func (p *process[T]) run() { defer p.system.wg.Done() p.state.Store(int32(StateRunning)) p.lifecycle(func(lc Lifecycle) { lc.Started() })
for { select { case c := <-p.ctrl: if p.handleCtrl(c) { p.terminate() return } case e := <-p.mailbox.user: if p.deliver(e) { p.terminate() return } case <-p.mailbox.done: p.terminate() return case <-p.system.ctx.Done(): p.terminate() return } }}它的核心是一个四路 select:ctrl 优先接住控制信令(交给 handleCtrl);mailbox.user 接住业务消息(交给 deliver);mailbox.done 感知 mailbox 被关闭(通常由 terminate 触发,走向优雅收尾);system.ctx.Done() 感知系统级 Shutdown(根 context 被取消,所有 Actor 同时被唤醒退出)。四个分支中任意一个判定「应终止」,循环就调用 terminate() 收尾并返回——这也是 spawn 里 sys.wg.Add(1) 必须在起 goroutine 之前调用的原因:run 退出时 defer p.system.wg.Done(),Shutdown 正是靠这个计数器知道所有 Actor 都已经真正退出。
下面这张图和上面的交互演示对应同一件事:一条消息如何被 select 取出、经 deliver 交给 Receive,以及崩溃后如何经 recover → 监督决策 → 换新实例走回循环:
flowchart TD A[run 循环: select] -->|ctrl| B[handleCtrl] A -->|mailbox.user 业务消息| C[deliver: Receive] A -->|mailbox.done 已关闭| D[terminate] A -->|system.ctx.Done 系统关停| D B -->|ctrlStop| D B -->|其他信令| A C -->|正常处理| A C -->|panic| E[recover] E --> F[applyFailure 监督策略] F -->|Resume 保留状态| A F -->|Restart| G[producer 造新实例] --> A F -->|Stop| D F -->|Escalate| H[parent.escalate] --> D
关键在于:这四路事件全部在同一个 goroutine 里被 select 依次取出、依次处理——Receive 因此永不会被并发调用。这不是「性能优化」,而是整个模型免锁的根本原因:会话 Actor 能把计数器写成裸的 s.counter++,不是因为幸运,而是因为架构层面就排除了并发写的可能性。并发的全部复杂度,被「关」进了 mailbox channel 那一次 select 里。可观测字段(restarts、handled)之所以另外用 atomic 保护,是因为它们要被 /debug/actors 这类跨 goroutine 的调试接口读取——那不是业务状态,业务状态从来不需要锁。
Context:Actor 与外界交互的唯一入口
业务代码摸到的另一个核心类型是 Context[T]——Receive 方法的第一个参数,也是 Actor 与外界交互的唯一入口:
// Context[T] 是 Actor 与外界交互的唯一入口。// 在 Receive 内部,Actor 通过 ctx 发消息、创建子 Actor、访问自身地址、调度定时任务。//// Context 只在「当前这次 Receive 调用」期间有效,不要跨消息保存。type Context[T any] struct { system *Actor self *PID[T] proc *process[T] msg T traceID string // 当前正在处理消息的 traceID,向下游发送时自动传播}Context 只在「当前这次 Receive 调用」期间有效,不要跨消息保存——它的 msg 与 traceID 字段每次 deliver 前都会被刷新,复用同一个实例是为了避免每条消息都分配一个新对象。它的方法组精简却完整:Self() 拿到自己的强类型地址;Tell 向「同类型」目标发消息并自动传播当前 traceID,让因果链在整棵调用树上保持完整;Schedule / ScheduleRepeated 安排延迟或周期消息——注意这里不是在 Receive 里 time.Sleep,而是把「未来的一条消息」交给系统级的单 goroutine 调度器,不管有多少个 Actor 安排了多少个定时任务,都不会多起一个 goroutine,从根源上消灭「每个 tick 一个 goroutine」的经典泄漏,而且任务与本 Actor 的生命周期绑定——重启或停止时自动取消,不会出现「stop 之后定时器仍打到已经不存在的旧实例」的诡异 bug。
let-it-crash:recover 兜住业务 panic,applyFailure 转交监督决策
deliver 是 let-it-crash 的落地点——用 defer/recover 兜住业务 panic,交给监督策略决定:
// deliver 处理单条业务消息,用 recover 兜住业务 panic——这是 let-it-crash 的落地点:// 业务代码可以放心地「崩」,运行时捕获后交给监督策略决定 Restart/Stop/Resume/Escalate。// 返回 true 表示处理循环应退出。func (p *process[T]) deliver(e envelope[T]) (stop bool) { if e.traceID == "" { e.traceID = newTraceID() // 入口处若无 trace,则生成一个,保证全链路可追踪 } defer func() { if r := recover(); r != nil { stop = p.applyFailure(r) } }() p.rctx.msg = e.msg p.rctx.traceID = e.traceID p.actor.Receive(p.rctx, e.msg) p.handled.Add(1) p.system.metrics.Counter(metricProcessed).Inc() return false}这里有两层含义值得慢下来读。第一,recover() 只发生在 defer 里,而这个 defer 早在 p.actor.Receive 被调用之前就已经挂好——不管业务代码在 Receive 内部炸得多深(哪怕嵌套调用了三层业务方法),这一层 recover 都能兜住,因为 Go 的 panic 会沿调用栈向上展开,直到遇到第一个 defer/recover。第二,recover() 返回值的类型是 any——这正是前面说的第一个正当边界:语言本身不知道你会 panic 出个 string 还是 error 还是自定义类型,类型擦除在这里是必要的,而不是偷懒。
崩溃后,applyFailure 把这个 any 值连同失败历史一起打包成 FailureContext,交给 p.props.supervisor.Decide(...) 决策:
// applyFailure 按监督策略处置一次失败。返回 true 表示本 Actor 应终止。func (p *process[T]) applyFailure(reason any) (stop bool) { now := time.Now() p.failures = append(p.failures, now) if len(p.failures) > 128 { // 有界,避免长命 Actor 的失败历史无限增长 p.failures = p.failures[len(p.failures)-128:] } p.system.metrics.Counter(metricProcessErrors).Inc()
decision := p.props.supervisor.Decide(FailureContext{ Reason: reason, RestartCount: int(p.restarts.Load()), Failures: p.failures, Now: now, }) p.system.logger.Warn("actor panic", "actor", p.pid.id, "reason", fmt.Sprint(reason), "directive", decision.Directive.String(), "restart_count", p.restarts.Load(), )
switch decision.Directive { case DirectiveResume: return false // 保留状态,忽略错误 case DirectiveRestart: if !p.restartWithBackoff(reason, decision.Backoff) { return true // 退避期间系统关停,直接终止 } if decision.Scope == AllForOne && p.parent != nil { p.parent.restartChildrenExcept(p.pid.id, reason) // 兄弟一起重启 } return false case DirectiveStop: return true case DirectiveEscalate: if p.parent != nil { p.parent.escalate(reason) // 上报给父,由父的策略处置 } return true default: return false }}监督策略会给出四种指令中的一种:Resume(保留当前实例与状态,只是当作这次错误没发生过);Restart(用工厂函数造一个全新实例顶替崩溃实例,可选叠加指数退避;若策略范围是 AllForOne,还会联动兄弟 Actor 一起重启);Stop(终止本 Actor);Escalate(把错误原封不动上报给父 Actor,由父的监督策略再决策一次——错误因此沿监督树向上传播,而不是原地终结)。restart 的精髓在于绝不复用可能已经处于脏状态的旧实例:
// Receiver[T] 是所有 Actor 必须实现的行为接口,也是模型的心脏。//// 运行时保证:同一个 Actor 的 Receive 永远不会被并发调用(一次一条消息,单线程语义)。// 因此在 Receive 内部你可以像写单线程代码一样自由读写自己的字段,完全不需要锁。// 并发的复杂度被「关」进了 mailbox channel。//// 接口只有一个方法——遵循「接口越小,抽象越强」。type Receiver[T any] interface { Receive(ctx *Context[T], msg T)}这就是为什么 Props 存的是一个 producer func() Receiver[T] 而不是现成实例——每次重启都重新跑一遍工厂函数,新实例从零初始化,旧实例连同它可能损坏的私有字段一起被整体丢弃。旧实例的子 Actor 也会被一并回收,避免留下没有父亲的孤儿 Actor。
spawn 与父子监督树
SpawnChild 在父 Actor 之下创建子 Actor 并建立监督关系。它被设计成一个包级泛型函数,而不是 Context 的方法——因为 Go 的方法不能引入新的类型参数,而子 Actor 的消息类型 U 通常不同于父的 T,SpawnChild[T, U] 让 T、U 都能从实参自动推导出来:
func SpawnChild[T, U any](parent *Context[T], name string, props *Props[U]) *PID[U] { parentProc := parent.self.proc child := spawn(parent.system, parentProc, parent.self.id+"/"+name, props) parentProc.addChild(name, child.proc) return child}真正干活的是内部的 spawn:建 mailbox → 造 PID → 造首个实例(调用一次 producer())→ 注册进系统的全局 map → 起 goroutine。
⚠️ 当心
spawn里sys.wg.Add(1)必须写在go p.run()之前,而不是放进run内部再Add——否则存在一个真实的竞态窗口:如果Shutdown恰好在这个 goroutine 被调度起来之前就执行到了wg.Wait(),计数器还是零,Wait会直接返回,而这个 Actor 的run循环随后才真正启动,永远不会被这次Shutdown等到。Add与go之间的先后顺序,是并发代码里典型的「看似无关紧要,实则决定正确性」的一行。
子 Actor 崩溃时,由父 Actor 的监督策略决定其生死;父 Actor 自己被停止或重启时,也会级联地停掉整棵子树(stopTree 沿 children 表递归下发),不留孤儿 Actor——地址路径本身也体现着这棵树,子 Actor 的 ID 是「父 ID + / + name」。
Ask:在 Tell 之上搭出类型安全的请求-响应
Actor 默认是 fire-and-forget 的 Tell。但很多时候你需要「问一句、等回答」(比如取一次会话计数):
// Ask 向 to 发送一个「携带回信通道」的请求消息,阻塞等待类型为 R 的回应或超时。//// build 回调负责把运行时提供的 reply 通道塞进你的请求消息里;服务端处理时// 通过该通道回信(reply <- resp)。回信通道带缓冲,服务端不会因回信而阻塞。// 整个过程完全类型安全:请求类型 U、响应类型 R 都由编译器检查。func Ask[U, R any](to *PID[U], build func(reply chan<- R) U, timeout time.Duration) (R, error) { var zero R if to == nil || to.proc == nil { return zero, ErrMailboxFull } to.proc.system.metrics.Counter(metricAsks).Inc() reply := make(chan R, 1) if !to.Tell(build(reply)) { return zero, ErrMailboxFull } timer := time.NewTimer(timeout) defer timer.Stop() select { case r := <-reply: return r, nil case <-timer.C: to.proc.system.metrics.Counter(metricAskTimeouts).Inc() return zero, ErrAskTimeout }}Ask 用一个带缓冲的回信 channel,在纯异步的 Tell 之上封装出请求-响应——build 回调负责把运行时提供的 reply 通道塞进你的请求消息里,服务端处理时通过这个通道回信(reply <- resp);回信通道带缓冲,服务端不会因为回信而阻塞,即便调用方已经超时放弃。请求类型 U、响应类型 R 全程由编译器检查,超时是一等公民——ErrAskTimeout 直接写在函数签名的返回值里,不是散落各处的补丁,这正回应了本书开篇说的「错误与超时建模的缺失」。
💡 小贴士
一个真实的泛型踩坑:
Tell最初被写成包级函数Tell[U](to *PID[U], msg U),结果Tell(sess, StartGen{})会让 Go 同时从两个实参回推同一个类型参数U——一个推出SessionMsg,一个推出StartGen,两者冲突直接报错。改成PID[T]的方法Tell(msg T)之后,接收者先把T绑定成SessionMsg,具体消息只需「可赋值给这个接口」即可,推导链就通了。这类「多个实参回推同一类型参数会打架」的坑,是写泛型 API 时最容易踩、也最容易被忽略的一类问题。
小结
- 骨架一句话:一个 Actor = 一个 goroutine + 一个有界 channel;其余全是围绕它展开的类型安全与容错封装。
- 泛型优先,
any死死关在两个正当边界:recover捕获的 panic 值,与系统内部的异构注册表procHandle。 - 密封接口消息 + 有界 mailbox(
post/postBlocking对应丢弃/等待两种背压姿态)+ 单 goroutine 串行run循环,把「状态私有、串行处理、地址寻址」全部落成了可编译、可运行的代码。 - let-it-crash 靠
deliver里的recover与applyFailure里的supervisor.Decide变成可执行的容错;restart靠工厂函数保证「永远是干净实例」。 spawn/SpawnChild把父子关系落成一棵可级联停止、可级联重启的监督树;Ask把「问一句等回答」和超时都变成了类型签名的一部分。- 下一章把这套运行时真正用起来:用一棵三层监督树建模一次 AI 流式会话的 demo,看
PID、mailbox、supervisor 这些抽象名词如何变成能扛住真实崩溃的系统。