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 被發現時,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 按註冊順序串行: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:_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 報錯不阻塞其他清理: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