1. 本章目标 #

第 22 章 §4 把七种 stream_mode 各演示了一遍,那是「知道有这回事」。真要做成能上线的流式界面,还会撞上一堆那一节没讲过的问题:

这一章把这些问题逐个解决,最后拼成一条前端可以直接消费的统一事件流。

学完你应能:

参考文档:

1.1 「能流式」和「流式能上线」差多远 #

很多人对流式的理解停在一行代码上:把 invoke 换成 stream,字就一个一个出来了。demo 里够用,产品里不够。差别在于:stream() 吐出的是图的内部动静,前端要的是界面语义,两者不是一回事:

stream() 给你的 前端真正想要的
「某个模型产生了一个片段,它来自 draft 节点,带着 draft-model 标签」 「往回答气泡里追加两个字」
「某个节点返回了一个 dict,键是 docs」 「把这 3 篇引用渲染到答案上方」
「__interrupt__ 出现在 updates 流里」 「弹一个审批对话框,别再转圈了」
一个 RuntimeError 从 for 循环里抛出来 「把步骤条标红,收起加载动画」

左边是实现细节:加节点、换模型、把某步收成子图,它就会变。右边是产品契约:定下来之后,不该跟着图的内部改动一起改。

所以本章只做一件事:在这两列之间加一层翻译。代码量不大,却是流式功能里唯一稳定的地方。省掉它,图每改一次前端就得跟着改;改错了多半还不报错,只是界面少了点东西。

1.2 三种流式 API,先认个脸 #

后面会反复碰到三个名字,先在这里对齐,免得混:

API 视角 形状 什么时候用
stream(stream_mode=...) 图视角:哪个节点动了、状态怎么变 随模式组合变化(§7.3) 默认选它,九成需求够用
stream(..., version="v2") 同上,但形状统一 一律 {"type", "ns", "data"} 写翻译层时更省心(§7.5)
astream_events() 调用栈视角:每个 Runnable 的起止 七字段的事件 dict 只有需要工具参数时才值得(§8)

三者不是互相替代。version="v2" 只是同一批数据的另一种包装;astream_events 才是另一套数据源。本章实战用的是第一种(也给出了第二种写法):同步异步都有,又和图的概念一一对应。

2. 和第 22 章的分工 #

同一批 API,两章切入点不同。第 22 章关心「怎么把执行过程打印出来看」,本章关心「怎么把执行过程变成产品」。也就是说:第 22 章的读者是正在调试图的你,本章的读者是将来要消费这条流的前端。

主题 第 22 章讲到哪 本章补什么
updates / values 完整覆盖,够用了 只在组合模式里作为配角出现
messages 会取 token、要过滤空片段 metadata 全字段、按节点/标签筛选、并行交错、四个类型陷阱
custom 有这个东西,能推字符串 writer 的两种拿法、进度协议怎么设计、四个边界行为
tasks / debug / checkpoints 排障用,本章不再展开 —
组合多模式 返回 (模式名, 内容) 元组 事件到达顺序、配 subgraphs=True 后元组会变三元素
子图 第 25 章用 subgraphs=True 看内部 不加这个参数,子图的 token 和进度全部丢失且不报错(§7.4)
version="v2" 没提 统一形状的流式协议,翻译层的首选(§7.5)
astream_events 提了一句「需要区分模型和工具时再上」 事件清单、结构、三个过滤参数、tool_call_chunks
失败与暂停 invoke 视角 流式视角:已产出的 chunk 怎么办、__interrupt__ 从哪露出来

一句话概括本章:在 LangGraph 吐出的原始流和前端需要的事件之间,加一层翻译。 这层翻译不复杂,但少了它,§1 开头那五个问题一个都躲不过。

3. messages 模式:token 是怎么出来的 #

这一节要打好本章所有筛选的地基。messages 模式的关键不在「能吐 token」,那是第 22 章讲过的;关键在于每个 token 都自带一张身份证,写明来自哪个节点、哪个模型、第几个超步。后面 §3.3、§3.4、§4.3 的三种筛法,读的都是这张身份证上的不同字段。

3.1 chunk 的形状 #

messages 模式流出来的不是字符串,而是二元组:第一个元素是消息(片段),第二个是那张身份证。第 22 章提过这一点,这里要看清第二个元素里到底有什么。

先搭一张后面反复用的两节点图:outline 列提纲(中间产物,不该给用户看),draft 写正文(要给用户看)。两个节点各调一次模型,正好演示「怎么只要后者」。

# operator 提供 add 函数,下面用它作为 log 字段的 reducer(追加而非覆盖)
import operator
# Annotated 用来给状态字段挂 reducer,TypedDict 用来声明状态结构
from typing import Annotated, TypedDict

# 从 .env 文件读取 DEEPSEEK_API_KEY 等环境变量
from dotenv import load_dotenv
# init_chat_model 用「供应商:模型名」一行创建聊天模型
from langchain.chat_models import init_chat_model
# START/END 是图的入口与出口哨兵节点,StateGraph 是图的构建器
from langgraph.graph import END, START, StateGraph
from rich import print
# override=True 让 .env 里的值覆盖已存在的同名环境变量,避免读到旧 key
load_dotenv(override=True)


# 图的状态结构:四个字段
class ArticleState(TypedDict):
    # 用户给定的主题,全程只读
    topic: str
    # outline 节点写入的提纲(中间产物)
    outline: str
    # draft 节点写入的正文(最终产物)
    draft: str
    # 执行日志,Annotated 里的 operator.add 让两个节点的日志相加而不是互相覆盖
    log: Annotated[list[str], operator.add]


# 第一个模型:负责列提纲。tags 在「创建时」指定,这个模型的每次调用都会带上它
outline_model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0, tags=["outline-model"])
# 第二个模型:负责写正文,标签不同,§3.4 会用这个差异来做筛选
draft_model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0, tags=["draft-model"])


# 第一个节点:把主题变成三条提纲
def make_outline(state: ArticleState) -> dict:
    # 调用提纲模型;注意这里用的是 invoke,节点内部并不需要写流式代码
    r = outline_model.invoke(f"给「{state['topic']}」列一个三条要点的提纲,每条不超过十个字,直接给要点。")
    # 返回局部更新:写 outline 字段,同时往 log 追加一条
    return {"outline": r.content, "log": ["提纲完成"]}


# 第二个节点:把提纲扩写成正文
def write_draft(state: ArticleState) -> dict:
    # 读上一个节点写进状态的 outline,交给正文模型扩写
    r = draft_model.invoke(f"根据提纲写一段 40 字左右的介绍:\n{state['outline']}")
    # 写 draft 字段,并追加日志
    return {"draft": r.content, "log": ["草稿完成"]}


# 用状态结构 ArticleState 新建一个图构建器
builder = StateGraph(ArticleState)
# 注册提纲节点,节点名 "outline" 后面会出现在 metadata 的 langgraph_node 里
builder.add_node("outline", make_outline)
# 注册正文节点,节点名 "draft"
builder.add_node("draft", write_draft)
# 入口指向 outline
builder.add_edge(START, "outline")
# outline 跑完接 draft,两者是串行的(§4 会改成并行看效果)
builder.add_edge("outline", "draft")
# draft 跑完结束
builder.add_edge("draft", END)
# 编译成可执行图;这张图不需要 checkpointer
graph = builder.compile()

# 统一的输入:TypedDict 的字段要给全,带 reducer 的 log 给空列表
INPUT = {"topic": "向量数据库", "outline": "", "draft": "", "log": []}

# 只看第一个 chunk 长什么样,看完立刻退出循环
for chunk in graph.stream(INPUT, stream_mode="messages"):
    # 解包:messages 模式的每一项固定是两个元素,第一个是消息,第二个是元数据
    token, meta = chunk
    print(token, meta)
    # 打印容器类型和长度,确认它就是个长度为 2 的 tuple
    print(f"类型: {type(chunk)}, 长度: {len(chunk)}")
    # 打印消息的具体类型和内容;注意 !r 会带引号,方便看出空字符串
    print(f"token: {type(token).__name__} content={token.content!r}")
    # 只看一条就够,跳出(生成器会被回收,剩下的模型调用不会继续跑完)
    break

输出:

类型: <class 'tuple'>, 长度: 2
token: AIMessageChunk content=''

三点要注意:

  1. 它是 tuple 不是对象,所以只能按位置解包,chunk["token"] 这种写法是错的。
  2. 第一个片段的 content 是空字符串。这不是 bug,模型的第一个片段通常只是在宣告「接下来是一条 AI 消息」,正文还没开始。§5.1 会讲清空片段一共有几种来源。
  3. 类型是 AIMessageChunk 而不是 AIMessage。这个区别看起来无所谓,实际是本章最隐蔽的一个坑,§5.3 专门讲。

break 的副作用: 提前跳出会让底层生成器被回收,后面的节点不再执行。调试时这样能省时间和 token,但如果你的节点里有副作用(写库、扣款),提前跳出等于让流程停在半路。这一点和第 26 章的暂停完全不同,它没有 checkpoint 可以恢复。

AIMessageChunk

字段 含义
content 本片段的正文。第一个 chunk 经常是空字符串,真正的字从后面的片段才开始出现
additional_kwargs 模型特有的附加字段。这里 DeepSeek 带了 reasoning_content(推理过程),普通对话模型通常是空 dict
response_metadata 响应侧元数据,例如 model_provider;流式过程中往往还不完整,结束时才会补全 token 用量等
id 这条消息的唯一 id,同一轮生成的各个 chunk 共享同一个 id,方便把碎片拼回一条完整消息
tool_calls 已经拼好的工具调用列表;纯文本生成时为空
invalid_tool_calls 解析失败的工具调用;正常路径下为空
tool_call_chunks 工具调用的流式碎片(名称、参数按 token 陆续到达);没有工具调用时为空列表

日常读流式正文,真正要盯的通常只有 content(以及需要推理过程时的 additional_kwargs["reasoning_content"])。tool_calls / tool_call_chunks 要到模型开始调工具时才有内容,§6 再展开。

3.2 metadata 里有什么 #

第二个元素是个 dict,字段比想象的多。与其猜有哪些,不如全打出来看一眼。碰到陌生的流式 API,这往往是最值得做的第一件事:

import operator
from typing import Annotated, TypedDict

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langgraph.graph import END, START, StateGraph

load_dotenv(override=True)


class ArticleState(TypedDict):
    topic: str
    outline: str
    draft: str
    log: Annotated[list[str], operator.add]


outline_model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0, tags=["outline-model"])
draft_model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0, tags=["draft-model"])


def make_outline(state: ArticleState) -> dict:
    r = outline_model.invoke(f"给「{state['topic']}」列一个三条要点的提纲,每条不超过十个字,直接给要点。")
    return {"outline": r.content, "log": ["提纲完成"]}


def write_draft(state: ArticleState) -> dict:
    r = draft_model.invoke(f"根据提纲写一段 40 字左右的介绍:\n{state['outline']}")
    return {"draft": r.content, "log": ["草稿完成"]}


builder = StateGraph(ArticleState)
builder.add_node("outline", make_outline)
builder.add_node("draft", write_draft)
builder.add_edge(START, "outline")
builder.add_edge("outline", "draft")
builder.add_edge("draft", END)
graph = builder.compile()

INPUT = {"topic": "向量数据库", "outline": "", "draft": "", "log": []}


# 重新跑一遍图,这次关心的是元数据而不是内容
for token, meta in graph.stream(INPUT, stream_mode="messages"):
    # sorted 让字段按字母序输出,方便和书上的清单对照
    for k, v in sorted(meta.items()):
        # 值可能很长(比如 lc_versions),截断到 70 个字符
        print(f"{k} = {str(v)[:70]}")
    # 同样只看第一条
    break

输出:

checkpoint_ns = outline:ce328eac-866f-c85e-64af-fed9f91a92ff
langgraph_checkpoint_ns = outline:ce328eac-866f-c85e-64af-fed9f91a92ff
langgraph_node = outline
langgraph_path = ('__pregel_pull', 'outline')
langgraph_step = 1
langgraph_triggers = ('branch:to:outline',)
lc_versions = {'langchain-core': '1.6.1', 'langchain': '1.3.18', 'langchain-openai':
ls_integration = langchain_chat_model
ls_model_name = deepseek-v4-flash
ls_model_type = chat
ls_provider = deepseek
ls_temperature = 0.0
tags = ['outline-model']

十三个字段,真正常用的就三个:

字段 含义 典型用途
langgraph_node 这个 token 来自哪个节点 按节点分流(§3.3),并行时的救命字段(§4)
tags 模型上挂的标签 按用途分流(§3.4),一个节点里有多次模型调用时用
langgraph_step 当前是第几个超步 前端显示进度、排查执行顺序

剩下的按用途分成三组,知道它们存在就行:

为什么不用 langgraph_step 做筛选? 它是超步序号,会随着图的结构变化而变(加一个前置节点,所有序号都往后挪)。节点名和标签是你自己起的,稳定得多。筛选用名字,展示用序号。

metadata

字段 示例值 含义 典型用途
langgraph_node 'outline' 产生该 token 的图节点名 按节点分流(§3.3);并行时的救命字段(§4)
tags ['outline-model'] 模型创建时挂上的标签 按用途分流(§3.4);同一节点多次调模型时区分
langgraph_step 1 当前超步序号 前端进度条、排查执行顺序(不要拿来做筛选)
ls_provider 'deepseek' 模型供应商 多模型监控、成本归因
ls_model_name 'deepseek-v4-flash' 具体模型名 记进日志,方便追查「换模型后流式静默退化」(§5.3)
ls_model_type 'chat' 模型类型 确认是聊天模型而非 embedding 等
ls_temperature 0.0 采样温度 复现与排障时对照调用参数
ls_integration 'langchain_chat_model' 集成入口标识 确认走的是哪条 LangChain 集成路径
lc_versions {...} 相关包版本号 版本差异导致行为不一致时对照
langgraph_triggers ('branch:to:outline',) 触发本节点的边/分支 排查「为什么这个节点跑了/没跑」
langgraph_path ('__pregel_pull', 'outline') 图内调度路径 偶发排障用,日常可忽略
checkpoint_ns 'outline:5769…' 本次节点调用的命名空间(节点名:uuid) 区分同名节点的不同调用实例
langgraph_checkpoint_ns 同上 与 checkpoint_ns 同义的 LangGraph 侧字段 父图 / 子图命名空间不同,§7.4 会用到

3.3 按节点过滤:只要某一步的输出 #

上面那张图有两个模型调用:outline 节点在打草稿提纲,draft 节点在写最终文案。给用户看的只有后者。靠 langgraph_node 一行就能筛出来:

import operator
from typing import Annotated, TypedDict

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langgraph.graph import END, START, StateGraph

load_dotenv(override=True)


class ArticleState(TypedDict):
    topic: str
    outline: str
    draft: str
    log: Annotated[list[str], operator.add]


outline_model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0, tags=["outline-model"])
draft_model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0, tags=["draft-model"])


def make_outline(state: ArticleState) -> dict:
    r = outline_model.invoke(f"给「{state['topic']}」列一个三条要点的提纲,每条不超过十个字,直接给要点。")
    return {"outline": r.content, "log": ["提纲完成"]}


def write_draft(state: ArticleState) -> dict:
    r = draft_model.invoke(f"根据提纲写一段 40 字左右的介绍:\n{state['outline']}")
    return {"draft": r.content, "log": ["草稿完成"]}


builder = StateGraph(ArticleState)
builder.add_node("outline", make_outline)
builder.add_node("draft", write_draft)
builder.add_edge(START, "outline")
builder.add_edge("outline", "draft")
builder.add_edge("draft", END)
graph = builder.compile()

INPUT = {"topic": "向量数据库", "outline": "", "draft": "", "log": []}


# 用一个列表攒 draft 节点的片段,最后一次性拼起来
buf = []
# 顺便统计一下总片段数和空片段数,为 §5.1 做准备
total = 0
# 单独数一下空片段,看看占比
empty = 0
# 遍历整条 messages 流
for token, meta in graph.stream(INPUT, stream_mode="messages"):
    # 每来一条就计数
    total += 1
    # content 为空的片段单独计一笔
    if not token.content:
        # 空片段计数加一
        empty += 1
    # 关键的一行:只留下来自 draft 节点的片段,用 get 避免 KeyError
    if meta.get("langgraph_node") == "draft":
        # 攒起来而不是直接 print,避免输出被换行打散
        buf.append(token.content)
# 拼接成完整答案
print("draft 节点输出:", "".join(buf))
# 打印片段构成,注意空片段占了大头
print(f"总片段 {total} 个,其中空内容 {empty} 个")

输出:

draft 节点输出: 高维向量索引与相似度搜索,为语义检索提供核心支撑,让精准匹配触手可及。
总片段 136 个,其中空内容 94 个

outline 节点的 token 一个都没混进来。这是流式界面最基础的过滤:图里的中间步骤不该给用户看。

顺便多看一个数字:136 个片段里有 94 个是空的,接近七成。这个比例波动很大:同一份代码换个时段再跑,出现过 627 个片段里 591 个是空的(九成四)。波动本身不重要,重要的是空片段永远占大头,绝不是「偶尔有一两个」。§5.1 会解释它们从哪来。

3.4 按标签过滤:更细的粒度 #

节点过滤有个局限:一个节点里可能调用模型好几次。比较常见的是「先让小模型判断意图,再让大模型正式回答」,两次调用挤在同一个节点里,langgraph_node 完全相同,分不开。

给模型打标签就能分开。标签有两种打法,区别在于作用范围。

第一种在创建模型时打,作用于这个模型的所有调用:

# tags 写在 init_chat_model 里,等于给这个模型对象贴了个永久标签
# 之后无论在哪个节点、调用多少次,流出来的 metadata 里都带着它
outline_model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0, tags=["outline-model"])

于是筛选条件从「节点名」换成「标签」,其他都不用动:

import operator
from typing import Annotated, TypedDict

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langgraph.graph import END, START, StateGraph

load_dotenv(override=True)


class ArticleState(TypedDict):
    topic: str
    outline: str
    draft: str
    log: Annotated[list[str], operator.add]


outline_model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0, tags=["outline-model"])
draft_model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0, tags=["draft-model"])


def make_outline(state: ArticleState) -> dict:
    r = outline_model.invoke(f"给「{state['topic']}」列一个三条要点的提纲,每条不超过十个字,直接给要点。")
    return {"outline": r.content, "log": ["提纲完成"]}


def write_draft(state: ArticleState) -> dict:
    r = draft_model.invoke(f"根据提纲写一段 40 字左右的介绍:\n{state['outline']}")
    return {"draft": r.content, "log": ["草稿完成"]}


builder = StateGraph(ArticleState)
builder.add_node("outline", make_outline)
builder.add_node("draft", write_draft)
builder.add_edge(START, "outline")
builder.add_edge("outline", "draft")
builder.add_edge("draft", END)
graph = builder.compile()

INPUT = {"topic": "向量数据库", "outline": "", "draft": "", "log": []}


# 攒 outline-model 这个标签的片段
buf = []
# 遍历整条流
for token, meta in graph.stream(INPUT, stream_mode="messages"):
    # tags 字段有可能不存在或为 None,用 or [] 兜一下,避免 in 操作报 TypeError
    if "outline-model" in (meta.get("tags") or []):
        # 命中标签就收下
        buf.append(token.content)
# 拼接输出,这次拿到的是提纲而不是正文
print("outline-model 输出:", "".join(buf))

输出:

outline-model 输出: 1. 高维向量索引
2. 相似度检索
3. 支持ANN查询

第二种用 with_config,给同一个模型的不同调用打不同标签,这才是解决「一个节点里调两次模型」的办法:

# 读取环境变量
from dotenv import load_dotenv
# 创建模型
from langchain.chat_models import init_chat_model
# MessagesState 是内置状态,只有一个带 add_messages reducer 的 messages 字段
from langgraph.graph import END, START, MessagesState, StateGraph

# 加载 .env
load_dotenv(override=True)

# 这个模型创建时没有任何标签
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)
# with_config 返回一个「配置好的副本」,原来的 model 不受影响
# 同一个模型可以派生出多个带不同标签的副本,分别用在不同的调用点上
tagged = model.with_config(tags=["final-answer"])


# 节点:用带标签的那个副本去调用
def node_tagged(state: MessagesState) -> dict:
    # 返回一条新的 AI 消息,add_messages reducer 会把它追加到历史里
    return {"messages": [tagged.invoke(state["messages"])]}


# 用内置的 MessagesState 建图
builder = StateGraph(MessagesState)
# 只有一个节点
builder.add_node("t", node_tagged)
# 入口连到它
builder.add_edge(START, "t")
# 它跑完就结束
builder.add_edge("t", END)

# 编译后直接流式执行,输入是一条用户消息
for token, meta in builder.compile().stream(
    {"messages": [{"role": "user", "content": "说三个字"}]}, stream_mode="messages"
):
    # 过滤掉空片段(§5.1),只打印有内容的
    if token.content:
        # 确认标签确实跟着 token 一起流了出来
        print(f"tags={meta.get('tags')} token={token.content!r}")

输出:

tags=['final-answer'] token='我爱你'

两种打法的取舍:

打法 作用范围 适合
init_chat_model(tags=[...]) 这个模型对象的所有调用 一个模型只干一件事(提纲模型、回答模型各建一个)
model.with_config(tags=[...]) 派生出的那个副本 同一个模型在不同位置扮演不同角色

with_config 还有一点值得记住:它返回副本,不改原对象。所以 model 本身仍然没有标签,可以继续派生出 tagged_a、tagged_b,互不干扰。

约定: 给「要给用户看的那次模型调用」固定打一个标签,比如 final-answer。前端只认这个标签。以后图里加多少个中间模型调用,前端代码都不用改,这就是 §1.1 说的「产品契约不随实现细节变」。本章实战(§10)就是这么做的。

4. 头号坑:并行节点的 token 会交错 #

这是本章最需要盯紧的一节。两个节点并行执行时,它们的 token 会挤在同一条流里,而且交错程度每次运行都不一样。

先说清为什么会这样,后面的现象就都能预期。第 20 章讲过,LangGraph 的并行分支是在同一个超步里由线程池并发执行的。两个分支各自调模型,两条 HTTP 流同时往回吐字,谁的字先到就先进 stream() 的输出队列。这里没有任何排序逻辑:LangGraph 不知道、也不该知道你打算怎么渲染这些字。所以「token 的到达顺序」从设计上就没有保证。

4.1 直接拼接得到的是两段答案粘在一起 #

构造一张有并行分支的图,两个分支各自调模型,输出故意选得完全无关(数字 vs 水果),混在一起一眼就能看出来:

# reducer 用的 add
import operator
# 状态声明用的两个类型工具
from typing import Annotated, TypedDict

# 读取 .env
from dotenv import load_dotenv
# 创建模型
from langchain.chat_models import init_chat_model
# 建图三件套
from langgraph.graph import END, START, StateGraph

# 加载环境变量
load_dotenv(override=True)


# 状态:两个分支各写一个字段,互不干扰
class BranchState(TypedDict):
    # 占位的问题字段,这个例子里用不到
    q: str
    # 分支 a 的结果
    a: str
    # 分支 b 的结果
    b: str
    # 两个分支都要往这里追加,所以必须挂 reducer,否则并行写同一个键会报冲突
    log: Annotated[list[str], operator.add]


# 这个例子里两个分支共用同一个模型对象,也不打标签,重点是节点名的差异
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)


# 分支 a:让模型数数字
def branch_a(state: BranchState) -> dict:
    # 输出短且规整,方便肉眼判断是否被切碎
    r = model.invoke("从 1 数到 8,用中文逗号分隔,只给数字。")
    # 写自己的字段,并追加日志
    return {"a": r.content, "log": ["a done"]}


# 分支 b:让模型列水果
def branch_b(state: BranchState) -> dict:
    # 和 a 的内容毫无关系,混在一起会非常明显
    r = model.invoke("列出八个水果名,用中文逗号分隔,只给名字。")
    # 写自己的字段,并追加日志
    return {"b": r.content, "log": ["b done"]}


# 建图
builder = StateGraph(BranchState)
# 注册分支 a
builder.add_node("a", branch_a)
# 注册分支 b
builder.add_node("b", branch_b)
# 汇合节点:什么都不做,只是给并行分支一个共同的下游
builder.add_node("join", lambda s: {"log": ["joined"]})
# 从入口连到 a
builder.add_edge(START, "a")
# 也从入口连到 b,两条边从同一个点出发,这就是并行的写法
builder.add_edge(START, "b")
# a 完成后进 join
builder.add_edge("a", "join")
# b 完成后也进 join,LangGraph 会等两个都完成才执行 join
builder.add_edge("b", "join")
# join 之后结束
builder.add_edge("join", END)
# 编译
graph = builder.compile()

# 输入:两个结果字段先给空串
INPUT = {"q": "x", "a": "", "b": "", "log": []}

# 常见的错误写法:不管来源,见到 token 就往界面上拼
naive = []
# 遍历整条流
for token, meta in graph.stream(INPUT, stream_mode="messages"):
    # 只做了「过滤空片段」这一道处理,完全没看来源
    if token.content:
        # 一股脑塞进同一个列表,等价于前端不分来源地追加到同一个气泡
        naive.append(token.content)
# 拼接后打印,看看用户会看到什么
print("".join(naive))

输出:

苹果,香蕉,橙子,葡萄,西瓜,草莓,桃子,梨1,2,3,4,5,6,7,8

两段完全无关的答案首尾相接。用户看到的就是这个。

这次运行还算「客气」:两段各自完整,只是接在一起。但同一份代码换一次运行,出现过这样的结果:

1,苹果2,3,4,,香蕉,橙子,5,6,7,葡萄,西瓜,草莓,8蓝莓,桃子

同一份代码、同一个输入,两种完全不同的破坏程度。 下一节把这件事量化。

4.2 同样的代码连跑八次,结果分成两类 #

上面已经看到同一份代码能跑出两种破坏程度。要判断这个问题有多严重,得把「多久出现一次」量化出来。

为了把不确定性放大,做两件事:把两个分支的输出都加长到 20 项(拉长两条流重叠的时间窗口),并且记录节点切换的轨迹:每当 token 的来源从 a 变成 b 或反过来,就记一笔。切换次数越多,说明咬合得越碎。

import operator
from typing import Annotated, TypedDict

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langgraph.graph import END, START, StateGraph

load_dotenv(override=True)


class BranchState(TypedDict):
    q: str
    a: str
    b: str
    log: Annotated[list[str], operator.add]


model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)


def branch_a(state: BranchState) -> dict:
    r = model.invoke("从 1 数到 8,用中文逗号分隔,只给数字。")
    return {"a": r.content, "log": ["a done"]}


def branch_b(state: BranchState) -> dict:
    r = model.invoke("列出八个水果名,用中文逗号分隔,只给名字。")
    return {"b": r.content, "log": ["b done"]}


builder = StateGraph(BranchState)
builder.add_node("a", branch_a)
builder.add_node("b", branch_b)
builder.add_node("join", lambda s: {"log": ["joined"]})
builder.add_edge(START, "a")
builder.add_edge(START, "b")
builder.add_edge("a", "join")
builder.add_edge("b", "join")
builder.add_edge("join", END)
graph = builder.compile()

INPUT = {"q": "x", "a": "", "b": "", "log": []}


def branch_a(state: BranchState) -> dict:
    # 从 8 个数字改成 20 个,模型要吐更久
    r = model.invoke("从 1 数到 20,用中文逗号分隔,只给数字。")
    # 返回值结构不变
    return {"a": r.content, "log": ["a"]}


# 分支 b 同样加长
def branch_b(state: BranchState) -> dict:
    # 从 8 个水果改成 20 个
    r = model.invoke("列出 20 个水果名,用中文逗号分隔,只给名字。")
    # 返回值结构不变
    return {"b": r.content, "log": ["b"]}


# 节点函数换了,图要重新建一遍(compile 之后的图是不可变的)
builder = StateGraph(BranchState)
# 重新注册加长版的 a
builder.add_node("a", branch_a)
# 重新注册加长版的 b
builder.add_node("b", branch_b)
# 汇合节点照旧
builder.add_node("join", lambda s: {"log": ["j"]})
# 从入口连到 a
builder.add_edge(START, "a")
# 也从入口连到 b,构成并行
builder.add_edge(START, "b")
# a 汇合到 join
builder.add_edge("a", "join")
# b 也汇合到 join
builder.add_edge("b", "join")
# 结束边
builder.add_edge("join", END)
# 重新编译
graph = builder.compile()

# 记录每一轮的切换次数,最后汇总
switches = []
# 同一份代码、同一份输入,连跑八次
for run in range(8):
    # seq 记录来源节点的变化轨迹,只在切换时追加
    seq = []
    # counts 记录两个节点各自产出了多少个非空片段
    counts = {"a": 0, "b": 0}
    # 跑一遍图
    for token, meta in graph.stream(INPUT, stream_mode="messages"):
        # 空片段不参与统计,否则噪音太大
        if token.content:
            # 取出这个片段的来源节点名
            n = meta["langgraph_node"]
            # 该节点的计数加一
            counts[n] += 1
            # 只有「第一条」或者「来源和上一条不同」时才记一笔,这样 seq 就是切换轨迹
            if not seq or seq[-1] != n:
                # 记一次切换
                seq.append(n)
    # 轨迹长度减一就是切换次数
    switches.append(len(seq) - 1)
    # 打印这一轮的轨迹、切换次数和 token 分布
    print(f"第 {run + 1} 次: 轨迹 {'->'.join(seq)} (切换 {len(seq) - 1} 次) token 数 {counts}")
# 汇总八次的切换次数,一眼看出分布
print(f"8 次的切换次数: {switches}")

输出:

第 1 次: 轨迹 b->a->b->a (切换 3 次) token 数 {'a': 39, 'b': 42}
第 2 次: 轨迹 b->a->b->a->b->a->b->a (切换 7 次) token 数 {'a': 39, 'b': 44}
第 3 次: 轨迹 b->a (切换 1 次) token 数 {'a': 39, 'b': 44}
第 4 次: 轨迹 a->b (切换 1 次) token 数 {'a': 39, 'b': 45}
第 5 次: 轨迹 b->a (切换 1 次) token 数 {'a': 39, 'b': 43}
第 6 次: 轨迹 a->b (切换 1 次) token 数 {'a': 39, 'b': 43}
第 7 次: 轨迹 b->a (切换 1 次) token 数 {'a': 39, 'b': 42}
第 8 次: 轨迹 b->a (切换 1 次) token 数 {'a': 39, 'b': 43}

八次运行清楚地分成两类:

这组数字正是这个 bug 最麻烦的地方:这一轮里,8 次有 6 次「看起来正常」。 换个时段再跑,同一份代码给出的是完全不同的分布:先前两轮实测分别是 [14, 1, 1, 1, 1, 1, 9, 13](3 次严重咬合,切换最多 14 次)和 [1, 7, 1, 1, 1, 5, 1, 1]。三轮合计 24 次,其中 17 次不会暴露问题。

这个比例就是你本地测试的漏检概率。本地测几次都正常,你会以为代码没问题;上线之后偶发乱码,每次复现条件都不同,日志里也没有任何异常。

你自己跑出来的分布大概率和这里不一样,这恰恰是问题所在。同一份代码、同一个模型、同样的输入,交错程度由线程调度和网络抖动共同决定,没有任何规律可循,也不可能靠多测几次测出来。

注意别被 token 数误导。 上面 a 每次都是 39 个片段(数字 1 到 20 的输出很稳定),b 在 42~45 之间波动(模型选的水果名长度不同,先前一轮里还出现过 61)。片段数波动是模型行为,和交错无关;决定要不要修的是切换次数,不是片段数。

4.3 正确做法:按节点分流 #

修法很简单,也很彻底:按 langgraph_node 分桶,每个来源各自攒。

import operator
from typing import Annotated, TypedDict

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langgraph.graph import END, START, StateGraph

load_dotenv(override=True)


class BranchState(TypedDict):
    q: str
    a: str
    b: str
    log: Annotated[list[str], operator.add]


model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)


def branch_a(state: BranchState) -> dict:
    r = model.invoke("从 1 数到 8,用中文逗号分隔,只给数字。")
    return {"a": r.content, "log": ["a done"]}


def branch_b(state: BranchState) -> dict:
    r = model.invoke("列出八个水果名,用中文逗号分隔,只给名字。")
    return {"b": r.content, "log": ["b done"]}


builder = StateGraph(BranchState)
builder.add_node("a", branch_a)
builder.add_node("b", branch_b)
builder.add_node("join", lambda s: {"log": ["joined"]})
builder.add_edge(START, "a")
builder.add_edge(START, "b")
builder.add_edge("a", "join")
builder.add_edge("b", "join")
builder.add_edge("join", END)
graph = builder.compile()

INPUT = {"q": "x", "a": "", "b": "", "log": []}


# 用 dict 做分桶容器:键是节点名,值是该节点的片段列表
bucket: dict[str, list[str]] = {}
# 遍历整条流
for token, meta in graph.stream(INPUT, stream_mode="messages"):
    # 依旧要过滤空片段
    if token.content:
        # setdefault 保证键不存在时自动建一个空列表,然后按来源节点分桶
        # 这一行就是修复的全部内容,不再往同一个列表里塞
        bucket.setdefault(meta["langgraph_node"], []).append(token.content)

# 每个桶各自拼接,互不污染
for node, parts in bucket.items():
    # 打印时带上来源节点名,对应前端的「渲染到哪个区域」
    print(f"[{node}] {''.join(parts)}")

输出:

[b] 苹果,香蕉,橙子,葡萄,西瓜,草莓,桃子,梨
[a] 1,2,3,4,5,6,7,8

两段各自完整。注意输出顺序仍然不确定(这次 b 在前),但这已经无害了:分桶之后,顺序只影响「哪个区域先开始填充」,不影响每个区域的内容正确性。这是关键区别:不确定性没有被消除,只是被隔离到了无害的地方。

推给前端时,事件里也要带上来源:

# 事件里带 node 字段,前端据此决定渲染到哪个区域
{"type": "token", "text": "苹果", "node": "b"}

前端拿到后按 node 分区渲染,而不是不分来源地追加到同一个气泡里。本章实战(§10)的 token 事件就是这个形状。

规则:只要图里存在并行分支,token 事件就必须带来源标识。 更严格一点的说法是:即使现在没有并行分支,也应该带上。 单链路的图确实可以偷懒,但一旦某天有人加了个并行节点,偷懒的代码会静默出错:没有报错、没有异常,只是用户偶尔看到乱码。带上 node 字段的成本是一个键,代价完全不对等。本章实战用「标签 + 节点」双重过滤,就是为了从结构上排除这个问题。

5. messages 模式的四个类型陷阱 #

上一节的坑是「顺序」,这一节的坑是「类型」。共同点是都不报错,所以值得单独用一节讲。

名字叫 messages,很容易让人以为流里全是「模型吐出的字」。其实流的是消息对象:种类不止一种,内容还可能是空的。四个陷阱按踩到的概率从高到低排:

# 陷阱 后果 会报错吗
5.1 大量 content 为空的片段 界面出现空白抖动,或误判「模型没输出」 不会
5.2 ToolMessage 也走这条流 工具返回值印在用户的回答里 不会
5.3 非流式模型退化成完整消息 流式效果消失;配合类型判断会让回答整条消失 不会
5.4 图里根本没有模型调用 token 流为空,前端一直转圈 不会

5.1 大量 content 为空的片段 #

§3.1 里第一个 chunk 的 content 就是 '',§3.3 的统计更直接:136 个片段里 94 个是空的,接近七成(另一轮实测是九成四)。

空片段一共有四种来源,知道来源才知道过滤掉它们的代价:

来源 说明
角色起始片段 模型先宣告「接下来是一条 AI 消息」,正文还没开始
推理内容(reasoning) 思考型模型把思维链放在 additional_kwargs 里,content 是空的
工具调用阶段 模型在生成函数名和参数 JSON,正文自然是空的(§8.4 会把这部分挖出来)
收尾片段 携带 finish_reason、token 用量等元信息,不带正文

处理方式就一行:

# 空内容的片段直接跳过,不要往界面上追加
# 少了这一行,前端会收到一串空字符串,导致光标闪动或占位符抖动
if not token.content:
    continue

但要清楚过滤掉的是什么。极端情况:某个模型回合全程只在调工具(Agent 里很常见),整轮 token 事件会被滤得一个不剩,界面一动不动,用户以为卡住了。这时要用 §6 的 custom 事件或 §8.4 的 tool_call_chunks 补位:过滤是对的,但滤完得有别的东西顶上。

5.2 ToolMessage 也走这条流 #

messages 模式不只流模型的输出,工具的返回值也在里面。搭一张最小的 Agent 图(模型 + 工具 + 回环),把消息类型统计一下:

# Counter 用来统计各类型出现的次数
from collections import Counter

# 读取 .env
from dotenv import load_dotenv
# 创建模型
from langchain.chat_models import init_chat_model
# @tool 装饰器把普通函数变成模型可调用的工具
from langchain_core.tools import tool
# MessagesState 自带 messages 字段,省得自己声明
from langgraph.graph import END, START, MessagesState, StateGraph
# ToolNode 是预置节点,负责执行 AI 消息里的 tool_calls
from langgraph.prebuilt import ToolNode

# 加载环境变量
load_dotenv(override=True)


# 定义一个查库存的工具
@tool
def get_stock(sku: str) -> str:
    """查询某个 SKU 的库存数量。"""
    # 用字典模拟数据库,查不到就返回兜底文案
    return {"A-100": "库存 12 件"}.get(sku, "查无此 SKU")


# 工具清单,绑定和 ToolNode 都要用,抽成变量避免两处不一致
TOOLS = [get_stock]
# bind_tools 把工具的 JSON Schema 告诉模型,模型才知道可以调用它
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0).bind_tools(TOOLS)


# 模型节点:把完整历史交给模型,返回一条新消息
def call_model(state: MessagesState) -> dict:
    # add_messages reducer 会把这条消息追加到历史末尾
    return {"messages": [model.invoke(state["messages"])]}


# 路由函数:看最后一条消息里有没有 tool_calls,决定去工具还是结束
def should_continue(state: MessagesState):
    # getattr 兜一下,因为 HumanMessage 之类没有 tool_calls 属性
    return "tools" if getattr(state["messages"][-1], "tool_calls", None) else END


# 建图
builder = StateGraph(MessagesState)
# 模型节点
builder.add_node("model", call_model)
# 工具节点,ToolNode 会自动匹配名字并执行
builder.add_node("tools", ToolNode(TOOLS))
# 入口进模型
builder.add_edge(START, "model")
# 模型之后走条件边,第三个参数声明可能的去向(供画图和校验用)
builder.add_conditional_edges("model", should_continue, ["tools", END])
# 工具执行完回到模型,让模型根据工具结果组织最终回答,这就是 Agent 的回环
builder.add_edge("tools", "model")
# 编译
graph = builder.compile()

# 一个会触发工具调用的问题
Q = {"messages": [{"role": "user", "content": "A-100 有货吗"}]}

# 用 (类型名, 节点名) 作为键统计,这样能同时看出「是什么」和「从哪来」
types = Counter()
# 遍历整条 messages 流
for token, meta in graph.stream(Q, stream_mode="messages"):
    # type(token).__name__ 拿到类名字符串,比 isinstance 更适合做统计
    types[(type(token).__name__, meta.get("langgraph_node"))] += 1
# 打印统计结果
for k, v in types.items():
    # 键是 (类型名, 节点名) 元组,值是出现次数
    print(f"{k} x{v}")

输出:

('AIMessageChunk', 'model') x71
('ToolMessage', 'tools') x1

两个关键区别:

所以如果你的渲染逻辑是「拿到 content 就追加到当前气泡」,那句「库存 12 件」会直接印在用户的回答正文里。修法是按类型分派:

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langchain_core.tools import tool
from langgraph.graph import END, START, MessagesState, StateGraph
from langgraph.prebuilt import ToolNode

load_dotenv(override=True)


@tool
def get_stock(sku: str) -> str:
    """查询某个 SKU 的库存数量。"""
    return {"A-100": "库存 12 件"}.get(sku, "查无此 SKU")


TOOLS = [get_stock]
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0).bind_tools(TOOLS)


def call_model(state: MessagesState) -> dict:
    return {"messages": [model.invoke(state["messages"])]}


def should_continue(state: MessagesState):
    return "tools" if getattr(state["messages"][-1], "tool_calls", None) else END


builder = StateGraph(MessagesState)
builder.add_node("model", call_model)
builder.add_node("tools", ToolNode(TOOLS))
builder.add_edge(START, "model")
builder.add_conditional_edges("model", should_continue, ["tools", END])
builder.add_edge("tools", "model")
graph = builder.compile()

Q = {"messages": [{"role": "user", "content": "A-100 有货吗"}]}


# 重新跑一遍,这次按消息类型分别处理
for token, meta in graph.stream(Q, stream_mode="messages"):
    # 先拿到类型名字符串
    name = type(token).__name__
    # 分支一:模型片段且内容非空,这才是要追加到回答气泡里的
    if name == "AIMessageChunk" and token.content:
        # 用 !r 打印,能看清模型的切分粒度
        print(f"[token] {token.content!r}")
    # 分支二:工具返回,单独渲染成「工具调用卡片」,别混进正文
    elif name == "ToolMessage":
        # 加个前缀,模拟前端把它渲染成独立的工具卡片
        print(f"[工具返回] {token.content!r}")

输出:

[工具返回] '库存 12 件'
[token] 'A'
[token] '-'
[token] '100'
[token] ' '
[token] '有'
[token] '货'
[token] ','
[token] '当前'
[token] '库存'
[token] '为'
[token] ' **'
[token] '12'
[token] ' '
[token] '件'
[token] '**'
[token] '。'

工具返回排在所有 token 之前,这个顺序值得停下来想一秒:第一个模型回合全程在生成函数调用(content 是空的,被 if token.content 滤掉了),所以真正流出正文的只有工具执行完之后的第二个回合。这就是 §5.1 说的那个极端情况的真实样子:如果你只监听 token 事件,从提问到工具返回这段时间里界面上一个字都不会动。

5.3 非流式模型会静默退化 #

这个陷阱最隐蔽,因为从头到尾不报错。

触发条件比想象的多:显式传 disable_streaming=True、某些供应商的特定模型不支持 SSE、开启了某些参数组合(比如部分供应商的 JSON 模式或 logprobs)、或者走了不支持流式的代理网关。共同点是:代码一行都不用改就会退化。

# 读取 .env
from dotenv import load_dotenv
# 创建模型
from langchain.chat_models import init_chat_model
# 建图
from langgraph.graph import END, START, MessagesState, StateGraph

# 加载环境变量
load_dotenv(override=True)

# disable_streaming=True 显式关掉流式,用来复现「换了个不支持流式的模型」的情形
nostream = init_chat_model("deepseek:deepseek-v4-flash", temperature=0, disable_streaming=True)


# 节点代码和普通节点完全一样,这正是问题所在,节点看不出任何区别
def plain_node(state: MessagesState) -> dict:
    # 依旧是 invoke,依旧返回一条消息
    return {"messages": [nostream.invoke(state["messages"])]}


# 建一张只有一个节点的图
builder = StateGraph(MessagesState)
# 注册节点
builder.add_node("plain", plain_node)
# 入口
builder.add_edge(START, "plain")
# 出口
builder.add_edge("plain", END)

# 数一下到底流出了几个 chunk
n = 0
# 调用方的代码也完全没变,还是 stream_mode="messages"
for token, meta in builder.compile().stream(
    {"messages": [{"role": "user", "content": "A-100 有货吗"}]}, stream_mode="messages"
):
    # 计数
    n += 1
    # 打印类型和内容前 40 个字符,重点看类型名
    print(f"chunk {n}: {type(token).__name__} content={token.content[:40]!r}")
# 总数
print(f"一共 {n} 个 chunk")

输出:

chunk 1: AIMessage content='抱歉,我无法查询实时库存或供货情况。如果您是指 **NVIDIA A100 GP'
一共 1 个 chunk

只有一个 chunk,类型是 AIMessage(不是 AIMessageChunk),里面是完整的回复。代码照跑、没有任何警告,界面上的文字一次性整段出现,流式效果消失了,但你只会以为是模型今天特别快。

更糟的是配合类型判断。看一下这两个类的继承关系:

# 从 langchain_core 导入两个消息类
from langchain_core.messages import AIMessage, AIMessageChunk

# 构造一条「完整消息」,模拟非流式模型的输出
full = AIMessage(content="完整回复")
# 构造一条「片段消息」,模拟流式模型的输出
chunk = AIMessageChunk(content="片")
# 完整消息不是片段的实例,用 AIMessageChunk 判断会漏掉它
print(f"isinstance(AIMessage,      AIMessageChunk) = {isinstance(full, AIMessageChunk)}")
# 片段却是完整消息的实例,用 AIMessage 判断两种都能接住
print(f"isinstance(AIMessageChunk, AIMessage)      = {isinstance(chunk, AIMessage)}")
# 打印继承链前三层,直接看出 AIMessageChunk 是 AIMessage 的子类
print(f"AIMessageChunk 的继承链: {[c.__name__ for c in AIMessageChunk.__mro__[:3]]}")

输出:

isinstance(AIMessage,      AIMessageChunk) = False
isinstance(AIMessageChunk, AIMessage)      = True
AIMessageChunk 的继承链: ['AIMessageChunk', 'AIMessage', 'BaseMessageChunk']

AIMessageChunk 是 AIMessage 的子类,反过来不成立。于是:

# 错的:非流式模型的整条回复会被这个条件直接丢掉,界面上一个字都没有,日志里也没有异常
if isinstance(token, AIMessageChunk):
    ...

# 对的:父类能同时接住片段和完整消息
if isinstance(token, AIMessage):
    ...

用父类判断之后,如果还想区分两种情况(比如打点统计有多少请求退化了),再补一次类名判断:

# 先用父类兜住所有模型输出
if isinstance(token, AIMessage):
    # 再用类名区分:拿到完整消息说明这次请求没有走流式,值得记一笔
    if type(token).__name__ == "AIMessage":
        # 把模型名一起记下来,这样才知道是哪个模型退化了(§3.2 的 ls_model_name)
        print(f"[警告] 非流式返回,模型={meta.get('ls_model_name')}")

这个坑的排查成本极高,因为现象是「界面上没字」,而代码、日志、异常栈全都干净。上线前值得专门加一条断言:如果一次请求里 AIMessage(而非 AIMessageChunk)的数量大于 0,就告警。 §12 的练习 5 会让你亲手复现一次「回答整条消失」。

5.4 图里没有模型调用 = 空流 #

最后一个很直白,但恰恰因为直白而容易忘:

# 提供 add 作为 reducer
import operator
# 状态类型
from typing import Annotated, TypedDict

# 建图;注意这段完全不需要模型,也不需要 .env
from langgraph.graph import END, START, StateGraph


# 状态结构
class PlainState(TypedDict):
    # 问题
    q: str
    # 答案
    a: str
    # 日志
    log: Annotated[list[str], operator.add]


# 节点:直接返回硬编码答案,全程没有任何模型调用
def no_model(state: PlainState) -> dict:
    # 现实里对应「命中缓存」「命中兜底话术」「被护栏拦下」这几条路径
    return {"a": "硬编码答案"}


# 建图
builder = StateGraph(PlainState)
# 注册这个不调模型的节点
builder.add_node("n", no_model)
# 入口
builder.add_edge(START, "n")
# 出口
builder.add_edge("n", END)

# 用 list() 把整条流收集起来,直接看长度
res = list(builder.compile().stream({"q": "x", "a": "", "log": []}, stream_mode="messages"))
# 数量是 0,但图其实是正常跑完的
print(f"chunk 数量: {len(res)}")

输出:

chunk 数量: 0

图正常跑完了,答案也写进状态了,但 messages 流是空的,而且没有任何异常。现实里以下路径都会这样:

结论:token 流不能作为唯一的事件来源。 前端不能靠「有没有收到 token」判断状态,必须搭配一个明确的终止事件。本章实战用三个终止事件覆盖所有情况:done(正常跑完)、paused(暂停等人)、error(失败):任何一次请求,最后一定恰好收到其中一个。 这是前端能安全收起加载动画的唯一依据。

6. custom 模式:让不可流式的步骤也有进度 #

前面三节都在处理「模型吐字」,但一次请求里真正花时间的往往不是模型。模型能一个字一个字吐,检索不能:它要么还在查,要么查完了,中间没有可以流出来的东西。第 15 章那条 RAG 链,检索那一两秒前端是完全静默的;如果再串一个重排(rerank),静默时间还要翻倍。

在用户看来,这段静默和「卡住了」没什么区别。custom 模式就是为此准备的:在节点里拿到一个 writer,想推什么推什么,不依赖任何模型。

6.1 拿 writer 的两种方式 #

get_stream_writer() 最常用:在节点函数体里调用,拿到一个可调用对象,调它就等于往 custom 流里推一条。

# 提供 add 作为 reducer
import operator
# 用 sleep 模拟检索耗时,好让进度事件肉眼可见
import time
# 状态类型
from typing import Annotated, TypedDict

# 读取 .env
from dotenv import load_dotenv
# 创建模型
from langchain.chat_models import init_chat_model
# get_stream_writer 从运行时上下文里取出当前这次执行的 writer
from langgraph.config import get_stream_writer
# 建图
from langgraph.graph import END, START, StateGraph

# 加载环境变量
load_dotenv(override=True)


# 一个简化的 RAG 状态
class RagState(TypedDict):
    # 用户问题
    q: str
    # 检索到的文档,用 add 累积
    docs: Annotated[list[str], operator.add]
    # 最终答案
    answer: str


# 生成用的模型
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)


# 检索节点:这是典型的「不可流式」步骤
def retrieve(state: RagState) -> dict:
    # 在节点内部取 writer;它和当前的 thread、节点绑定,不需要显式传参
    writer = get_stream_writer()
    # 一进节点就先推一条「开始」,这样前端立刻有反馈,不用等检索完
    writer({"stage": "retrieve", "status": "start", "text": "正在检索知识库"})
    # 模拟建立连接、编码 query 的耗时
    time.sleep(0.1)
    # 每命中一篇就推一条进度,让用户看到事情在推进
    for i in range(3):
        # 同一个 stage 下的 progress 事件,text 各不相同
        writer({"stage": "retrieve", "status": "progress", "text": f"命中第 {i + 1} 篇"})
    # 收尾事件带上结构化数据(命中数),前端可以显示「找到 3 篇」
    writer({"stage": "retrieve", "status": "done", "hits": 3})
    # 真正的返回值仍然走状态,和 custom 事件是两条独立的通道
    return {"docs": [f"文档{i}" for i in range(3)]}


# 生成节点:这一步是可流式的,但仍然值得推阶段事件
def generate(state: RagState) -> dict:
    # 同样先取 writer
    writer = get_stream_writer()
    # 推「开始生成」,前端可以把加载动画从「检索中」切到「生成中」
    writer({"stage": "generate", "status": "start", "text": "开始生成"})
    # 这里用 invoke,token 会通过 messages 模式单独流出去(§7 会两条一起收)
    r = model.invoke(f"用一句话回答:{state['q']}")
    # 推「生成完成」,前端收起打字光标
    writer({"stage": "generate", "status": "done"})
    # 把答案写进状态
    return {"answer": r.content}


# 建图
builder = StateGraph(RagState)
# 检索节点
builder.add_node("retrieve", retrieve)
# 生成节点
builder.add_node("generate", generate)
# 入口进检索
builder.add_edge(START, "retrieve")
# 检索完进生成
builder.add_edge("retrieve", "generate")
# 生成完结束
builder.add_edge("generate", END)
# 编译;这张图后面 §7 还要复用
graph = builder.compile()

# 统一输入
INPUT = {"q": "向量数据库是什么", "docs": [], "answer": ""}

# 只订阅 custom 模式:拿到的就是 writer 推的原始内容,没有任何包装
for chunk in graph.stream(INPUT, stream_mode="custom"):
    # 直接打印,形状就是你推进去的那个 dict
    print(chunk)

输出:

{'stage': 'retrieve', 'status': 'start', 'text': '正在检索知识库'}
{'stage': 'retrieve', 'status': 'progress', 'text': '命中第 1 篇'}
{'stage': 'retrieve', 'status': 'progress', 'text': '命中第 2 篇'}
{'stage': 'retrieve', 'status': 'progress', 'text': '命中第 3 篇'}
{'stage': 'retrieve', 'status': 'done', 'hits': 3}
{'stage': 'generate', 'status': 'start', 'text': '开始生成'}
{'stage': 'generate', 'status': 'done'}

注意 custom 流里没有任何包装:推进去什么,流出来就是什么。这和 messages 模式(强制二元组)不同,自由度全在你手上,代价是你得自己定协议(§6.2)。

第二种方式是把 writer 声明成节点的参数,LangGraph 会自动注入:

import operator
from typing import Annotated, TypedDict

from langgraph.graph import END, START, StateGraph


class RagState(TypedDict):
    q: str
    docs: Annotated[list[str], operator.add]
    answer: str


INPUT = {"q": "向量数据库是什么", "docs": [], "answer": ""}


# StreamWriter 是类型标注,LangGraph 靠它识别「这个参数要注入 writer」
from langgraph.types import StreamWriter


# 节点签名里多一个参数,名字必须是 writer,类型标注必须是 StreamWriter
def retrieve2(state: RagState, writer: StreamWriter) -> dict:
    # 直接用,不需要再调 get_stream_writer()
    writer({"stage": "retrieve", "status": "start", "text": "注入式 writer"})
    # 返回值写法不变
    return {"docs": ["x"]}


# 建一张只有这个节点的图来验证
builder = StateGraph(RagState)
# 注册节点,LangGraph 在注册时就会检查签名并准备注入
builder.add_node("r", retrieve2)
# 入口
builder.add_edge(START, "r")
# 出口
builder.add_edge("r", END)
# 跑起来看 custom 流
for chunk in builder.compile().stream(INPUT, stream_mode="custom"):
    # 输出和上面那种写法完全一致
    print(chunk)

输出:

{'stage': 'retrieve', 'status': 'start', 'text': '注入式 writer'}

两种方式功能完全一样,选哪种看场景:

方式 优点 缺点
get_stream_writer() 任何函数里都能用,包括节点调用的工具、子函数 依赖隐式上下文,单元测试时不好替换
writer: StreamWriter 参数 依赖显式,测试时直接传个 lambda x: None 就行 只有节点函数能用,深层调用的函数拿不到

建议: 节点自己推进度用参数注入(好测试);节点调用的工具或工具库内部推进度用 get_stream_writer()(没别的选择)。

6.2 工具内部也能推进度 #

这一点很实用但容易被忽略:get_stream_writer() 不要求调用者是节点函数,只要求在图的执行过程中被调用。所以工具函数内部可以直接推进度,这对「工具里有个慢查询」的场景非常有用。

import operator
import time
from typing import Annotated, TypedDict

from langgraph.config import get_stream_writer
from langgraph.graph import END, START, StateGraph


class RagState(TypedDict):
    q: str
    docs: Annotated[list[str], operator.add]
    answer: str


INPUT = {"q": "向量数据库是什么", "docs": [], "answer": ""}


# @tool 装饰器
from langchain_core.tools import tool


# 一个模拟耗时的检索工具
@tool
def slow_search(kw: str) -> str:
    """按关键词检索资料。"""
    # 在工具函数体内部取 writer,此时的运行时上下文是「调用这个工具的那个节点」
    w = get_stream_writer()
    # 工具一开始执行就推一条,前端能显示「正在查资料」
    w({"stage": "tool", "status": "start", "text": f"工具正在查 {kw}"})
    # 模拟慢查询
    time.sleep(0.05)
    # 查完推一条完成
    w({"stage": "tool", "status": "done"})
    # 正常返回工具结果
    return f"{kw} 的资料若干"


# 节点:显式调用这个工具(不用 ToolNode,看得更清楚)
def tool_node(state: RagState) -> dict:
    # invoke 工具时传的是参数字典
    r = slow_search.invoke({"kw": state["q"]})
    # 把工具结果写进状态
    return {"docs": [r]}


# 建图验证
builder = StateGraph(RagState)
# 注册节点
builder.add_node("t", tool_node)
# 入口
builder.add_edge(START, "t")
# 出口
builder.add_edge("t", END)
# 订阅 custom 流,看工具内部推的事件有没有出来
for chunk in builder.compile().stream(INPUT, stream_mode="custom"):
    # 工具内部推的两条事件都到了
    print(chunk)

输出:

{'stage': 'tool', 'status': 'start', 'text': '工具正在查 向量数据库是什么'}
{'stage': 'tool', 'status': 'done'}

这就补上了 §5.2 留下的缺口:Agent 调工具那段时间,token 流是空的,但 custom 流可以不空。 让工具自己汇报进度,界面就不会僵住。

6.3 进度协议怎么设计 #

writer 收什么都行(字符串、dict、列表、Pydantic 对象),所以必须自己定协议,否则前端会退化成一堆靠字符串内容分支的 if。

第 22 章的例子推的是字符串:demo 可以,产品不行:字符串前端只能原样显示,dict 才能驱动进度条、图标、状态机。一个够用的约定:

字段 作用
stage 哪个阶段(retrieve / answer / guard),对应前端哪个步骤条目
status start / progress / done,驱动这个条目的图标
text 给人看的一句话,可空
其他 阶段特有的数据,比如 hits 命中数、docs 引用列表

设计时有两条硬要求:

  1. stage + status 要能唯一确定一个「事件语义」。 前端才能做幂等处理:收到两次同样的 retrieve/start 只当一次。§9.3 会证明这个去重不是可选项而是必需项(恢复执行时进度事件真的会重发)。
  2. stage 的取值要独立于节点名。 听起来多余,其实很重要:节点是实现细节(可能拆成两个、可能改名),stage 是产品概念(界面上那一行步骤)。两者一对一时可以偷懒同名,但心里要清楚它们是两层东西。

本章实战里还多加了一个 kind 字段,用来把「阶段进度」和「引用列表」这两种结构不同的 custom 事件区分开,这是因为 custom 通道是共享的,一条流里可能混着好几种业务事件。

6.4 四个边界行为 #

custom 模式的边界行为有四个,其中三个是「不报错」,一个是「响亮报错」。

一、图外调用会报错。 get_stream_writer() 依赖运行时上下文:

# 单独导入,不建图、不执行任何节点
from langgraph.config import get_stream_writer

# 直接在图外面调用,看会发生什么
try:
    # 此时没有任何运行时上下文
    w = get_stream_writer()
# 捕获所有异常,打印类型和消息
except Exception as e:
    # 报错信息提到的是 get_config,说明 writer 是从运行时配置里取的
    print(f"{type(e).__name__}: {e}")

输出:

RuntimeError: Called get_config outside of a runnable context

这是本节唯一会响亮报错的行为,反而是好事。它的含义是:想推进度的函数必须在图的执行过程中被调用(节点里、或节点调用的工具里)。你不能在图外面单独测这个函数:单元测试时要么用参数注入版本(§6.1),要么把 writer 作为可选参数传进去。

二、用 invoke 跑时,writer 静默丢弃。 不报错,推出去的东西也没人接:

import operator
import time
from typing import Annotated, TypedDict

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langgraph.config import get_stream_writer
from langgraph.graph import END, START, StateGraph

load_dotenv(override=True)


class RagState(TypedDict):
    q: str
    docs: Annotated[list[str], operator.add]
    answer: str


model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)


def retrieve(state: RagState) -> dict:
    writer = get_stream_writer()
    writer({"stage": "retrieve", "status": "start", "text": "正在检索知识库"})
    time.sleep(0.1)
    for i in range(3):
        writer({"stage": "retrieve", "status": "progress", "text": f"命中第 {i + 1} 篇"})
    writer({"stage": "retrieve", "status": "done", "hits": 3})
    return {"docs": [f"文档{i}" for i in range(3)]}


def generate(state: RagState) -> dict:
    writer = get_stream_writer()
    writer({"stage": "generate", "status": "start", "text": "开始生成"})
    r = model.invoke(f"用一句话回答:{state['q']}")
    writer({"stage": "generate", "status": "done"})
    return {"answer": r.content}


builder = StateGraph(RagState)
builder.add_node("retrieve", retrieve)
builder.add_node("generate", generate)
builder.add_edge(START, "retrieve")
builder.add_edge("retrieve", "generate")
builder.add_edge("generate", END)
graph = builder.compile()

INPUT = {"q": "向量数据库是什么", "docs": [], "answer": ""}


# 同一张图,这次用 invoke 而不是 stream
out = graph.invoke(INPUT)
# 图正常跑完,返回最终状态;writer 推的七条事件全部被丢掉,没有任何警告
print(f"invoke 正常返回: docs={len(out['docs'])} 篇")

输出:

invoke 正常返回: docs=3 篇

这是好事:同一份节点代码,invoke 和 stream 都能跑,不用写两套。 批处理用 invoke、在线接口用 stream,节点代码完全复用。

三、没有任何 writer 调用时,custom 流是空的,也不报错:

# 提供 add 作为 reducer
import operator
# 状态类型
from typing import Annotated, TypedDict

# 建图
from langgraph.graph import END, START, StateGraph


# 状态结构
class PlainState(TypedDict):
    # 问题
    q: str
    # 答案
    answer: str
    # 日志
    log: Annotated[list[str], operator.add]


# 建一张什么都不推的图
builder = StateGraph(PlainState)
# 节点直接返回答案,没有 writer 调用
builder.add_node("plain", lambda s: {"answer": "直接返回"})
# 入口
builder.add_edge(START, "plain")
# 出口
builder.add_edge("plain", END)

# 订阅 custom 流并收集全部内容
out = list(builder.compile().stream({"q": "x", "answer": "", "log": []}, stream_mode="custom"))
# 空列表,但图是正常跑完的
print(f"custom 流的内容: {out}")

输出:

custom 流的内容: []

和 §5.4 一个道理:空流不代表出错,前端不能靠「有没有事件」判断状态。

四、子图里推的事件,父图默认收不到。 这个最容易造成线上事故,放到 §7.4 单独讲。

7. 组合多个 stream_mode #

7.1 一次拿到进度、token 和状态变更 #

到这里三种模式都单独用过了,但真实界面需要它们同时存在:进度靠 custom、正文靠 messages、节点产出靠 updates。

不需要开三条流(那会把图跑三遍),把模式传成列表就行。代价是每一条要多解一层包装:内容变成 (模式名, 内容) 二元组。

import operator
import time
from typing import Annotated, TypedDict

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langgraph.config import get_stream_writer
from langgraph.graph import END, START, StateGraph

load_dotenv(override=True)


class RagState(TypedDict):
    q: str
    docs: Annotated[list[str], operator.add]
    answer: str


model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)


def retrieve(state: RagState) -> dict:
    writer = get_stream_writer()
    writer({"stage": "retrieve", "status": "start", "text": "正在检索知识库"})
    time.sleep(0.1)
    for i in range(3):
        writer({"stage": "retrieve", "status": "progress", "text": f"命中第 {i + 1} 篇"})
    writer({"stage": "retrieve", "status": "done", "hits": 3})
    return {"docs": [f"文档{i}" for i in range(3)]}


def generate(state: RagState) -> dict:
    writer = get_stream_writer()
    writer({"stage": "generate", "status": "start", "text": "开始生成"})
    r = model.invoke(f"用一句话回答:{state['q']}")
    writer({"stage": "generate", "status": "done"})
    return {"answer": r.content}


builder = StateGraph(RagState)
builder.add_node("retrieve", retrieve)
builder.add_node("generate", generate)
builder.add_edge(START, "retrieve")
builder.add_edge("retrieve", "generate")
builder.add_edge("generate", END)
graph = builder.compile()

INPUT = {"q": "向量数据库是什么", "docs": [], "answer": ""}


# stream_mode 传列表,一次订阅三种模式;顺序不影响事件的到达顺序
for mode, chunk in graph.stream(INPUT, stream_mode=["custom", "messages", "updates"]):
    # 先按模式名分派,这是多模式流的第一件事
    if mode == "messages":
        # 注意这里要再解一层:messages 的内容本身就是个二元组
        token, meta = chunk
        # 依旧要过滤空片段(§5.1)
        if token.content:
            # 打印时带上来源节点,对应 §4.3 的规则
            print(f"[messages/{meta['langgraph_node']}] {token.content!r}")
    # custom 模式的内容就是 writer 推进去的原始 dict,不用再解包
    elif mode == "custom":
        # 直接打印原始 dict,形状就是 writer 推进去的那个
        print(f"[custom] {chunk}")
    # 剩下的就是 updates:形状是 {节点名: 该节点的返回值}
    else:
        # 一条 updates 通常只有一个键(一个节点的产出),取第一个键就是节点名
        node = list(chunk.keys())[0]
        # 只打印产出的字段名,避免文档输出太长
        print(f"[updates] {node} -> {list(chunk[node].keys())}")

输出(token 部分省略中间):

[custom] {'stage': 'retrieve', 'status': 'start', 'text': '正在检索知识库'}
[custom] {'stage': 'retrieve', 'status': 'progress', 'text': '命中第 1 篇'}
[custom] {'stage': 'retrieve', 'status': 'progress', 'text': '命中第 2 篇'}
[custom] {'stage': 'retrieve', 'status': 'progress', 'text': '命中第 3 篇'}
[custom] {'stage': 'retrieve', 'status': 'done', 'hits': 3}
[updates] retrieve -> ['docs']
[custom] {'stage': 'generate', 'status': 'start', 'text': '开始生成'}
[messages/generate] '向量'
[messages/generate] '数据库'
[messages/generate] '是'
[messages/generate] '专门'
...
[messages/generate] '。'
[custom] {'stage': 'generate', 'status': 'done'}
[updates] generate -> ['answer']

7.2 到达顺序:custom 在节点内,updates 在节点后 #

上面的输出顺序不是随机的,有规律,而且这个规律直接决定了界面该怎么写:

custom(retrieve 的三条进度)  ← 节点执行中,writer 一调就推出来
updates(retrieve)            ← 节点执行完,返回值合并进状态才有
custom(generate start)
messages(token...)           ← 模型吐字的过程中
custom(generate done)
updates(generate)

custom 和 messages 是「过程中」的,updates 是「完成后」的。 差别的根源在于它们的产生时机:writer(...) 是节点函数体里的一行代码,调到就推;而 updates 必须等节点函数 return、返回值经过 reducer 合并进状态之后才有。所以:

有意思的是,本章实战没有用 updates 发引用列表,而是用 custom 推的。原因就是这个时机差:引用列表在检索节点内部算出来的那一刻就能发,不用等节点里剩下的逻辑(打分、去重、日志)跑完。对界面来说,早半秒把来源渲染出来,用户等正文时就有东西可看。

一个推论: 如果某个节点很慢且不推 custom 事件,那么从它开始执行到它结束,这条流上什么都不会有:updates 要等它结束,messages 要等它调模型。这就是 §6 存在的意义。

7.3 加上 subgraphs=True,元组会变三元素 #

这是个容易写错的细节,而且写错的方式有两种,一种响亮一种安静。

第 25 章用过 subgraphs=True 看子图内部。它的作用是给每条事件再加一层来源标识(命名空间),于是形状又变了。先搭一个父图套子图的结构:

# 提供 add 作为 reducer
import operator
# 状态类型
from typing import Annotated, TypedDict

# 读取 .env(这段其实不调模型,留着是为了和前面的代码块保持一致)
from dotenv import load_dotenv
# 子图节点里要推 custom 事件
from langgraph.config import get_stream_writer
# 建图
from langgraph.graph import END, START, StateGraph

# 加载环境变量
load_dotenv(override=True)


# 父图和子图共用同一个状态结构(这样子图可以直接当节点用)
class DocState(TypedDict):
    # 问题
    q: str
    # 文档,带 reducer
    docs: Annotated[list[str], operator.add]
    # 答案
    answer: str


# 子图里的节点:推一条 custom 事件,方便观察它的命名空间
def inner_step(state: DocState) -> dict:
    # 子图节点里同样能取 writer
    writer = get_stream_writer()
    # 推一条带 stage=subgraph 的事件,一眼能认出是子图发的
    writer({"stage": "inner", "status": "running"})
    # 往共享的 docs 字段追加一条
    return {"docs": ["子图产出"]}


# 子图:只有一个节点
inner_builder = StateGraph(DocState)
# 注册子图内部的节点
inner_builder.add_node("step", inner_step)
# 子图的入口
inner_builder.add_edge(START, "step")
# 子图的出口
inner_builder.add_edge("step", END)
# 子图也要 compile,编译后的图才能当节点用
subgraph = inner_builder.compile()

# 父图
builder = StateGraph(DocState)
# 父图的普通节点
builder.add_node("pre", lambda s: {"answer": "预处理完成"})
# 把编译好的子图直接注册成一个节点,节点名叫 sub
builder.add_node("sub", subgraph)
# 父图入口进 pre
builder.add_edge(START, "pre")
# pre 之后进子图
builder.add_edge("pre", "sub")
# 子图跑完结束
builder.add_edge("sub", END)
# 编译父图
graph = builder.compile()

# 统一输入
INPUT = {"q": "x", "docs": [], "answer": ""}

# 组合一:单模式 + subgraphs=True,每条是 (命名空间, 内容)
print("单模式 + subgraphs=True:")
# 两个变量:命名空间 + 内容
for ns, chunk in graph.stream(INPUT, stream_mode="custom", subgraphs=True):
    # ns 是元组,子图事件里带着「子图节点名:任务 id」
    print(f"  ns={ns} chunk={chunk}")

# 组合二:列表模式 + subgraphs=True,每条变成三个元素
print("\n多模式 + subgraphs=True:")
# 不解包,直接收整条,才能看出它有几个元素
for item in graph.stream(INPUT, stream_mode=["custom", "updates"], subgraphs=True):
    # 故意不解包,直接打印长度,确认它是 3 而不是 2
    print(f"  len={len(item)} item={item}")
    # 看一条就够
    break

输出:

单模式 + subgraphs=True:
  ns=('sub:c9c6f836-3e10-2d38-d0f6-b507b41bd275',) chunk={'stage': 'inner', 'status': 'running'}

多模式 + subgraphs=True:
  len=3 item=((), 'updates', {'pre': {'answer': '预处理完成'}})

把四种组合列齐,解包方式各不相同:

stream_mode subgraphs 每条的形状
单个 False chunk
列表 False (mode, chunk)
单个 True (ns, chunk)
列表 True (ns, mode, chunk)

踩法有两种,第二种是安静的:

响亮版:元素个数变了,解包直接失败:代码本来是 for mode, chunk in ...,某天为了排查子图问题加了个 subgraphs=True,三个元素塞不进两个变量,立刻 ValueError: too many values to unpack (expected 2)。看到报错就改,损失有限。

安静版:元素个数没变,但含义换了位置。两个变量的形状之间可以互相「假装兼容」:

# 原来的代码:单模式 + subgraphs=True,ns 是命名空间元组
for ns, chunk in graph.stream(INPUT, stream_mode="custom", subgraphs=True):
    # 靠 len(ns) == 0 判断这条来自主图还是子图
    if len(ns) == 0:
        ...

# 某天有人改成了「列表模式 + 不要 subgraphs」,解包照样成功:
# 此时第一个变量拿到的是模式名字符串 "custom",不再是元组
for ns, chunk in graph.stream(INPUT, stream_mode=["custom", "updates"]):
    # len("custom") 是 6,永远不等于 0,于是所有事件都被当成子图事件丢掉
    # 没有异常、没有警告,只是界面上少了东西
    if len(ns) == 0:
        ...

str 和 tuple 都支持 len() 和 in,所以这类错位能一路跑下去不报错。这也是 §7.5 那个方案真正的价值所在:形状固定之后,这种错位在语法层面就不成立了。

主图的事件里 ns 是空元组 (),子图的是 ('子图节点名:uuid',)。要判断「这条来自主图还是子图」,用 ns == () 比 len(ns) == 0 更安全:前者对字符串永远是 False,能让上面那种错位暴露出来。

觉得这张表太啰嗦?§7.5 有一个一劳永逸的办法。

7.4 子图里的 token 和进度,默认全部丢失 #

这是本章第二个「不报错但要命」的坑,而且比 §4 的交错更隐蔽,因为它是重构引入的。

场景很常见:检索逻辑越写越长,于是把它抽成一个子图。图的行为、状态、返回值一切正常,测试全过。但前端的进度条突然空了。

import operator
from typing import Annotated, TypedDict

from langgraph.config import get_stream_writer
from langgraph.graph import END, START, StateGraph


class DocState(TypedDict):
    q: str
    docs: Annotated[list[str], operator.add]
    answer: str


def inner_step(state: DocState) -> dict:
    writer = get_stream_writer()
    writer({"stage": "inner", "status": "running"})
    return {"docs": ["子图产出"]}


inner_builder = StateGraph(DocState)
inner_builder.add_node("step", inner_step)
inner_builder.add_edge(START, "step")
inner_builder.add_edge("step", END)
subgraph = inner_builder.compile()

builder = StateGraph(DocState)
builder.add_node("pre", lambda s: {"answer": "预处理完成"})
builder.add_node("sub", subgraph)
builder.add_edge(START, "pre")
builder.add_edge("pre", "sub")
builder.add_edge("sub", END)
graph = builder.compile()

INPUT = {"q": "x", "docs": [], "answer": ""}


# 对照实验:同一张图、同一个模式,只有 subgraphs 参数不同
without_ns = list(graph.stream(INPUT, stream_mode="custom"))
# 加上 subgraphs=True
with_ns = list(graph.stream(INPUT, stream_mode="custom", subgraphs=True))
# 不加参数:子图里 writer 推的那条事件根本没出现
print(f"subgraphs=False 收到 {len(without_ns)} 条: {without_ns}")
# 加了参数:事件出现了,还带上了命名空间
print(f"subgraphs=True  收到 {len(with_ns)} 条: {with_ns}")

# 再验证一遍 messages 模式是不是同样的问题
c = list(graph.stream(INPUT, stream_mode="messages"))
# 子图里如果有模型调用,token 也一样收不到(这张图的子图没调模型,所以是 0)
print(f"(对照)messages 模式 subgraphs=False 收到 {len(c)} 条")

输出:

subgraphs=False 收到 0 条: []
subgraphs=True  收到 1 条: [(('sub:28cbd64c-ce40-f8e3-afdb-c7d77d54c9e2',), {'stage': 'inner', 'status': 'running'})]
(对照)messages 模式 subgraphs=False 收到 0 条

0 条。不是少了几条,是一条都没有,而且没有任何警告。

这个行为本身是合理的:subgraphs=False 的语义是「只订阅这一层的事件」,子图对父图来说就是一个普通节点,它内部的动静默认不外泄。问题在于它和「把节点重构成子图」这个动作组合起来时是静默的:

你做的事 图的行为 流式的行为
把 retrieve 节点拆成三个函数 不变 不变
把 retrieve 节点改成子图 不变 进度事件和 token 全部消失

后果按严重程度排:

修法只有一个字:加 subgraphs=True,然后接受多出来的 ns 元素。

约定: 如果你的图现在或将来可能有子图,从第一天就写 subgraphs=True。没有子图时它的成本只是每条事件多一个空元组 (),而事后补的成本是一次线上事故。本章实战那张图没有子图,所以没加;但只要引入子图,翻译层的第一行就该改。

7.5 version="v2":一种形状打天下 #

§7.3 那张四行的表格,本质上是历史包袱:为了向后兼容,stream() 的返回形状随参数变化。LangGraph 提供了一个显式的出路:version="v2",让所有情况都返回同一种形状。

import operator
from typing import Annotated, TypedDict

from langgraph.config import get_stream_writer
from langgraph.graph import END, START, StateGraph


class DocState(TypedDict):
    q: str
    docs: Annotated[list[str], operator.add]
    answer: str


def inner_step(state: DocState) -> dict:
    writer = get_stream_writer()
    writer({"stage": "inner", "status": "running"})
    return {"docs": ["子图产出"]}


inner_builder = StateGraph(DocState)
inner_builder.add_node("step", inner_step)
inner_builder.add_edge(START, "step")
inner_builder.add_edge("step", END)
subgraph = inner_builder.compile()

builder = StateGraph(DocState)
builder.add_node("pre", lambda s: {"answer": "预处理完成"})
builder.add_node("sub", subgraph)
builder.add_edge(START, "pre")
builder.add_edge("pre", "sub")
builder.add_edge("sub", END)
graph = builder.compile()

INPUT = {"q": "x", "docs": [], "answer": ""}


# 加一个 version="v2",其他参数不动
for part in graph.stream(INPUT, stream_mode=["custom", "updates"], subgraphs=True, version="v2"):
    # 不再是元组,而是一个带固定键的 dict:type / ns / data
    print(f"  type={part['type']:8s} ns={part['ns']} data={str(part['data'])[:60]}")

输出:

  type=updates  ns=() data={'pre': {'answer': '预处理完成'}}
  type=custom   ns=('sub:b10080d9-c577-2345-49ef-cea3a093f09d',) data={'stage': 'inner', 'status': 'running'}
  type=updates  ns=('sub:b10080d9-c577-2345-49ef-cea3a093f09d',) data={'step': {'docs': ['子图产出']}}
  type=updates  ns=() data={'sub': {'q': 'x', 'docs': ['子图产出'], 'answer': '预处理完成'}}

三个键,含义一目了然:

键 含义 对应 v1 里的
type 哪种模式 元组里的 mode
ns 来自哪一层(空元组 = 主图) 元组里的 ns
data 内容本身 元组里的 chunk

关键在于这三个键永远都在,不管你传的是单个模式还是列表、subgraphs 是 True 还是 False。§7.3 那张四行表格直接作废:

# v1:形状随参数变,翻译层要按四种情况分别写解包逻辑
for mode, chunk in graph.stream(INPUT, stream_mode=["custom", "updates"]):
    ...

# v2:形状固定,加模式、加 subgraphs 都不用改这一行
for part in graph.stream(INPUT, stream_mode=["custom", "updates"], version="v2"):
    ...

还有一个额外好处,和第 26 章的 interrupt 有关。v1 里 __interrupt__ 是塞进 values 内容里的,所以状态字典会被污染;v2 把它提成了一个独立的字段。搭一张最小的可暂停图来看(§9.2 会用同一张图讲 v1 的表现):

import operator
from typing import Annotated, TypedDict

from langgraph.config import get_stream_writer
from langgraph.graph import END, START, StateGraph


# 内存版 checkpointer,interrupt 要能恢复必须有它
from langgraph.checkpoint.memory import InMemorySaver
# interrupt 用来暂停
from langgraph.types import interrupt


# 一个需要人工确认的状态
class GateState(TypedDict):
    # 问题
    q: str
    # 答案
    answer: str
    # 日志
    log: Annotated[list[str], operator.add]


# 会暂停的节点
def gate(state: GateState) -> dict:
    # 先推一条进度,说明正在等人
    get_stream_writer()({"stage": "gate", "status": "waiting"})
    # 暂停,把问题交给人
    d = interrupt("要继续吗")
    # 恢复后用人的答复组织返回值
    return {"answer": f"人说 {d}"}


# 建图
builder = StateGraph(GateState)
# 注册节点
builder.add_node("gate", gate)
# 入口
builder.add_edge(START, "gate")
# 出口
builder.add_edge("gate", END)
# 编译时带上 checkpointer
graph = builder.compile(checkpointer=InMemorySaver())

# 这张图的输入
GINPUT = {"q": "x", "answer": "", "log": []}
# 会话配置
cfg2 = {"configurable": {"thread_id": "v2demo"}}

# 用 v2 订阅,注意 values 类型的 part 多了一个 interrupts 键
for part in graph.stream(GINPUT, cfg2, stream_mode=["updates", "custom", "values"], version="v2"):
    # 只有 values 类型带 interrupts 字段,其他类型没有这个键,所以要先判断
    extra = f" interrupts={part.get('interrupts')}" if "interrupts" in part else ""
    # 打印类型、数据和中断信息
    print(f"[{part['type']}] {str(part['data'])[:70]}{extra}")

输出:

[values] {'q': 'x', 'answer': '', 'log': []} interrupts=()
[custom] {'stage': 'gate', 'status': 'waiting'}
[updates] {'__interrupt__': (Interrupt(value='要继续吗', id='41d019ca252bfa043bf536a
[values] {'q': 'x', 'answer': '', 'log': []} interrupts=(Interrupt(value='要继续吗', id='41d019ca252bfa043bf536a2e5ff76f9'),)

对比 §9.2 的 v1 输出就能看出差别:v1 的最后一条是 {'q': 'x', 'answer': '', 'log': [], '__interrupt__': (...)}:中断信息混在业务状态里,如果你把 values 直接反序列化成前端模型就会多出一个字段。v2 的 data 保持干净,中断在 interrupts 里。

顺便一提,invoke() 也支持 version="v2",返回一个带两个字段的容器:

import operator
from typing import Annotated, TypedDict

from langgraph.config import get_stream_writer
from langgraph.graph import END, START, StateGraph


class DocState(TypedDict):
    q: str
    docs: Annotated[list[str], operator.add]
    answer: str


def inner_step(state: DocState) -> dict:
    writer = get_stream_writer()
    writer({"stage": "inner", "status": "running"})
    return {"docs": ["子图产出"]}


inner_builder = StateGraph(DocState)
inner_builder.add_node("step", inner_step)
inner_builder.add_edge(START, "step")
inner_builder.add_edge("step", END)
subgraph = inner_builder.compile()

builder = StateGraph(DocState)
builder.add_node("pre", lambda s: {"answer": "预处理完成"})
builder.add_node("sub", subgraph)
builder.add_edge(START, "pre")
builder.add_edge("pre", "sub")
builder.add_edge("sub", END)
graph = builder.compile()

INPUT = {"q": "x", "docs": [], "answer": ""}


# invoke 也可以要求 v2 协议
out = graph.invoke(INPUT, version="v2")
# 返回的不再是 dict,而是 GraphOutput 对象
print(f"类型={type(out).__name__}")
# .value 是最终状态,内容和 v1 的返回值一致
print(f"value 的键={list(out.value.keys())}")
# .interrupts 是这次执行产生的中断,没有中断就是空元组
print(f"interrupts={out.interrupts}")

输出:

类型=GraphOutput
value 的键=['q', 'docs', 'answer']
interrupts=()

这比第 26 章那种「检查返回值里有没有 __interrupt__ 键」的写法干净:out.interrupts 一定存在,不用 in 判断。

那为什么本章实战没用 v2? 三个理由,按重要性排:

  1. 现有的示例、Stack Overflow 答案绝大多数是 v1 形状,你读别人的代码时躲不开,所以必须先看懂 §7.3 那张表。
  2. v1 是默认值(version="v1"),团队里只要有一个人漏写参数,形状就变了。
  3. 本章实战只有一层图、两个模式,v1 的形状已经足够简单。

建议:新写的翻译层用 v2。 特别是当你的图有子图、或者模式组合会变化时,v2 能省掉一整类 bug。已有的代码不用急着改:两者可以在同一个项目里共存,因为它是逐次调用的参数,不是全局开关。

还有个 v3。 stream_events(version="v3") 和 astream_events(version="v3") 是更新的实验协议,事件形状是 {"type": "event", "method": ..., "params": {...}, "seq": ...}(类似 JSON-RPC,带序号,天生适合走 WebSocket)。它被 @beta 标记,调用时会打印 LangChainBetaWarning,形状随时可能变。知道有这回事就行,生产代码现在别用。

8. astream_events:最细的一层 #

第 22 章的建议是「先用 stream_mode,九成需求够了」。现在依然成立。这一节讲剩下那一成。

两者的区别是视角:

stream(stream_mode=...) astream_events()
视角 图视角:哪个节点动了、状态怎么变 调用栈视角:每个 Runnable 的开始、流动、结束
事件密度 一个节点几条 一个节点十几条,加上模型每个 token 一条
能拿到 节点产出、token、自定义进度 上面全部,外加工具的解析后参数、tool_call_chunks
同步版本 有(stream) 没有稳定的同步版(见下)

同步版本有个容易误判的地方。 Pregel 上确实有 stream_events(),看名字像是 astream_events 的同步版,但它只接受实验性的 v3:

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langchain_core.tools import tool
from langgraph.graph import END, START, MessagesState, StateGraph
from langgraph.prebuilt import ToolNode

load_dotenv(override=True)


@tool
def get_stock(sku: str) -> str:
    """查询某个 SKU 的库存数量。"""
    return {"A-100": "库存 12 件"}.get(sku, "查无此 SKU")


TOOLS = [get_stock]
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0).bind_tools(TOOLS)


def call_model(state: MessagesState) -> dict:
    return {"messages": [model.invoke(state["messages"])]}


def should_continue(state: MessagesState):
    return "tools" if getattr(state["messages"][-1], "tool_calls", None) else END


builder = StateGraph(MessagesState)
builder.add_node("model", call_model)
builder.add_node("tools", ToolNode(TOOLS))
builder.add_edge(START, "model")
builder.add_conditional_edges("model", should_continue, ["tools", END])
builder.add_edge("tools", "model")
graph = builder.compile()

Q = {"messages": [{"role": "user", "content": "A-100 有货吗"}]}


# 试着用同步的 stream_events 要 v2 协议
try:
    # 看起来天经地义的写法
    for ev in graph.stream_events(Q, version="v2"):
        pass
# 捕获并打印异常,重点看提示里说了什么
except NotImplementedError as e:
    # 报错信息直接给出了替代方案
    print(f"{type(e).__name__}: {e}")

输出:

NotImplementedError: stream_events(version='v2') is not supported. Use astream_events() for v1/v2, or stream_events(version='v3') on a supported subclass.

报错信息把结论说得很清楚:v1/v2 只有异步版,同步版只有实验性的 v3。 如果你的服务是同步框架(Flask、同步 Django),用 astream_events 就得自己套一层 asyncio.run 或线程池;这也是本章实战选 stream_mode 的现实原因之一。

8.1 事件清单 #

拿 §5.2 那张带工具的图,把事件类型统计出来。先看看一次「查个库存」到底产生多少事件:

# 异步入口需要 asyncio
import asyncio
# 统计各类型的数量
from collections import Counter

# 读取 .env
from dotenv import load_dotenv
# 创建模型
from langchain.chat_models import init_chat_model
# 工具装饰器
from langchain_core.tools import tool
# 建图
from langgraph.graph import END, START, MessagesState, StateGraph
# 预置工具节点
from langgraph.prebuilt import ToolNode

# 加载环境变量
load_dotenv(override=True)


# 和 §5.2 完全相同的工具
@tool
def get_stock(sku: str) -> str:
    """查询某个 SKU 的库存数量。"""
    # 字典模拟数据库
    return {"A-100": "库存 12 件"}.get(sku, "查无此 SKU")


# 工具清单
TOOLS = [get_stock]
# 绑定工具
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0).bind_tools(TOOLS)


# 模型节点
def call_model(state: MessagesState) -> dict:
    # 返回一条新消息
    return {"messages": [model.invoke(state["messages"])]}


# 路由:有 tool_calls 就去工具,否则结束
def should_continue(state: MessagesState):
    # getattr 兜住没有该属性的消息类型
    return "tools" if getattr(state["messages"][-1], "tool_calls", None) else END


# 建图
builder = StateGraph(MessagesState)
# 模型节点
builder.add_node("model", call_model)
# 工具节点
builder.add_node("tools", ToolNode(TOOLS))
# 入口
builder.add_edge(START, "model")
# 条件边
builder.add_conditional_edges("model", should_continue, ["tools", END])
# 工具回到模型,形成回环
builder.add_edge("tools", "model")
# 编译
graph = builder.compile()

# 会触发工具调用的问题
Q = {"messages": [{"role": "user", "content": "A-100 有货吗"}]}


# astream_events 是异步生成器,必须在协程里消费
async def main():
    # 按事件类型计数
    kinds = Counter()
    # async for 遍历事件流;每个 ev 是一个 dict
    async for ev in graph.astream_events(Q, version="v2"):
        # ev["event"] 就是事件类型名,比如 on_chat_model_stream
        kinds[ev["event"]] += 1
    # 先打印总数,直观感受事件密度
    print(f"总事件数 {sum(kinds.values())}")
    # 按出现次数从多到少打印
    for k, v in kinds.most_common():
        # 左对齐到 28 列,方便对照
        print(f"{k:28s} x{v}")


# 启动事件循环
asyncio.run(main())

输出:

总事件数 81
on_chat_model_stream         x57
on_chain_start               x6
on_chain_end                 x6
on_chain_stream              x6
on_chat_model_start          x2
on_chat_model_end            x2
on_tool_start                x1
on_tool_end                  x1

事件类型按前缀分三组:

前缀 来自 这次的数量
on_chat_model_* 模型调用(start / stream / end) 61 条,其中 57 条是 token
on_chain_* 图本身和每个节点的包装(LangGraph 把它们都算 chain) 18 条
on_tool_* 工具执行 2 条

一次「查个库存」产生 81 个事件,其中你可能真正关心的只有 2 个(on_tool_start 和 on_tool_end)。这就是用 astream_events 的第一个代价,所以 §8.3 的过滤参数几乎是必须配的,而不是可选优化。

顺便看一眼 on_chain_* 的 name 字段,就知道那 18 条是什么:图本身的 name 是 LangGraph,每个节点的 name 是节点名(model、tools)。也就是说,节点的进出都会各产生一条 chain 事件,这些对界面基本没用,exclude_types=["chain"] 一次滤掉。

事件总数每次运行都不同(这次 81,先前两轮记录的是 93 和 116),因为模型输出长度会变;看数量级就行,别对着数字调代码。

关于 version 参数:astream_events 的默认值就是 "v2",写不写都一样(可以用 inspect.signature(graph.astream_events).parameters['version'].default 自己确认)。version="v1" 是遗留版本,事件命名和结构都和 v2 不同,新代码不要用。注意这个 version 和 §7.5 那个 stream(version=...) 是两套完全独立的东西,只是名字撞了:前者管事件协议的代号,后者管流式返回值的形状。

8.2 事件结构 #

每个事件的字段是固定的七个。用「每种类型留一个样本」的办法把结构打出来:

import asyncio

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langchain_core.tools import tool
from langgraph.graph import END, START, MessagesState, StateGraph
from langgraph.prebuilt import ToolNode

load_dotenv(override=True)


@tool
def get_stock(sku: str) -> str:
    """查询某个 SKU 的库存数量。"""
    return {"A-100": "库存 12 件"}.get(sku, "查无此 SKU")


TOOLS = [get_stock]
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0).bind_tools(TOOLS)


def call_model(state: MessagesState) -> dict:
    return {"messages": [model.invoke(state["messages"])]}


def should_continue(state: MessagesState):
    return "tools" if getattr(state["messages"][-1], "tool_calls", None) else END


builder = StateGraph(MessagesState)
builder.add_node("model", call_model)
builder.add_node("tools", ToolNode(TOOLS))
builder.add_edge(START, "model")
builder.add_conditional_edges("model", should_continue, ["tools", END])
builder.add_edge("tools", "model")
graph = builder.compile()

Q = {"messages": [{"role": "user", "content": "A-100 有货吗"}]}


async def main():
    # 用 dict 给每种事件类型留第一个样本
    samples = {}
    # 跑一遍,setdefault 保证只留第一次出现的那条
    async for ev in graph.astream_events(Q, version="v2"):
        # 已存在的键不会被覆盖,所以留下的是每种类型的第一条
        samples.setdefault(ev["event"], ev)

    # 挑三种最有代表性的类型来看结构
    for k in ["on_chat_model_stream", "on_tool_start", "on_tool_end"]:
        # 取出样本
        ev = samples[k]
        # 打印类型名作为小标题
        print(f"\n{k}:")
        # name 是组件名:模型是类名,工具是工具名
        print(f"  name  = {ev['name']}")
        # metadata 里同样有 langgraph_node,和 messages 模式一致
        print(f"  node  = {(ev.get('metadata') or {}).get('langgraph_node')}")
        # 顶层七个键,所有事件类型都一样
        print(f"  keys  = {sorted(ev.keys())}")
        # data 里的键随事件类型变化,这是唯一不固定的部分
        print(f"  data  = {sorted(ev['data'].keys())}")


# 启动
asyncio.run(main())

输出:

on_chat_model_stream:
  name  = ChatDeepSeek
  node  = model
  keys  = ['data', 'event', 'metadata', 'name', 'parent_ids', 'run_id', 'tags']
  data  = ['chunk']

on_tool_start:
  name  = get_stock
  node  = tools
  keys  = ['data', 'event', 'metadata', 'name', 'parent_ids', 'run_id', 'tags']
  data  = ['input']

on_tool_end:
  name  = get_stock
  node  = tools
  keys  = ['data', 'event', 'metadata', 'name', 'parent_ids', 'run_id', 'tags']
  data  = ['input', 'output']

七个顶层字段,逐个说清用途:

字段 内容 用途
event 事件类型名 分派的第一依据
name 组件名:模型类名 ChatDeepSeek、工具名 get_stock、节点名 按工具名筛选靠它(§8.3 的 include_names)
data 载荷,唯一随类型变化的字段 _stream 给 chunk,_start 给 input,_end 给 output
metadata 和 messages 模式同一份元数据 里面有 langgraph_node、ls_model_name 等
tags 组件上的标签 配合 §3.4 做最精确的筛选
run_id 这次组件调用的唯一 id 关联同一次调用的 start / stream / end
parent_ids 祖先调用的 id 列表 拼调用树,做自定义追踪(这活儿 LangSmith 更擅长,第 29 章讲)

工具事件比 messages 模式好用的地方就在 data 上:on_tool_start 的 data["input"] 是解析好的参数字典(不是 JSON 字符串),on_tool_end 的 data["output"] 是完整的 ToolMessage 对象。想在界面上显示「正在查询 A-100 的库存」,需要的就是这个参数,而 messages 模式给不了,那边只有一条 ToolMessage,参数已经不在了。

import asyncio

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langchain_core.tools import tool
from langgraph.graph import END, START, MessagesState, StateGraph
from langgraph.prebuilt import ToolNode

load_dotenv(override=True)


@tool
def get_stock(sku: str) -> str:
    """查询某个 SKU 的库存数量。"""
    return {"A-100": "库存 12 件"}.get(sku, "查无此 SKU")


TOOLS = [get_stock]
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0).bind_tools(TOOLS)


def call_model(state: MessagesState) -> dict:
    return {"messages": [model.invoke(state["messages"])]}


def should_continue(state: MessagesState):
    return "tools" if getattr(state["messages"][-1], "tool_calls", None) else END


builder = StateGraph(MessagesState)
builder.add_node("model", call_model)
builder.add_node("tools", ToolNode(TOOLS))
builder.add_edge(START, "model")
builder.add_conditional_edges("model", should_continue, ["tools", END])
builder.add_edge("tools", "model")
graph = builder.compile()

Q = {"messages": [{"role": "user", "content": "A-100 有货吗"}]}


async def main():
    # 这次不统计,直接按类型分派打印,模拟真实的翻译层
    async for ev in graph.astream_events(Q, version="v2"):
        # 先取事件类型
        kind = ev["event"]
        # 分支一:模型片段
        if kind == "on_chat_model_stream":
            # chunk 是 AIMessageChunk,取它的正文
            c = ev["data"]["chunk"].content
            # 依旧要过滤空片段(§5.1),工具调用阶段的片段就是空的
            if c:
                # 打印正文片段
                print(f"token: {c!r}")
        # 分支二:工具开始执行,这里能拿到解析好的参数字典
        elif kind == "on_tool_start":
            # ev["name"] 是工具名,data["input"] 是 {'sku': 'A-100'}
            print(f">> 调用 {ev['name']} 参数={ev['data']['input']}")
        # 分支三:工具执行完毕
        elif kind == "on_tool_end":
            # data["output"] 是完整的 ToolMessage 对象,取 content 拿返回值
            print(f"<< 返回 {ev['data']['output'].content!r}")


# 启动
asyncio.run(main())

输出(token 部分省略):

>> 调用 get_stock 参数={'sku': 'A-100'}
<< 返回 '库存 12 件'
token: 'A'
token: '-'
token: '100'
token: ' '
token: '有'
token: '货'
...
token: '。'

注意工具事件出现在所有 token 之前,原因和 §5.2 一样:第一个模型回合全程在生成函数调用,正文是空的,被 if c 过滤掉了(§5.1 说的情况)。

8.3 三个过滤参数 #

事件太多,astream_events 提供了六个过滤参数:include_types / include_names / include_tags 和对应的三个 exclude_*。三个维度分别对应「哪一类组件」「哪一个组件」「哪一批组件」。

按类型筛:include_types=["chat_model"] 只要模型事件:

import asyncio
from collections import Counter

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langchain_core.tools import tool
from langgraph.graph import END, START, MessagesState, StateGraph
from langgraph.prebuilt import ToolNode

load_dotenv(override=True)


@tool
def get_stock(sku: str) -> str:
    """查询某个 SKU 的库存数量。"""
    return {"A-100": "库存 12 件"}.get(sku, "查无此 SKU")


TOOLS = [get_stock]
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0).bind_tools(TOOLS)


def call_model(state: MessagesState) -> dict:
    return {"messages": [model.invoke(state["messages"])]}


def should_continue(state: MessagesState):
    return "tools" if getattr(state["messages"][-1], "tool_calls", None) else END


builder = StateGraph(MessagesState)
builder.add_node("model", call_model)
builder.add_node("tools", ToolNode(TOOLS))
builder.add_edge(START, "model")
builder.add_conditional_edges("model", should_continue, ["tools", END])
builder.add_edge("tools", "model")
graph = builder.compile()

Q = {"messages": [{"role": "user", "content": "A-100 有货吗"}]}


async def main():
    # 计数器
    kinds = Counter()
    # include_types 收的是组件类型,注意值是 chat_model 而不是事件名 on_chat_model_stream
    async for ev in graph.astream_events(Q, version="v2", include_types=["chat_model"]):
        # 统计剩下的事件
        kinds[ev["event"]] += 1
    # 打印分布,chain 和 tool 事件都不见了
    print(dict(kinds))


# 启动
asyncio.run(main())

输出:

{'on_chat_model_start': 2, 'on_chat_model_stream': 68, 'on_chat_model_end': 2}

按名字筛:include_names=["get_stock"] 只要某个工具的事件,这是「显示正在调用哪个工具」最直接的写法:

import asyncio

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langchain_core.tools import tool
from langgraph.graph import END, START, MessagesState, StateGraph
from langgraph.prebuilt import ToolNode

load_dotenv(override=True)


@tool
def get_stock(sku: str) -> str:
    """查询某个 SKU 的库存数量。"""
    return {"A-100": "库存 12 件"}.get(sku, "查无此 SKU")


TOOLS = [get_stock]
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0).bind_tools(TOOLS)


def call_model(state: MessagesState) -> dict:
    return {"messages": [model.invoke(state["messages"])]}


def should_continue(state: MessagesState):
    return "tools" if getattr(state["messages"][-1], "tool_calls", None) else END


builder = StateGraph(MessagesState)
builder.add_node("model", call_model)
builder.add_node("tools", ToolNode(TOOLS))
builder.add_edge(START, "model")
builder.add_conditional_edges("model", should_continue, ["tools", END])
builder.add_edge("tools", "model")
graph = builder.compile()

Q = {"messages": [{"role": "user", "content": "A-100 有货吗"}]}


async def main():
    # include_names 匹配的是事件里的 name 字段(工具名 / 模型类名 / 节点名)
    async for ev in graph.astream_events(Q, version="v2", include_names=["get_stock"]):
        # 81 个事件里只剩下这个工具的两条
        print(f"{ev['event']:16s} name={ev['name']}")


# 启动
asyncio.run(main())

输出:

on_tool_start    name=get_stock
on_tool_end      name=get_stock

反向排除:exclude_types=["chain"] 去掉图和节点包装产生的噪音,保留模型和工具:

import asyncio
from collections import Counter

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langchain_core.tools import tool
from langgraph.graph import END, START, MessagesState, StateGraph
from langgraph.prebuilt import ToolNode

load_dotenv(override=True)


@tool
def get_stock(sku: str) -> str:
    """查询某个 SKU 的库存数量。"""
    return {"A-100": "库存 12 件"}.get(sku, "查无此 SKU")


TOOLS = [get_stock]
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0).bind_tools(TOOLS)


def call_model(state: MessagesState) -> dict:
    return {"messages": [model.invoke(state["messages"])]}


def should_continue(state: MessagesState):
    return "tools" if getattr(state["messages"][-1], "tool_calls", None) else END


builder = StateGraph(MessagesState)
builder.add_node("model", call_model)
builder.add_node("tools", ToolNode(TOOLS))
builder.add_edge(START, "model")
builder.add_conditional_edges("model", should_continue, ["tools", END])
builder.add_edge("tools", "model")
graph = builder.compile()

Q = {"messages": [{"role": "user", "content": "A-100 有货吗"}]}


async def main():
    # 计数器
    kinds = Counter()
    # exclude_ 系列是黑名单,和 include_ 系列相反;这里去掉 18 条 chain 事件
    async for ev in graph.astream_events(Q, version="v2", exclude_types=["chain"]):
        # 统计
        kinds[ev["event"]] += 1
    # 模型和工具的事件都留着,噪音少了一大半
    print(dict(kinds))


# 启动
asyncio.run(main())

输出:

{'on_chat_model_start': 2, 'on_chat_model_stream': 55, 'on_chat_model_end': 2, 'on_tool_start': 1, 'on_tool_end': 1}

三个维度的取舍:

参数 匹配的字段 适合
include_types 组件类型(chat_model / tool / chain / retriever / prompt) 粗筛,先把一大类噪音去掉
include_names name 字段 精确到某个工具或某个模型
include_tags tags 字段 最推荐:配合 §3.4 的标签,跨组件类型筛「同一个用途」

include_tags 是最精确也最稳定的一种:给最终回答的模型打 final-answer,然后 include_tags=["final-answer"],一步到位。它比 include_names 好的地方在于换模型不用改筛选条件:name 是模型类名(ChatDeepSeek),换成别的供应商就变了;标签是你自己起的,不会变。

注意 include_* 和 exclude_* 不要混用同一个维度。 同时传 include_types 和 exclude_types 时的行为容易反直觉(先按白名单选,再按黑名单剔),调试成本高。一个维度只用一个方向。

8.4 tool_call_chunks:把「正在准备调用工具」也流出来 #

这是 astream_events 相比 messages 模式唯一真正独有的能力,值得单独一节。

模型生成函数调用参数的过程,其实也是流式的:先吐工具名,再逐字吐参数 JSON。这些片段的 content 是空的(所以在 §5.1 被过滤了),但它们的 tool_call_chunks 属性里有东西:

import asyncio

from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
from langchain_core.tools import tool
from langgraph.graph import END, START, MessagesState, StateGraph
from langgraph.prebuilt import ToolNode

load_dotenv(override=True)


@tool
def get_stock(sku: str) -> str:
    """查询某个 SKU 的库存数量。"""
    return {"A-100": "库存 12 件"}.get(sku, "查无此 SKU")


TOOLS = [get_stock]
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0).bind_tools(TOOLS)


def call_model(state: MessagesState) -> dict:
    return {"messages": [model.invoke(state["messages"])]}


def should_continue(state: MessagesState):
    return "tools" if getattr(state["messages"][-1], "tool_calls", None) else END


builder = StateGraph(MessagesState)
builder.add_node("model", call_model)
builder.add_node("tools", ToolNode(TOOLS))
builder.add_edge(START, "model")
builder.add_conditional_edges("model", should_continue, ["tools", END])
builder.add_edge("tools", "model")
graph = builder.compile()

Q = {"messages": [{"role": "user", "content": "A-100 有货吗"}]}


async def main():
    # 收集所有带 tool_call_chunks 的片段
    seen = []
    # 遍历事件流
    async for ev in graph.astream_events(Q, version="v2"):
        # 只看模型片段事件
        if ev["event"] == "on_chat_model_stream":
            # 取出 AIMessageChunk
            ch = ev["data"]["chunk"]
            # getattr 兜一下:不是所有片段都有这个属性,有的话也可能是空列表
            if getattr(ch, "tool_call_chunks", None):
                # 收下这一批(一个片段里可能有多个并行的工具调用)
                seen.append(ch.tool_call_chunks)
    # 先看总数,感受一下参数是分多少片流出来的
    print(f"带 tool_call_chunks 的 chunk 数: {len(seen)}")
    # 只看前四片,重点是第一片和后面的区别
    for s in seen[:4]:
        # 每条是一个列表,元素是 tool_call_chunk 字典
        print(s)


# 启动
asyncio.run(main())

输出:

带 tool_call_chunks 的 chunk 数: 12
[{'name': 'get_stock', 'args': '', 'id': 'call_00_wtH4aoRRxa6eeCelVJVF7376', 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '{', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': 'sku', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]

注意第一片和后面的区别:

这个结构意味着:模型刚决定要调工具的那一刻(第一片),你就能在界面上显示「正在查询库存…」,不用等参数拼完(省下 11 个片段的时间)。如果参数很长(比如让模型生成一段 SQL 或一大段 JSON),这个提前量能到几百毫秒,用户是能感觉到的。

要拼出完整参数的话,把 args 按 index 分组累加成字符串,再 json.loads。但大多数情况不用自己拼:on_tool_start 的 data["input"](§8.2)已经是解析好的字典了。tool_call_chunks 的价值只在「提前」这两个字上。

8.5 选型 #

需求 用什么
聊天界面流字 stream_mode="messages"
步骤进度条 stream_mode="custom"
前两者一起 stream_mode=["custom", "messages"]
上面这些,但想少写解包逻辑 加 version="v2"(§7.5)
图里有子图 加 subgraphs=True(§7.4),别忘
显示「正在调用某工具(参数是什么)」 astream_events 的 on_tool_start
抢在参数生成完之前提示 astream_events 的 tool_call_chunks(§8.4)
自己搭调用树 / 做追踪 astream_events 的 run_id + parent_ids,或者直接上 LangSmith
排查图的执行顺序 stream_mode="tasks"(第 22 章 §4.6)
服务是同步框架 只能 stream_mode,astream_events 没有稳定的同步版(§8 开头)

默认选 stream_mode。 三个理由:同步异步都有、形状简单、和图的概念一一对应。只有上表下半部分那几种需求,才值得付「81 个事件里挑 2 个」的代价。本章实战用的就是 stream_mode。

有人会想「两个一起用」:主流程走 stream_mode,再开一条 astream_events 只订阅工具事件。别这么做:一次 stream() 或 astream_events() 就是一次完整的图执行,两条流意味着图跑两遍、模型收两次费、副作用做两次。

真正需要「工具执行进度」时,正确的做法是 §6.2:让工具自己用 writer 推一条 custom 事件。一行代码,走的还是同一条流。只有当你需要的是模型还在拼参数时就提示(§8.4 的 tool_call_chunks),才不得不换成 astream_events,因为那个时刻工具函数还没开始执行,writer 无从下手。

9. 流式下的失败与暂停 #

前面八节讲的都是「顺利跑完」。真实服务里还有两种结局:挂了,以及停下来等人。流式下它们多出一个共同问题:已经推给前端的内容收不回来了。invoke 视角下不存在这事(要么返回结果要么抛异常);流式下必须正面处理。

9.1 节点抛异常:已产出的 chunk 先送到,异常从迭代器抛出 #

第 22 章 §7 讲的是 invoke 视角。流式下要先搞清两件事:异常前推出去的内容还在不在,以及异常到底从哪一行抛出来。

# 提供 add 作为 reducer
import operator
# 状态类型
from typing import Annotated, TypedDict

# 失败节点里要先推一条进度,才能验证「异常前的 chunk 是否送达」
from langgraph.config import get_stream_writer
# 建图;这段不需要模型,也不需要 .env
from langgraph.graph import END, START, StateGraph


# 状态结构
class FlowState(TypedDict):
    # 文档列表,带 reducer
    docs: Annotated[list[str], operator.add]
    # 答案
    answer: str


# 会抛异常的节点
def boom(state: FlowState) -> dict:
    # 先取 writer
    writer = get_stream_writer()
    # 关键:异常之前先推一条事件,模拟「前端已经渲染了半句话」
    writer({"stage": "boom", "status": "start"})
    # 模拟下游服务挂了
    raise RuntimeError("下游服务 503")


# 建图
builder = StateGraph(FlowState)
# 第一个节点正常完成,用来确认它的 updates 事件也送达了
builder.add_node("ok", lambda s: {"docs": ["a"]})
# 第二个节点炸掉
builder.add_node("boom", boom)
# 入口进 ok
builder.add_edge(START, "ok")
# ok 之后进 boom
builder.add_edge("ok", "boom")
# boom 之后结束(实际到不了)
builder.add_edge("boom", END)
# 编译
graph = builder.compile()

# 记录异常前收到的所有 chunk
got = []
# try 包住整个 for 循环,这个位置很重要,下面会解释
try:
    # 同时订阅 updates 和 custom
    for mode, chunk in graph.stream({"docs": [], "answer": ""}, stream_mode=["updates", "custom"]):
        # 收下并打印,模拟推给前端
        got.append((mode, chunk))
        # 打印出来,确认它在异常之前确实送达了
        print(f"收到 [{mode}] {chunk}")
# 捕获节点里抛出的异常
except Exception as e:
    # 打印异常类型和消息,确认它原样传上来了(没有被包装)
    print(f"异常在迭代中抛出: {type(e).__name__}: {e}")
# 统计异常前到底送达了几条
print(f"异常前收到 {len(got)} 个 chunk")

输出:

收到 [updates] {'ok': {'docs': ['a']}}
收到 [custom] {'stage': 'boom', 'status': 'start'}
异常在迭代中抛出: RuntimeError: 下游服务 503
异常前收到 2 个 chunk

两个结论:

  1. 异常前的 chunk 都正常送达了,包括失败节点自己推的那条 boom/start。前端已经渲染了这些内容,不能假装没发生过,所以收尾方式不是「清空重来」,而是「补一个失败状态」。
  2. 异常是从 for 循环里抛出来的,不是在 stream() 调用处。

第二点值得单独验证一下,因为它决定 try 该写在哪:

import operator
from typing import Annotated, TypedDict

from langgraph.config import get_stream_writer
from langgraph.graph import END, START, StateGraph


class FlowState(TypedDict):
    docs: Annotated[list[str], operator.add]
    answer: str


def boom(state: FlowState) -> dict:
    writer = get_stream_writer()
    writer({"stage": "boom", "status": "start"})
    raise RuntimeError("下游服务 503")


builder = StateGraph(FlowState)
builder.add_node("ok", lambda s: {"docs": ["a"]})
builder.add_node("boom", boom)
builder.add_edge(START, "ok")
builder.add_edge("ok", "boom")
builder.add_edge("boom", END)
graph = builder.compile()


# 只调用 stream(),不消费它
try:
    # stream() 返回的是生成器,此时图还一行都没跑
    it = graph.stream({"docs": [], "answer": ""}, stream_mode=["updates", "custom"])
    # 所以这里不会有任何异常
    print("stream() 调用本身没有抛异常,拿到迭代器: " + type(it).__name__)
# 这个 except 永远不会命中
except Exception as e:
    # 如果真打印了这一行,说明前面的结论是错的
    print(f"stream() 处就抛了: {type(e).__name__}")
# 现在才开始真正消费迭代器
try:
    # 消费的过程中图才执行,异常在这里才出现
    for x in it:
        pass
# 这个 except 才是能拦住它的那个
except Exception as e:
    # 打印异常,证明它出现在消费阶段
    print(f"异常确实在消费迭代器时才抛: {type(e).__name__}: {e}")

输出:

stream() 调用本身没有抛异常,拿到迭代器: generator
异常确实在消费迭代器时才抛: RuntimeError: 下游服务 503

stream() 是生成器函数,调用它只是造了个迭代器,图一行都没跑。所以 try 必须包住整个循环:

# 错的:图还没开始跑,这个 try 什么都拦不到
try:
    # 这一行只是造了个生成器
    it = graph.stream(payload, cfg, stream_mode=["custom", "messages"])
except Exception:
    ...
for mode, chunk in it:      # 异常在这里抛,没人管
    ...

# 对的:try 包住消费的整个过程
try:
    # 迭代放在 try 里面,节点抛的异常才拦得住
    for mode, chunk in graph.stream(payload, cfg, stream_mode=["custom", "messages"]):
        ...
except Exception:
    ...

正确的收尾方式是补一个错误事件,让前端从「加载中」切到「失败」,而不是把异常直接往上扔:那样 HTTP 连接会断开,前端的 EventSource 只看到连接中断,界面永远转圈:

# 翻译层的骨架:异常兜底 + 明确终止
def to_events(graph, payload, cfg):
    # try 包住整个循环(上面刚验证过的位置)
    try:
        # 正常路径:把图的原始流翻译成前端事件
        for mode, chunk in graph.stream(payload, cfg, stream_mode=["custom", "messages"]):
            # 这里是各模式的翻译逻辑,§10.2 有完整版
            ...
    # 任何异常都转成一个 error 事件,而不是继续往上抛
    except Exception as exc:
        # 带上异常类型和消息,前端可以显示也可以只记日志
        yield {"type": "error", "message": f"{type(exc).__name__}: {exc}"}
        # 发完 error 就结束,别再 yield done,否则前端会收到两个终止事件
        return
    # 只有正常跑完才发 done
    yield {"type": "done"}

注意那个 return。 少了它,异常路径会同时发出 error 和 done,前端的状态机就乱了:先标红又收起加载态,用户看到的是一闪而过的错误。任何一次请求,终止事件必须恰好一个。

9.2 interrupt 在流里长什么样 #

第 26 章用 invoke 看 interrupt,返回值里有个 __interrupt__ 键。流式下它出现在哪条流里,是个必须确认的问题,因为前端要靠它弹审批框。

# 提供 add 作为 reducer
import operator
# 状态类型
from typing import Annotated, TypedDict

# interrupt 需要 checkpointer 才能恢复,内存版够演示(第 26 章 §3.1)
from langgraph.checkpoint.memory import InMemorySaver
# 节点里推进度
from langgraph.config import get_stream_writer
# 建图
from langgraph.graph import END, START, StateGraph
# Command 用来携带人的答复,interrupt 用来暂停
from langgraph.types import Command, interrupt


# 状态结构
class GateState(TypedDict):
    # 问题
    q: str
    # 答案
    answer: str
    # 日志
    log: Annotated[list[str], operator.add]


# 需要人工确认的节点
def gate(state: GateState) -> dict:
    # 取 writer
    writer = get_stream_writer()
    # 在 interrupt 之前推一条进度,§9.3 会看到这一行的副作用
    writer({"stage": "gate", "status": "waiting"})
    # 暂停并把问题抛给人;恢复后 d 就是人给的答复
    d = interrupt("要继续吗")
    # 用人的答复组织返回值
    return {"answer": f"人说 {d}"}


# 建图
builder = StateGraph(GateState)
# 只有一个节点
builder.add_node("gate", gate)
# 入口
builder.add_edge(START, "gate")
# 出口
builder.add_edge("gate", END)
# 编译时必须带 checkpointer,否则 resume 会失败
graph = builder.compile(checkpointer=InMemorySaver())

# thread_id 标识这次会话,resume 时要用同一个
cfg = {"configurable": {"thread_id": "s1"}}
# 输入
INPUT = {"q": "x", "answer": "", "log": []}

# 三种模式一起订阅,看 __interrupt__ 出现在哪几条里
for mode, chunk in graph.stream(INPUT, cfg, stream_mode=["updates", "custom", "values"]):
    # 截断打印,避免输出过长
    print(f"[{mode}] {str(chunk)[:80]}")

# 分隔线
print("--- resume ---")
# 用 Command(resume=...) 恢复;输入位置传 Command 而不是新的状态
for mode, chunk in graph.stream(Command(resume="好"), cfg, stream_mode=["updates", "custom"]):
    # 同样截断打印
    print(f"[{mode}] {str(chunk)[:80]}")

输出:

[values] {'q': 'x', 'answer': '', 'log': []}
[custom] {'stage': 'gate', 'status': 'waiting'}
[updates] {'__interrupt__': (Interrupt(value='要继续吗', id='32ea75d1bc61c33657700917de347b37'
[values] {'q': 'x', 'answer': '', 'log': [], '__interrupt__': (Interrupt(value='要继续吗', id='32ea75d1
--- resume ---
[custom] {'stage': 'gate', 'status': 'waiting'}
[updates] {'gate': {'answer': '人说 好'}}

三点:

一、__interrupt__ 只出现在 updates 和 values 里,不在 messages 里。 顺手验证一下,免得有人指望从 token 流里发现暂停:

import operator
from typing import Annotated, TypedDict

from langgraph.checkpoint.memory import InMemorySaver
from langgraph.config import get_stream_writer
from langgraph.graph import END, START, StateGraph
from langgraph.types import interrupt


class GateState(TypedDict):
    q: str
    answer: str
    log: Annotated[list[str], operator.add]


def gate(state: GateState) -> dict:
    writer = get_stream_writer()
    writer({"stage": "gate", "status": "waiting"})
    d = interrupt("要继续吗")
    return {"answer": f"人说 {d}"}


builder = StateGraph(GateState)
builder.add_node("gate", gate)
builder.add_edge(START, "gate")
builder.add_edge("gate", END)
graph = builder.compile(checkpointer=InMemorySaver())

INPUT = {"q": "x", "answer": "", "log": []}


cfg3 = {"configurable": {"thread_id": "s3"}}
# 数一下 messages 流里有几条
n = 0
# 这张图的节点根本不调模型,所以本来就是空的;但重点是暂停信息也不在这里
for chunk in graph.stream(INPUT, cfg3, stream_mode="messages"):
    # 有多少条就加多少次
    n += 1
# 结论:只订阅 messages 的前端永远不知道图停下来了
print(f"messages 流收到 {n} 条(interrupt 不会出现在这里)")

输出:

messages 流收到 0 条(interrupt 不会出现在这里)

只订阅 messages 的前端永远不会知道图停下来了,它只会一直等下一个 token。所以翻译层必须同时订阅 updates。

二、在 updates 流里,__interrupt__ 占的是「节点名」的位置(和 {'gate': {...}} 同一层),所以判断方式是查键:

# updates 的形状是 {节点名: 更新内容},而 __interrupt__ 混在同一层,用 in 判断
if "__interrupt__" in chunk:
    # 值是 Interrupt 元组;并行中断时可能有多个,这里取第一个(第 26 章 §7.1)
    itr = chunk["__interrupt__"][0]
    # 翻译成前端事件:id 用于 resume 时对应,value 是给人看的内容
    yield {"type": "interrupt", "id": itr.id, "value": itr.value}

如果嫌这种「特殊键」的写法不干净,§7.5 的 version="v2" 把它提成了独立字段。

三、「暂停了」和「跑完了」是两种不同的结束。 前端对它们的处理完全相反:一个要弹审批框继续等,一个要收起加载状态。所以事件协议里得有两个终止事件,不能都叫 done。本章实战用 paused 和 done 区分,加上 §9.1 的 error,一共三种终止。

9.3 恢复时进度事件会重发 #

上面的输出里有个细节容易被忽略:resume 之后,{'stage': 'gate', 'status': 'waiting'} 又出现了一次。

这就是第 26 章 §4 那个坑在流式场景下的表现:恢复执行时,被中断的节点会从头重跑,interrupt() 之前的所有代码都会再执行一遍,包括那句 writer(...)。第 26 章担心的后果是重复扣款,这里的后果轻得多但一样烦人:进度条上会多出一个已经完成过的步骤,用户看到「正在等待确认」又出现了一次,以为要再批一遍。

修法是去重,但去重集合放在哪里是关键。它必须活得比单次 stream 更长:因为第一次和第二次是两次独立的 stream 调用(中间隔了几分钟甚至几小时,可能还换了进程,见第 26 章 §8)。

# 错的:seen 是函数局部变量,第二次调用时又是个全新的空集合,去重完全没生效
def to_events(graph, payload, cfg):
    # 这一行每次调用都重新执行
    seen = set()
    # 于是恢复时那条重发的进度事件照样会被 yield 出去
    ...


# 对的:让调用方传进来,由会话持有
def to_events(graph, payload, cfg, seen):
    # seen 的生命周期由外面决定,跨两次 stream 调用依然是同一个集合
    ...


# 会话级的去重集合;真实服务里放 Redis 或会话对象,别用进程内的字典
SEEN: dict[str, set[str]] = {}
# key 用 thread_id:同一个会话共享一个集合,不同会话互不影响
seen = SEEN.setdefault(thread_id, set())

去重的 key 就用 §6.3 约定的 stage + status + text 拼起来,这也解释了为什么 §6.3 要求「stage + status 能唯一确定一个事件语义」:没有这个约定,就没法可靠去重。

两个实践细节:

§10.3 的实战会给出去重生效的实际输出。

10. 实战:可调试的流式图 #

10.1 设计 #

把前面九节拼起来:一张 RAG 图,外面套一层事件翻译,输出前端能直接消费的流。

START → retrieve → guard → answer → END
         │          │        │
         │          │        └─ 模型 token(打 final-answer 标签)
         │          └─ 敏感问题 interrupt 挂人工
         └─ custom 进度 + 引用事件

三个节点各自代表一类流式难题,不是随便挑的:

节点 代表的难题 靠什么解决
retrieve 不可流式的慢步骤,前端会空转 custom 进度事件(§6)
guard 中途停下来等人,且恢复时会重跑 interrupt + 去重(§9.2、§9.3)
answer 要流字,但不能把中间模型的输出漏出去 标签过滤(§3.4)

对外的事件协议只有六种 type,前端认这六个就够:

type 何时发出 前端该做什么
stage 节点内 writer 推的进度 更新步骤条
citation 检索完成,引用就位 先把来源渲染出来
token 带 final-answer 标签的模型片段 追加到回答气泡
interrupt 图暂停等人 弹审批框
error 节点抛异常 切失败态
paused / done 两种终止 等人操作 / 收起加载态

这份协议有一条不变量,也是它能称为「契约」的关键:任何一次请求,最后恰好收到一个终止事件(done、paused、error 三者之一)。有了它,前端收起加载动画只需一行逻辑,不必靠超时兜底;§5.4 那种「零 token」路径也不会让界面卡住。

翻译层做的四件事,各自对应前面的一节:

  1. 三重过滤 token:标签(§3.4)+ 非空(§5.1)+ 带上来源节点(§4.3)
  2. 进度去重:跨调用持有 seen(§9.3)
  3. 异常兜底:try 包住整个循环,补 error 事件(§9.1)
  4. 区分终止:paused 和 done(§9.2)

几处刻意的简化,避免偏离主题:

10.2 stream_server.py #

"""第 27 章产出:可调试的流式图

把一张 RAG 图的内部动静,翻译成前端能直接消费的统一事件流。
覆盖:阶段进度(custom)、token 分流(messages)、引用、暂停(interrupt)、失败(error)。

跑法:
    python stream_server.py              # 正常问答,看完整事件流
    python stream_server.py --sse        # 输出 SSE 报文格式
    python stream_server.py --sensitive  # 触发人工确认,演示暂停与恢复
    python stream_server.py --fail       # 演示节点失败时事件流怎么收尾
"""

# 解析命令行参数,让四种演示共用一个入口
import argparse
# 把事件序列化成 SSE 报文里的 data 字段
import json
# 提供 add,作为 log 字段的 reducer
import operator
# SSE 模式直接写 stdout,避免 print 自动加换行破坏报文格式
import sys
# 模拟检索耗时,让进度事件肉眼可见
import time
# Annotated 挂 reducer,Any/Iterator 给翻译层做类型标注,TypedDict 声明状态
from typing import Annotated, Any, Iterator, TypedDict

# 从 .env 读取模型 key
from dotenv import load_dotenv
# 创建聊天模型
from langchain.chat_models import init_chat_model
# 内存版 checkpointer:interrupt 要恢复必须有它(第 26 章 §3.1)
from langgraph.checkpoint.memory import InMemorySaver
# 在节点内部取 writer,推自定义进度事件
from langgraph.config import get_stream_writer
# 建图三件套
from langgraph.graph import END, START, StateGraph
# Command 携带人的决策,interrupt 用来暂停
from langgraph.types import Command, interrupt

# 加载环境变量,override=True 让 .env 覆盖同名的旧变量
load_dotenv(override=True)

# ---------------------------------------------------------------- 知识库
# 用内存字典代替向量库,让本章聚焦在「流式」而不是检索质量上。
KB = [
    # 每条知识带 id(引用用)、title(前端展示用)、text(喂模型用)、kw(关键词匹配用)
    {"id": "faq-01", "title": "退货时限", "text": "签收后 7 天内可无理由退货,需保持商品完好。", "kw": ["退货", "7 天", "无理由"]},
    # 第二条:退款到账时间
    {"id": "faq-02", "title": "退款到账", "text": "退款审核通过后 3 个工作日内退回原支付渠道。", "kw": ["退款", "到账", "几天"]},
    # 第三条:运费规则
    {"id": "faq-03", "title": "运费规则", "text": "订单满 99 元包邮,未满收取 10 元运费。", "kw": ["运费", "包邮", "邮费"]},
    # 第四条:账号注销,命中它会触发人工确认
    {"id": "faq-04", "title": "账号注销", "text": "注销账号需联系客服人工处理,注销后数据不可恢复。", "kw": ["注销", "删除账号"]},
]

# 命中这些词就需要人工确认,用来演示流式过程中的暂停
SENSITIVE = ["注销", "删除账号"]


# 图的状态:一次问答的全部上下文
class RagState(TypedDict):
    # 用户问题
    question: str
    # 检索命中的文档(原始 dict,含 id/title/text)
    docs: list[dict]
    # 最终答案
    answer: str
    # 人工是否批准;非敏感问题自动置 True
    approved: bool
    # 执行轨迹,带 reducer 所以每个节点都能安全追加
    log: Annotated[list[str], operator.add]


# 给最终回答的模型打标签,前端只把带这个标签的 token 显示给用户
FINAL_TAG = "final-answer"
# 创建回答模型,tags 在创建时就固定下来(§3.4 的第一种打法)
answer_model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0, tags=[FINAL_TAG])


# ---------------------------------------------------------------- 节点
def retrieve(state: RagState) -> dict:
    """检索:不可流式,所以用 custom 事件汇报进度,避免前端空转。"""
    # 取 writer,后面所有进度都靠它推
    writer = get_stream_writer()
    # 一进节点就推「开始」,用户立刻有反馈,不用等检索结束(§7.2 的时机差)
    writer({"kind": "stage", "stage": "retrieve", "status": "start", "text": "正在检索知识库"})

    # 取出问题,下面反复用
    q = state["question"]
    # 特殊哨兵值:用来演示 --fail 场景下事件流怎么收尾
    if q == "__fail__":  # 仅用于演示失败路径
        # 抛异常,会被翻译层的 except 兜住并转成 error 事件(§9.1)
        raise RuntimeError("知识库连接超时")

    # 命中的文档
    hits = []
    # 逐条扫知识库
    for doc in KB:
        time.sleep(0.05)  # 模拟检索耗时,让进度事件肉眼可见
        # 关键词有一个出现在问题里就算命中
        if any(k in q for k in doc["kw"]):
            # 收下这篇
            hits.append(doc)
            # 每命中一篇推一条进度,让用户看到事情在推进
            writer({"kind": "stage", "stage": "retrieve", "status": "progress", "text": f"命中《{doc['title']}》"})

    # 检索结束,推一条带命中数的完成事件
    writer({"kind": "stage", "stage": "retrieve", "status": "done", "hits": len(hits)})
    # 引用单独发一个事件,前端可以先把来源渲染出来,再等正文
    writer({"kind": "citation", "docs": [{"id": d["id"], "title": d["title"]} for d in hits]})
    # 把命中的文档写进状态,供 answer 节点组装 prompt
    return {"docs": hits, "log": [f"retrieve:{len(hits)}"]}


def guard(state: RagState) -> dict:
    """敏感问题挂人工确认。演示 interrupt 在流式里的表现。"""
    # 取 writer
    writer = get_stream_writer()
    # 不敏感就直接放行,连进度事件都不推(少一个噪音)
    if not any(k in state["question"] for k in SENSITIVE):
        # approved=True 让 answer 节点走正常分支
        return {"approved": True, "log": ["guard:auto"]}

    # 敏感问题:先告诉前端「要等人了」
    writer({"kind": "stage", "stage": "guard", "status": "waiting", "text": "该问题需人工确认"})
    # 暂停,把结构化信息交给人;恢复后 decision 就是 Command(resume=...) 里的值
    # 注意这一行之前的 writer 调用在恢复时会重跑一遍(§9.3),靠翻译层的 seen 去重
    decision = interrupt({"question": state["question"], "reason": "涉及账号注销"})
    # 只有明确的 approve 才算批准,其他一律当拒绝
    ok = decision == "approve"
    # 把决策结果和原始决策都记下来,便于审计
    return {"approved": ok, "log": [f"guard:{decision}"]}


def answer(state: RagState) -> dict:
    """生成:token 会通过 messages 模式流出去。"""
    # 取 writer
    writer = get_stream_writer()
    # 推「开始生成」,前端把加载状态从「检索中」切到「生成中」
    writer({"kind": "stage", "stage": "answer", "status": "start"})

    # 没批准就走固定话术,全程不调模型,这正是 §5.4 说的「零 token」路径
    if not state["approved"]:
        # 直接返回,注意这条路径下前端一个 token 都收不到,只能靠终止事件收尾
        return {"answer": "该操作需人工处理,已转交客服。", "log": ["answer:blocked"]}

    # 把命中的文档拼成上下文;一篇都没命中时给个占位符,避免 prompt 出现空段
    ctx = "\n".join(f"[{d['id']}] {d['text']}" for d in state["docs"]) or "(无资料)"
    # 组 prompt:强调只依据资料回答,并限制长度,让演示输出短一些
    prompt = (
        "只依据资料回答,资料里没有就说不知道。用完整的一两句话回答,不超过 50 字。\n"
        f"资料:\n{ctx}\n\n问题:{state['question']}"
    )
    # 用带 final-answer 标签的模型;节点里仍然是 invoke,token 由外层的 stream 负责流出
    reply = answer_model.invoke(prompt)
    # 推「生成完成」,前端收起打字光标
    writer({"kind": "stage", "stage": "answer", "status": "done"})
    # 把答案写进状态
    return {"answer": reply.content, "log": ["answer:ok"]}


# 组装并编译这张图;抽成函数是为了让每种演示都能拿到一张干净的图
def build_graph():
    # 用 RagState 建图
    builder = StateGraph(RagState)
    # 检索节点
    builder.add_node("retrieve", retrieve)
    # 人工确认节点
    builder.add_node("guard", guard)
    # 生成节点
    builder.add_node("answer", answer)
    # 入口进检索
    builder.add_edge(START, "retrieve")
    # 检索完进 guard
    builder.add_edge("retrieve", "guard")
    # guard 通过后进生成(拒绝的情况由 answer 内部分支处理,图结构保持线性)
    builder.add_edge("guard", "answer")
    # 生成完结束
    builder.add_edge("answer", END)
    # 必须带 checkpointer,否则 guard 里的 interrupt 没法恢复
    return builder.compile(checkpointer=InMemorySaver())


# ---------------------------------------------------------------- 事件翻译层
def to_events(graph, payload: Any, cfg: dict, seen: set[str] | None = None) -> Iterator[dict]:
    """把图的多路流,翻译成一条扁平、自洽的事件流。

    前端只需要认这几个 type:stage / citation / token / interrupt / error / paused / done。

    seen 由调用方按会话持有。恢复执行时被中断的节点会从头重跑(第 26 章 §4),
    进度事件会重发一遍;只有让去重集合活得比单次 stream 长,才能真正去掉重复。
    """
    # 允许不传 seen(比如一次性脚本),但这样就没有跨调用去重能力
    if seen is None:
        # 退化成单次去重:只能挡住同一次 stream 内部的重复
        seen = set()
    # 标记这次执行是不是以暂停结束,决定最后发 paused 还是 done
    paused = False

    # try 必须包住整个 for 循环,否则拦不住节点抛出的异常(§9.1)
    try:
        # 一次订阅三种模式:进度、token、节点产出(含 __interrupt__)
        for mode, chunk in graph.stream(payload, cfg, stream_mode=["custom", "messages", "updates"]):
            # 分支一:自定义进度事件
            if mode == "custom":
                # 引用事件结构和阶段事件不同,先分出去;它不参与去重(同一会话只会发一次)
                if chunk.get("kind") == "citation":
                    yield {"type": "citation", "docs": chunk["docs"]}
                    continue
                # 同一个阶段状态只发一次,避免恢复时重复
                key = f"{chunk.get('stage')}:{chunk.get('status')}:{chunk.get('text', '')}"
                # 见过就丢掉,这一行就是 §9.3 那个坑的修法
                if key in seen:
                    continue
                # 没见过就记下来
                seen.add(key)
                # 翻译成前端的 stage 事件,字段名和内部协议解耦
                yield {
                    "type": "stage",
                    "stage": chunk.get("stage"),
                    "status": chunk.get("status"),
                    "text": chunk.get("text", ""),
                }

            # 分支二:模型 token
            elif mode == "messages":
                # messages 的内容本身是二元组,要再解一层
                token, meta = chunk
                # 三重过滤:只要带标签的最终回答、只要有内容的 chunk、按节点确认来源
                if FINAL_TAG not in (meta.get("tags") or []):
                    continue
                # 空片段占七成,必须滤掉(§5.1)
                if not token.content:
                    continue
                # 事件里带上来源节点,前端才能分区渲染(§4.3 的规则)
                yield {"type": "token", "text": token.content, "node": meta.get("langgraph_node")}

            # 分支三:节点产出,这里只关心暂停信号
            elif mode == "updates":
                # __interrupt__ 占的是「节点名」的位置,用 in 判断(§9.2)
                if "__interrupt__" in chunk:
                    # 并行中断时可能有多个,这个流程只有一个,取第一个
                    itr = chunk["__interrupt__"][0]
                    # 记下来,等循环结束后发 paused 而不是 done
                    paused = True
                    # id 供前端 resume 时回传,value 是给人看的内容
                    yield {"type": "interrupt", "id": itr.id, "value": itr.value}

    # 任何异常都转成 error 事件,不让它继续往上抛断掉 HTTP 连接
    except Exception as exc:
        # 已产出的事件前端都收到了,这里补一个终止事件,不让界面卡在「加载中」
        yield {"type": "error", "message": f"{type(exc).__name__}: {exc}"}
        # 发完就返回,保证终止事件恰好一个(§9.1 强调的那个 return)
        return

    # 暂停和跑完是两种不同的结束,前端的处理也不同:一个等人操作,一个收尾
    yield {"type": "paused"} if paused else {"type": "done"}


def to_sse(ev: dict) -> str:
    """按 SSE 报文格式包一层,浏览器 EventSource 可直接消费。"""
    # event: 行让前端按类型分派,data: 行是 JSON 负载,末尾空行表示一条报文结束
    # ensure_ascii=False 保留中文原文,否则会变成 \uXXXX 转义
    return f"event: {ev['type']}\ndata: {json.dumps(ev, ensure_ascii=False)}\n\n"


# ---------------------------------------------------------------- 演示
# 会话级的去重集合。真实服务里放 Redis 或会话对象,key 就是 thread_id。
SEEN: dict[str, set[str]] = {}


# 消费事件流并渲染;sse=True 输出报文,否则输出人类可读格式
def render(graph, payload, cfg, sse: bool) -> list[dict]:
    # 从配置里取出 thread_id,作为去重集合的 key
    tid = cfg["configurable"]["thread_id"]
    # 同一个会话复用同一个集合,这样第二次 stream 才能查到第一次的记录
    seen = SEEN.setdefault(tid, set())
    # 收集全部事件,便于调用方做断言
    collected = []
    # 消费翻译层产出的事件流
    for ev in to_events(graph, payload, cfg, seen):
        # 先存一份
        collected.append(ev)
        # SSE 模式:直接写报文,不做人类可读的美化
        if sse:
            # 用 sys.stdout.write 而不是 print,避免多出一个换行破坏报文格式
            sys.stdout.write(to_sse(ev))
            continue
        # 以下是给人看的渲染,模拟前端对六种 type 的处理
        t = ev["type"]
        # 阶段进度:更新步骤条
        if t == "stage":
            # rstrip 去掉 text 为空时多出的尾部空格
            print(f"  [stage] {ev['stage']}/{ev['status']} {ev['text']}".rstrip())
        # 引用:先把来源列出来
        elif t == "citation":
            # 一篇都没命中时显示「无」
            titles = "、".join(d["title"] for d in ev["docs"]) or "无"
            # 只打印标题,真实前端这里是可点击的来源链接
            print(f"  [cite ] {titles}")
        # token:追加到回答气泡
        elif t == "token":
            # 用 !r 带引号打印,能看清模型的切分粒度和空格
            print(f"  [token] {ev['text']!r}")
        # 暂停:弹审批框
        elif t == "interrupt":
            # 打印 interrupt 的 value,也就是要给审批人看的内容
            print(f"  [pause] {ev['value']}")
        # 失败:切失败态
        elif t == "error":
            # 打印异常类型和消息,真实前端会渲染成一条错误提示
            print(f"  [error] {ev['message']}")
        # 剩下的是 paused / done 两个终止事件,只打个标记
        else:
            # 固定宽度对齐,让 done 和 paused 在输出里一眼可辨
            print(f"  [{t:5s}]")
    # 返回收集到的事件,方便在测试里断言
    return collected


# 造一份初始状态;每个新会话都要一份全新的
def new_payload(q: str) -> dict:
    # 每次新会话的初始状态;TypedDict 的字段要给全
    return {"question": q, "docs": [], "answer": "", "approved": False, "log": []}


# 命令行入口:按参数分派到三种演示场景
def main():
    # 用 argparse 做子场景开关
    ap = argparse.ArgumentParser()
    # --sse:输出 SSE 报文而不是人类可读格式
    ap.add_argument("--sse", action="store_true", help="输出 SSE 报文格式")
    # --sensitive:走人工确认路径,演示 paused 和恢复去重
    ap.add_argument("--sensitive", action="store_true", help="演示暂停与恢复")
    # --fail:让检索节点抛异常,演示 error 收尾
    ap.add_argument("--fail", action="store_true", help="演示节点失败")
    # 解析
    args = ap.parse_args()

    # 每次运行都新建图(InMemorySaver 的数据不跨进程)
    graph = build_graph()

    # 场景一:节点失败
    if args.fail:
        # 打印场景标题
        print("=== 节点失败 ===")
        # 用哨兵问题触发异常,thread_id 随便给一个
        render(graph, new_payload("__fail__"), {"configurable": {"thread_id": "f1"}}, args.sse)
        return

    # 场景二:暂停与恢复
    if args.sensitive:
        # 两段必须用同一个 thread_id,否则恢复不到同一个会话
        cfg = {"configurable": {"thread_id": "s1"}}
        # 打印第一段的标题
        print("=== 第一段:跑到人工确认处暂停 ===")
        # 第一次调用:正常传初始状态,会停在 guard
        render(graph, new_payload("我要注销账号"), cfg, args.sse)
        # 打印第二段的标题,前面留一个空行分隔
        print("\n=== 第二段:批准后恢复(注意进度事件没有重复)===")
        # 第二次调用:传 Command 而不是状态;SEEN 里已有第一段的记录,重发的进度会被挡掉
        render(graph, Command(resume="approve"), cfg, args.sse)
        return

    # 场景三(默认):两个正常问题,各自独立会话
    for i, q in enumerate(["多少天内可以退货?", "满多少包邮?"]):
        # 打印当前问题作为标题
        print(f"=== 问题:{q} ===")
        # thread_id 用序号区分,避免两个问题共享去重集合
        render(graph, new_payload(q), {"configurable": {"thread_id": f"t{i}"}}, args.sse)
        # 空行分隔两个问题的输出
        print()


# 作为脚本运行时才执行 main,被 import 时不执行
if __name__ == "__main__":
    main()

10.3 跑起来 #

正常问答:

python stream_server.py
=== 问题:多少天内可以退货? ===
  [stage] retrieve/start 正在检索知识库
  [stage] retrieve/progress 命中《退货时限》
  [stage] retrieve/done
  [cite ] 退货时限
  [stage] answer/start
  [token] '7'
  [token] '天内'
  [token] '可'
  [token] '无'
  [token] '理由'
  [token] '退货'
  [token] ','
  [token] '需'
  [token] '保持'
  [token] '商品'
  [token] '完好'
  [token] '。'
  [stage] answer/done
  [done ]

=== 问题:满多少包邮? ===
  [stage] retrieve/start 正在检索知识库
  [stage] retrieve/progress 命中《运费规则》
  [stage] retrieve/done
  [cite ] 运费规则
  [stage] answer/start
  [token] '满'
  [token] '99'
  [token] '元'
  [token] '包'
  [token] '邮'
  [token] '。'
  [stage] answer/done
  [done ]

三个细节值得对照着看:

暂停与恢复:

python stream_server.py --sensitive
=== 第一段:跑到人工确认处暂停 ===
  [stage] retrieve/start 正在检索知识库
  [stage] retrieve/progress 命中《账号注销》
  [stage] retrieve/done
  [cite ] 账号注销
  [stage] guard/waiting 该问题需人工确认
  [pause] {'question': '我要注销账号', 'reason': '涉及账号注销'}
  [paused]

=== 第二段:批准后恢复(注意进度事件没有重复)===
  [stage] answer/start
  [token] '请联系'
  [token] '客服'
  [token] '人工'
  [token] '处理'
  [token] '注销'
  [token] ','
  [token] '注销'
  [token] '后'
  [token] '数据'
  [token] '不可'
  [token] '恢复'
  [token] '。'
  [stage] answer/done
  [done ]

这段输出里有三处直接印证了前面各节:

节点失败:

python stream_server.py --fail
=== 节点失败 ===
  [stage] retrieve/start 正在检索知识库
  [error] RuntimeError: 知识库连接超时

retrieve/start 已经发给前端了(§9.1 验证过:异常前的 chunk 都会送达),然后是明确的 error。前端把步骤条标红、收起加载态,不会一直转圈。

注意没有 done:error 就是这次请求的终止事件。这就是 §9.1 那个 return 的作用;少了它这里会多出一个 done,前端先标红又收起,用户只看到错误一闪而过。

SSE 报文格式:

python stream_server.py --sse
event: stage
data: {"type": "stage", "stage": "retrieve", "status": "start", "text": "正在检索知识库"}

event: stage
data: {"type": "stage", "stage": "retrieve", "status": "progress", "text": "命中《退货时限》"}

event: stage
data: {"type": "stage", "stage": "retrieve", "status": "done", "text": ""}

event: citation
data: {"type": "citation", "docs": [{"id": "faq-01", "title": "退货时限"}]}

event: stage
data: {"type": "stage", "stage": "answer", "status": "start", "text": ""}

event: token
data: {"type": "token", "text": "签", "node": "answer"}

event: done
data: {"type": "done"}

(上面截取了几条,实际输出包含全部事件。)

SSE 的报文格式只有三条规则,to_sse 那三行就是全部:

  1. event: 类型:前端 EventSource 靠它分派到不同的处理函数
  2. data: JSON:负载,必须是单行(所以不能用带缩进的 json.dumps)
  3. 末尾一个空行:表示一条报文结束,少了它浏览器会一直等下一行

还有两个实践要点:

10.4 验收清单 #

检查项 怎么验
中间步骤的 token 不外泄 图里再加一个不打 final-answer 标签的模型节点,确认它的 token 不出现在事件流里
并行安全 token 事件带 node 字段;加一个并行分支后前端仍能正确分区渲染
检索期间有反馈 retrieve/start 在任何 token 之前到达
引用早于正文 citation 事件在第一个 token 之前
失败不卡界面 --fail 能看到 error 事件而不是异常栈,且没有 done
终止事件恰好一个 三种跑法各数一遍 done / paused / error 的总数,必须都是 1
暂停与跑完可区分 --sensitive 第一段以 paused 结束,正常问答以 done 结束
恢复不重复进度 --sensitive 第二段没有 guard/waiting;把 seen 改成局部变量后它应该出现
零 token 路径也能收尾 把 --sensitive 的 resume 值改成 "reject",走固定话术不调模型,确认仍有 done
子图安全 把 retrieve 改成子图,确认进度事件消失;加上 subgraphs=True 后恢复(§7.4)
事件类型封闭 前端只需处理六种 type,加新节点不用改前端

11. 实用约定与坑 #

约定:

约定 为什么
给「给用户看」的那次模型调用固定打一个标签 前端只认标签,图里加多少中间模型都不用改前端(§3.4)
token 事件一律带来源 node 并行分支一出现就必须有;现在没有也要带(§4.3)
custom 推 dict 而不是字符串 字符串只能原样显示,dict 才能驱动 UI 状态(§6.1)
stage + status 要能唯一确定一个事件语义 没有这个约定就没法可靠去重(§6.3、§9.3)
图里有子图就一定要写 subgraphs=True 不写会静默丢掉子图的 token 和进度(§7.4)
新写的翻译层用 version="v2" 形状固定,加模式、加子图都不用改解包逻辑(§7.5)
事件协议里必须有明确的终止事件 token 流可能是空的,前端不能靠「没消息了」判断结束(§5.4)
任何一次请求,终止事件恰好一个 少了前端一直转圈,多了状态机会乱(§9.1)
暂停和跑完用不同的终止事件 一个等人操作,一个收起加载态(§9.2)
默认用 stream_mode,不用 astream_events 形状简单、有同步版;后者噪音比信号多(§8.5)
用 isinstance(token, AIMessage) 而不是 AIMessageChunk 非流式模型的整条回复也能接住(§5.3)
把 ls_model_name 记进日志 「换了个模型导致流式失效」唯一的线索(§3.2、§5.3)

响亮的错误(会抛异常,看到就能改):

报错 原因 怎么办
RuntimeError: Called get_config outside of a runnable context 在图外调用 get_stream_writer() 只在节点或节点调用的工具里用(§6.4)
ValueError: too many values to unpack (expected 2) 加了 subgraphs=True,元组从 2 个变 3 个 按四种组合对照解包,或改用 version="v2"(§7.3、§7.5)
NotImplementedError: stream_events(version='v2') is not supported 想用同步版的事件流 v1/v2 只有异步版;同步版只有实验性的 v3(§8 开头)

安静的 bug(不报错,只是结果不对,这是本章的重点):

现象 原因 怎么办
界面上两段答案混在一起,且偶发 并行节点的 token 交错,交错程度不确定(实测 8 次里 3 次严重交错) 按 langgraph_node 分流(§4.3)
换了个模型后流式效果消失 模型不支持流式,退化成一个完整 AIMessage 用 isinstance(..., AIMessage) 判断(§5.3)
回答整条消失,日志干净 用 isinstance(token, AIMessageChunk) 筛选,非流式模型的整条回复被丢掉 改用父类 AIMessage(§5.3)
工具返回值印在了用户的回答里 ToolMessage 也走 messages 流 按 type(token).__name__ 分派(§5.2)
大量空内容的 chunk,界面抖动 角色片段、推理内容、工具调用阶段、收尾片段 if not token.content: continue(§5.1)
Agent 调工具那几秒界面完全不动 那一轮的片段 content 全是空的,被过滤干净了 工具内部用 writer 推进度(§6.2),或用 tool_call_chunks(§8.4)
界面一直转圈但图已经跑完 该路径没有模型调用,token 流是空的 加明确的终止事件(§5.4)
把某个节点重构成子图后,进度条空了 / 正文没了 不加 subgraphs=True 时子图事件不外泄 加 subgraphs=True(§7.4)
加了第二个 stream_mode 后条件判断全不匹配 (ns, chunk) 和 (mode, chunk) 位置错位,str 和 tuple 都支持 len/in,不报错 用 ns == () 而非 len(ns) == 0;或改用 version="v2"(§7.3、§7.5)
恢复后进度条多出已完成的步骤 被中断的节点从头重跑,writer 又推了一遍 去重集合按会话持有,别放在函数局部(§9.3)
去重在单机上生效,多进程部署后失效 seen 放在进程内字典里 放 Redis,key 用 thread_id(§9.3)
正文缺字 对 token 也做了去重,重复的字被吃掉 只对 custom 事件去重(§9.3)
try 包了 graph.stream(...) 但异常没拦住 stream() 只是造迭代器,图还没跑 try 要包住整个 for 循环(§9.1)
前端渲染了半句话然后连接断了 异常直接往上抛,没有兜底事件 except 里补 error 事件并 return(§9.1)
错误一闪而过又变成成功态 异常路径同时发了 error 和 done except 里补完 error 记得 return(§9.1)
本地流式正常,上线后变成一次性出现 Nginx 缓冲了响应 加 X-Accel-Buffering: no(§10.3)

这张表里没有一条会在测试里自己暴露出来,这就是为什么流式功能值得专门写一层翻译,并且照着 §10.4 的清单验收一遍。

12. 练习 #

  1. 测出你自己环境的漏检概率。 跑 §4.2 的脚本至少十次,记录每次的切换次数,算出「切换 1 次」(看起来正常)的比例。书里三轮 24 次是 17 次正常,也就是近七成的情况下这个 bug 不会暴露。这个数字就是你在本地测试时的漏检概率,把它记住,下次有人说「我测了几次没问题」时拿出来。
  2. 给本章实战加一个不该外泄的节点。 在 retrieve 前面加一个 rewrite 节点,用另一个模型改写用户问题(不打 final-answer 标签)。确认它的 token 一个都没进事件流;然后故意把标签加上,看会漏出什么。这一步验证的是 §3.4 那个「产品契约」是否真的成立。
  3. 亲手制造子图事故。 把本章实战的 retrieve 节点改成一个子图(内部两个节点:embed 和 search,各推一条进度),不改翻译层。观察进度条变成什么样。然后把生成节点也挪进子图,看正文是不是整段消失。最后加 subgraphs=True + 处理多出来的 ns 元素修好它(§7.4)。
  4. 把翻译层改成 version="v2"。 对照 §7.5,把 to_events 里的 for mode, chunk in ... 改成 for part in ...,四种模式组合都试一遍,确认解包逻辑再也不用跟着变。顺便把 __interrupt__ 的判断换成 part["interrupts"]。
  5. 把 messages 换成 astream_events。 用 include_tags=["final-answer"] 拿最终回答的 token,再用 on_tool_start 加一个「正在调用工具」的事件。给图加一个工具调用节点来验证。对比两种实现的代码量,以及需要处理的事件数量(提示:会多一个数量级)。
  6. 接进 FastAPI。 把 to_events 包成一个 StreamingResponse(media_type="text/event-stream") 的接口,写个最小 HTML 页面用 EventSource 消费,六种事件各自渲染。别忘了 §10.3 提到的 Nginx 缓冲问题。这一步做完,前面十节的东西才算真正落地。
  7. 测试非流式退化。 把 answer_model 加上 disable_streaming=True,观察事件流的变化(提示:token 事件会一条不剩,因为标签过滤仍然通过但类型变了)。然后把翻译层的类型判断改成 isinstance(token, AIMessageChunk),看回答是怎么整条消失的,注意整个过程没有任何异常。
  8. 验证零 token 路径的收尾。 把 --sensitive 的 Command(resume="approve") 改成 "reject",这条路径走固定话术、完全不调模型。确认事件流里没有任何 token,但仍然恰好有一个 done。这就是 §5.4 那条结论的实际价值。

13. 本章小结 #

下一章是阶段四的收尾:把第 19~27 章的东西合起来,做一条完整的审批流:提交、风控检查、Agent 起草、人工审批、发送。