目录 · 第 13 / 16 章
ActorPart IV · 流式 Markdown 解析器

13四大机制:增量 · 重排 · 自愈 · 背压

冻结已定稿只重渲染活动块、按 seq 重排、let-it-crash + 缺口自愈、全链路背压。

mdparser/block.go:BlockStatemdparser/actors.go:assemblermdparser/actors.go:advance

上一章画完了拓扑图:一颗串行心脏(DocActor)接住 LLM 吐出的 chunk,一池无状态 worker 并行算行内,一个 Assembler 在出口重排序并推送。拓扑只是管道的形状——这一章要证明这根管道值这个价钱。从管子里流出来的是四条具体机制:冻结已定稿只重渲染活动块、并行算但串行呈、let-it-crash 把崩溃损失降到一个块、以及一路传导回生产者的背压。逐条对照代码看下来你会发现同一件事:没有一条是我们为 Markdown 发明的新招,它们全部是 Actor 运行时本来就有的能力,只是被投影到了「解析」这个具体领域上。

🔑 设计钥匙

把这四条机制拆开看,每一条都能在 Actor 运行时的通用能力表里找到原型:冻结增量对应「私有状态只对自己可见,从不共享」;并行重排对应「消息传递不保证到达顺序,顺序要在某个边界被重新收拢」;let-it-crash 对应「监督树把故障半径缩到一个可重建的子树」;背压对应「有界 mailbox 是涌现性质,不是外挂上去的限流器」。这些不是我们为 Markdown 发明的新招——它们是运行时本来就有的通用能力,这一章只是逐条对照 mdparser 的代码,把它们在「解析」这个具体领域里指认出来。

增量:冻结已定稿,只重渲染活动块

块级状态机的全部状态小到能画在纸上——容器栈,加上至多一个打开的叶子:

// BlockState 是块级状态机的可变状态。规则通过它读写容器栈、产出事件。
type BlockState struct {
cfg *Config
stack []container
leaf *openLeaf // 指向 leafStore(打开时)或 nil
leafStore openLeaf // 唯一的叶子存储:任一时刻至多一个叶子打开,故整篇复用一个,不每块新分配
seq int
events []BlockEvent
}

feedLine 每来一行分三步走:第一步 matchPrefix 自底向上匹配既有容器栈前缀,一旦某层匹配失败,栈顶到失败层之间的所有容器立即闭合;第二步反复尝试打开新容器(`> - text` 这类一行内嵌套多层的写法就是在这一步的循环里建立的);第三步把剩下的文本交给叶子规则归类,没有规则接住就落进段落规则兜底。三步全部只触碰栈顶,已经闭合的块永不回头重解析——这正是 CommonMark 规范附录 A 说的「处理完这一行就可以扔掉」。增量成本因此只正比于新增的这一行,不正比于文档已经写了多长。

流式输出的形状由此自然长成已定稿前缀 + 暂定尾巴:定稿片段按序号拼接,只增不改,这个 append-only 的形状与 streaming-markdown(自报数据,未独立核实)的乐观提交策略同形——它见到开分隔符就先落笔,不等闭合符出现。打开的容器链与打开的叶子则相反:每来一个 chunk 都要重渲染一次「暂定视图」,靠的正是这个方法——

// provisional 渲染「当前打开但未闭合」的部分:容器开标签 + 叶子临时渲染 + 容器闭标签。
// 这是 AI 流式场景的核心能力:给未完结的尾巴一个暂定视图。
func (s *BlockState) provisional(w Writer) {
if len(s.stack) == 0 && s.leaf == nil {
return
}
for i := range s.stack {
c := s.stack[i]
s.cfg.writeContainer(w, BlockEvent{Kind: EvOpen, Container: c.kind, Ordered: c.ordered, Start: c.start})
}
if s.leaf != nil {
s.cfg.writeLeaf(w, s.buildLeaf())
}
for i := len(s.stack) - 1; i >= 0; i-- {
c := s.stack[i]
s.cfg.writeContainer(w, BlockEvent{Kind: EvClose, Container: c.kind, Ordered: c.ordered, Start: c.start})
}
}

它做的事情很朴素:把容器栈从外到内挨个写开标签,把当前打开的叶子渲染一遍,再把容器栈从内到外挨个写闭标签——相当于「假装这一刻整个文档就在这里结束」渲染一次。**加粗 没写完时按字面输出,闭合的瞬间升级成 <strong>;这个渲染每个 chunk 都要重新算一遍,但代价很小,因为暂定尾巴通常只有几百字节。连没写完的半行都参与这次渲染——docActor.provisional() 把跨 chunk 的行缓冲拷贝进一份状态机的影子副本,追加进当前叶子后再渲染,不需要真的修改状态机本身。这就是 token 级即时反馈的来源,也是「对 AI 流式场景极度友好」这句话的真正兑现点——phoenix_streamdown(自报数据,未独立核实)在 DOM 层做的是同一件事:冻结已完成块,只重渲染最后一个活动块;我们把它挪到了解析器层,少的不是 DOM 更新次数,而是重解析本身。

对照朴素做法——每来一个 chunk 把累计的全部文本重新丢进 parse() 一遍——Chrome 官方的 LLM 渲染指南(自报数据,未独立核实)把这种写法明确列为反面教材:文档越长,每一次追加就越贵。增量状态机把这条曲线从 O(n²) 砍成了 O(n)——每个 chunk 的成本只正比于它自己带来的新内容,与文档已有多长无关。

重排:并行算,串行呈

CommonMark 规范原话:第一阶段必须逐行顺序处理,但第二阶段「可以并行化,因为一个块的行内解析不影响任何其他块的行内解析」。DocActor 把每个闭合的叶子块连同一个单调递增的文档序号 seq 打包成 leafJob{seq, leaf},按 d.rr % len(d.workers) 轮转派发给 N 个无状态 worker——这就是把规范写明的并行缝隙真正用起来。

📝 注意

诚实的限定:单篇文档的行内解析通常算不上瓶颈——批量基线是 52 µs/4.4KB,轮不到并行出场。并行池真正的意义在多会话服务器形态:一个系统同时承载几百篇文档流式生成时,worker 池是天然跨文档共享的计算资源,而每篇文档自己的 DocActor 依然保持串行、互不加锁——这和 Flink 按 key 分区、组内串行组间并行是同一个形状。

并行带来一个必然的副作用:worker 算完的顺序不等于文档顺序,谁先算完谁先到 Assembler。顺序因此要在出口边界重建——Assembler 维护 pending map[int]stringnext int,新到的片段先按 seq 存进 pending;只有 next 号片段就位,才把它定稿、推送、next 前进一位,如此循环直到缺口再次出现。这是一个经典的 resequencer:并行只活在实现细节里,一个整数 next 就是对外承诺的全部——seq 是承诺,并行只是细节。

// 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
}
// advance 把连续就绪的片段依次定稿并推送。
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()
}

「并行计算,串行呈现」不是 Markdown 解析独有的技巧,而是任何消息系统一旦引入并行都要面对的账:算力可以在多个无状态 Actor 间自由分摊,但顺序必须在某个边界被重新收拢。这里选的边界是 Assembler 的 advance()——不到二十行代码,把并行留在实现细节里,顺序作为对外承诺原封不动地兑现。

自愈:let-it-crash 与断线重连

worker 是无状态、可随时重建的执行体:一个工厂函数(PropsFromProducer)配 RestartStrategy(-1)(不限重启次数)。它 panic 时,运行时 recover、监督者换一个新实例顶上——但块级状态机的进度从没住在 worker 里,它一直安然待在 DocActor 中,一次 worker 崩溃碰不到解析状态的一根汗毛。

真正的问题在别处:panic 吞掉的那条消息永远不会重新出现——这是运行时「至多一次投递」的诚实代价,不是「恰好一次」。seq=N 的片段可能就此消失,Assembler 若什么都不做,会在头部缺口上永远等下去。

对策是缺口自愈:一个定时的 gapTick(懒启动,间隔取 gapTimeout/2)检查头部缺口滞留了多久——maxSeen >= next 或者 total > next 说明确实缺了一块,time.Since(lastProgress) 一旦超过阈值,就插入一个占位符继续前进,而不是死等一条再也不会到来的消息:

// healGap 缺口自愈:头部片段缺失超时(worker 崩溃吞掉了在途消息),
// 插入占位符跳过。这是「至多一次投递 + let-it-crash」组合拳的最后一环。
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
}
ctx.Logger().Warn("gap healed with placeholder", "seq", a.next)
placeholder := `<p class="md-error">[block ` + strconv.Itoa(a.next) + ` unavailable]</p>` + "\n"
a.pending[a.next] = placeholder
a.advance()
}

一个毒块只损失它自己——两侧的段落完好,文档照常读完。对照批量解析器:一个 panicunwrap 失败等于整次 parse() 报废,调用方只能选全文档重试或彻底放弃。测试里用 WithPanicOn 故障注入验证这条闭环:命中触发条件的内容必定让处理它的 worker panic,占位符必定出现在正确的位置,监督树的重启计数必定大于零。

断线重连用的是同一套家底,几乎不需要新代码:Assembler 手里已定稿的 frags 切片本身就是一份「活的检查点」,Subscribe(fromSeq) 能从任意序号重放已定稿片段加最新暂定视图。SSE 断线重连就是带着 Last-Event-ID 调一次 Subscribe(lastSeq+1)——零外部存储,状态从未离开进程。

💡 小贴士

断线重连没有为它单独写一段「恢复逻辑」——它复用的是 Assembler 本来就在维护的同一份状态。这是个可以推广的经验:一旦一个系统的状态天然长成「只增不改的日志 + 一个游标」,checkpoint 与 resume 几乎是免费拿到的,不必额外设计一套快照协议。

背压:全链路,慢消费者拖不垮

推送与派发全部走阻塞式投递(宁慢勿丢)加有界 mailbox:

// postBlocking 阻塞投递:等到有空位则入队返回 true;mailbox 关闭则返回 false。
func (m *mailbox[T]) postBlocking(e envelope[T]) bool {
select {
case m.user <- e:
return true
case <-m.done:
return false
}
}

它的实现只是一个 select:能塞进 channel 就塞进去返回 true,mailbox 已关闭就返回 false——没有轮询,没有额外的排队或限流代码,阻塞本身就是全部机制。这个原语在 Assembler 的推送边界上再现了一次:

// send 阻塞推送(带 done 逃生阀)——背压源头与扇出版的 Assembler.send 相同。
func (d *fusedDocActor) send(sub *subscription, ev Event) bool {
select {
case sub.ch <- ev:
return true
case <-sub.done:
return false
}
}

慢消费者的减速由此一路传回生产者:

订阅者读得慢 → Assembler 阻塞在推送 → 它的 mailbox 堆满
→ worker 的阻塞投递也开始等 → worker mailbox 堆满
→ DocActor 派发阻塞 → DocActor mailbox 堆满 → Feed() 阻塞
flowchart RL
  sub["订阅者读得慢"] -->|"send() 阻塞"| asm["Assembler mailbox 堆满"]
  asm -->|"TellBlocking 阻塞"| worker["worker mailbox 堆满"]
  worker -->|"TellBlocking 阻塞"| doc["DocActor mailbox 堆满"]
  doc -->|"TellBlocking 阻塞"| feed["Feed(chunk) 阻塞"]

背压如何一路传回 Feed()

没有丢消息,没有 OOM,也没有为此写过一行专门代码——背压是有界 mailbox 的涌现性质,不是额外挂上去的机制。内存占用有硬上界(各级 mailbox 容量之和),因为拓扑是无环的(doc → worker → assembler → 订阅者),这条阻塞链不会绕回来死锁自己。TestBackpressureSlowSubscriber 把订阅缓冲压到 1、每个事件人为拖 1ms,跑一篇 30 块的文档:零丢失、严格有序——这不是运气,是有界系统本该有的样子。

延伸阅读

以下三条均为对方项目自报数据,未经本书独立核实,仅作同类思路的参照:

  • streaming-markdown——append-only 输出加乐观提交解析,与本章「已定稿前缀 + 暂定尾巴」是同一个形状。
  • Chrome:Render LLM responses——官方指南点名「每 chunk 全量重解析」的朴素做法,论证动机与本章 O(n²) 对比 O(n) 一致。
  • phoenix_streamdown——冻结已完成块、只重渲染活动块,与本章「暂定视图」同思路,区别是他们做在 DOM 层,我们做在解析器层。

小结

  • 四条机制——冻结增量、并行重排、let-it-crash 自愈、全链路背压——没有一条是为 Markdown 定制的新发明,全部来自 Actor 运行时本身的通用能力:私有状态的串行心脏、可监督重建的无状态池、有界 mailbox。解析器这层壳很薄,薄到几乎只是规则表 + 状态机。
  • 但这层壳目前还很粗糙:渲染靠字符串拼接手写转义,解析规则焊死在一个 switch 里——这些短板妨碍它真正被用于生产,与四条机制的架构价值无关,是纯粹的工程债。
  • 下一章把这层壳打磨成生产可用:渲染改成 Renderer 接口写 io.Writer(基线快 3.5×),解析改成可注册的规则表(ruler),再用一个删除线扩展证明——加语法、改输出,都不必再动核心代码。
源码

正在读取完整文件…