Skip to content

Stream Transformers:把 raw Pregel 输出变成可用投影

源码版本1.2.9

职责

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)。
  • SubgraphTransformerSubgraphRunStream 而非裸事件:发现一个子图调用时,_on_startedSubgraphRunStream 包装 mini-mux (SubgraphTransformer:670-705),用户拿到的不是事件流而是句柄对象,可以在句柄上递归 handle.values / handle.messages / handle.subgraphs,跟根 run 是同构接口。

关键文件

数据流

整个 v3 链路从 PregelLoop._emit 开始,经过 SyncQueue、GraphRunStream._pump_next、convert_to_protocol_event,最后到 mux。convert_to_protocol_event 这一段把 v2 的 {type, ns, data, interrupts} 直接转成 v3 ProtocolEvent:

python
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):

python
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_nextmux.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 只能被一个 for loop 消费 (single-consumer note:38-40);要 fan-out 得 projection.tee(n) (tee:245),开销是缓冲。
  • schedule(on_error="raise") 会把 close 路径变 fail:scheduleon_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