9. 流式与实时交互

9.流式与实时交互

LangGraph 的流式与实时交互机制,本质上是为了解决一个核心问题:如何让 AI 应用的响应过程对用户可见。传统的请求-响应模式像在黑箱中操作,用户只能等待最终结果。而流式交互则像打开了一扇窗,让整个过程变得透明、可感知。

这种透明性在现代 AI 应用中至关重要。当 LLM 生成一段长文本时,用户希望看到文字逐字出现;当工具在后台执行耗时操作时,用户需要进度提示;当多步骤工作流运行时,开发者需要调试信息。LangGraph 的流式系统正是为满足这些需求而设计。

流模式类型配置

LangGraph 提供了六种流模式,每种模式对应不同的数据粒度。理解它们的区别,是合理使用流式功能的第一步。

六种流模式详解

values 模式会推送完整的图状态。每次节点执行后,整个状态对象都会被发送。这种模式适合需要完整上下文的场景,比如调试时查看所有变量。但缺点是数据量大,网络开销高。

updates 模式只推送状态的变化部分。它返回节点名称和该节点产生的更新,格式为 {'node_name': {'field': 'new_value'}}。这是大多数应用的首选模式,既保持了信息完整性,又避免了冗余数据传输。

messages-tuple(或 messages)模式专门用于流式传输 LLM 生成的 Token。它返回一个元组 (message_chunk, metadata),其中 message_chunk 包含 Token 内容,metadata 包含节点信息、标签等上下文。构建聊天界面时,这个模式几乎是必选项。

custom 模式允许从图的任意位置发送自定义数据。通过 get_stream_writer() 获取写入器,可以推送任意 JSON 可序列化的数据。这在工具执行进度、业务状态通知等场景中非常有用。

debug 模式会推送尽可能多的执行信息,包括节点名称、完整状态、执行时间等。它主要用于开发和调试阶段,生产环境应谨慎使用,避免性能损耗。

events 模式会流式传输所有事件,包括图状态变化、LLM 调用、工具执行等。这个模式主要用于迁移旧版 LCEL 应用,新应用通常不需要。

模式组合使用

实际应用中,往往需要同时获取多种类型的数据。LangGraph 支持传入模式列表来实现这一点:

# Python 示例:同时获取状态更新和 LLM Token
for mode, chunk in agent.stream(
    {"messages": [{"role": "user", "content": "讲个笑话"}]},
    stream_mode=["updates", "messages-tuple"]
):
    if mode == "updates":
        print(f"状态更新: {chunk}")
    elif mode == "messages-tuple":
        message_chunk, metadata = chunk
        print(f"Token: {message_chunk.content}")
// JavaScript 示例:组合模式
for await (const [mode, chunk] of await agent.stream(
  { messages: [{ role: "user", content: "讲个笑话" }] },
  { streamMode: ["updates", "messages-tuple"] }
)) {
  if (mode === "updates") {
    console.log(`状态更新: ${JSON.stringify(chunk)}`);
  } else if (mode === "messages-tuple") {
    const [messageChunk, metadata] = chunk;
    console.log(`Token: ${messageChunk.content}`);
  }
}

代码中,stream_mode 参数接受列表形式,流式输出会返回 (mode, chunk) 元组。通过判断 mode 值,可以区分不同类型的数据。这种模式组合让前端能够同时渲染对话历史和实时生成的文本。

LLM Token 实时推送

Token 流式传输是聊天应用的核心功能。LangGraph 的实现有几个关键点值得注意。

基本用法

即使使用 .invoke() 而非 .stream() 调用 LLM,只要设置了 stream_mode="messages-tuple",Token 依然会被流式传输。这是因为 LangGraph 在运行时层面拦截了 LLM 调用,自动处理流式逻辑。

# Python 示例:Token 流式传输
from langchain.chat_models import init_chat_model
from langgraph.graph import StateGraph, START

llm = init_chat_model(model="openai:gpt-4o-mini")

def call_model(state):
    # 使用 .invoke() 而非 .stream()
    response = llm.invoke([{"role": "user", "content": f"讲个关于{state['topic']}的笑话"}])
    return {"joke": response.content}

graph = StateGraph(State).add_node(call_model).add_edge(START, "call_model").compile()

# 但依然可以流式接收 Token
async for message_chunk, metadata in graph.astream(
    {"topic": "猫"},
    stream_mode="messages-tuple"
):
    if message_chunk.content:
        print(message_chunk.content, end="|", flush=True)
// JavaScript 示例
import { ChatOpenAI } from "@langchain/openai";

const llm = new ChatOpenAI({ model: "gpt-4o-mini" });

const callModel = async (state: any) => {
  const response = await llm.invoke([
    { role: "user", content: `讲个关于${state.topic}的笑话` }
  ]);
  return { joke: response.content };
};

// 流式接收 Token
for await (const [messageChunk, metadata] of await graph.stream(
  { topic: "猫" },
  { streamMode: "messages-tuple" }
)) {
  if (messageChunk.content) {
    process.stdout.write(messageChunk.content + "|");
  }
}

这种设计简化了开发流程。无需在业务逻辑中考虑流式细节,只需在调用图时指定模式即可。.invoke() 和 .stream() 的区别仅在于调用方式,流式能力由运行时统一提供。

过滤特定 LLM 调用

复杂应用中可能包含多个 LLM 调用。通过 tags 参数可以区分它们:

# Python 示例:为不同 LLM 添加标签
joke_model = init_chat_model(model="openai:gpt-4o-mini", tags=['joke'])
poem_model = init_chat_model(model="openai:gpt-4o-mini", tags=['poem'])

async def call_model(state, config):
    topic = state["topic"]
    # 生成笑话
    joke_response = await joke_model.ainvoke(
        [{"role": "user", "content": f"写个关于{topic}的笑话"}],
        config
    )
    # 生成诗歌
    poem_response = await poem_model.ainvoke(
        [{"role": "user", "content": f"写首关于{topic}的诗"}],
        config
    )
    return {"joke": joke_response.content, "poem": poem_response.content}

# 只接收笑话模型的 Token
async for msg, metadata in graph.astream(
    {"topic": "猫"},
    stream_mode="messages-tuple"
):
    if metadata["tags"] == ["joke"]:
        print(msg.content, end="|", flush=True)
// JavaScript 示例
const jokeModel = new ChatOpenAI({ model: "gpt-4o-mini", tags: ["joke"] });
const poemModel = new ChatOpenAI({ model: "gpt-4o-mini", tags: ["poem"] });

// 只接收笑话模型的 Token
for await (const [msg, metadata] of await graph.stream(
  { topic: "猫" },
  { streamMode: "messages-tuple" }
)) {
  if (metadata.tags?.includes("joke")) {
    process.stdout.write(msg.content + "|");
  }
}

tags 参数在初始化模型时设置,会包含在元数据中。通过过滤 metadata["tags"],可以精确控制显示哪个 LLM 的输出。这在多模型协作的场景中特别有用,比如一个模型生成内容,另一个模型评估质量。

按节点过滤

除了按标签过滤,还可以按节点名称过滤:

# Python 示例:只接收特定节点的 Token
for msg, metadata in graph.stream(
    inputs,
    stream_mode="messages-tuple"
):
    # 检查节点名称
    if metadata["langgraph_node"] == "write_poem":
        print(msg.content, end="|", flush=True)
// JavaScript 示例
for await (const [msg, metadata] of await graph.stream(
  inputs,
  { streamMode: "messages-tuple" }
)) {
  if (metadata.langgraph_node === "writePoem") {
    process.stdout.write(msg.content + "|");
  }
}

metadata["langgraph_node"] 字段记录了产生该 Token 的节点名称。这在并行执行多个 LLM 调用时非常实用,可以分别渲染不同节点的输出。

工具执行进度流

工具执行可能耗时较长,比如查询数据库、调用外部 API。让用户感知到进度,能显著提升体验。

在工具中发送进度

Python 中通过 get_stream_writer() 获取写入器,JavaScript 中通过 config.writer 实现:

# Python 示例:工具进度流
from langgraph.config import get_stream_writer
from langchain_core.tools import tool

@tool
def query_database(query: str) -> str:
    """查询数据库"""
    writer = get_stream_writer()
    
    # 发送进度更新
    writer({"type": "progress", "data": "已检索 0/100 条记录"})
    
    # 模拟耗时操作
    # ... 执行查询 ...
    
    writer({"type": "progress", "data": "已检索 50/100 条记录"})
    # ... 更多操作 ...
    
    writer({"type": "progress", "data": "已检索 100/100 条记录"})
    
    return "查询完成"

# 使用 custom 模式接收进度
for chunk in graph.stream(
    {"query": "SELECT * FROM users"},
    stream_mode="custom"
):
    if chunk.get("type") == "progress":
        print(f"进度: {chunk['data']}")
// JavaScript 示例
import { tool } from "@langchain/core/tools";

const queryDatabase = tool(
  async (input: any, config: any) => {
    // 发送进度更新
    config.writer({ type: "progress", data: "已检索 0/100 条记录" });
    
    // 模拟耗时操作
    // ... 执行查询 ...
    
    config.writer({ type: "progress", data: "已检索 50/100 条记录" });
    // ... 更多操作 ...
    
    config.writer({ type: "progress", data: "已检索 100/100 条记录" });
    
    return "查询完成";
  },
  {
    name: "query_database",
    description: "查询数据库",
    schema: z.object({ query: z.string() })
  }
);

// 接收进度
for await (const chunk of await graph.stream(
  { query: "SELECT * FROM users" },
  { streamMode: "custom" }
)) {
  if (chunk.type === "progress") {
    console.log(`进度: ${chunk.data}`);
  }
}

进度数据可以是任意结构。示例中使用 {"type": "progress", "data": "..."} 格式,前端可以根据 type 字段区分消息类型,渲染不同的 UI 组件。

注意事项

在 Python 异步代码中(Python < 3.11),get_stream_writer() 无法自动获取上下文。需要显式传递 writer 参数:

# Python 3.10 及以下版本需要这样做
from langgraph.types import StreamWriter

async def query_database(query: str, writer: StreamWriter) -> str:
    writer({"type": "progress", "data": "开始查询"})
    # ... 异步操作 ...
    return "结果"

# 图定义中不需要特殊处理
graph = StateGraph(State).add_node(query_database).compile()

LangGraph 会自动检测函数签名中的 writer 参数并注入。这个限制在 Python 3.11+ 中已解决,因为 asyncio 支持了上下文传播。

自定义事件数据流

除了进度,还可以发送任意业务相关的事件数据。比如,在电商客服机器人中,可以推送订单状态更新;在数据分析工具中,可以发送中间计算结果。

节点内发送自定义数据

# Python 示例:节点内自定义事件
from langgraph.config import get_stream_writer

def process_order(state):
    writer = get_stream_writer()
    
    # 发送订单确认
    writer({"event": "order_confirmed", "order_id": state["order_id"]})
    
    # 模拟处理
    # ... 处理逻辑 ...
    
    # 发送发货通知
    writer({"event": "order_shipped", "tracking_number": "12345"})
    
    return {"status": "completed"}

# 接收自定义事件
for chunk in graph.stream(
    {"order_id": "ORD-001"},
    stream_mode="custom"
):
    if chunk.get("event") == "order_confirmed":
        print(f"订单已确认: {chunk['order_id']}")
    elif chunk.get("event") == "order_shipped":
        print(f"订单已发货,追踪号: {chunk['tracking_number']}")
// JavaScript 示例
import { LangGraphRunnableConfig } from "@langchain/langgraph";

const processOrder = async (state: any, config: LangGraphRunnableConfig) => {
  config.writer({ event: "order_confirmed", order_id: state.order_id });
  
  // ... 处理逻辑 ...
  
  config.writer({ event: "order_shipped", tracking_number: "12345" });
  
  return { status: "completed" };
};

// 接收自定义事件
for await (const chunk of await graph.stream(
  { order_id: "ORD-001" },
  { streamMode: "custom" }
)) {
  if (chunk.event === "order_confirmed") {
    console.log(`订单已确认: ${chunk.order_id}`);
  } else if (chunk.event === "order_shipped") {
    console.log(`订单已发货,追踪号: ${chunk.tracking_number}`);
  }
}

自定义事件让图的状态机具备了事件驱动能力。前端可以订阅特定事件,触发相应的 UI 更新,比如显示通知、更新进度条、切换页面等。

与非 LangChain LLM 集成

如果使用的 LLM 没有 LangChain 集成,可以通过 custom 模式手动流式传输:

# Python 示例:集成自定义 LLM
from langgraph.config import get_stream_writer
from openai import AsyncOpenAI

client = AsyncOpenAI()

async def stream_custom_llm(state):
    writer = get_stream_writer()
    
    response = await client.chat.completions.create(
        model="gpt-4o-mini",
        messages=[{"role": "user", "content": state["query"]}],
        stream=True
    )
    
    full_content = ""
    async for chunk in response:
        if chunk.choices[0].delta.content:
            content = chunk.choices[0].delta.content
            full_content += content
            writer({"custom_llm_chunk": content})  # 手动发送 Token
    
    return {"response": full_content}

# 接收自定义 LLM 的 Token
for chunk in graph.stream(
    {"query": "讲个笑话"},
    stream_mode="custom"
):
    if "custom_llm_chunk" in chunk:
        print(chunk["custom_llm_chunk"], end="|", flush=True)

这种方式提供了最大灵活性。无论 LLM 提供何种接口,只要支持流式响应,就可以集成到 LangGraph 的流式体系中。

React 应用流集成

前端集成是流式功能的最终落地点。LangGraph 提供了 useStream Hook,大幅简化了 React 应用的开发。

基础用法

// React 组件示例
"use client";

import { useStream } from "@langchain/langgraph-sdk/react";
import type { Message } from "@langchain/langgraph-sdk";

export default function ChatApp() {
  const thread = useStream<{ messages: Message[] }>({
    apiUrl: "http://localhost:2024",
    assistantId: "agent",
    messagesKey: "messages",
  });

  return (
    <div>
      {/* 渲染消息历史 */}
      <div>
        {thread.messages.map((message) => (
          <div key={message.id}>
            <strong>{message.type}:</strong> {message.content as string}
          </div>
        ))}
      </div>

      {/* 输入表单 */}
      <form
        onSubmit={(e) => {
          e.preventDefault();
          const form = e.target as HTMLFormElement;
          const message = new FormData(form).get("message") as string;
          form.reset();
          
          // 提交消息并自动处理流式响应
          thread.submit({ messages: [{ type: "human", content: message }] });
        }}
      >
        <input type="text" name="message" placeholder="输入消息..." />
        
        {/* 根据加载状态显示不同按钮 */}
        {thread.isLoading ? (
          <button type="button" onClick={() => thread.stop()}>
            停止
          </button>
        ) : (
          <button type="submit">发送</button>
        )}
      </form>
    </div>
  );
}

useStream Hook 封装了所有复杂逻辑:创建线程、管理状态、处理流式响应、更新消息列表。开发者只需关注 UI 渲染。

关键特性

自动消息管理:thread.messages 会自动更新。当 LLM 生成新 Token 时,Hook 会合并到最新消息中,触发组件重新渲染。

加载状态:thread.isLoading 指示是否有正在进行的流式请求。可以用来禁用输入框、显示加载动画。

中断支持:thread.stop() 可以取消当前流式请求。这在用户想停止生成时很有用。

分支管理:thread.branch 和 thread.setBranch() 支持对话分支。用户可以回到历史消息,从该点开启新对话分支。

错误处理:thread.error 包含任何发生的错误。可以在 UI 中显示错误提示。

事件回调

对于更细粒度的控制,可以使用回调函数:

// 高级事件处理
const thread = useStream<{ messages: Message[] }>({
  apiUrl: "http://localhost:2024",
  assistantId: "agent",
  messagesKey: "messages",
  
  // 自定义事件
  onCustomEvent: (data) => {
    if (data.type === "progress") {
      console.log(`进度: ${data.data}`);
    }
  },
  
  // 错误处理
  onError: (error) => {
    console.error("流式错误:", error);
    alert("发生错误,请重试");
  },
  
  // 完成回调
  onFinish: (state) => {
    console.log("流式完成:", state);
  },
  
  // 元数据事件(包含 runId 和 threadId)
  onMetadataEvent: (data) => {
    console.log("运行元数据:", data);
  }
});

这些回调让应用能够响应各种事件,实现进度条、错误提示、日志记录等功能。

断线重连

useStream 支持自动重连,这在页面刷新后恢复对话时特别有用:

// 自动重连配置
const thread = useStream<{ messages: Message[] }>({
  apiUrl: "http://localhost:2024",
  assistantId: "agent",
  reconnectOnMount: true,  // 启用自动重连
});

// 自定义存储
const thread = useStream<{ messages: Message[] }>({
  apiUrl: "http://localhost:2024",
  assistantId: "agent",
  reconnectOnMount: () => window.localStorage,  // 使用 localStorage
});

默认使用 sessionStorage 存储运行 ID。启用后,组件挂载时会自动尝试恢复未完成的流式请求,确保不丢失任何消息。

TypeScript 类型支持

对于 LangGraph.js 用户,可以复用状态类型定义:

import {
  Annotation,
  MessagesAnnotation,
  type StateType,
  type UpdateType,
} from "@langchain/langgraph/web";

const AgentState = Annotation.Root({
  ...MessagesAnnotation.spec,
  context: Annotation<string>(),
});

const thread = useStream<
  StateType<typeof AgentState.spec>,
  { UpdateType: UpdateType<typeof AgentState.spec> }
>({
  apiUrl: "http://localhost:2024",
  assistantId: "agent",
  messagesKey: "messages",
});

这种方式保证了类型安全,避免手动定义接口的繁琐和错误。

总结

LangGraph 的流式系统提供了从底层到前端的全栈解决方案。六种流模式覆盖了从调试到生产的各种需求,Token 流式传输让 LLM 响应实时可见,工具进度流提升了长操作的体验,自定义事件为业务逻辑提供了事件驱动能力,React 集成则让前端开发变得简单高效。

流式交互不仅是技术优化,更是用户体验的升级。它让 AI 应用从"黑箱"变为"白箱",从"等待"变为"参与"。在下一章中,我们将探讨 LangGraph 的持久化与时间旅行功能,看看如何让对话状态可追溯、可回退,进一步提升应用的可靠性和可控性。