跳到主要內容

AG Day 19 串流輸出:stream 與事件流

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 部署篇章預作準備。

留言

這個網誌中的熱門文章

Day 2 變數與資料型別

Day 2 變數與資料型別 引言 寫程式的過程中,變數與資料型別是處理資料的基礎。變數是存放資料的容器,資料型別則決定這筆資料有哪些特性、可以進行哪些操作。學會定義變數、認識各種資料型別,是學好 Python 的關鍵一步。 這篇文章會帶你了解 Python 中變數的觀念、如何定義變數,以及常見的資料型別,包括整數、浮點數、字串、布林值,還有串列、元組、字典與集合等容器型別。我們也會介紹變數的命名規則與撰寫風格建議,以及如何用 type() 檢查資料型別。 什麼是變數?如何在 Python 中定義變數 變數是在程式執行時用來存放資料的名稱。透過定義變數,我們可以給一筆資料一個名字,並在程式的其他地方用這個名字取用該筆資料。在 Python 中,變數不需要事先宣告型別,因為 Python 是動態型別語言,變數的型別由指定給它的值決定。 定義變數的基本語法 在 Python 中定義變數非常簡單,只要用賦值符號 = 把值指定給變數即可。例如: x = 5 # 定義變數 x,並把整數 5 賦值給它 name = "Alice" # 定義變數 name,並把字串 "Alice" 賦值給它 在這裡,x 是一個變數,被賦予整數 5;name 是另一個變數,被賦予字串 "Alice"。 變數的更新與覆寫 變數的值可以修改,也就是說,我們可以在程式的不同地方給同一個變數新的值。例如: x = 10 # x 最初被賦予 10 x = 15 # x 的值現在被更新為 15 這樣就能依照需求,在程式執行過程中靈活調整變數的值。 Python 的動態型別系統 Python 和某些靜態型別語言不同,定義變數時不需要宣告型別。賦值時,Python 會根據值自動判斷變數的型別。例如: x = 5 # x 是整數 x = 3.14 # x 變成浮點數 x = "Hi" # x 變成字串 同一個變數在程式執行過程中可以存放不同型別的值,這是 Python 的彈性之一。 常見資料型別 在 Python 中,資料型別決定我們可以對變數進行哪些操作...

Day 1 Python 簡介與環境設定

Day 1 Python 簡介與環境設定 引言 在現在的科技環境裡,程式設計已經是一項重要技能。無論你是對資料科學有興趣、想成為開發者,或是想踏入人工智慧(AI)領域,學會寫程式都能明顯提升你的競爭力。在眾多程式語言中,Python 因為語法簡單、功能強大、應用範圍廣泛,成為許多人進入程式世界的第一選擇。這篇文章會帶你認識 Python 的背景與優勢,並一步步教你在不同系統上安裝與設定 Python 開發環境,最後寫出第一支 Python 程式。 為什麼選擇 Python? Python 是一種高階程式語言,由 Guido van Rossum 在 1991 年發布。Python 的設計哲學強調程式碼的可讀性,並用縮排來定義程式區塊,這點和許多使用大括號的語言不同。簡潔的語法讓它成為初學者的理想選擇;就算是經驗豐富的開發者,也能用它完成複雜的專案。 Python 的優勢如下: 簡單易學 :Python 的語法清楚、結構簡潔,初學者很快就能上手。和其他語言相比,學習曲線相對平緩,不需要先弄懂一堆複雜觀念,就能開始寫程式。 應用範圍廣泛 :從資料科學、網頁開發、人工智慧、機器學習、自動化測試到網路爬蟲,Python 都有大量開源函式庫與工具支援,而且在這些領域都扮演關鍵角色。 豐富的函式庫與框架 :Python 的函式庫生態系非常龐大。做資料分析有 NumPy、Pandas;開發網站有 Django、Flask;做深度學習有 TensorFlow、PyTorch。各種需求幾乎都能找到對應的套件,讓開發更有效率。 跨平台支援 :Python 支援 Windows、macOS、Linux 等作業系統,程式通常不需要太多修改就能跨平台執行,讓開發與部署更有彈性。 活躍的社群 :Python 擁有龐大的開發者社群。學習或開發上遇到問題,幾乎都能在社群與論壇(例如 Stack Overflow)找到答案,對初學者來說是很強的後盾,也能減少卡關時的挫折感。 Python 的應用領域 Python 的流行與強大功能,讓許多領域都開始大量使用它。以下是幾個常見的應用方向: 資料科學 :隨著大數據與人工智慧興起,資料科學大量使用 Python。NumPy、Pandas 與 Matplotlib 等工具能處理和分析龐...

Python 從入門到 PyTorch 深度學習:開啟 AI 世界的大門

Python 從入門到 PyTorch 深度學習:開啟 AI 世界的大門 隨著人工智慧(AI)與深度學習(Deep Learning)快速發展,越來越多人對這些技術產生興趣。不論你是想踏入 AI 領域的初學者,還是已經有程式基礎的開發者,學好 Python 與深度學習框架(例如 PyTorch),都能為你打開更多可能。 為什麼選擇 Python? Python 已經是資料科學與人工智慧領域的首選語言。它的語法簡潔、容易上手,而且擁有龐大的生態系與大量開源函式庫。無論是資料處理、資料視覺化,還是建立機器學習與深度學習模型,Python 都能勝任。對想進入 AI 或資料科學領域的人來說,它幾乎是必備工具。 PyTorch 是什麼? PyTorch 是由 Meta(原 Facebook)AI 研究團隊開發的開源深度學習框架,以易用、靈活和動態計算圖著稱,是許多 AI 研究人員與開發者的首選。相較於其他框架,PyTorch 的寫法更貼近原生 Python,對初學者相對友善。無論是簡單的實驗,還是複雜的深度學習模型,PyTorch 都能提供強大的支援。 這個系列能帶給你什麼? 這個系列會從 Python 的基礎開始,帶你一步一步學習,最後能自己用 PyTorch 建立深度學習模型。即使你完全沒有寫過程式,也能跟著文章的節奏累積技能,理解 AI 與深度學習的核心觀念。 本系列涵蓋的主題 Python 基礎:從變數、條件判斷到函式與模組。 資料處理工具:用 NumPy 與 Pandas 有效率地操作資料。 資料視覺化:用 Matplotlib 與 Seaborn 把資料畫成圖表。 深度學習的數學基礎:線性代數、微積分與機率。 PyTorch 入門:理解張量、模型建構與 GPU 加速。 基礎深度學習模型:CNN 與 RNN 的實作應用。 深度學習專案實戰:從資料前處理到模型部署的端到端流程。 誰適合這個系列? 程式初學者 :如果你對 AI 充滿好奇,卻還沒寫過程式,系列的第一部分會帶你快速上手 Python,並幫助你理解深度學習的基本觀念。 資料科學愛好者 :如果你已經熟悉一些資料處理方法,進階部分會教你如何用 PyTorch 建構深度學習模型。 開發者與研究人員 :想更深入了...