12拓扑设计:三种 Actor 与消息协议
DocActor(串行心脏)→ InlineWorker 池(并行)→ Assembler(重排序 + 推送);单一入口门面。
mdparser/actors.go:docActormdparser/actors.go:Parser前一章看过 CommonMark 官方策略里画出的那道切缝:块级解析必须串行,行内解析可以并行。这一章把这道切缝真正浇筑成代码——mdparser 用三种 Actor 搭出一棵监督树,外加一份密封接口消息协议,把「解析」拆成串行与并行两半,再在边界上把并行算出来的碎片重新缝合成一条有序的输出流。三种 Actor 各司其职、消息协议把「谁能给谁发什么」焊死在类型系统里、三个刻意的设计决策把「并行之后怎么办」的难题一一接住——这三层加在一起,才是这套拓扑真正站得住的原因,而不只是一张好看的流程图。
一张图:三个 Actor,一条流水线
流水线的形状是这样的:输入从 Feed(chunk) / Close() 进来,先落到 DocActor——它是整条流水线唯一的「串行心脏」,持有行缓冲与块级状态机(容器栈 + 打开的叶子块)。每条 feedChunk 到达后,先按换行符切出完整行,逐行喂给状态机(bp.feedLine(line)),状态机再吐出零到多个 BlockEvent。块级解析天然是顺序的:一行会不会开新块、该挂到哪个容器下,完全取决于「之前已经打开了哪些块」——这是一条绕不开的因果链,所以 DocActor 从不并行处理消息,它只做一件事:把唯一可并行的部分——行内解析——切出去,自己保持轻快。
// docActor 持有块级状态机与未完行缓冲,是整条流水线唯一的「串行心脏」。// 块级解析天然是顺序的(每一行都依赖之前的打开块结构),所以它不并行;// 它把唯一可并行的部分(行内解析)切出去,自己保持轻快。type docActor struct { cfg parserConfig bp *BlockState buf string // 跨 chunk 的未完行缓冲 asm *actor.PID[AsmMsg] workers []*actor.PID[InlineMsg] rr int // round-robin 派发指针 inited bool closed bool}被切出去的部分交给一组 InlineWorker。它们是无状态执行体:领到一个 leafJob{seq, leaf},拿 leaf.Content 跑一遍行内解析(分隔符栈配对、链接递归、转义),交给配置好的渲染器,吐出一个 fragment{seq, html},仅此而已——不持有跨消息的状态,也不记得上一个 job 是谁。无状态意味着它们可以放心并行、放心崩溃:崩了就换一个新实例顶上,不影响其他 worker,更碰不到 DocActor 的状态机(状态机本来就没有放在 worker 里)。
// inlineWorker 是无状态执行体:收叶子块、行内解析、渲染、上报片段。// 它是「let it crash」的落地对象——panicOn 命中时直接 panic,// 监督者换一个新实例顶上,被吞掉的那条消息由 Assembler 的缺口自愈兜底。type inlineWorker struct { asm *actor.PID[AsmMsg] parse *Config panicOn string}并行算出来的碎片必然乱序到达,于是需要第三个角色:Assembler。它是流水线里另一个「有状态」节点,但状态的性质和 DocActor 完全不同——DocActor 的状态是「解析到哪了」,Assembler 的状态是「谁先谁后」。它按 DocActor 分配的文档序号 seq 把碎片重新排好,只有连续就位的前缀才会定稿并推送给订阅者;容器标签片段、暂定尾巴更新(provisionalUpdate)、文档终结信号(endOfDoc)则由 DocActor 直接旁路发给它,不经过 worker——这三类消息本来就不带需要并行计算的行内内容,插队反而更快。
// assembler 是重排序器 + 推送边界:// - pending 暂存乱序到达的片段,next 指向待定稿序号——经典 resequencer;// - 订阅推送是阻塞式的(带 done 逃生阀),慢订阅者 → 本 Actor 停摆 →// mailbox 堆积 → worker/doc 的 TellBlocking 依次变慢 → Feed 变慢。// 背压就是这样一路传回生产者的;// - 缺口自愈:头部序号迟迟不来(worker 崩溃吞了消息),超时后插入// 占位符继续前进——单个坏块不能扣住整篇文档。type assembler struct { cfg parserConfig pending map[int]string next int frags []string // 已定稿片段(下标 == Seq),订阅重放的数据源 provisional string total int // -1 = 输入尚未终结 maxSeen int subs []*subscription done bool lastProgress time.Time tickCancel actor.CancelFunc ticking bool}flowchart LR
feed["Feed(chunk) / Close()"] --> doc["DocActor(有状态)<br/>行缓冲 + 块级状态机"]
doc -->|"leafJob{seq,leaf}<br/>轮转派发"| w0["worker-0<br/>无状态"]
doc -->|"leafJob{seq,leaf}"| w1["worker-N<br/>无状态"]
w0 -->|"fragment{seq,html}"| asm["Assembler(有状态)<br/>按 seq 重排序"]
w1 -->|"fragment{seq,html}"| asm
doc -.->|"fragment{容器}<br/>provisionalUpdate<br/>endOfDoc"| asm
asm -->|"阻塞推送 = 背压源头"| sub["Subscription.C<br/>SSE / WebSocket / UI"]
mdparser 拓扑:DocActor → InlineWorker 池 → Assembler
这张图和前面 demo 里的会话三层树几乎是同一张图:DocActor 对应 SessionActor(有状态、按 key 隔离),InlineWorker 对应 GenWorker(无状态、易崩、可重建),Assembler 对应 SSE 边界(订阅、重放、背压)。同一个运行时、同一套模式,换一个领域照样成立——这正是验证学习成果想要的效果。
消息协议:密封接口 + 强类型 PID
三个 Actor 各自持有一族密封接口消息,和 demo 里的 SessionMsg 是同一个惯用法——发错类型编译期就报错,不必等到运行时:
type DocMsg interface{ isDocMsg() } // feedChunk / closeInput / docSubscribe / docSnapshottype InlineMsg interface{ isInlineMsg() } // leafJob{seq, leaf}type AsmMsg interface{ isAsmMsg() } // fragment / provisionalUpdate / endOfDoc / subscribeMsg / snapshotReq / gapTick密封接口只是第一层保险。第二层是强类型 PID——docActor 手里握着的不是一个泛泛的「某个 Actor 地址」,而是 asm *actor.PID[AsmMsg] 和 workers []*actor.PID[InlineMsg]:泛型参数把「这个地址只收 AsmMsg」焊死在类型里。想把一个 leafJob 错发给 Assembler、或者把一个 fragment 错发给 worker,代码根本编译不过——这条防线比密封接口生效得更早,连「构造出一条消息」这一步都不给你机会走错地址。demo 里的 SessionActor / GenWorker 用的是同一个模式,mdparser 只是把它原样搬到了新领域。
值得留意的是,所有载荷都是值类型——Leaf 只有 Kind、Level、Info、Content、Tight 五个字段,没有指针共享。「消息不可变」这条金律,在这里不是靠约定,而是靠类型系统兑现:worker 拿到的 leaf 是它自己的一份拷贝,怎么改都动不到 DocActor 那边的状态机。ast.go 里 BlockEvent 的注释还点出了一层更深的对照:pulldown-cmark 的 Event 是消费者主动拉出来的(pull),这里的 BlockEvent 是推给下游 Actor 的(push)——pull 与 push 的对偶,正是两种架构分野的起点,下一节会具体展开这份「push」到底要付出什么代价。
Feed / Close / Subscribe / Snapshot 各自落在哪条消息上
门面的四个方法,每一个都对应一条确定的消息路径,没有例外:
| 门面方法 | 发给 DocActor 的消息 | 后续流向 |
|---|---|---|
Feed(chunk) | feedChunk{chunk} | DocActor 就地驱动状态机,派发 leafJob / fragment |
Close() | closeInput{} | 状态机收口:补上最后的容器闭合,再发 endOfDoc{total} |
Subscribe(fromSeq) | docSubscribe{sub} | 原样转发成 subscribeMsg{sub} 给 Assembler |
Snapshot(timeout) | docSnapshot{reply}(经 actor.Ask) | 原样转发成 snapshotReq{reply} 给 Assembler |
Subscribe 与 Snapshot 都要先经过 DocActor 才能到 Assembler,而不是直连 Assembler 的地址——这一手多绕的设计不是偷懒,是刻意的:Feed、Subscribe、Snapshot 全部先落进 DocActor 的 mailbox,天然共享同一条 FIFO 队列。所以「Feed 之后立刻 Snapshot」一定能看到这次 Feed 产生的暂定视图——不需要轮询、不需要等待,顺序由 mailbox 的先进先出语义直接保证,这也正是 TestProvisionalRendering 不必轮询就能稳定通过的原因。
三个刻意的设计决策
拓扑图好画,难的是图背后要回答的三个问题:并行算出来的顺序怎么保证?worker 崩溃丢的消息怎么办?对外到底暴露几个方法?mdparser 对这三个问题给出的答案,都不是默认选项。
🔑 设计钥匙
三个刻意的设计决策撑起了这套拓扑:顺序在边界重建(Assembler 用
seq重排,「并行计算,串行呈现」)、接受至多一次投递的代价并用缺口自愈兜底(worker 崩溃丢的消息靠超时占位符续上,一个坏块不拖垮全文档)、以及单一入口门面(Feed/Close/Subscribe/Snapshot四个方法共享同一 mailbox 的 FIFO,把 push 的状态杂耍全部关进 Assembler 内部)。
决策一:顺序在边界重建。 序号从哪来?块级状态机每吐出一个 BlockEvent 就带着一个单调递增的 seq,不管这个事件最终是 EvLeaf(要过 worker)还是 EvOpen / EvClose(直接旁路),seq 都从同一个计数器上分配——这保证了两条路径产出的碎片仍然共享同一套坐标系,Assembler 完全不用关心某个 seq 到底是从 worker 绕了一圈回来的,还是 DocActor 自己直接发的。
// dispatch 把块级事件派发下游:叶子块轮转给 worker 池(并行行内解析),// 容器标签直接渲染发给 Assembler(无行内内容,不值得过一遍 worker)。func (d *docActor) dispatch(events []BlockEvent) { for _, ev := range events { if ev.Kind == EvLeaf { w := d.workers[d.rr%len(d.workers)] d.rr++ w.TellBlocking(leafJob{seq: ev.Seq, leaf: ev.Leaf}) } else { d.asm.TellBlocking(fragment{seq: ev.Seq, html: d.cfg.parse.ContainerHTML(ev)}) } }}dispatch 就是这个分岔口:EvLeaf 轮转派给 worker 池,其余的调用 d.cfg.parse.ContainerHTML(ev) 同步渲染后直接发给 Assembler。到了 Assembler 这一侧,重排逻辑只有一个 pending map[int]string 加一个 next int:
func (a *assembler) advance() { for { html, ok := a.pending[a.next] if !ok { break } delete(a.pending, a.next) a.frags = append(a.frags, html) a.broadcast(Event{Kind: EventFragment, Seq: a.next, HTML: html}) a.next++ a.lastProgress = time.Now() } a.maybeFinish()}只有 next 号片段就位,才定稿推送、next 才前进;后到的碎片先趴在 pending 里,等前面的坑填上才轮到它。并行是实现细节,顺序是对外承诺——这笔「顺序税」被明明白白记在 advance() 里,不到二十行,没有藏进锁或者原子操作背后。
决策二:至多一次投递的代价,用缺口自愈兜底。 worker 是无状态执行体,监督策略是 RestartStrategy(-1)——不限重启次数,崩了立刻换新实例顶上。但重启解决不了一件事:worker panic 那一刻正在处理的那条 leafJob 随崩溃实例一起消失了,这是运行时的 at-most-once 语义,没有重投递这回事。换上的新 worker 只能收到它启动之后发来的消息,seq=N 那个 fragment 永远不会来了——如果 Assembler 什么都不做,它会在这个缺口上永远等下去:advance() 里的循环第一次查 pending[N] 就会 miss、break,后面所有已经就位的碎片全都被堵在原地,一起陪葬。
对策是一条心跳消息 gapTick{},由 ensureTick 懒启动的 ctx.ScheduleRepeated 定时投递:
func (a *assembler) healGap(ctx *actor.Context[AsmMsg]) { if a.done { return } headMissing := a.maxSeen >= a.next || (a.total >= 0 && a.total > a.next) if !headMissing || time.Since(a.lastProgress) < a.cfg.gapTimeout { return } placeholder := `<p class="md-error">[block ` + strconv.Itoa(a.next) + ` unavailable]</p>` + "\n" a.pending[a.next] = placeholder a.advance()}headMissing 判断「后面确实还有东西没到」(maxSeen 已经超过 next,或者文档已经收口但总数比 next 大);lastProgress 判断「卡住多久了」。两个条件同时成立、且超过 gapTimeout,才伪造一个占位片段塞进 pending[a.next],把 advance() 硬顶开。一个毒块只损失它自己那一段 HTML,文档照常收尾;换成传统的 pull 解析器,一个 panic 掀翻的是整篇文档的调用栈。教学用的 WithPanicOn 选项可以在指定内容里现场触发这条路径,TestCrashIsolationAndGapHealing 验证的正是「毒块两侧的段落完好、占位符如期出现、监督树的重启计数大于零」。
💡 小贴士
把
WithGapTimeout调到几十毫秒,再配合WithPanicOn,单元测试里就能把「worker 崩溃 → 出现缺口 → 占位符愈合」整条链路几乎零延迟地复现一遍,不必真的等满两秒的默认超时。
决策三:单一入口门面,把 push 关进系统内部。 还记得前几章提到的那条警告吗——pulldown-cmark 的 README 明确点破了 push 式接口的病根:靠一连串事件回调驱动解析,调用方就得在回调之间手忙脚乱地拼出「当前处于哪个块」这类状态,一步错了就全盘乱。Actor 消息骨子里就是这种 push:上游把消息扔过来,下游只能被动接住、被迫维护自己的状态机。mdparser 的应对不是绕开 push——拓扑内部从头到尾都是 push——而是把这份复杂度整个焊死在系统内部,只留一张干净的对外脸:
p := mdparser.New(sys, "doc-1")p.Feed("# 标题\n正在打") // 任意切分追加,无需对齐行/块边界p.Close() // 输入终结sub := p.Subscribe(0) // 有序事件流:Fragment(seq) / Provisional / Donesnap, _ := p.Snapshot(time.Second) // Ask:一致视图 = Stable + Provisional// Parser 是流式 Markdown 解析器的对外门面。一个 Parser 对应一篇文档// (一个 DocActor 子树),Feed 追加输入、Close 终结输入。type Parser struct { doc *actor.PID[DocMsg] cfg parserConfig}调用方看到的只有四个方法。内部三族消息、一个重排序器、一套缺口自愈定时器,全都不需要它操心。Feed 把任意切分的输入喂给 DocActor,阻塞语义保证下游拥塞时是「变慢」而不是「丢数据」;Subscribe 与 Snapshot 也都先发给 DocActor 再转发 Assembler,与 Feed 共享同一 mailbox 的 FIFO 顺序。
// Feed 追加一段输入(任意切分,无需对齐行/块边界——LLM 的 token 流就是这样)。// 阻塞语义:下游拥塞时 Feed 变慢而不是丢数据。返回 false 表示解析器已停止。func (p *Parser) Feed(chunk string) bool { return p.doc.TellBlocking(feedChunk{chunk: chunk}) }延伸阅读
- pulldown-cmark — GitHub README:作者说明了选择 pull 而非 push 作为对外接口的理由——push 式解析接口出了名地难用、易出错,因为使用者要在一连串回调之间小心维护状态。
mdparser走的是相反的路径:内部彻头彻尾是 push(Actor 消息),但用本章「单一入口门面」把这份复杂度关起来,不泄漏给调用方——算是对同一个问题的另一种回答,而不是回避这个问题。
小结
mdparser的拓扑就是三个 Actor:DocActor(串行心脏)派发行内工作给一组 InlineWorker(并行、易崩),它们的产出在 Assembler 里重排、推送。- 三族密封接口消息 + 值类型载荷 + 强类型 PID,把「消息不可变」「发错地址编译不过」从约定升级为编译期保证;
Feed/Close/Subscribe/Snapshot各自对应一条确定的消息路径,且共享同一 mailbox 的 FIFO 顺序。 - 三个刻意的设计决策——顺序在边界重建、至多一次投递配缺口自愈、单一入口门面——是这套拓扑能稳定跑起来的真正原因,而不是拓扑图本身。
- 下一章把这四大机制逐个拆开验收:增量解析、并行重排、let-it-crash 自愈、全链路背压,并配一个可拖拽的数据流可视化,让这套拓扑在你眼前真正流动起来。