Skip to content

Stream Transformers:把 raw Pregel 輸出變成可用投影

源码版本1.2.9

職責

PregelLoop 在每個超步 (superstep) 裡透過 _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