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

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 的 ActorRefAny,发错消息只能运行时炸;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() 收尾并返回——这也是 spawnsys.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 里。可观测字段(restartshandled)之所以另外用 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 调用」期间有效,不要跨消息保存——它的 msgtraceID 字段每次 deliver 前都会被刷新,复用同一个实例是为了避免每条消息都分配一个新对象。它的方法组精简却完整:Self() 拿到自己的强类型地址;Tell 向「同类型」目标发消息并自动传播当前 traceID,让因果链在整棵调用树上保持完整;Schedule / ScheduleRepeated 安排延迟或周期消息——注意这里不是在 Receivetime.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]TU 都能从实参自动推导出来:

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。

⚠️ 当心

spawnsys.wg.Add(1) 必须写在 go p.run() 之前,而不是放进 run 内部再 Add——否则存在一个真实的竞态窗口:如果 Shutdown 恰好在这个 goroutine 被调度起来之前就执行到了 wg.Wait(),计数器还是零,Wait 会直接返回,而这个 Actor 的 run 循环随后才真正启动,永远不会被这次 Shutdown 等到。Addgo 之间的先后顺序,是并发代码里典型的「看似无关紧要,实则决定正确性」的一行。

子 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 里的 recoverapplyFailure 里的 supervisor.Decide 变成可执行的容错;restart 靠工厂函数保证「永远是干净实例」。
  • spawn/SpawnChild 把父子关系落成一棵可级联停止、可级联重启的监督树;Ask 把「问一句等回答」和超时都变成了类型签名的一部分。
  • 下一章把这套运行时真正用起来:用一棵三层监督树建模一次 AI 流式会话的 demo,看 PID、mailbox、supervisor 这些抽象名词如何变成能扛住真实崩溃的系统。
源码

正在读取完整文件…