StreamMux:多模式流的事件分发器
职责
在 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 被发现时,SubgraphTransformer 调 mux._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 按注册顺序串行:
push用for 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.extensions。native_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_child传False(_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,全过完还活着就塞主日志:
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:
_register调transformer_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 报错不阻塞其他清理:
close用first_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 在aclose时gather所有 schedule 的 task,async 跑完才afinalize;sync 模式调 schedule 会失败。
小结
StreamMux 是 v3 流式协议的枢纽:它不实现任何模式逻辑,只按顺序调度 transformer、按 scope 建子 mux、把 raw 事件和 transformer 投影都收进主日志。具体每个 mode 怎么把事件转成用户友好的形状见 stream/transformers;raw 事件怎么从 PregelLoop 喂进来见 Pregel 引擎。
对照官方资料:LangGraph 文档 · README。