Jean's Blog

一个专注软件测试开发技术的个人博客

0%

LangGraph之流式机制

介绍

LangGraph 实现了一个流式传输系统,用于实时更新信息。对于基于大型语言模型(LLMs)构建的应用程序来说,流式传输对于增强其响应性至关重要。这是因为大型语言模型在处理复杂的自然语言任务时,可能会存在一定的延迟。而通过流式传输,可以在模型还在处理数据的过程中,就将部分结果逐步展示给用户,而不是让用户等待一个完整的响应,从而让用户能够更快地获得反馈,提高了应用程序的交互性。

官方文档:https://docs.langchain.com/oss/python/langgraph/streaming

  • Stream graph state(流式传输图状态):可以通过“updates”和“values”模式获取状态更新/值。其中,“updates”模式会流式传输图在每一步之后的状态更新,如果在同一个步骤中进行了多次更新(例如,运行了多个节点),这些更新将分别流式传输;而“values”模式则会流式传输图在每一步之后的完整状态值。
  • Stream subgraph outputs(流式传输子图输出):包括父图和任何嵌套子图的输出。通过在父图的.stream()方法中设置subgraphs=True,可以实现这一点。输出将以(namespace, data)元组的形式流式传输,其中namespace是一个包含调用子图的节点路径的元组,例如(“parent_node:“, “child_node:“)。
  • Stream LLM tokens(流式传输LLM令牌):可以从任何地方捕获令牌流,包括节点内部、子图内部或工具内部。这使得能够实时获取大型语言模型(LLM)的输出令牌,从而提升应用的响应性和用户体验。
  • Stream custom data(流式传输自定义数据):可以直接从工具函数发送自定义更新或进度信号。这为开发者提供了更大的灵活性,可以根据需要自定义流式传输的数据内容,以满足特定的应用需求。
  • Use multiple streaming modes(使用多种流式传输模式):可以选择以下多种流式传输模式:
    • values(完整状态):流式传输每一步之后图的完整状态值。
    • updates(状态增量):流式传输每一步之后图的状态更新。
    • messages(LLM令牌+元数据):流式传输2元组(LLM令牌,元数据),这些元数据包含有关图节点和LLM调用的详细信息。
    • custom(任意用户数据):流式传输任意用户定义的数据。
    • debug(详细跟踪):在图执行过程中尽可能多地流式传输信息,包括节点名称和完整状态等详细信息。

流式处理如何重新定义 AI 交互体验

核心是围绕「流式处理(Streaming)」,讲它对 AI 用户体验的革命性影响,以及 LangGraph 在这方面的创新方案。

流式处理的革命性意义

核心观点

  • 流式处理,已经成为 AI 应用开发中提升用户体验的关键技术
  • 传统的「全量输出」模式,是让用户等待整个响应完全生成后,再一次性展示结果
  • 这种模式在处理复杂任务(比如长文本生成、多步骤推理)时,会产生明显的延迟,用户体验非常差。

通俗理解

你可以把它想象成:

  • 传统模式:点一份外卖,必须等所有菜做好、打包、骑手送到,你才能一口吃到饭。
  • 流式处理:厨师做好一道菜就先给你上一道,边做边吃,不用干等着。

LangGraph 的创新解决方案

  1. 它解决了什么?

LangGraph 通过创新的多模式流式处理机制,解决了传统 AI 应用的延迟痛点,让开发者能构建「真正实时、响应迅速」的 AI 系统。

  1. 核心价值是什么?

LangGraph 的流式处理,核心在于把应用的运行状态,抽象成了一个「可观测的图结构」

  • 图里的每一个节点(比如一个工具调用、一次 LLM 生成、一次分支判断),执行结果都能实时反馈给客户端
  • 开发者可以根据场景,自由选择 3 种流式传输方案:
    1. 完整状态快照:每次节点执行完,把整个应用的最新状态一次性推给前端。
    2. 增量更新:只把节点执行后变化的那部分状态,推送给前端。
    3. Token 级输出:直接捕获 LLM 生成文本时的每一个 token,实现打字机效果。

通俗理解

LangGraph 就像一个「智能厨房管理系统」:

  • 它把做菜的每个步骤(节点)都可视化了,厨师做完一步,系统就立刻通知你。
  • 你可以自己选怎么接收信息:
    • 要么每次都看完整的进度报告(完整快照);
    • 要么只看这次又做了哪道菜(增量更新);
    • 甚至直接看厨师切菜、炒菜的每一个动作细节(Token 级输出)。

LangGraph 把流式处理玩出了新高度,它让 AI 应用不再是 “等半天再出结果”,而是可以边算边反馈,用户体验和交互方式都被彻底改变了

messages模式:实现打字机效果的核心

LangGraph 里最常用的流式处理模式 ——messages 模式,它是实现 AI 对话里 “打字机逐字输出” 效果的关键。

模式定义

  • 核心定位messages 模式是 LangGraph 流式处理中最直观、最常用的模式。

  • 核心用途:专门用于流式传输大语言模型(LLM)生成的令牌(Token),也就是让文本一个字一个字地 “打出来”。

  • 和其他模式的区别:它关注的是最终内容的逐字输出,而不是像 values/updates 模式那样,关注中间状态的变更。

代码示例

1
2
3
4
5
# ✅ 正确写法:正确解包元组
for message_chunk, metadata in agent.stream(initial_state, stream_mode="messages"):
# 检查对象是否是预期的消息类型并有内容属性
if message_chunk.content:
print(message_chunk.content, end='', flush=True)
  1. agent.stream(initial_state, stream_mode="messages")
    • 这是开启流式处理的核心方法,指定 stream_mode="messages" 就进入了 messages 模式。
    • 它会返回一个可迭代的流,每次迭代会给你一对数据。
  2. for message_chunk, metadata in ...
    • 这是解包元组的关键操作。每次迭代返回的是一个包含两个元素的元组:
      • message_chunk:LLM 生成的单块消息内容(也就是单个 / 多个 Token 组成的片段)。
      • metadata:和这条消息相关的元数据(比如是哪个节点生成的、时间戳等)。
  3. if message_chunk.content:
    • 这是一个安全判断,避免处理没有文本内容的空消息块(比如纯工具调用的消息)。
  4. print(message_chunk.content, end='', flush=True)
    • end='':让打印不换行,实现连续输出。
    • flush=True:强制把内容立刻输出到控制台,而不是等缓冲区满了再一次性输出,这样才能看到实时的打字机效果。

关键点

  • stream_mode="messages" 模式下,每次迭代返回的是一个包含两个元素的元组 (message_chunk, metadata)

  • 必须正确解包这两个元素,才能访问到 message_chunk 里的 .content 属性,否则会直接报错。

  • 很多新手的坑就是忘了解包,直接用 for chunk in agent.stream(...),结果拿到的是整个元组,访问 .content 就会失败。

stream_mode="messages" 开启逐字流,然后通过解包 (message_chunk, metadata) 拿到每个 Token 片段,再配合 print(..., flush=True) 实现打字机效果。

异步流式处理优化

  • 适用场景:Web 应用等高并发场景(比如同时给多个用户提供对话服务)。

  • 推荐方案:使用异步流式处理方法,也就是用 Python 的 async/await 语法来实现。

  • 为什么要用异步?

​ 同步的 agent.stream() 在处理请求时会阻塞线程,高并发下很容易卡死;而异步的 agent.astream() 可以让单个线程同时处理多个请求,大幅提升服务的并发能力和响应速度。

1
2
3
4
5
6
7
8
async def stream_agent_response():
"""异步流式处理AI代理的响应"""
initial_state = {"messages": [{"role": "user", "content": "请帮我写一首诗"}]}

async for message_chunk, metadata in agent.astream(initial_state, stream_mode="messages"):
if message_chunk.content:
# 在实际应用中,这里会将内容实时推送到前端
yield message_chunk.content
  1. async def stream_agent_response()

    • 定义了一个异步生成器函数,函数里可以使用 awaitasync for

    • 这种函数通常会被 FastAPI、Starlette 等 Web 框架用来实现 SSE(Server-Sent Events)流式接口。

  2. initial_state = {"messages": [...]}

    • 初始化对话状态,这里模拟用户发送了一句 “请帮我写一首诗” 的请求。
  3. async for message_chunk, metadata in agent.astream(...)

    • agent.astream() 是 LangGraph 提供的异步流式方法,对应同步的 agent.stream()

    • async for 用来异步迭代流数据,在等待 LLM 生成 Token 的间隙,线程可以去处理其他请求,不会被阻塞。

    • 和同步模式一样,每次迭代返回的是 (message_chunk, metadata) 元组,需要解包才能拿到消息内容。

  4. if message_chunk.content:

    • 安全判断,过滤掉没有文本内容的消息块(比如纯工具调用的消息)。
  5. yield message_chunk.content

    • yield 把生成的文本片段一个个 “吐出来”,形成一个异步生成器。

    • 在实际 Web 应用中,框架会把这些片段实时推送给前端,实现打字机效果。

在 Web 高并发场景下,用 agent.astream() + async for 实现异步流式处理,既能让用户看到打字机效果,又能保证服务同时处理大量请求不卡顿。

同步 vs 异步流

同步流 异步流
agent.stream() agent.astream()
for ... in ... 迭代 async for ... in ... 迭代
阻塞线程,不适合高并发 非阻塞,高并发性能好
适合脚本、命令行工具 适合 Web 服务、API 接口

在 FastAPI 里,你可以直接把这个异步生成器返回成一个 StreamingResponse,前端就能收到实时的文本流了。

支持的流模式(stream_mode)

将以下一种或多种流模式作为列表传递给 stream()astream() 方法:

模式 核心概念 适用场景
values 模式完整状态快照 在每个节点执行后返回完整的图状态,提供全局视角的状态监控。 调试和开发阶段查看完整状态;需要全局上下文进行审计追踪;状态规模较小或网络带宽充足的情况。
updates 模式增量状态更新 只传输状态发生变化的部分,大幅减少数据传输量。 生产环境下的状态同步;需要高效网络传输的场景;实时监控特定字段的变化。
messages 模式LLM 令牌流 实时流式传输大语言模型的生成过程,实现打字机效果 聊天机器人界面;需要实时显示 AI 思考过程;提升用户交互体验。
custom 模式自定义数据流 允许开发者发送自定义的业务数据,与核心状态分离。 进度条和状态更新;业务特定的通知消息;调试和监控信息。
debug 模式详细调试信息 提供最详细的执行信息,包括节点生命周期、执行时间等。 深度调试和性能分析;排查复杂的执行问题;学习和理解 LangGraph 内部机制。

支持将多种模式作为列表同时传递,如:

1
stream_mode=["updates", "messages", "custom"]

定义一个常见的图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
from typing import TypedDict
from langgraph.graph import StateGraph, START


class State(TypedDict):
topic: str
joke: str


def refine_topic(state: State):
return {"topic": state["topic"] + "和小狗"}


def generate_joke(state: State):
return {"joke": f"这是一个关于{state['topic']}的笑话"}


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

流式传输stream_mode=”values”

在图表的每个步骤之后传输状态的完整值

1
2
3
4
5
for chunk in graph.stream(
{"topic": "冰激凌"},
stream_mode="values",
):
print(chunk)

执行结果

1
2
3
{'topic': '冰激凌'}
{'topic': '冰激凌和小狗'}
{'topic': '冰激凌和小狗', 'joke': '这是一个关于冰激凌和小狗的笑话'}

说明:values 则输出完整状态值。

流式传输stream_mode=”updates”

将图表每一步之后的更新流式传输到状态

1
2
3
4
5
for chunk in graph.stream(
{"topic": "冰激凌"},
stream_mode="updates",
):
print(chunk)

执行结果

1
2
{'refine_topic': {'topic': '冰激凌和小狗'}}
{'generate_joke': {'joke': '这是一个关于冰激凌和小狗的笑话'}}

说明:updates 模式输出节点的状态更新

流式传输stream_mode=”debug”

在图表的整个执行过程中传输尽可能多的信息

1
2
3
4
5
for chunk in graph.stream(
{"topic": "冰激凌"},
stream_mode="debug",
):
print(chunk)

执行结果

1
2
3
4
{'step': 1, 'timestamp': '2025-09-18T08:51:31.478123+00:00', 'type': 'task', 'payload': {'id': '671c3949-69b8-8e90-b997-e76a1f711f5f', 'name': 'refine_topic', 'input': {'topic': '冰激凌'}, 'triggers': ('branch:to:refine_topic',)}}
{'step': 1, 'timestamp': '2025-09-18T08:51:31.478680+00:00', 'type': 'task_result', 'payload': {'id': '671c3949-69b8-8e90-b997-e76a1f711f5f', 'name': 'refine_topic', 'error': None, 'result': [('topic', '冰激凌和小狗')], 'interrupts': []}}
{'step': 2, 'timestamp': '2025-09-18T08:51:31.478828+00:00', 'type': 'task', 'payload': {'id': 'b7f1816a-7dab-4847-e8e6-78ae4d691a8e', 'name': 'generate_joke', 'input': {'topic': '冰激凌和小狗'}, 'triggers': ('branch:to:generate_joke',)}}
{'step': 2, 'timestamp': '2025-09-18T08:51:31.479006+00:00', 'type': 'task_result', 'payload': {'id': 'b7f1816a-7dab-4847-e8e6-78ae4d691a8e', 'name': 'generate_joke', 'error': None, 'result': [('joke', '这是一个关于冰激凌和小狗的笑话')], 'interrupts': []}}

从以上结果来看,可以看打执行的详细信息,每一个步骤执行

LLM token 流式传输stream_mode=”messages”

为调用LLM的图形节点传输LLM令牌和元数据

需要接入大模型

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
from langchain_deepseek import ChatDeepSeek
import os

llm = ChatDeepSeek(
model="deepseek-chat",
temperature=0,
api_key=os.environ.get("DEEPSEEK_API_KEY"),
base_url=os.environ.get("DEEPSEEK_API_BASE"),
)


def generate_joke(state: State):
llm_response = llm.invoke(
[
{"role": "user", "content": f"生成一个关于 {state['topic']}的笑话"}
]
)
return {"joke": llm_response.content}


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

for message_chunk, metadata in graph.stream(
{"topic": "冰激凌"},
stream_mode="messages",
):
if message_chunk.content:
print(message_chunk.content, end="|", flush=True)

执行结果(流式输出的)

1
2
3
4
5
6
小狗|走进|一家|冰|激|凌|店|,|店员|问|:“|想要|什么|口|味的|?”
|小狗|说|:“|汪|草|味的|!”
|店员|愣|住|:“|抱歉|…|我们没有|汪|草|口味|。”
|小狗|叹气|:“|那|好吧|,|给我|一个|‘|爪子|’|饼干|筒|装|香|草|味|——|但|记住|,|这次|别|再把|我的|球|藏|进|冰|激|凌|里|了|,|上次|我|挖|了|半小时|!”|🍦|🐶|

|(|注|:|谐|音|梗|:|汪|草|=|香|草|,|爪子|=|甜|筒|品牌|“|可爱|多|”|的|经典|筒|身|设计|)|

说明:每生成一个 token 就会立即输出,并附带上下文信息。

工具中自定义数据流式输出stream_mode=”custom”

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
# @Time:2025/9/30 14:48
# @Author:jinglv
from dataclasses import dataclass
from typing import TypedDict

from langchain_core.runnables import RunnableConfig
from langgraph.config import get_stream_writer
from langgraph.constants import START, END
from langgraph.graph import StateGraph
from langgraph.runtime import Runtime

from app.agent.model.llms import qv_llm


# 定义状态
class State(TypedDict):
"""状态"""
user_input: str
test_cases: str
result: str


# 定义运行时的上下文参数
@dataclass
class RuntimeContext:
"""运行时上下文参数"""
test_env: str # 测试环境
tester_name: str # 测试人员名称


def generator_test_case(state: State):
"""生成测试用例"""
writer = get_stream_writer()
writer(f"开始执行生成测试用例的节点")
prompt = """
请帮我生成用户5条登录的测试用例,登录账号密码的长度限制为8到16位
"""
cases = qv_llm.invoke(prompt)
return {"test_cases": cases}


def run_test_cases(state: State, runtime: Runtime[RuntimeContext]):
"""执行测试用例"""
writer = get_stream_writer()
writer(f"开始执行【分析测试用例】的节点")
cases = state['test_cases']
prompt = f"""
请分析当前的五条测试用例,是否有缺陷,
用例数据如下:{cases}
"""
result = qv_llm.invoke(prompt)
return {"result": result}


def generator_test_report(state: State, config: RunnableConfig):
"""生成测试报告"""
writer = get_stream_writer()
writer(f"开始【生成测试报告】的节点运行")
return {"report": "这个是一个测试报告"}


# =================工作流的创建个编排===================
graph = StateGraph(State, context_schema=RuntimeContext)
graph.add_node("生成测试用例", generator_test_case)
graph.add_node("执行测试用例", run_test_cases)
graph.add_node("生成测试报告", generator_test_report)

# 设置起点
# graph.set_entry_point("生成测试用例")
graph.add_edge(START, "生成测试用例")
graph.add_edge("生成测试用例", "执行测试用例")
graph.add_edge("执行测试用例", "生成测试报告")
# graph.set_finish_point("生成测试报告")
graph.add_edge("生成测试报告", END)

app = graph.compile()

response = app.stream({"user_input": "测试项目A"},
config={"recursion_limit": 5},
context=RuntimeContext(test_env="测试环境A", tester_name="张三"),
stream_mode=['messages', 'custom']
)

for input_type, chunk in response:
if input_type == "messages":
# ai的输出内容
print(chunk[0].content, end="", flush=True)
elif input_type == "custom":
# 工具执行的输出内容
print(chunk)

注意:必须在 LangGraph 执行上下文中调用 get_stream_writer(),否则无效。

messages模式与其他模式的对比

特性维度 messages 模式 updates 模式 values 模式 custom 模式 debug 模式
流式单位 LLM 令牌 / 消息块 状态增量更新 完整状态快照 自定义业务数据 详细调试信息
数据粒度 最细(字符级) 中等(字段级) 最粗(状态级) 灵活可定义 最详细(节点级)
主要用途 聊天界面、逐字输出 状态监控、进度跟踪 调试、审计 业务事件通知 深度调试和性能分析
性能开销 中等 最低 最高 非常高
网络传输 适合实时传输 高效,适合频繁更新 数据量大,适合内网 按需定义 不适合生产环境
  1. 用户交互首选:messages 模式

    它是实现打字机效果的核心,数据粒度最细,专门为聊天界面的实时逐字输出设计,是面向用户的 AI 对话应用的标配。

  2. 生产环境高效之选:updates 模式

    只传输状态变化的部分,性能开销最低、传输效率最高,适合生产环境的状态同步和频繁更新场景。

  3. 开发调试专用:valuesdebug 模式

    • values 模式返回完整状态快照,适合开发调试和审计追踪,但数据量大、开销高,不适合高并发生产环境。
    • debug 模式提供最详细的节点级执行信息,适合排查复杂问题和性能分析,但开销非常高,绝对不建议在生产环境使用。
  4. 灵活扩展方案:custom 模式

    可以自定义业务数据,与核心状态分离,适合发送进度条更新、业务通知等场景,兼顾灵活性和低开销。

多模式组合应用

LangGraph 中多种流式处理模式的组合使用,解决了实际项目中 “单一模式无法同时满足多种需求” 的问题。

核心概念:多模式组合

  • 在实际项目中,单一的流式模式往往无法同时满足用户交互、状态监控、业务事件、调试等多种需求。
  • LangGraph 支持通过 stream_mode 参数传入一个模式列表(比如 ["messages", "updates"]),让你同时接收多种类型的流数据。
  • 流数据会以 (mode, chunk) 的形式返回,你可以根据 mode 字段来区分处理不同类型的数据。

代码示例

1
2
3
4
5
6
7
8
9
10
11
12
13
# 同时接收LLM令牌和状态更新
async for mode, chunk in agent.astream(
initial_state,
stream_mode=["messages", "updates"] # 多模式组合
):
if mode == "messages":
msg_chunk, meta = chunk
if msg_chunk.content:
# 实时显示到前端界面
update_ui_with_text(msg_chunk.content)
elif mode == "updates":
# 更新后台状态监控面板
update_status_panel(chunk)
  1. stream_mode=["messages", "updates"]
    • 这是多模式组合的核心写法,同时开启了两种流模式。
    • 之后每次迭代返回的 (mode, chunk) 元组中,mode 会告诉你当前 chunk 属于哪种模式。
  2. if mode == "messages":
    • modemessages 时,chunk(message_chunk, metadata) 元组,和你之前学的 messages 模式用法完全一样。
    • 这里的逻辑是:提取文本内容,调用 update_ui_with_text 推送到前端,实现打字机效果。
  3. elif mode == "updates":
    • modeupdates 时,chunk 是状态的增量更新数据。
    • 这里的逻辑是:调用 update_status_panel 更新后台的状态监控面板,让开发者看到节点执行过程中的状态变化。

组合优势

  • 实时用户体验messages 模式提供打字机效果,满足前端用户交互需求。

  • 后台状态监控updates 模式高效跟踪状态变化,适合后台监控和调试。

  • 灵活的业务集成:如果加入 custom 模式,还可以支持自定义业务事件(比如进度条、通知消息)。

  • 完整的调试能力:如果加入 debug 模式,可以同时获得深度调试信息,排查问题更方便。

LangGraph 支持同时开启多种流式模式,通过 (mode, chunk) 的形式返回不同类型的数据,让你可以同时满足前端用户体验、后台状态监控、业务事件通知和调试分析等多种需求。