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 + AllForOne | Erlang/OTP、Flink restart-strategy |
| 可观测性 | 散落的 log.Printf | slog 结构化日志 + 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 只回答「做什么」——重启、停止、恢复还是上报,配多大退避、多大范围;「怎么做」——真正销毁旧实例、造出新实例,交给 process。TestAllForOneRestartsSiblings 端到端验证了这一行为:崩溃一个子,兄弟确实一起重启了。
可观测性: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,从不是 []byte 或 any。
// 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 提供便捷封装——SaveSnapshot、LoadSnapshot、DeleteSnapshot——统一用调用方自己的 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 解析从「批量吞吐」逼成了「流处理」问题,而本部分打磨出的运行时,恰好是解这道题最趁手的工具。