Streaming 流式输出
流式输出(Streaming)允许 Agent 在执行过程中实时输出结果,而不是等待完整执行完毕后再一次性返回。这大大提升了用户体验,让响应看起来更快速、更自然。
概述
Streaming 的核心优势:
- 即时反馈:用户不必等待完整的响应
- 更好的体验:类似打字效果的实时输出
- 可见进度:Agent 思考过程的可见性
- 早期中断:用户可以在看到部分结果后中断
LangChain 支持多种流式模式:
stream_mode="values" → 每次发送所有状态值
stream_mode="messages" → 仅发送新增消息
stream_mode="updates" → 仅发送状态更新
stream_mode="events" → 发送详细事件(逐步推理等)
stream_mode="debug" → 发送调试信息基础流式输出
使用 .stream() 方法
python
from langgraph.prebuilt import create_agent
from langchain_openai import ChatOpenAI
agent = create_agent(
model=ChatOpenAI(model="gpt-4o-mini", streaming=True),
tools=[]
)
# 流式输出
for event in agent.stream(
{"messages": [("human", "请用三段话介绍人工智能")]},
stream_mode="values"
):
messages = event.get("messages", [])
if messages:
# 输出最后一条消息的最新内容
message = messages[-1]
if hasattr(message, 'content') and message.content:
print(message.content, end="", flush=True)使用 create_agent 的流式支持
python
agent = create_agent(
model=ChatOpenAI(model="gpt-4o-mini", streaming=True),
tools=tools,
system_prompt="你是一个实时助手。"
)
# 逐块输出
async for event in agent.astream_events(
{"messages": [("human", "写一首关于春天的诗")]},
version="v2"
):
kind = event["event"]
if kind == "on_chat_model_stream":
content = event["data"]["chunk"].content
if content:
print(content, end="", flush=True)Stream Mode 详解
stream_mode="values"
每次流式输出都包含完整的当前状态值:
python
print("=== stream_mode='values' ===")
for event in agent.stream(
{"messages": [("human", "你好")]},
stream_mode="values"
):
# event 包含完整的状态
messages = event["messages"]
print(f"当前消息数: {len(messages)}")
if messages:
last_msg = messages[-1]
if hasattr(last_msg, 'content') and last_msg.content:
print(f"内容: {last_msg.content}")
print("---")stream_mode="messages"
只输出新增的消息,减少冗余:
python
print("=== stream_mode='messages' ===")
for msg_chunk, metadata in agent.stream(
{"messages": [("human", "你好")]},
stream_mode="messages"
):
if hasattr(msg_chunk, 'content') and msg_chunk.content:
print(msg_chunk.content, end="", flush=True)
print()stream_mode="updates"
仅发送各节点的新增更新:
python
print("=== stream_mode='updates' ===")
for update in agent.stream(
{"messages": [("human", "你好")]},
stream_mode="updates"
):
print(f"更新: {update}")事件类型
当使用 stream_mode="events" 时,可以收到丰富的事件:
python
# 常用事件类型
ON_CHAIN_START # 链开始执行
ON_CHAIN_END # 链执行完成
ON_CHAT_MODEL_START # 模型开始生成
ON_CHAT_MODEL_STREAM # 模型流式生成中
ON_CHAT_MODEL_END # 模型生成完成
ON_TOOL_START # 工具开始调用
ON_TOOL_END # 工具调用完成
ON_RETRIEVER_START # 检索器开始检索
ON_RETRIEVER_END # 检索器检索完成python
# 监听详细事件
async for event in agent.astream_events(
{"messages": [("human", "搜索关于 Python 的信息")]},
version="v2"
):
event_type = event["event"]
name = event.get("name", "")
if event_type == "on_tool_start":
print(f"🔧 开始调用工具: {name}")
print(f" 输入: {event['data'].get('input')}")
elif event_type == "on_tool_end":
print(f"✅ 工具完成: {name}")
print(f" 输出: {str(event['data'].get('output', ''))[:100]}")
elif event_type == "on_chat_model_stream":
chunk = event["data"]["chunk"]
if hasattr(chunk, 'content') and chunk.content:
print(chunk.content, end="", flush=True)
elif event_type == "on_chat_model_end":
print("\n📝 模型生成完成")带工具的流式输出
python
from langchain.tools import tool
from langgraph.prebuilt import create_agent
from langchain_openai import ChatOpenAI
@tool
def get_weather(city: str) -> str:
"""获取城市天气"""
return f"{city} 的天气:晴朗,25°C"
@tool
def search_info(query: str) -> str:
"""搜索信息"""
return f"关于'{query}'的搜索结果..."
tools = [get_weather, search_info]
agent = create_agent(
model=ChatOpenAI(model="gpt-4o-mini", streaming=True),
tools=tools,
system_prompt="你是助手,使用工具回答问题。"
)
# 流式输出含工具调用的响应
async for event in agent.astream_events(
{"messages": [("human", "北京的天气怎么样?顺便搜索一下北京的景点")]},
version="v2"
):
kind = event["event"]
if kind == "on_tool_start":
print(f"\n⚡ 调用工具: {event['name']}")
print(f" 参数: {event['data'].get('input')}")
elif kind == "on_tool_end":
print(f"✓ 工具返回: {str(event['data'].get('output', ''))[:80]}")
elif kind == "on_chat_model_stream":
chunk = event["data"]["chunk"]
if hasattr(chunk, 'content') and chunk.content:
print(chunk.content, end="", flush=True)构建流式聊天应用
结合 Stream 实现实时聊天界面:
python
import asyncio
from langgraph.prebuilt import create_agent
from langgraph.checkpoint import InMemorySaver
from langchain_openai import ChatOpenAI
from langchain.tools import tool
@tool
def search_knowledge(query: str) -> str:
"""搜索知识库"""
return f"关于 {query} 的知识库内容"
@tool
def calculate(expression: str) -> str:
"""计算数学表达式"""
try:
return str(eval(expression))
except:
return "无法计算"
tools = [search_knowledge, calculate]
checkpointer = InMemorySaver()
agent = create_agent(
model=ChatOpenAI(model="gpt-4o-mini", streaming=True),
tools=tools,
checkpointer=checkpointer,
system_prompt="你是一个全能助手。"
)
async def stream_chat(message: str, thread_id: str = "chat"):
"""流式聊天函数"""
print(f"用户: {message}")
print("助手: ", end="", flush=True)
async for event in agent.astream_events(
{"messages": [("human", message)]},
config={"configurable": {"thread_id": thread_id}},
version="v2"
):
kind = event["event"]
if kind == "on_chat_model_stream":
chunk = event["data"]["chunk"]
if hasattr(chunk, 'content') and chunk.content:
print(chunk.content, end="", flush=True)
elif kind == "on_tool_start":
print(f"\n\n[调用工具: {event['name']}]")
elif kind == "on_tool_end":
print(f"[工具完成]")
print("\n")
# 交互示例
async def main():
await stream_chat("你好,你是谁?")
await stream_chat("100 + 200 等于多少?")
await stream_chat("你知道什么是机器学习吗?")
# 运行
asyncio.run(main())Web 应用的流式集成
python
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from langgraph.prebuilt import create_agent
from langchain_openai import ChatOpenAI
app = FastAPI()
agent = create_agent(
model=ChatOpenAI(model="gpt-4o-mini", streaming=True),
tools=[]
)
@app.post("/chat")
async def chat_endpoint(message: str):
async def generate():
async for event in agent.astream_events(
{"messages": [("human", message)]},
version="v2"
):
if event["event"] == "on_chat_model_stream":
chunk = event["data"]["chunk"]
if hasattr(chunk, 'content') and chunk.content:
yield f"data: {chunk.content}\n\n"
yield "data: [DONE]\n\n"
return StreamingResponse(
generate(),
media_type="text/event-stream"
)
# uvicorn server:app --reload前端集成示例(JavaScript)
javascript
// 浏览器端代码
const eventSource = new EventSource('/chat?message=你好');
eventSource.onmessage = (event) => {
if (event.data === '[DONE]') {
eventSource.close();
return;
}
// 增量显示内容
document.getElementById('response').textContent += event.data;
};
eventSource.onerror = (error) => {
console.error('Stream error:', error);
eventSource.close();
};Stream 模式对比
| stream_mode | 输出内容 | 适用场景 |
|---|---|---|
values | 完整状态值 | 需要完整上下文 |
messages | 新增消息 | 简单的聊天应用 |
updates | 状态更新 | 监控执行流程 |
events | 详细事件 | 完整追踪、UI 更新 |
debug | 调试信息 | 开发和调试 |
最佳实践
- 启用 streaming:创建 LLM 时设置
streaming=True - 选择合适的 stream_mode:聊天用
messages,监控用events - 处理工具调用流:在事件中区分工具调用和文本生成
- 前端适配:使用 Server-Sent Events(SSE)或 WebSocket
- 错误处理:流式传输中可能出现网络中断,需做重连处理
- 结合 checkpointer:流式输出 + 记忆保存实现完整的聊天体验
下一步
- Memory 记忆:结合流式输出实现持久化对话
- Callbacks 与追踪:使用 LangSmith 追踪流式执行
- RAG 应用设计:为 RAG 系统添加流式输出
- 测试与评估:测试流式输出的正确性