AG Day 19 串流輸出:stream 與事件流
執行需求:CPU+API key。在昨天 AG Day 18(原文連結)中,我們把研究流程拆成檢索子圖與摘要子圖,讓 research-agent 的主圖維持精簡、每個子圖都能獨立測試。但到目前為止,不管圖的內部結構多乾淨,我們呼叫它的方式一直是 app.invoke(...):送出問題後畫面完全靜止,直到整條圖跑完,才一次看到最終結果。如果一次研究任務要花十幾秒、甚至更久的時間,使用者只會看到一個沒有任何回饋的空白畫面,體感上跟當機沒兩樣。今天我們要改用 LangGraph 的串流介面,讓 CLI 能即時顯示代理正在做什麼、正在思考些什麼。
引言
「串流」在使用者體驗上的價值,遠比表面看起來的「即時顯示文字」更重要。心理學上有個常被引用的現象:使用者對「等待中、有進度」的容忍度,遠高於「完全靜止、不知道發生什麼事」的等待。一個會逐字吐出回應的聊天介面,即使總耗時跟一次性回傳完全一樣,使用者的主觀體感也會覺得快得多,因為畫面持續在告訴他們「系統正在工作,沒有卡住」。對研究助理這種可能需要呼叫多次工具、思考好幾輪才能給出答案的系統來說,串流幾乎是必要而非加分項。
LangGraph 提供的串流能力不只侷限於「把模型生成的文字逐字吐出來」,它其實有好幾種不同粒度的串流模式,分別對應「使用者想看到模型打字的過程」「開發者想看到每個節點執行完的完整狀態」「想看到節點內部發生的細節事件(例如工具呼叫開始、結束)」等不同需求。今天我們會依序介紹這幾種模式,並示範怎麼在 research-agent 的 CLI 裡組合使用,讓一般使用者看到自然語言的即時回覆,同時讓開發模式下能看到更細節的除錯資訊。
原理與觀念
stream_mode="values":每個節點執行完的完整狀態
最簡單的串流模式是 values:圖每執行完一個節點,就把當下完整的狀態物件回傳一次。這種模式很適合用來觀察「代理現在進行到哪一步」,例如印出目前訊息歷史的長度、或最新一則訊息的內容,但它不會給你模型生成過程中的逐字內容,因為它是以「節點」為單位回傳,而不是以「token」為單位。對於只想知道「查詢進度」而不需要打字機效果的場景(例如背景排程任務的日誌),values 模式已經足夠。
stream_mode="updates":只回傳這個節點新增或改變的部分
updates 模式跟 values 很像,差別在於它只回傳「這個節點造成了什麼變化」,而不是整個狀態。當狀態物件本身很龐大(例如訊息歷史已經累積幾十則)時,updates 能大幅減少每次串流事件的資料量,因為你只關心「剛剛哪個節點跑完、它做了什麼」,不需要每次都拿到完整的歷史紀錄。這在我們之後要做的觀測與日誌記錄(AG Day 9 已經打過基礎)特別實用,可以直接把每次 updates 的內容摘要寫進 events 資料表。
stream_mode="messages":token 級的逐字輸出
如果目標是做出聊天介面常見的「打字機效果」,需要用 messages 模式,它會在底層模型支援串流回應的前提下,把模型逐步生成的訊息片段(token 或字元區塊)即時吐出來,搭配 metadata 還能知道這個片段來自圖裡的哪個節點。要注意這個模式依賴底層模型 API 本身有沒有提供串流端點;我們選用的模型透過環境變數 RESEARCH_AGENT_MODEL 指定,實際是否支援串流、串流的介面細節,請以對應廠商官方文件為準,不同模型、不同版本可能有落差。
astream_events:更細粒度的事件流,涵蓋工具呼叫的開始與結束
前面三種模式關注的都是「狀態怎麼變化」或「文字怎麼生成」,如果想知道更細節的事件,例如「工具呼叫從什麼時候開始、什麼時候結束、中間輸出了什麼」,LangGraph 提供了 astream_events 這個非同步方法,會依序吐出一連串帶有事件類型(例如節點開始、節點結束、工具開始執行、模型開始生成)的事件物件。這對建置更豐富的除錯介面或可觀測性儀表板很有用,但事件種類與資料結構相對複雜,一般 CLI 場景不一定需要用到這麼細的粒度,我們今天會示範基本用法,實務上先從 updates 或 messages 開始,有明確需求再往 astream_events 進階。
完整實作
我們先在 src/research_agent/cli.py 加一個串流版的執行函式,用 stream_mode="updates" 即時印出每個節點的執行進度,這是最容易上手、也最實用的一種模式:
# research-agent/src/research_agent/cli.py(新增串流函式)
from langchain_core.messages import HumanMessage
from research_agent.graph import build_graph
from research_agent.memory import get_checkpointer
def run_with_progress(query: str, thread_id: str) -> None:
"""用 updates 模式即時顯示代理目前跑到哪個節點。"""
app = build_graph(checkpointer=get_checkpointer())
config = {"configurable": {"thread_id": thread_id}}
for update in app.stream({"messages": [HumanMessage(content=query)]}, config=config, stream_mode="updates"):
for node_name, node_output in update.items():
print(f"[進度] 節點「{node_name}」執行完畢")
messages = node_output.get("messages") if isinstance(node_output, dict) else None
if messages:
last = messages[-1]
preview = str(getattr(last, "content", ""))[:60]
print(f" 最新訊息預覽:{preview}")
離線模式下的示範輸出:
[進度] 節點「agent」執行完畢
最新訊息預覽:(模型決定呼叫 search_web 的示範內容)
[進度] 節點「approval_gate」執行完畢
[進度] 節點「retrieval_subgraph」執行完畢
最新訊息預覽:[離線模擬結果] 針對「...」的模擬搜尋摘要...
[進度] 節點「agent」執行完畢
最新訊息預覽:根據搜尋結果(示範資料),整理如下...
接著示範 messages 模式的逐字輸出,這個模式需要非同步呼叫,並且需要有真實的 API 金鑰才能觀察到逐字生成的效果(離線模擬模式下不會有逐字輸出,會直接一次性回傳模擬字串):
# research-agent/src/research_agent/cli.py(新增逐字輸出函式)
import asyncio
from langchain_core.messages import HumanMessage
from research_agent.graph import build_graph
from research_agent.memory import get_checkpointer
async def run_with_typing_effect(query: str, thread_id: str) -> None:
"""用 messages 模式模擬打字機效果,需模型 API 支援串流回應。"""
app = build_graph(checkpointer=get_checkpointer())
config = {"configurable": {"thread_id": thread_id}}
print("代理:", end="", flush=True)
async for chunk, metadata in app.astream(
{"messages": [HumanMessage(content=query)]}, config=config, stream_mode="messages"
):
if metadata.get("langgraph_node") == "agent" and getattr(chunk, "content", None):
print(chunk.content, end="", flush=True)
print() # 換行,結束這一輪輸出
把兩種模式接進 CLI 的命令列參數,讓使用者可以選擇要看進度模式還是逐字模式:
# research-agent/src/research_agent/cli.py(整合命令列參數)
import argparse
import asyncio
def main() -> None:
parser = argparse.ArgumentParser(description="research-agent 命令列入口")
parser.add_argument("query", help="要研究的問題")
parser.add_argument("--thread-id", default="cli-default", help="對話串識別碼")
parser.add_argument(
"--stream", choices=["progress", "typing", "none"], default="progress",
help="選擇串流顯示模式:progress(節點進度)、typing(逐字打字機)、none(一次性輸出)",
)
args = parser.parse_args()
if args.stream == "progress":
run_with_progress(args.query, args.thread_id)
elif args.stream == "typing":
asyncio.run(run_with_typing_effect(args.query, args.thread_id))
else:
from research_agent.graph import build_graph
from research_agent.memory import get_checkpointer
from langchain_core.messages import HumanMessage
app = build_graph(checkpointer=get_checkpointer())
config = {"configurable": {"thread_id": args.thread_id}}
result = app.invoke({"messages": [HumanMessage(content=args.query)]}, config=config)
print(result["messages"][-1].content)
if __name__ == "__main__":
main()
最後示範最細粒度的 astream_events,讓開發者能看到工具呼叫的起訖時間,這在後面 AG Day 41 效能總檢章節排查延遲來源時會很有幫助:
# research-agent/debug_events.py
import asyncio
from langchain_core.messages import HumanMessage
from research_agent.graph import build_graph
from research_agent.memory import get_checkpointer
async def debug_run(query: str) -> None:
app = build_graph(checkpointer=get_checkpointer())
config = {"configurable": {"thread_id": "debug-events-001"}}
async for event in app.astream_events(
{"messages": [HumanMessage(content=query)]}, config=config, version="v2"
):
kind = event["event"]
if kind in ("on_tool_start", "on_tool_end"):
name = event.get("name", "未知工具")
print(f"[{kind}] 工具:{name}")
if __name__ == "__main__":
asyncio.run(debug_run("幫我查一下台灣的太陽能發電趨勢"))
這裡的 version="v2" 對應目前 LangGraph 事件結構的主要版本,實際參數名稱與事件種類(on_tool_start、on_tool_end、on_chat_model_stream 等)請以官方文件為準,不同版本可能新增或調整事件類型。
我們在 AG Day 9(原文連結)已經替代理建立了觀測基礎,把每一輪思考與工具呼叫寫進 events 資料表。串流介面剛好提供了一個很自然的切入點,把 updates 模式吐出的每一次節點結果,順手寫進資料庫,不需要額外的攔截邏輯:
# research-agent/src/research_agent/cli.py(串流結果同步寫入 events)
from research_agent.storage import record_event # 沿用 AG Day 9 建立的事件寫入函式
from research_agent.config import load_settings
def run_with_progress_and_logging(query: str, thread_id: str) -> None:
"""在顯示進度的同時,把每個節點的輸出摘要寫進 events 資料表,供事後查詢。"""
settings = load_settings()
app = build_graph(checkpointer=get_checkpointer())
config = {"configurable": {"thread_id": thread_id}}
step = 0
for update in app.stream({"messages": [HumanMessage(content=query)]}, config=config, stream_mode="updates"):
step += 1
for node_name, node_output in update.items():
print(f"[進度] 節點「{node_name}」執行完畢")
messages = node_output.get("messages") if isinstance(node_output, dict) else None
summary = str(getattr(messages[-1], "content", ""))[:200] if messages else ""
record_event(
db_path=settings.database_path,
run_id=thread_id,
step_number=step,
event_type=node_name,
tool_name=None,
tool_args=None,
tool_result=summary,
)
這樣一來,即使使用者只看到終端機上的即時進度,我們在背後也同步留下了一份結構化、之後可以直接下 SQL 查詢的完整執行紀錄,串流與觀測性不會互相取捨。
最後,替 run_with_progress 補一個簡單的測試,確認它至少會走過每一個節點、不會在串流過程中意外中斷。因為離線模擬模式下所有工具呼叫都是同步、可預期的,這個測試完全不需要任何 API 金鑰:
# research-agent/tests/test_cli_streaming.py
from research_agent.cli import run_with_progress
def test_run_with_progress_completes_without_error(capsys):
run_with_progress("測試用查詢字串", thread_id="pytest-stream-001")
output = capsys.readouterr().out
assert "[進度]" in output
assert "agent" in output
常見錯誤與踩雷
第一個常見錯誤是把同步的 stream() 跟非同步的 astream() 搞混。messages 模式與 astream_events 通常需要用 async for 搭配 astream/astream_events,如果誤用同步的 for ... in app.stream(...) 呼叫非同步介面,會直接拿到一個協程物件而不是預期的資料,或是丟出型別相關的錯誤。動手前先確認自己用的是同步版還是非同步版的 API,並搭配對應的呼叫語法。
第二個是誤以為只要換成串流模式,總耗時就會變短。串流改變的是「使用者何時開始看到內容」,而不是「整個流程總共花多少時間」;如果代理需要呼叫三次工具才能給出答案,不管用不用串流,這三次工具呼叫的總時間都是省不掉的。把串流當成使用者體驗上的改善,而不是效能上的最佳化,才不會對它有錯誤的期待。
第三個是在 messages 模式下,忘記過濾 metadata 裡的節點來源,導致把子圖內部節點、或非對外顯示用途的中間訊息也一併印給使用者看,畫面變得雜亂。我們在範例裡明確篩選 metadata.get("langgraph_node") == "agent",只顯示主要回覆節點的內容,這種篩選邏輯建議明確寫出來,而不是依賴「反正只有一個地方會生成文字」的僥倖假設。
第四個雷是離線模擬模式(--dry-run)下,工具直接回傳一整段固定字串,並不會產生逐字生成的效果,如果沒有在說明文字或程式碼註解裡提到這一點,團隊裡沒有金鑰的成員可能會誤以為 messages 模式的串流功能沒有正常運作,實際上只是離線模式本來就沒有逐字生成這件事可言。
效能與實務提醒
串流輸出對後端資源的消耗,通常不會比一次性回傳更高,因為底層呼叫模型 API 的次數是一樣的,差別只在於資料是分批送達還是一次送達。但如果是透過 HTTP 對外提供服務(我們會在 AG Day 38 用 FastAPI 包裝),串流回應需要搭配伺服器發送事件(Server-Sent Events)或分塊傳輸編碼(chunked transfer encoding),前端與反向代理設定要正確處理長連線與逐塊資料,否則使用者可能看不到任何串流效果,或是連線在中途被逾時機制切斷。
另外,串流過程中如果使用者中途取消請求(例如關掉瀏覽器分頁),伺服器端需要有機制正確釋放對應的資源、停止繼續呼叫模型 API,避免產生使用者已經不在乎、卻仍持續計費的呼叫。這部分的完整處理會在部署篇章具體示範,今天先在 CLI 場景熟悉串流介面本身的用法。
如果同時要顯示串流進度、又要把每一步寫進資料庫,記得評估這份額外的資料庫寫入會不會拖慢使用者感受到的即時性。今天的 record_event 是同步呼叫,寫入 SQLite 通常在毫秒等級,對整體串流體驗影響不大;但如果之後改成寫進需要走網路的觀測服務(例如 AG Day 35 要介紹的 Langfuse),建議改成非同步、非阻塞的方式送出,避免觀測本身反而變成拖慢使用者體驗的瓶頸。
最後,開發與除錯階段建議預設用 progress(updates 模式)而不是逐字打字機效果,因為前者的輸出內容更適合直接讀、複製貼上做除錯,逐字模式的輸出在終端機裡不容易複製成乾淨的文字,比較適合留給真正要面向一般使用者展示的介面場合。
小結
今天我們讓 research-agent 從「整條圖跑完才看到結果」進化成能即時回饋執行進度:stream_mode="updates" 用來顯示節點層級的進度、stream_mode="messages" 用來做逐字打字機效果、astream_events 用來取得更細粒度的除錯事件。三種模式各有適合的場景,不必為了追求最細節的資訊,而在一般使用情境也硬套 astream_events。
今天新增的關鍵詞:串流(stream)——資料分批、即時送達,而不是等全部處理完才一次回傳;事件流(event stream)——描述執行過程中發生了哪些具體事件(節點開始、工具呼叫結束等)的細粒度資料流;打字機效果——把模型逐步生成的文字片段即時顯示出來的使用者介面模式。我們也順手把串流事件接回了 events 資料表,證明「使用者體驗」與「可觀測性」這兩個目標,其實可以用同一份資料流一次滿足,不需要分別維護兩套機制。
結語
有了即時的執行進度回饋,research-agent 在使用體驗上跨出了重要一步。但目前代理的能力還侷限在「呼叫外部搜尋工具、拿回一段文字摘要」,它並沒有一份真正屬於自己、能反覆查詢的知識庫,每次都得重新向外查詢,也無法針對已經蒐集過的資料做更精準的比對。
明天,我們會進入「AG Day 20 RAG 基礎:embedding 與向量檢索」,開始建立檢索增強生成(RAG)的基礎能力,介紹文字如何被轉換成向量、向量之間的相似度怎麼計算,並讓 research-agent 第一次擁有能查詢的本機向量知識庫,而不是每次都只能仰賴即時的網路搜尋結果。
延伸資源
- LangGraph 官方文件:Streaming 概念頁面,涵蓋
values、updates、messages等stream_mode的完整說明。 - LangGraph 官方文件:
astream_events的事件類型與版本差異說明。 - MDN Web Docs:Server-Sent Events 介紹,理解串流回應在 HTTP 層的實作方式,供 AG Day 38 部署篇章預作準備。
留言
張貼留言