Skip to content

StreamMux:多模式流的事件分发器

源码版本1.2.9

职责

在 LangGraph v3 流式协议里,用户调 graph.stream_events(version="v3") 拿到的不再是一串 raw (namespace, mode, payload) 元组,而是一个 GraphRunStream 对象,上面挂着 run.values / run.messages / run.updates / run.subgraphs / run.lifecycle 这些「投影 (projection)」(GraphRunStream:31-49)。把这些 raw 事件按模式拆到不同投影上,再让用户用普通 for 循环拉动的,就是 StreamMux——事件分发中心。代码在 libs/langgraph/langgraph/stream/_mux.py (StreamMux:26)。

StreamMux 本身不实现任何具体模式的逻辑,它只做三件事:注册 transformer、push 事件、关闭。push 时按注册顺序依次调每个 StreamTransformer.process(event) (push:269-296),transformer 决定要不要把事件塞进自己的 projection channel;然后 mux 自己再把事件追加到 _events 这个主日志,主日志是 raw 流的「完整审计副本」。transformer 们开的 projection 都是 StreamChannel 实例,带名字的 channel 还会被 mux 自动转发到主日志,所以 transformer 的副作用既体现在自己的 channel 上,也保留在主日志里。

StreamMux 还是子图嵌套的基础。每个 subgraph 被发现时,SubgraphTransformermux._make_child(scope) (_make_child:193-225) 给子图建一个 mini-mux,这个 mini-mux 跟父 mux 共享同一组 factory(所以 transformer 实例是每个 scope 一份),也共享父 mux 的 pump 绑定——根 pump 一动,所有子图的事件都往前走一格。

设计动机

  • caller-driven pump,不要后台线程:GraphRunStream 的迭代器就是 pump (_pump_next:107-130)。用户 for event in run.values 时,每次迭代触发 _pump_next 拉一个 raw 事件喂给 mux;后台不开线程,内存由消费速率决定。bind_pump 把这个 pump 函数登记到 mux (bind_pump:162-180),然后 mini-mux 子图通过 _make_child 自动继承。
  • transformer 按注册顺序串行:pushfor transformer in self._transformers: if not transformer.process(event): keep = False (push loop:288-292),前一个 transformer 的副作用先发生,后一个能看见;某个 transformer 返回 False 只是抑制事件进主日志,前面的 transformer 还是吃到了。
  • native vs extension 投影区分:_native = True 的 transformer 投影会被 mux 当成 run.values 这种直接属性挂到 GraphRunStream 上 (native attrs:78-79);非 native 的只进 run.extensionsnative_keys 集合 (native_keys:109) 是 v3 协议给 SDK 暴露的「直接属性」白名单。
  • before_builtins 给内容改写 transformer 让位:像 PII 过滤、内容审查这类 transformer 必须在 MessagesTransformer 等内置 transformer 之前看到原始 text 字段 (before_builtins:94-109);mux 注册时按 before_builtins = True 分两 lane 重组 (partition:131-148),但 lane 内顺序保留。
  • seq 只在 root mux 分配:_assign_seq 默认 True,child mux 通过 _make_childFalse (_assign_seq:215-220)。这样子图事件直接转发到根日志时不被改 envelope,保持 seq 单调对应到 root 的写入顺序。

关键文件

  • StreamMux 类:26-46 — 中心调度器属性和 docstring。
  • StreamMux.__init__:48-148 — 接 transformers 列表或 factories,按 before_builtins 分 lane 注册。
  • _make_child:193-225 — 子图发现时建 mini-mux,继承 pump 绑定。
  • _register:227-267 — 调 transformer.init() 拿 projection dict,检查 key 冲突,绑定带名字的 StreamChannel。
  • push:269-296 — 同步路径:按顺序调 process,再决定是否进主日志。
  • apush:351-378 — 异步路径:aprocess 串行 await,保证后面的 transformer 看到前面的 async 结果。
  • close:298-324 — 同步关闭,调每个 transformer 的 finalize,出错也保证全部 channel 关掉。
  • aclose:380-404 — 异步关闭,先 gather 所有 schedule 出来的 task,再 afinalize
  • _collect_stream_modes:398-414 — 从 mux 已注册 transformer 的 required_stream_modes 反推要给图请求哪些模式。
  • _pregel_stream_v3:3533-3558 — v3 同步入口,实例化 StreamMux + GraphRunStream,把 mux 的 stream_mode 反推给底层 stream()

数据流

StreamMux.push 的核心就这么几行——按注册顺序过一遍 transformer,全过完还活着就塞主日志:

python
def push(self, event: ProtocolEvent) -> None:
    keep = True
    for transformer in self._transformers:
        if not transformer.process(event):
            keep = False
    if keep:
        if self._assign_seq:
            self._seq += 1
            event["seq"] = self._seq
        self._events.push(event)

(push body:288-296) 看似朴素,但 transformer.process 里可能往自己的 StreamChannel push,带名字的 channel 还会被 mux 自动 wire 进主日志(_register:227-267),所以一个 raw values 事件可能同时被 ValuesTransformer 装进 run.values,又被某个用户自定义 transformer 转换后塞回主日志。事件在 transformer 之间走串行链路,顺序严格保持。

边界与失败

  • async transformer 在 sync 模式下立刻 raise:_registertransformer_requires_async 检查 (requires_async check:234-239),如果 transformer 重写了 aprocess 但没设 supports_sync=True,注册到 sync mux 会报错,提示用 astream() 而不是 stream()
  • projection key 冲突直接 ValueError:_register 检查 set(projection) & set(self.extensions) (conflict check:246-256),报错时列出冲突的 key 和归属的 transformer 名字;自定义 transformer 取 projection key 不能跟内置重名。
  • _make_child 拒绝 pre-built transformers:mux 用 transformers= 参数构造时不能建子 mux (factory required:209-214),因为 pre-built 实例无法在新的 scope 下重新初始化。SubgraphTransformer 只能跑在用 factories= 构造的 mux 上。
  • finalize 报错不阻塞其他清理:closefirst_error 收第一个异常,继续关剩下的 transformer 和 channel,最后再 raise (close error handling:312-324),保证资源不泄漏。
  • bind_pump 必须在 channel 读取前调:GraphRunStream.__init__wire_pump=True 默认就走 _wire_request_more (_wire_request_more:80-92),子 mux 走 _make_child 时父 binding 自动传下去;但用户自己手搓 mux 实例时,忘 bind pump 会让 StreamChannel.__iter___request_more 永远是 None(_request_more:191-192),迭代就直接结束。
  • schedule(coro) 只能在 async 模式用:StreamTransformer.schedule 需要事件循环 (schedule:233-263),mux 在 aclosegather 所有 schedule 的 task,async 跑完才 afinalize;sync 模式调 schedule 会失败。

小结

StreamMux 是 v3 流式协议的枢纽:它不实现任何模式逻辑,只按顺序调度 transformer、按 scope 建子 mux、把 raw 事件和 transformer 投影都收进主日志。具体每个 mode 怎么把事件转成用户友好的形状见 stream/transformers;raw 事件怎么从 PregelLoop 喂进来见 Pregel 引擎

对照官方资料:LangGraph 文档 · README