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