Skip to content

串流輸出:astream 與 stream_mode

源码版本1.2.9

職責

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:450Pregel 類別定義,所有串流入口都在這裡。
  • Pregel.stream:2655 — 同步串流入口,簽名裡列出全部 stream_mode 選項。
  • Pregel.astream:3063 — 非同步串流入口,真正的非同步主迴圈寫在它的實作體裡。
  • Pregel.invoke:3836 — 同步收斂入口,內部 for chunk in self.stream(...) 取最後一個 values chunk。
  • sync main loop:2964 — 同步主迴圈:while loop.tick()runner.tickloop.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:709Pregel 預設 stream_mode="values",可被 stream(stream_mode=...) 覆蓋。

資料流

非同步主迴圈骨架在 astream 實作體裡(async main loop:3437),三段式和 Pregel 引擎頁講的一樣,差別在每段中間夾了 _output 把佇列裡的東西吐給呼叫方:

python
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_modeprint 一份,如果 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_stepsGraphRecursionError,提示調高 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_modesubgraphs 參數,這兩個被 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