在构建基于大语言模型(LLM)的应用程序时,响应速度和用户体验至关重要。流式输出(Streaming)技术允许应用程序在 LLM 生成文本的每个 Token 完成后立即将其呈现给用户,而不是等待整个响应完成。LangGraph,作为构建有状态、多参与者 LLM 应用的强大框架,对流式输出提供了全面且深入的支持。
本文将详细介绍 LangGraph 的流式输出机制,涵盖其核心概念、API 使用方法,并深入探讨如“流式冒泡”等高级应用场景。

引言:为什么流式输出如此重要

想象一下,你正在使用一个 AI 助手来撰写一封长邮件或分析一份复杂的文档。如果 AI 在几秒钟内“沉默不语”,然后一次性吐出整个结果,体验会很糟糕。你会怀疑它是否在正常工作,并感到长时间的等待。
流式输出解决了这个问题。它就像一个打字员,在键盘上敲击每个字后立即显示在屏幕上。这带来了以下好处:

  • 更好的用户体验 (UX):用户可以即时看到 AI 的“思考”和输出过程,减少了等待的焦虑感。
  • 更快的感知响应速度:即使用户获得完整结果的时间没有变化,但看到第一个字符的时间大大缩短,让系统感觉更“快”。
  • 实时反馈与交互:流式输出是构建真正交互式应用的基础。例如,用户可以在 AI 生成过程中打断它,或者根据已生成的内容进行后续操作。
    LangGraph 的优势在于它不仅能流式传输 LLM 生成的单个 Token,还能流式传输整个图的状态更新,这对于构建复杂的、多步骤的智能体应用来说是无价的。

LangGraph 流式输出的核心概念

在深入了解 API 之前,我们需要区分 LangGraph 中两种不同的流式输出维度。

Token 级流式输出

这是最常见也最直观的流式方式。它指的是从 LLM 模型(如 OpenAI 的 GPT-4、Anthropic 的 Claude)中逐个 Token 地获取输出。

  • 特点:细粒度,逐字符/逐词。
  • 适用场景:聊天界面、代码生成、文本创作等任何需要实时显示生成内容的场景。
  • LangGraph 中的实现:LangGraph 通过 astream_events API,可以捕获来自图中任何节点的 LLM 模型生成的 on_chat_model_stream 事件。

状态更新流式输出

LangGraph 的核心是一个有向图,每个节点执行后会更新图的共享状态。状态更新流式输出是指在图执行的每个步骤(或每个节点完成后)将状态的变化流式传输出来。

  • 特点:粗粒度,以“步骤”或“节点”为单位。
  • 适用场景
    • 调试和监控:实时查看 Agent 的决策过程、使用了哪个工具、工具的输入输出是什么。
    • 复杂工作流:在一个多步骤的分析流程中,每完成一个分析步骤就向用户报告进度和中间结果。
  • LangGraph 中的实现:主要通过 astream API 实现,它可以配置为流式传输 updates(增量更新)或 values(当前完整状态)。

两者的差异与结合

特性 Token 级流式输出 状态更新流式输出
数据单元 单个 Token 节点的输出/状态的变化
实时性 极高 中等(取决于节点执行时间)
信息量 仅包含生成的文本内容 包含整个图的状态、工具调用结果、决策逻辑等
主要用途 终端用户体验 (UX) 开发者调试、复杂工作流状态展示

关键点:LangGraph 的强大之处在于你可以同时使用这两种流式输出。例如,在一个聊天应用中,你可以用 astream_events 来流式显示 LLM 正在生成的文本(为了 UX),同时用 astream 来在后台监控 Agent 的整个执行流程(为了调试和日志记录)。

LangGraph 流式输出 API

LangGraph 提供了几个核心方法来实现流式输出,其中 astream_events 是最新、最强大的。

astream_events: 最全面的事件流

astream_events 是 LangGraph 0.2 版本后引入的、推荐的流式事件 API。它提供了一个统一的、基于生成器的接口,可以访问图执行过程中发生的所有事件。

  • 方法签名async def astream_events(input, config, version="v1", *, interrupt_before=None, interrupt_after=None)
  • 核心参数
    • input: 图的初始输入。
    • config: 一个包含配置信息的字典,其中最重要的是 configurable 字段,可以用来传递 thread_id 等信息。
    • version: 协议版本,当前推荐使用 "v1"
  • 返回值:一个异步生成器,它会按顺序产生一系列事件对象。
    事件结构:每个事件都是一个字典,包含以下关键字段:
  • event: 事件类型,例如 on_chat_model_start, on_chat_model_stream, on_tool_start, on_tool_end
  • name: 触发事件的组件名称(如节点名称、模型名称)。
  • run_id: 当前运行的唯一 ID,用于追踪一次完整的图执行。
  • tags: 与事件关联的标签,可以用来过滤事件。
  • metadata: 事件的元数据,包含诸如 langgraph_step(当前步骤数)、langgraph_node(当前节点名)等信息。
  • data: 事件的具体数据,例如 on_chat_model_stream 事件中,data 包含一个 chunk 对象,其中就有生成的 Token。

astream: 简化的状态更新流

astream 是一个更高级的 API,专注于流式传输图的状态更新。它比 astream_events 更简单,当你只关心状态如何变化,而不关心底层的 Token 或工具调用细节时,它非常有用。

  • 方法签名async def astream(input, config, *, stream_mode="updates", interrupt_before=None, interrupt_after=None)
  • 核心参数 stream_mode:这是最重要的参数,它决定了流式输出的内容。
    • "updates" (默认):流式传输每个节点执行后产生的状态更新。这是一个字典,只包含被节点修改的状态键值对。
    • "values":流式传输每个节点执行后的完整状态。这是一个包含当前所有状态键值的字典。
    • "debug":流式传输用于调试的内部信息。
  • 返回值:一个异步生成器,产生 (node_name, state_update_or_value) 元组。

stream: astream 的同步版本

如果你在同步环境中工作,可以使用 stream,它是 astream 的同步包装器,用法和参数完全一致。


代码实战:构建一个支持流式输出的图

让我们通过一个简单的例子来演示如何使用这些 API。我们将构建一个简单的两节点图:一个节点生成一个故事大纲,另一个节点根据大纲扩展成一个完整的故事。

步骤 1: 定义状态和模型

首先,我们定义图的共享状态,并初始化一个 ChatModel。

import operator
from typing import Annotated, List, TypedDict
from langchain_core.messages import BaseMessage, HumanMessage
from langchain_openai import ChatOpenAI
from langgraph.graph import StateGraph, END
# 1. 定义状态
class AgentState(TypedDict):
    messages: Annotated[List[BaseMessage], operator.add]  # 消息列表,使用add操作符合并
# 2. 初始化模型
llm = ChatOpenAI(model="gpt-4o-mini", temperature=0)

步骤 2: 定义节点和边

现在,我们定义两个节点函数,并构建图的结构。

# 3. 定义节点函数
def outline_node(state: AgentState):
    """第一个节点:生成故事大纲"""
    print("---outline_node 正在执行---")
    messages = [HumanMessage(content="请为一个关于AI的故事生成一个简短的大纲。")]
    response = llm.invoke(messages)
    return {"messages": [response]} # 返回一个字典,更新状态
def story_node(state: AgentState):
    """第二个节点:根据大纲扩展成完整故事"""
    print("---story_node 正在执行---")
    # 获取上一个节点的大纲
    outline = state["messages"][-1].content
    messages = [
        HumanMessage(content=f"请根据以下大纲,写一个完整的AI故事:\n\n{outline}")
    ]
    response = llm.invoke(messages)
    return {"messages": [response]}
# 4. 构建图
workflow = StateGraph(AgentState)
# 添加节点
workflow.add_node("outline_agent", outline_node)
workflow.add_node("story_agent", story_node)
# 设置边
workflow.set_entry_point("outline_agent")
workflow.add_edge("outline_agent", "story_agent")
workflow.add_edge("story_agent", END)
# 5. 编译图
app = workflow.compile()

步骤 3: 编译图并实现流式输出

现在,我们使用 astreamastream_events 来流式执行这个图。

使用 astream 流式传输状态更新
import asyncio
async def main():
    print("====== 使用 astream (stream_mode='updates') ======")
    # 使用 'updates' 模式,只流式传输每个节点的状态更新
    async for event in app.astream({"messages": []}, stream_mode="updates"):
        for node_name, node_output in event.items():
            print(f"\n[节点 {node_name} 完成执行,产生更新]")
            print(f"更新内容: {node_output}")
    print("\n====== 使用 astream (stream_mode='values') ======")
    # 使用 'values' 模式,流式传输每个节点执行后的完整状态
    async for event in app.astream({"messages": []}, stream_mode="values"):
        print(f"\n[当前步骤完成,完整状态中的消息数量: {len(event['messages'])}]")
        # print(event) # 如果取消注释,将打印整个状态
if __name__ == "__main__":
    asyncio.run(main())

预期输出(简化版)

====== 使用 astream (stream_mode='updates') ======
[节点 outline_agent 完成执行,产生更新]
更新内容: {'messages': [AIMessage(content='一个关于AI的觉醒与自我探索的故事...')]}
[节点 story_agent 完成执行,产生更新]
更新内容: {'messages': [AIMessage(content='在未来的某一天,一个名为“艾瑞安”的AI...')]}
====== 使用 astream (stream_mode='values') ======
[当前步骤完成,完整状态中的消息数量: 1]
[当前步骤完成,完整状态中的消息数量: 2]
使用 astream_events 流式传输 Token

这是实现聊天界面的关键。

async def print_token_stream():
    print("\n====== 使用 astream_events 流式传输 Token ======")
    # 使用 astream_events 来获取所有事件
    async for event in app.astream_events({"messages": []}, version="v1"):
        kind = event["event"]
        # 过滤出 chat_model 的流式输出事件
        if kind == "on_chat_model_stream":
            content = event["data"]["chunk"].content
            if content:
                # 实时打印每个 Token,用 end="" 防止换行
                print(content, end="", flush=True)
        # 也可以选择性地打印工具调用等事件
        elif kind == "on_tool_start":
            print(f"\n[工具调用开始]: {event['name']}")
        elif kind == "on_tool_end":
            print(f"\n[工具调用结束]: {event['name']}")
if __name__ == "__main__":
    asyncio.run(print_token_stream())

预期输出(模拟)

====== 使用 astream_events 流式传输 Token ======
一
个
关
于
AI
的
觉
醒
...
[这里会逐字打印出完整的故事]
...
。

进阶场景:流式冒泡机制

在一个复杂的图中,你可能有多个分支并行执行,或者一个节点的输出需要被另一个节点处理。在这些情况下,如何优雅地将不同来源的流式输出呈现给用户,是一个挑战。这就是“流式冒泡”机制要解决的问题。

什么是流式冒泡?

“流式冒泡”不是一个 LangGraph 的官方术语,而是一种设计模式。它指的是将来自图中不同节点、不同层级的流式输出(通常是 Token)有序地、聚合地“冒泡”到顶层,最终呈现给用户。
想象一个气泡从水底升起:它可能来自不同的源头,但最终都会冒出水面。在我们的场景中,“水面”就是用户界面。

如何实现流式冒泡

实现流式冒泡的核心思路是:

  1. 使用 astream_events 作为唯一数据源:它能捕获所有节点的所有事件。
  2. 在应用层进行聚合和排序:你需要编写逻辑来处理这些事件流。
  3. 使用元数据来区分来源:利用事件中的 metadata 字段(如 langgraph_node)来知道 Token 来自哪个节点。
  4. 管理输出顺序:这是最棘手的部分。如果节点 A 和节点 B 并行执行,它们的 Token 会交替到达。你需要决定是:
    • 按节点分组:先显示完节点 A 的所有输出,再显示节点 B 的。
    • 按时间顺序混合:就像两个人同时说话,交替显示它们的输出。
    • 优先级冒泡:某个节点的输出(如最终答案)优先级更高,优先显示。

代码示例:实现多分支流式冒泡

让我们修改之前的图,让它并行生成两个不同风格的故事开头,然后将这些流式输出“冒泡”到控制台。

from langgraph.graph import StateGraph, END
# 1. 定义状态 (与之前相同)
class AgentState(TypedDict):
    messages: Annotated[List[BaseMessage], operator.add]
llm = ChatOpenAI(model="gpt-4o-mini", temperature=0)
# 2. 定义两个并行的故事生成节点
def scifi_story_node(state: AgentState):
    print("---scifi_story_node 正在执行---")
    messages = [HumanMessage(content="写一个科幻风格的故事开头,大约50字。")]
    response = llm.invoke(messages)
    return {"messages": [response]}
def fairytale_story_node(state: AgentState):
    print("---fairytale_story_node 正在执行---")
    messages = [HumanMessage(content="写一个童话风格的故事开头,大约50字。")]
    response = llm.invoke(messages)
    return {"messages": [response]}
# 3. 构建一个并行图
workflow = StateGraph(AgentState)
workflow.add_node("scifi_agent", scifi_story_node)
workflow.add_node("fairytale_agent", fairytale_story_node)
workflow.set_entry_point("scifi_agent")
# 让两个节点并行执行。LangGraph 会从入口点开始,同时执行所有没有依赖的节点。
# 在这个简单例子中,我们需要一个虚拟的入口来触发并行。
# 更好的做法是使用 `Send` API 或者一个条件边。为了简单起见,我们用条件边。
from langgraph.graph import START
def route_to_parallel(state):
    # 这个函数只是简单地触发两个节点
    return ["scifi_agent", "fairytale_agent"]
# 更简单的并行触发方式是直接在 START 后连接两个节点
workflow.add_edge(START, "scifi_agent")
workflow.add_edge(START, "fairytale_agent")
workflow.add_edge("scifi_agent", END)
workflow.add_edge("fairytale_agent", END)
app = workflow.compile()
async def stream_with_bubbling():
    print("\n====== 并行节点的流式冒泡 ======")
    # 用于存储每个节点的输出字符串
    node_outputs = {
        "scifi_agent": "",
        "fairytale_agent": ""
    }
    async for event in app.astream_events({"messages": []}, version="v1"):
        kind = event["event"]
        if kind == "on_chat_model_stream":
            node_name = event["metadata"]["langgraph_node"]
            content = event["data"]["chunk"].content
            if content and node_name in node_outputs:
                # 将 Token 追加到对应节点的字符串中
                node_outputs[node_name] += content
                # --- 冒泡逻辑 ---
                # 这里我们选择“按时间顺序混合”的冒泡方式
                print(f"[{node_name}]: {content}", end="", flush=True)
                
    print("\n\n====== 最终聚合结果 ======")
    for node, text in node_outputs.items():
        print(f"{node} 的完整输出: {text}")
if __name__ == "__main__":
    asyncio.run(stream_with_bubbling())

预期输出(模拟)

====== 并行节点的流式冒泡 ======
[scifi_agent]: 在 [fairytale_agent]: 从 [scifi_agent]: 2 [fairytale_agent]: 前 [scifi_agent]: 0 [fairytale_agent]: , [scifi_agent]: 4 [fairytale_agent]: 有 [scifi_agent]: 2 [fairytale_agent]: 一 [scifi_agent]: 年 [fairytale_agent]: 个...
...
[scifi_agent]: 星际飞船。
[fairytale_agent]: 魔法森林。
====== 最终聚合结果 ======
scifi_agent 的完整输出: 在2042年,星际飞船...
fairytale_agent 的完整输出: 从前,有一个魔法森林...

这个例子展示了如何将来自两个并行节点的流式输出交错地“冒泡”到用户界面。你可以根据需要修改冒泡逻辑,例如,先收集完一个节点的所有输出再显示下一个,或者根据用户的选择优先显示某个节点的输出。


实践技巧

  1. 始终使用异步 API (astream, astream_events):流式输出是一个 I/O 密集型操作。在 Python 中,使用 async/await 可以显著提高性能,尤其是在处理多个并发流时。
  2. 理解 stream_mode 的选择
    • 如果你的目标是构建一个聊天 UI,请使用 astream_events 来获取 Token。
    • 如果你的目标是构建一个进度展示或调试工具,astream 配合 stream_mode="values" 是更好的选择。
  3. 善用 tagsmetadata:在创建图或调用 LLM 时,可以添加自定义 tags。在 astream_events 中,你可以根据这些 tags 来过滤事件,只处理你关心的部分。
  4. 处理流式输出的中断:LangGraph 支持在图的特定节点之前或之后中断执行 (interrupt_before, interrupt_after)。这对于需要人工介入的流程非常有用。在流式输出时,你可以在用户界面上提供一个“停止”按钮,当被点击时,可以设置一个标志来中止流式生成循环。
  5. 为不同的事件类型设计 UI:不要只把所有 Token 混在一起显示。
    • on_chat_model_stream: 显示为普通文本。
    • on_tool_start: 可以显示为一个“正在使用工具 X…”的加载动画。
    • on_tool_end: 可以显示工具调用的结果,例如,如果工具是搜索,可以展示搜索结果的摘要。
  6. 在生产环境中使用回调处理器astream_events 非常适合调试和原型开发。但在生产环境中,你可能会使用 LangChain 的回调处理器 (Callbacks),如 StreamingStdOutCallbackHandler 或自定义的回调处理器,它们可以更紧密地集成到你的日志和监控系统中。

LangGraph 正在快速发展,其流式输出能力也在不断增强。未来的改进可能包括:

  • 更精细的事件过滤:提供更强大的 API 来订阅特定类型、特定节点或特定标签的事件,而不是在客户端手动过滤。
  • 原生支持 Server-Sent Events (SSE) 和 WebSocket:LangGraph 可能会提供更直接的集成,使得将流式输出推送到前端变得更简单。
  • 与前端框架的深度集成:可能会出现与 React, Vue 等前端框架更紧密绑定的库,自动处理流的订阅、聚合和 UI 更新。

总结

流式输出是构建现代、响应式 LLM 应用的核心技术。LangGraph 通过 astream_eventsastream 这两个强大而灵活的 API,为开发者提供了对整个图执行过程的细粒度控制。
本文从基础概念入手,详细介绍了 Token 级和状态级流式输出的区别,并通过代码示例展示了如何使用核心 API。我们还深入探讨了“流式冒泡”这一高级模式,展示了如何在复杂的并行图中优雅地处理和呈现流式数据。
掌握 LangGraph 的流式输出机制,将帮助你打造出不仅智能,而且快速、流畅、用户友好的 AI 应用。随着 LangGraph 的不断演进,其流式能力必将变得更加强大和易用,为下一代 AI 体验奠定基础。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐