Skip to content

流式输出 Streaming

LLM 回答通常是逐字生成的,前端做"打字机"效果就要靠流式。LangGraph 的 stream 不仅能流式输出 LLM token,还能流式输出每步状态变化。本篇讲清楚各种 stream_mode 的区别。

一、为什么需要流式

场景非流式流式
用户体验等几秒一次性蹦出整段逐字出现,像 ChatGPT
长流程可见性全跑完才看到结果每步实时可见进度
调试只能看最终状态看到每个节点的增量
前端集成一次返回 JSONSSE/WebSocket 持续推送

LangGraph 的 stream 是同步流式,astream 是异步流式。核心参数是 stream_mode

二、stream_mode 参数详解

stream_mode 决定每次回调给你"什么形态的数据":

取值每次回调内容典型用途
values完整状态快照(每步一份)看状态全貌
updates每步的增量更新字典(默认推荐)调试单节点
messagesLLM 逐 token 输出 + 元数据打字机效果
debug最详细的调度日志排查图行为
custom节点内用 get_stream_writer 主动发的自定义事件业务自定义进度

可以同时订阅多个:stream_mode=["updates", "messages"],回调里按类型区分。

三、示例图:一个简单两节点工作流

python
from typing import TypedDict
from langgraph.graph import StateGraph, START, END

class State(TypedDict):
    topic: str
    draft: str
    final: str

def draft_node(state):
    # 假装逐步生成草稿(真实场景是 llm.stream)
    chunks = ["LangGraph", "是一个", "用于", "构建", "智能体", "的框架。"]
    text = ""
    for c in chunks:
        text += c
    return {"draft": text}

def polish_node(state):
    return {"final": state["draft"] + "(已润色)"}

graph = StateGraph(State)
graph.add_node("draft", draft_node)
graph.add_node("polish", polish_node)
graph.add_edge(START, "draft")
graph.add_edge("draft", "polish")
graph.add_edge("polish", END)
app = graph.compile()

四、同一次运行下不同 mode 的输出差异

mode="values":每步完整状态

python
for chunk in app.stream({"topic": "LangGraph"}, stream_mode="values"):
    print(chunk)
text
{'topic': 'LangGraph'}                          # 初始
{'topic': 'LangGraph', 'draft': 'LangGraph是一个...'}   # draft 后
{'topic': 'LangGraph', 'draft': '...', 'final': '...(已润色)'}  # polish 后

每步都是完整状态,字段越来越多。适合"想看状态演进全貌"。

mode="updates":每步增量

python
for chunk in app.stream({"topic": "LangGraph"}, stream_mode="updates"):
    print(chunk)
text
{'draft': {'draft': 'LangGraph是一个用于构建智能体的框架。'}}
{'polish': {'final': 'LangGraph是一个用于构建智能体的框架。(已润色)'}}

每个 chunk 只包含本节点改了什么,结构是 {节点名: 更新字典}。最适合调试。

mode="messages":逐 token

python
for chunk, metadata in app.stream({"topic": "LangGraph"}, stream_mode="messages"):
    print(chunk.content, end="", flush=True)

messages 模式需要图里有真正支持流式的 LLM 调用(llm.streamllm.astream)。上面假例子没有真 LLM,所以这个模式看不到 token。真实用法见下一节。

mode="debug":最详细

python
for chunk in app.stream({"topic": "LangGraph"}, stream_mode="debug"):
    print(chunk)

会输出每次调度的细节:执行了哪个节点、输入输出、耗时等。排查"图为什么这么走"时最有用。

五、真实 LLM 流式:messages 模式

把节点换成流式 LLM,messages 模式就能拿到逐 token:

python
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
from langgraph.graph import StateGraph, MessagesState, START, END

llm = ChatOpenAI(model="gpt-4o-mini", temperature=0, streaming=True)

def chat(state):
    # llm.invoke 在 messages 模式下会被自动流式
    resp = llm.invoke(state["messages"])
    return {"messages": [resp]}

g = StateGraph(MessagesState)
g.add_node("chat", chat)
g.add_edge(START, "chat")
g.add_edge("chat", END)
app = g.compile()

for chunk, meta in app.stream(
    {"messages": [HumanMessage(content="用三句话介绍 LangGraph")]},
    stream_mode="messages",
):
    # chunk 是 AIMessageChunk,content 是增量文本
    print(chunk.content, end="", flush=True)

要点:

  • LLM 要开 streaming=True
  • messages 模式的每个 chunk 是 AIMessageChunk.content 是这一小段文本。
  • meta 里有 langgraph_node 等字段,告诉你这个 token 来自哪个节点。

六、同时订阅多种 mode

python
for chunk in app.stream(
    {"messages": [HumanMessage(content="你好")]},
    stream_mode=["updates", "messages"],
):
    # chunk 是元组 (mode, data)
    mode, data = chunk
    if mode == "updates":
        print("节点更新:", data)
    elif mode == "messages":
        print("token:", data[0].content, end="")

注意:多模式时,回调项是 (mode_name, data) 元组,要按 mode 分流处理。

七、astream:异步流式

在异步框架(FastAPI、aiohttp)里用 astream,不阻塞事件循环:

python
async def run():
    async for chunk, meta in app.astream(
        {"messages": [HumanMessage(content="你好")]},
        stream_mode="messages",
    ):
        print(chunk.content, end="", flush=True)

import asyncio
asyncio.run(run())

ainvoke / astream 对应同步的 invoke / stream,节点里用 async def 定义异步节点即可。详见 创建与连接节点

八、前端 SSE 对接简介

后端把 astream 的输出转成 SSE(Server-Sent Events)推给浏览器,就能做打字机效果。大致结构:

python
from fastapi import FastAPI
from fastapi.responses import StreamingResponse

app_server = FastAPI()

@app_server.get("/chat")
async def chat_stream():
    async def event_gen():
        async for chunk, meta in app.astream(
            {"messages": [HumanMessage(content="你好")]},
            stream_mode="messages",
            config={"configurable": {"thread_id": "web-1"}},
        ):
            text = chunk.content
            if text:
                yield f"data: {text}\n\n"   # SSE 格式
    return StreamingResponse(event_gen(), media_type="text/event-stream")

前端用 EventSource 监听即可。完整的部署对接见 部署与运维

九、custom 模式:自定义进度

节点内用 get_stream_writer 主动发事件,配合 stream_mode="custom"

python
from langgraph.config import get_stream_writer

def long_task(state):
    writer = get_stream_writer()
    for i in range(3):
        writer({"progress": f"步骤 {i+1}/3 完成"})
    return {"done": True}

# 监听
for chunk in app.stream(inputs, stream_mode="custom"):
    print("自定义事件:", chunk)

适合上报业务语义的进度("已检索到 5 篇文档""正在生成第 2 段")。

十、常见踩坑

踩坑 1:messages 模式看不到 token

原因通常是 LLM 没开 streaming=True,或节点用的是 llm.invoke 但模型不支持流式。确保 LLM 支持 streaming 且节点确实在流式生成。

踩坑 2:流式与 checkpointer 配合出问题

带 checkpointer 的图在流式时,每个 step 仍会写检查点。如果节点很慢,可能看到"状态先更新、检查点后写完"的延迟。一般无碍,但 HITL 中断点判断要基于 get_state 而非流式 chunk。

踩坑 3:把 messages 模式的 chunk 当成完整消息

messages 模式的每个 chunk 只是一小段,要拼接才是完整回复。完整消息会通过 values/updates 模式给出(节点返回后)。

踩坑 4:同步流式阻塞了主线程

长流程同步 stream 会阻塞。在 Web 服务里一定用 astream,否则一个用户的长回复会卡住整个进程。

踩坑 5:多模式时没按 mode 分流

stream_mode=["a", "b"] 时每个 chunk 是 (mode, data) 元组,直接当 dict 用会报错。先解包判断 mode。

十一、小结

  • stream/astream + stream_mode 控制输出形态。
  • 调试首选 updates,看全貌用 values,打字机用 messages,排查调度用 debug,业务进度用 custom
  • LLM 流式要开 streaming=True,多模式时按 (mode, data) 分流。
  • Web 服务用 astream + SSE 推前端。