流式输出与订阅

第 11 章我们用普通 API 实现了终端逐字输出。LangGraph 把"过程"也做成了流的:graph.stream() 既能按节点产出状态更新,也能把模型回复按 token 逐个吐出来。前者适合观察智能体每一步在做什么,后者适合做"打字机"体验或对接服务端推送。本章沿用最小对话图做演示。

准备工作:一个最小对话图

import os
from typing import Annotated, TypedDict
from langchain_core.messages import HumanMessage
from langchain_openai import ChatOpenAI
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages

model = ChatOpenAI(
    model=os.getenv("LLM_MODEL", "deepseek-chat"),
    api_key=os.getenv("LLM_API_KEY"),
    base_url=os.getenv("LLM_BASE_URL"),
)

class State(TypedDict):
    messages: Annotated[list, add_messages]

def agent(state):
    reply = model.invoke(state["messages"])
    return {"messages": [reply]}

builder = StateGraph(State)
builder.add_node("agent", agent)
builder.add_edge(START, "agent")
builder.add_edge("agent", END)
graph = builder.compile()

invoke 是"等全部跑完拿最终结果";stream 是"边跑边看"。两者输入相同,都可以配 checkpointer 与 thread_id。

stream_mode="updates":按节点输出

每个产出是一份字典 {节点名: 该节点返回的部分更新},能看清执行了哪些节点:

for update in graph.stream(
    {"messages": [HumanMessage("一句话介绍大模型")]},
    stream_mode="updates",
):
    for node, value in update.items():
        # 本例只有 agent 节点,value 形如 {'messages': [AIMessage(...)]}
        print("运行节点:", node)
# 输出:运行节点: agent(出现几次就代表图走了几个节点/几轮)

多节点图(如第 18 章的 ReAct)会依次打出 agent、tools、agent……直到结束,非常适合观察循环轮次与分支走向。

stream_mode="messages":逐 token 输出

messages 模式产出 (消息块, 元信息) 二元组,消息块可用 .content 累积成完整回复:

full = ""
for chunk, meta in graph.stream(
    {"messages": [HumanMessage("一句话介绍大模型")]},
    stream_mode="messages",
):
    piece = chunk.content
    if piece:
        print(piece, end="", flush=True)   # 逐字打印,形成打字机效果
        full += piece
# 输出:大模型是通过海量文本训练出来的……(文字逐个出现)
print("\n[完整回复]", full)

meta 里带 langgraph_node、langgraph_checkpoint_ns 等字段,可用于判断消息来自哪个节点、是否需要跳过工具消息。

组合两种模式

想同时拿到"节点轨迹"和"逐字内容",传一个模式列表即可:

for kind, data in graph.stream(
    {"messages": [HumanMessage("一句话介绍大模型")]},
    stream_mode=["updates", "messages"],
):
    if kind == "updates":
        print("[节点更新]", list(data.keys()))
    else:
        chunk = data[0]
        if chunk.content:
            print(chunk.content, end="", flush=True)
# 输出:同时出现 [节点更新] ['agent'] 与逐字打印的回复内容(两者先后由框架调度决定)

订阅与推送

  • 异步服务端用 graph.astream(async 版本),在协程里逐块处理即可。
  • Web 服务把每个 chunk 包成 SSE(Server-Sent Events)推给浏览器,即可实现流式对话界面;完整 FastAPI + SSE 实现将在第 28 章给出。思路一句话:一个生成器里 yield f"data: {文本块}\n\n",浏览器用 EventSource 消费。

三种模式怎么选

模式产出典型用途
updates每执行一个节点给一份 {节点: 状态更新}观察执行轨迹、调试循环
messages每个 token 一个消息块打字机效果、SSE 推送
updates + messages两类事件按类型混排界面同时显示步骤与文字

常见问题

  • 为什么 messages 模式会吐出空内容块?工具调用、消息元数据等 chunk 不一定带文本,用 if piece 跳过即可。
  • stream 与 invoke 的最终结果一致吗?一致,stream 只是提前暴露中间过程;想核对终态可收集全部内容,或用第 20 章的 get_state 查询。
  • 想区分文字来自哪个节点?读消息块元信息 meta 里的 langgraph_node 字段。
  • 输出乱序或丢字?服务端推送建议按 SSE 事件逐块下发、客户端按序追加,避免并发连接各自拼接。

小结:graph.stream 按 stream_mode 输出不同粒度——"updates" 逐节点看执行轨迹,"messages" 逐 token 拿模型回复;两者可组合,服务端再配 astream + SSE 就能把打字机效果推给用户。

笔记加载中…