Stream Transformers:把 raw Pregel 输出变成可用投影
职责
PregelLoop 在每个超步里通过 _emit 把 (checkpoint_ns, mode, payload) 三元组推进一个 SyncQueue (_emit:1380-1414),里面是 raw 的「值」「更新」「任务」「检查点」这些事件。这是 v1 / v2 协议的输出形态——用户拿到的就是一串 chunk。但 v3 协议里,用户希望拿到的不是 chunk,而是「按模式组织的、可以单消费者迭代的投影 (projection)」:run.values 是一份完整状态快照流,run.messages 是一组 LLM token 流的 ChatModelStream 句柄,run.updates 是节点更新流,run.subgraphs 是子图入口句柄。
把 raw chunk 转成这些投影的逻辑就放在 libs/langgraph/langgraph/stream/transformers.py 里 (ValuesTransformer:28),每个模式一个 transformer,共同继承 StreamTransformer 抽象基类 (StreamTransformer:44)。它们被 StreamMux 注册、被 GraphRunStream 推送事件,各自维护内部 StreamChannel 作为投影。用户调 graph.stream_events(version="v3") 时,实际上拿到的是 GraphRunStream (GraphRunStream:31),它的 .values / .messages / .updates / .subgraphs / .custom / .lifecycle 这些属性就是 native transformer 的投影入口 (native attrs:78-79)。
整套 transformer 体系是 v3 协议的扩展点——除了内置这五个,用户可以在 Pregel.compile(stream_transformers=[...]) 加自己的 transformer (stream_transformers param:783),在 _pregel_stream_v3 里跟内置 transformer 一起以 factory 形式传给 StreamMux (factories list:3533-3544)。
设计动机
- native vs extension 投影双轨:
_native = True的 transformer 投影挂在run.<key>直接属性上 (_native:50-52);非 native 进run.extensions,只暴露在字典里。内置五个都是 native,用户自定义默认进 extensions,要直接属性就显式设_native = True。 required_stream_modes反推图要发哪些模式:transformer 自己声明要哪些 mode (required_stream_modes:88-93),_collect_stream_modes取并集 (_collect_stream_modes:398-414),v3 调用底层stream()时把这个并集当stream_mode传进去 (stream_mode=_collect_stream_modes:3549)。用户不直接指定stream_mode,而是「按需拉取」。- scope 过滤避免噪音:
ValuesTransformer.process检查params["namespace"] != self._scope_list(ValuesTransformer.process:70-82),只接受自己 scope 的 events;子图的事件留给子图自己的 mini-mux 处理,根run.values不会被嵌套图污染。 - sync / async 双轨 + 自动检测:transformer 默认
supports_sync = False,只要重写了aprocess/afinalize/afail就被transformer_requires_async判为 async-only (transformer_requires_async:308-330),sync mux 注册时直接报错。要 sync / async 都支持得显式supports_sync = True,比如SubgraphTransformer(supports_sync:689)。 SubgraphTransformer走SubgraphRunStream而非裸事件:发现一个子图调用时,_on_started建SubgraphRunStream包装 mini-mux (SubgraphTransformer:670-705),用户拿到的不是事件流而是句柄对象,可以在句柄上递归handle.values/handle.messages/handle.subgraphs,跟根run是同构接口。
关键文件
StreamTransformer 基类:44-115— 抽象init/process/aprocess/finalize/fail+ ClassVar 配置。init:128-138— 返回 projection dict,key 进extensions,_native=True时还挂run.<key>。process:148-164— sync lane,默认raise NotImplementedError,subclass 必须重写。aprocess:166-184— async lane,默认委托process;需要 async 工作重写。schedule:233-263— 跟 mux 生命周期绑定的asyncio.Task,aclose 时 gather,afail 时 cancel。transformer_requires_async:308-330— 检测aprocess/afinalize/afail是否重写,决定 transformer 是否能在 sync mux 跑。ValuesTransformer:28-82—run.values投影,required_stream_modes=("values",)。UpdatesTransformer:120-152—run.updates投影。MessagesTransformer:155-335—run.messages投影,把 LLM token 流打包成ChatModelStream句柄,处理 v2 协议的(payload, metadata)输入。LifecycleTransformer:608-667—run.lifecycle投影,把子图生命周期事件以 native protocol event 形式发回主日志,远程 SDK 可见。SubgraphTransformer:670-705—run.subgraphs投影,发现子图调用时建SubgraphRunStream+ mini-mux。GraphRunStream._pump_next:107-130— pump:拉一个 chunk,convert_to_protocol_event,丢给 mux。convert_to_protocol_event:10-32— v2StreamPart→ v3ProtocolEvent的字段映射。
数据流
整个 v3 链路从 PregelLoop._emit 开始,经过 SyncQueue、GraphRunStream._pump_next、convert_to_protocol_event,最后到 mux。convert_to_protocol_event 这一段把 v2 的 {type, ns, data, interrupts} 直接转成 v3 ProtocolEvent:
def convert_to_protocol_event(part: StreamPart) -> ProtocolEvent:
part_dict = cast(dict[str, Any], part)
params: _ProtocolEventParams = {
"namespace": list(part_dict["ns"]),
"timestamp": int(time.time() * 1000),
"data": part_dict["data"],
}
if "interrupts" in part_dict:
params["interrupts"] = part_dict["interrupts"]
return {
"type": "event",
"method": part_dict["type"],
"params": params,
}(convert body:20-32) 注意 method 就是 v2 的 type——values / updates / messages / custom / tasks / checkpoints / debug,这是 StreamTransformer.process 里 dispatch 的依据,比如 ValuesTransformer 只对 method == "values" 的事件动手 (method check:70-71)。
而 v1 / v2 协议在 Pregel.stream 里走的是另一条路径,raw chunk 直接从 SyncQueue 出来过 _output 函数 (_output:4184-4243),按 stream_mode 是字符串还是列表决定输出 payload 还是 (mode, payload) 还是 (ns, mode, payload):
if version == "v2":
...
yield {"type": mode, "ns": ns, "data": payload, "interrupts": ints}
elif stream_subgraphs and isinstance(stream_mode, list):
yield (ns, mode, payload)
elif isinstance(stream_mode, list):
yield (mode, payload)
elif stream_subgraphs:
yield (ns, payload)
else:
yield payload(_output shape branches:4220-4243) 这就是为什么 stream_mode=["updates", "values"] 会拿到三元组而单字符串会拿到单值——同一份 raw chunk 在 _output 这里被分形状。v3 不走 _output,而是走 GraphRunStream._pump_next → mux.push → transformer。
边界与失败
process必须实现一个,init必须 return dict:init默认是@abstractmethod,子类不实现会启动失败 (init abstract:127-138);process默认raise NotImplementedError(process default:162-164),不重写会在第一个事件 push 时炸。- transformer 之间顺序敏感:
before_builtins = True的 transformer 先跑 (before_builtins:94-109),但 docstring 警告它能看到tasks事件在LifecycleTransformer/SubgraphTransformer消费前;改event["params"]["namespace"]或id/result/error/interrupts字段会让后两个的记账错乱,只能观察、改冷字段。 MessagesTransformer跟 v2 协议耦合:stream_events(version="v3")在调用底层stream()时强制传version="v2"和CONFIG_KEY_STREAM_MESSAGES_V2=True(_pregel_stream_v3 stream call:3545-3551),否则 messages 投影拿不到正确 payload 格式。StreamChannel单消费者,迭代两次报错:run.values只能被一个forloop 消费 (single-consumer note:38-40);要 fan-out 得projection.tee(n)(tee:245),开销是缓冲。schedule(on_error="raise")会把 close 路径变 fail:schedule的on_error="raise"让 task 异常在aclose时把 mux 推到afail(on_error:254-257),默认"log"只记日志不传播——选错会让单个 transformer 失败拖垮整条流。SubgraphTransformer._should_track严格大于 scope:只追踪当前 scope 下一层的子图 (should_track:634-636),孙图通过子句柄的handle.subgraphs递归发现;如果用户跳过handle.subgraphs直接看根run.subgraphs,孙图事件不会出现在那里。
小结
StreamTransformer 是 v3 流式协议的扩展点:五个内置 transformer 把 raw (ns, mode, payload) 转成可迭代的 typed 投影,用户也可以自己加 transformer 进 stream_transformers 参数。整个分发机制见 StreamMux;raw chunk 怎么从 PregelLoop 喂进 SyncQueue 见 PregelLoop。
对照官方资料:LangGraph 文档 · README。