Appearance
流式输出 Streaming
LLM 回答通常是逐字生成的,前端做"打字机"效果就要靠流式。LangGraph 的
stream不仅能流式输出 LLM token,还能流式输出每步状态变化。本篇讲清楚各种stream_mode的区别。
一、为什么需要流式
| 场景 | 非流式 | 流式 |
|---|---|---|
| 用户体验 | 等几秒一次性蹦出整段 | 逐字出现,像 ChatGPT |
| 长流程可见性 | 全跑完才看到结果 | 每步实时可见进度 |
| 调试 | 只能看最终状态 | 看到每个节点的增量 |
| 前端集成 | 一次返回 JSON | SSE/WebSocket 持续推送 |
LangGraph 的 stream 是同步流式,astream 是异步流式。核心参数是 stream_mode。
二、stream_mode 参数详解
stream_mode 决定每次回调给你"什么形态的数据":
| 取值 | 每次回调内容 | 典型用途 |
|---|---|---|
values | 完整状态快照(每步一份) | 看状态全貌 |
updates | 每步的增量更新字典(默认推荐) | 调试单节点 |
messages | LLM 逐 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.stream 或 llm.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 推前端。