串流輸出:astream 與 stream_mode
職責
Pregel.astream / Pregel.stream 是編譯圖對外暴露的執行入口。你把輸入丟進來,它啟動 PregelLoop 一圈圈跑,過程中把每一超步 (superstep) 產生的資料按呼叫方指定的 stream_mode 投影成流事件,塞進一個佇列,再由呼叫方迭代消費。ainvoke / invoke 本質上就是 astream / stream 收斂成單值——Pregel.invoke 的實作體就是 for chunk in self.stream(...) 取最後一片(Pregel.invoke body:3891)。
這一層的責任是把「執行」和「輸出形狀」解耦:執行還是 Pregel 那套 BSP 迴圈,輸出形狀由 stream_mode 決定——同一個圖跑一次,你可以要 values(每步完整狀態)、updates(每步增量)、messages(LLM token 流)、custom(節點內自訂輸出)、checkpoints(檢查點事件)、tasks(任務起止)、debug(全部)等不同投影,甚至傳一個 list 同時拿多種模式。
設計動機
- 不阻塞執行:輸出走佇列,
runner.tick/atick每完成一個任務就往裡塞事件,外層for迴圈邊迭代邊拿,而不是跑完整張圖再一次性回傳。這對 LLM token 流(messages)是剛需——使用者得邊生成邊看。 - 多投影複用同一份執行:同一個
runner.atick(...)內部產生多種事件,_output函式(_output:4184)只做過濾/分發的動作,不重新跑圖。傳stream_mode=["values","updates"]不會跑兩次。 - 子圖命名空間統一:開
subgraphs=True後,事件帶namespace前綴((ns, mode, payload)),父圖能看見子圖內部事件,靠_output在出隊時打標籤,而不是子圖自己額外開一條流。 - v2 型別化:新版本透過
version="v2"把事件變成{"type": mode, "ns": ..., "data": ..., "interrupts": ...}字典,便於下游程式化處理;values模式還會把__interrupt__單獨抽出成interrupts欄位。
關鍵檔案
Pregel class:450—Pregel類別定義,所有串流入口都在這裡。Pregel.stream:2655— 同步串流入口,簽名裡列出全部stream_mode選項。Pregel.astream:3063— 非同步串流入口,真正的非同步主迴圈寫在它的實作體裡。Pregel.invoke:3836— 同步收斂入口,內部for chunk in self.stream(...)取最後一個valueschunk。sync main loop:2964— 同步主迴圈:while loop.tick()→runner.tick→loop.after_tick()。async main loop:3437— 非同步主迴圈:三段式同上,runner.atick替換同步版。_output:4184— 出隊過濾函式,按stream_mode/print_mode/subgraphs決定每個事件怎麼 yield。GraphRunStream:31— 同步呼叫方驅動的流包裝,for迴圈就是 pump,沒有背景執行緒。AsyncGraphRunStream:304— 非同步對應物,多個投影(run.values/run.messages)共用一個 pump。stream_mode attr:709—Pregel預設stream_mode="values",可被stream(stream_mode=...)覆蓋。
資料流
非同步主迴圈骨架在 astream 實作體裡(async main loop:3437),三段式和 Pregel 引擎頁講的一樣,差別在每段中間夾了 _output 把佇列裡的東西吐給呼叫方:
while loop.tick():
for task in await loop.amatch_cached_writes():
loop.output_writes(task.id, task.writes, cached=True)
async for _ in runner.atick(
[t for t in loop.tasks.values() if not t.writes],
timeout=self.step_timeout,
get_waiter=get_waiter,
schedule_task=loop.aaccept_push,
):
# emit output
for o in _output(
stream_mode,
print_mode,
subgraphs,
stream.get_nowait,
asyncio.QueueEmpty,
version,
_output_mapper,
_state_mapper,
):
yield o
loop.after_tick()
await aemit_graph_lifecycle_events(loop)
# wait for checkpoint
if durability_ == "sync":
await cast(asyncio.Future, loop._put_checkpoint_fut)runner.atick 內部把每個任務產生的 writes 寫回 loop,同時把 (ns, mode, payload) 三元組塞進 stream(asyncio.Queue);外層 async for 拿到執行權時,_output 呼叫 stream.get_nowait() 把佇列裡所有就緒事件抽乾,按 stream_mode 過濾後再 yield 給使用者。_output 的核心邏輯(_output:4184):拿到 (ns, mode, payload),如果 mode in print_mode 先 print 一份,如果 mode in stream_mode 才向外 yield——所以 print_mode 是「只看不發」的除錯開關,不影響真正 yield 出去的內容。
下面這張圖把 stream 呼叫從入口到事件出隊走一遍:
邊界與失敗
- 預設 stream_mode:
stream_mode=None時,如果是被作為子圖呼叫(CONFIG_KEY_TASK_ID在 config 裡),預設values,否則用self.stream_mode(stream_mode default:2740),避免被子圖預設updates模式污染父圖輸出。 - 步數超限:主迴圈退出後看
loop.status,如果是out_of_steps拋GraphRecursionError,提示調高recursion_limit(out_of_steps:3002);draining狀態拋GraphDrained,把控制權交還給外部RunControl。 - 多模式同時輸出:
stream_mode傳 list 時,每個事件變成(mode, payload)元組;再疊subgraphs=True則是(ns, mode, payload)(tuple mode:4240),呼叫方要按元組解構。 - messages 模式的繼承處理:
stream_mode="messages"且version="v1"時,會剝掉繼承自上層的 v2 messages handler,避免 v1 流被路由到 content-block 事件協議(strip v2 handler:2773),但保留 v1 handler 以支援subgraphs=True時內層messages事件被外層觀察。 - durability 與 checkpoint 同步:
durability="sync"時主迴圈每步結束都await loop._put_checkpoint_fut,保證 checkpoint 落盤後才進入下一步;"async"是預設值,持久化和下一步並行;"exit"只在退出時落盤。 - v3 實驗性 stream_events:
stream_events(version="v3")不接受stream_mode和subgraphs參數,這兩個被 mux 內部接管,顯式傳會被_reject_v3_invariant_kwargs拒掉(_reject_v3_invariant_kwargs:387)。
小結
astream / stream 把 Pregel 的 BSP 迴圈包成可迭代介面,stream_mode 選擇決定輸出投影,_output 是出隊過濾的單一咽喉。要追執行細節看 /pregel/pregel 與 /pregel/loop,要追多模式 mux 看 libs/langgraph/langgraph/stream/_mux.py。對照官方資料:LangGraph 文件 · README