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

09从教学玩具到生产可用

集中调度消灭 goroutine 泄漏、死信队列、优雅停止、退避+失败率监督、快照持久化、可观测性。

actor/scheduler.go:scheduleractor/supervisor.go:BackoffRestartStrategyactor/deadletter.go

前八章的运行时,足以把 Actor 的思想讲得通透,但它离「敢放上线」仍有一段肉眼可见的距离。本章不再讲新概念,而是逐条把「教学版」补齐到「生产可用」,并交代每一次加固为何必要、对齐了工业界哪一套成熟做法——这也是本项目 actor/ 目录从约 700 行膨胀到如今规模的缘由:多出来的代码几乎全花在了正确性、健壮性、可观测性与工程完整度上,而非新功能。

差距总账

维度教学版的问题加固后的做法对标
定时任务每个 tick 起一个 goroutine + time.Sleep,高负载下 goroutine 爆炸单 goroutine 最小堆调度器,任务与 Actor 生命周期绑定Netty HashedWheelTimer、Akka Scheduler
消息丢弃mailbox 满 / Actor 停止时静默返回 false死信队列 + 计数指标,绝不静默Akka DeadLetters、RabbitMQ DLX
停止语义stopTree 直接关 channel,排队消息全丢优雅 drain(可配)+ 带超时的系统级关停Akka CoordinatedShutdown
监督策略只有「固定次数重启」指数退避 + 抖动 / 失败率窗口 / OneForOne + AllForOneErlang/OTP、Flink restart-strategy
可观测性散落的 log.Printfslog 结构化日志 + traceID 贯穿 + Prometheus 指标OpenTelemetry、Prometheus
状态恢复只活在进程内存里泛型 Store[S] + 文件快照,跨进程重启恢复Flink state backend、Akka Persistence

定时任务:从「按 tick 起 goroutine」到集中调度器

教学版驱动逐 token 吐流是这么写的:每个 tick 单独 go func() { time.Sleep(...); self.Tell(tick{}) }()。单会话看不出问题,但设想「十万并发会话、每个每 60ms 一个 tick」——瞬时就是海量 goroutine 阻塞在 Sleep 上,Go 调度器和内存双重承压。更棘手的是:worker 崩溃重启后,旧实例遗留的那个 goroutine 依然会醒来,把 tick 投给已经换了新实例的 Actor,造成跨代消息错乱。

解法是引入一个单 goroutine 的集中式调度器,用最小堆按触发时间排序,全局只有一个调度循环 + 一个 time.Timer——这正是 Netty HashedWheelTimer、Akka Scheduler 的思路:定时是共享基础设施,不该随任务数线性膨胀 goroutine。

// scheduler 是一个单 goroutine 的定时任务调度器,用一个最小堆(按触发时间排序)
// 管理所有延迟/周期任务。
//
// 它解决了 demo 里「每个 tick 起一个 goroutine」的泄漏问题:无论有多少 Actor、
// 多少定时任务,全局只有一个调度 goroutine + 一个 time.Timer。这与 Netty 的
// HashedWheelTimer、Akka 的 Scheduler 是同一思路——集中式定时,避免 goroutine 爆炸。
//
// 任务的动作是一个 func(),通常是「向某个 Actor 投递一条消息」。
type scheduler struct {
mu sync.Mutex
pq timerHeap
timer *time.Timer
wake chan struct{}
closed chan struct{}
nextID atomic.Int64
metric *Counter
}
flowchart TD
  A["Context.Schedule(delay, msg)"] --> B["scheduler.add():压入最小堆"]
  B --> C{{"新任务是堆顶吗?"}}
  C -->|是| D["唤醒调度循环,重置 time.Timer"]
  C -->|否| E["循环继续等待"]
  D --> F["单一调度 goroutine"]
  E --> F
  F -->|定时器触发| G["弹出到期任务,执行 action(兜住 panic)"]
  G --> H["把信封投递给 Actor mailbox"]
  I["Actor 停止 / 重启"] --> J["process.cancelScheduled()"]
  J --> K["CancelFunc 标记任务已取消"]
  K -.->|下次弹出时跳过,不再投递| G

集中调度器:单 goroutine + 最小堆,任务随 Actor 停止/重启自动取消

真正干活的是几个实现细节。堆只在 push 和 pop 时重排——取消一个任务从不直接动堆,只是把 timerTask 上的原子 canceled 标志置位;循环下次轮到它触发时才惰性跳过,比堆内删除便宜得多。wake 是一个容量为 1 的信号 channel,只有「新任务成为堆顶」时才会被唤醒,其余时间循环高效阻塞在 time.Timer 或关停信号上——不忙轮询。而 runAction 把每次任务执行都包在自己的 recover 里,一个爆炸的定时回调永远不可能拖垮这个被全进程所有 Actor 共用的调度器。

worker 迁移后只剩一行,且不再持有 self——由 Context.Schedule 代为登记:

func (w *GenWorker) scheduleTick(ctx *actor.Context[WorkerMsg]) {
ctx.Schedule(tokenInterval, tick{})
}
func (c *Context[T]) Schedule(delay time.Duration, msg T) CancelFunc {
self := c.self
trace := c.traceID
cancel := c.system.scheduler.scheduleOnce(delay, func() {
self.deliver(envelope[T]{msg: msg, traceID: trace})
})
c.proc.trackCancel(cancel)
return cancel
}

关键就在最后一行:Context.Schedule 把返回的 CancelFunc 通过 trackCancel 登记到该 Actor 自己的 process 上,Actor 自己停止或重启时统一调用 cancelScheduled(),把它注册过的所有定时器一并取消——这同时堵死了两个洞。goroutine 不再泄漏,而「旧实例的 stale tick 打到新实例」也不可能发生,因为旧实例的定时器在它不复存在的那一刻就已被取消。TestContextScheduleCanceledOnStop 专门守护这一点:系统停止后,即便还有定时器在飞,死信数也不会持续增长。

消息丢弃:从「返回 false」到死信队列

教学版的 Tell 一旦碰上 mailbox 满或 Actor 已停止,返回个 false 便草草了事——而这个返回值极易被调用方忽略,消息就此无声无息地蒸发。这正是生产环境里最难排查的一类故障:「用户说没收到回复,日志里却什么都没有」。加固的原则只有一条,且处处适用:任何一条无法投递的消息都必须被显式记录。

func (p *PID[T]) deliver(e envelope[T]) bool {
if p == nil || p.proc == nil {
return false
}
if p.proc.mailbox.post(e) {
return true
}
// 投递失败:记录死信 + 指标,绝不静默丢弃。
reason := "mailbox full"
if p.proc.mailbox.closed() {
reason = "actor stopped"
}
p.proc.system.metrics.Counter(metricMailboxFull).Inc()
p.proc.system.deadLetter(p.id, e.msg, reason)
return false
}
// deadLetter 记录一条无法投递的消息:既进死信环形缓冲(供排查),也计数(供监控)。
func (s *Actor) deadLetter(target string, msg any, reason string) {
s.metrics.Counter(metricDeadLetters).Inc()
s.deadletters.record(DeadLetter{
Target: target,
Reason: reason,
Message: fmt.Sprintf("%T:%v", msg, msg),
Time: time.Now(),
})
s.logger.Warn("dead letter", "target", target, "reason", reason, "msg_type", fmt.Sprintf("%T", msg))
}

deadLetter 每次都做两件事:给一个计数器(metricDeadLetters)加一,供监控聚合;同时把这次丢弃写进一个有界环形缓冲,供人工排查。

package actor
import (
"sync"
"time"
)
// DeadLetter 记录一条「无法投递」的消息:目标 mailbox 满、目标已停止、或投递超时。
//
// 在生产系统里,静默丢消息是最难排查的故障之一。死信队列(Dead Letter Queue)
// 把这些丢弃事件显式记录下来,既能观测(有多少、丢给谁、为什么),也能作为
// 补偿/重投的依据。这是对齐 Akka DeadLetters、RabbitMQ DLX 的标准做法。
type DeadLetter struct {
Target string `json:"target"` // 目标 Actor 路径
Reason string `json:"reason"` // 丢弃原因(mailbox full / stopped / ...)
Message string `json:"message"` // 消息的字符串化摘要(避免持有原始引用)
Time time.Time `json:"time"`
}
// deadLetterQueue 是一个有界环形缓冲:只保留最近 N 条死信,避免无界增长。
// 累计计数由 metrics 记录,这里只留最近样本供人排查。
type deadLetterQueue struct {
mu sync.Mutex
buf []DeadLetter
next int
filled bool
cap int
}
func newDeadLetterQueue(capacity int) *deadLetterQueue {
if capacity <= 0 {
capacity = 256
}
return &deadLetterQueue{buf: make([]DeadLetter, capacity), cap: capacity}
}
func (q *deadLetterQueue) record(dl DeadLetter) {
q.mu.Lock()
q.buf[q.next] = dl
q.next = (q.next + 1) % q.cap
// … 这只是文件开头 40 行,并非完整声明;点击上方「浏览完整文件」

环形缓冲只留最近 N 条——既可观测(丢了多少、丢给了谁、因何而丢),又不会无界膨胀。运维时,一条 curl /debug/deadletters 就能回答「消息究竟去了哪」,不必翻遍日志。对那些真的无法容忍丢消息的调用方,Tell 还有个兄弟 TellBlocking:阻塞等到 mailbox 有空位或目标停止为止,用延迟换取确定送达——这与下一节「宁可慢、不可丢」的直觉一脉相承。

🔑 设计钥匙

「let it crash」这套哲学能立得住,前提是失败永远不会静默——死信队列、优雅停止、退避与失败率窗口,本质上都是同一件事:把「崩溃」从一次不可见的信息丢失,转化为一个可观测、有边界的事件。没有这层保证,「让它崩」只是「让它悄悄坏」的委婉说法。

优雅停止:不丢在途消息

教学版的 stopTree 直接关闭 mailbox channel,排队中还没处理的消息全部丢弃——对「宁可慢、不可丢」的可靠投递场景,这不可接受。加固后的 terminate 引入优雅 drain:收尾之前,先把 mailbox 里已排队的业务消息处理完,之后(drain 期间新到的)才记为死信。

func (p *process[T]) terminate() {
if p.props.drainOnStop {
p.drainRemaining() // 优雅停止:把已排队的消息处理完,而非直接丢弃
}
p.cancelScheduled()
p.stopChildren()
p.lifecycle(func(lc Lifecycle) { lc.Stopped() })
p.mailbox.close()
p.state.Store(int32(StateStopped))
p.system.deregister(p.pid.id)
// mailbox 关闭后仍残留的消息(drain 之后到达的),记为死信,绝不静默丢弃。
for _, e := range p.mailbox.drain() {
p.system.deadLetter(p.pid.id, e.msg, "actor stopped")
}
}

drainRemaining 对每条排队消息调用 safeDeliver——与正常处理走的是同一个 Receive,唯一区别是这里的 panic 会被吞掉而非交给监督者:清理阶段不该再触发一次重启风暴。是否 drain 是按 Actor 可配的(Props.WithDrainOnStop),因为不是所有排队消息都值得处理完——比如一条过期的心跳 tick,丢了也无妨。

系统级 Shutdown() 也相应升级:从「无返回、发完即忘」变成等待所有 Actor 退出、带超时、并报告是否干净退出:

// Shutdown 优雅停止所有 Actor:cancel 根 context 唤醒每个 run 循环(触发排空),
// 然后在 shutdownTimeout 内等待所有 goroutine 退出,并停掉定时调度器。
// 返回 false 表示超时仍有 Actor 未退出(生产中应告警)。
func (s *Actor) Shutdown() bool {
s.cancel()
done := make(chan struct{})
go func() {
s.wg.Wait()
close(done)
}()
select {
case <-done:
s.scheduler.stop()
s.logger.Info("actor system shut down cleanly")
return true
case <-time.After(s.shutdownTimeout):
s.scheduler.stop()
s.logger.Error("actor system shutdown timed out", "timeout", s.shutdownTimeout)
return false
}
}

超时未排空应当告警——「宁可等、不可丢」不只落在单个 Actor 上,也落在系统边界上。TestGracefulDrainOnShutdown 端到端验证了这一点:50 条排队消息,全部在 shutdown 返回前处理完。

监督策略:从「数次数」到「退避 + 失败率 + 范围」

教学版监督者只会「重启到上限就停」。生产里这远远不够,真正的容错需要三个正交能力,本项目把它们做成可组合的策略,而非一坨写死的判断:

指数退避 + 抖动——崩溃后不是立刻重启,而是越崩越慢地重启,给瞬时故障恢复时间,也防止「崩溃—重启—又崩」打满 CPU;抖动避免多个 Actor 同时重启的惊群。对齐 Flink 的 exponential-delay restart-strategy。

// BackoffRestartStrategy 返回带【指数退避 + 抖动】的重启策略。
//
// 第 n 次重启等待 min(BaseDelay * 2^n, MaxDelay) ± 抖动。这是生产级容错的标配:
// 崩溃后不立刻重启,而是越崩越慢地重启,给「瞬时故障」恢复的时间,也防止 CPU 空转风暴。
// 对齐 Flink 的 exponential-delay restart-strategy。
func BackoffRestartStrategy(cfg BackoffConfig) SupervisorStrategy {
if cfg.BaseDelay <= 0 {
cfg.BaseDelay = 100 * time.Millisecond
}
if cfg.MaxDelay <= 0 {
cfg.MaxDelay = 30 * time.Second
}
return StrategyFunc(func(fc FailureContext) Decision {
if cfg.MaxRestarts >= 0 && fc.RestartCount >= cfg.MaxRestarts {
return Decision{Directive: DirectiveStop}
}
return Decision{Directive: DirectiveRestart, Backoff: computeBackoff(cfg, fc.RestartCount)}
})
}

退避的计算本身与任何具体 Actor 的状态无关,这正是它能被系统里所有受监督 Actor 安全共用的原因:

func computeBackoff(cfg BackoffConfig, restartCount int) time.Duration {
delay := cfg.BaseDelay
for i := 0; i < restartCount && delay < cfg.MaxDelay; i++ {
delay *= 2
}
if delay > cfg.MaxDelay {
delay = cfg.MaxDelay
}
if cfg.Jitter > 0 {
factor := 1 - cfg.Jitter + rand.Float64()*(2*cfg.Jitter)
delay = time.Duration(float64(delay) * factor)
}
return delay
}

失败率窗口——比「累计次数」更贴近真实健康度:只要滚动窗口内失败不超阈值就一直重启,一旦窗口内失败太密集,判定为持续性故障、果断停止。一个每天崩两次的 Actor 和一个每分钟崩十次的 Actor 状况天差地别,固定次数分辨不出来,滚动窗口的失败率却能。对齐 Flink 的 failure-rate restart-strategy。

// FailureRateStrategy 返回「时间窗口内失败次数超限则停止」的策略。
//
// 只要 window 时间内的失败次数不超过 maxFailures,就一直重启(可选退避);
// 一旦窗口内失败太密集,说明是持续性故障而非瞬时抖动,果断停止。
// 对齐 Flink 的 failure-rate restart-strategy——比「累计次数」更贴近真实健康度。
func FailureRateStrategy(maxFailures int, window time.Duration, backoff time.Duration) SupervisorStrategy {
return StrategyFunc(func(fc FailureContext) Decision {
cutoff := fc.Now.Add(-window)
recent := 0
for _, t := range fc.Failures {
if t.After(cutoff) {
recent++
}
}
if recent > maxFailures {
return Decision{Directive: DirectiveStop}
}
return Decision{Directive: DirectiveRestart, Backoff: backoff}
})
}

重启范围 OneForOne / AllForOne——一个子崩溃,是只重启它(故障隔离最好,默认),还是把同一监督者下的所有兄弟一起重启(兄弟状态强相关、一个坏了其它也不可信时)——这是 Erlang/OTP 的经典监督语义。与其为「退避 × 范围」的每种组合各写一个策略,范围用装饰器叠加上去:

func WithScope(s SupervisorStrategy, scope RestartScope) SupervisorStrategy {
return StrategyFunc(func(fc FailureContext) Decision {
d := s.Decide(fc)
if d.Directive == DirectiveRestart {
d.Scope = scope
}
return d
})
}

三者通过一个统一接口解耦:SupervisorStrategy.Decide 只回答「做什么」——重启、停止、恢复还是上报,配多大退避、多大范围;「怎么做」——真正销毁旧实例、造出新实例,交给 processTestAllForOneRestartsSiblings 端到端验证了这一行为:崩溃一个子,兄弟确实一起重启了。

可观测性:slog、traceID 与 Prometheus

教学版靠 log.Printf 打印:字符串拼接、无法聚合,一旦请求跨越不止一个 Actor 就丢了因果链。生产可观测性是三件套协同工作。

结构化日志给每条日志打上「谁在处理、属于哪条请求」,可以按字段过滤聚合,而不必按子串 grep:

func (c *Context[T]) Logger() *slog.Logger {
return c.system.logger.With("actor", c.self.id, "trace_id", c.traceID)
}

traceID 不是裸值,而是包在消息信封里贯穿全程——入口无 trace 时自动生成,Context.Tell 向下游发送时自动传播。这样「一条请求经过 Manager → Session → Worker」的链路才能事后从日志里串起来。

指标是一个零依赖的进程内注册表:

package actor
import (
"fmt"
"io"
"sort"
"sync"
"sync/atomic"
)
// Metrics 是一个零依赖的进程内指标注册表,提供计数器(counter)与瞬时量(gauge)。
//
// 设计取舍:
// - 热路径(每条消息 +1)走【缓存好的指针】,只做一次 atomic.Add,不碰 map、不加锁。
// process 在 spawn 时把自己要用的 Counter 指针取好并缓存,避免每条消息都查表。
// - 冷路径(注册新指标、导出快照)才加锁。
// - 只用 int64,因为底层是 atomic;需要小数时按需换算(如 P99 延迟另行统计)。
type Metrics struct {
mu sync.RWMutex
counters map[string]*Counter
gauges map[string]*Gauge
}
// Counter 是只增计数器(如「已处理消息总数」)。
type Counter struct{ v atomic.Int64 }
// Add 增加 n(可为负,但语义上计数器应只增)。
func (c *Counter) Add(n int64) { c.v.Add(n) }
// Inc 自增 1,是最常见的热路径操作。
func (c *Counter) Inc() { c.v.Add(1) }
// Value 读取当前值。
func (c *Counter) Value() int64 { return c.v.Load() }
// Gauge 是可增可减、可直接设值的瞬时量(如「当前存活 Actor 数」)。
type Gauge struct{ v atomic.Int64 }
// Add 增减 n。
func (g *Gauge) Add(n int64) { g.v.Add(n) }
// … 这只是文件开头 40 行,并非完整声明;点击上方「浏览完整文件」

设计上的取舍是刻意的:热路径(每条消息处理、每次重启、每条死信)只做一次 atomic.Add,增的是一个缓存好的 *Counter 指针——不查 map、不加锁。process 在 spawn 时把要用的计数器解析一次并持有指针。只有冷路径——注册一个全新的指标名,或为 /metrics 导出快照——才会碰注册表的锁。正是这个拆分,让「加指标」永远不会拖累 Tell 的基准测试。注册表覆盖处理消息数、重启数、死信数、mailbox 已满次数、Ask 超时、存活 Actor 数、快照读写次数,一个 /metrics 端点直接以 Prometheus 文本格式被抓取。

状态持久化:泛型 Store[S] 与跨进程恢复

「Actor 私有状态即活的检查点」是这个运行时的核心理念之一,但教学版的检查点只活在进程内存里——进程一重启,全丢。生产需要这份状态活在进程之外,与 Flink 的 state backend、Akka Persistence 是同一件事。

持久化层刻意分成两层。底层是字节级、可插拔的存储后端契约 SnapshotStore——存、取、删一段字节,仅此而已——今天有内存实现和文件实现,未来 Redis、RocksDB 都能直接接入。上层是泛型的 Store[S]:业务代码读写的是具体状态类型 S,从不是 []byteany

// Store[S] 是持久化的【上层接口】:它在字节级 SnapshotStore 之上,提供针对具体状态
// 类型 S 的强类型存取。业务代码面向 Store[S] 编程,存进去、拿出来的都是 S,而不是 []byte
// 或 any——这就是「泛型优先」在持久化层的落地。
//
// 序列化用字节跳动的 sonic(高性能 JSON),而非标准库 encoding/json。
type Store[S any] struct {
backend SnapshotStore
}
func (s *Store[S]) Save(id string, state S) error {
if s.backend == nil {
return ErrNoSnapshotStore
}
data, err := sonic.Marshal(state)
if err != nil {
return err
}
return s.backend.Save(id, data)
}

序列化用字节跳动的 sonic(高性能 JSON)而非标准库;Context 提供便捷封装——SaveSnapshotLoadSnapshotDeleteSnapshot——统一用调用方自己的 Actor id 作为持久化键,业务代码从不用手拼键。demo 里的 SessionActor 据此每产出一个 token 就落一次快照,处理首条消息前先尝试恢复;把 SNAPSHOT_DIR 指向一个真实目录,kill 掉进程再拉起,客户端重连即从磁盘续传,接着上次的位置继续吐流。

明确不做的事:跨机器分布式

坦诚地划清边界,与努力补齐能力同等重要。

⚠️ 当心

本运行时刻意不碰跨机器分布式:远程透明寻址、集群成员管理、网络分区下的一致性——这些各自需要独立的传输层、序列化协议与共识/成员协议,复杂度比本章讲的一切都高出一个数量级,硬塞进「单进程学习运行时」只会两头不讨好。生产上真要分布式 Actor,应直接选型成熟框架:Go 有 proto-actor(含 cluster),.NET/跨语言有 Microsoft Orleans 的 Virtual Actor,JVM 有 Akka Cluster。本项目的价值,是把单机 Actor 的骨架与本章的加固思路讲透,让你日后翻阅这些框架源码时不再陌生。

小结

  • 教学版到生产版的每一步加固,补的都不是新概念,而是「失败被看见」的能力:goroutine 不再泄漏、消息不再静默丢弃、停止不再丢在途消息、重启策略考虑退避与真实健康度而非死数字、状态能跨进程存活。
  • 这些加固共同的方向,是把「让它崩」从一句口号变成一套可观测、有边界的工程实践——上面差距总账里的每一种失败模式,最终都走向同一个结局:计数器加一、留下记录,绝不是沉默。
  • 单机 Actor 运行时到此收尾。下一部分我们换一个战场:AI 把 Markdown 解析从「批量吞吐」逼成了「流处理」问题,而本部分打磨出的运行时,恰好是解这道题最趁手的工具。
源码

正在读取完整文件…