Skip to content

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调试信息开发和调试

最佳实践

  1. 启用 streaming:创建 LLM 时设置 streaming=True
  2. 选择合适的 stream_mode:聊天用 messages,监控用 events
  3. 处理工具调用流:在事件中区分工具调用和文本生成
  4. 前端适配:使用 Server-Sent Events(SSE)或 WebSocket
  5. 错误处理:流式传输中可能出现网络中断,需做重连处理
  6. 结合 checkpointer:流式输出 + 记忆保存实现完整的聊天体验

下一步

本站为非官方中文学习站点,不代表 LangChain 官方。部分内容参考官方文档并重新整理为中文学习笔记。