1. 本章目标 #
前两章把图「搭」出来了:第 20 章讲节点和边,第 21 章讲状态怎么合并。但一直有个动作被一笔带过:.compile(),以及编译之后到底能怎么用它。
这一章专门讲运行这一侧:
compile()做了什么,invoke/stream/get_state各自解决什么问题,出错了怎么查、怎么续跑。
学完这一章,你应该能做到:
- 说清
StateGraph和CompiledStateGraph的分工,以及compile()会检查什么、不会检查什么 - 在
invoke(要结果)、stream(要过程)、get_state(要现场)之间正确选择 - 掌握
stream的七种stream_mode,知道哪种用来调试、哪种用来做 UI,以及为什么打字机效果开头会蹦出一堆空片段 - 用
config控制thread_id、recursion_limit、tags和metadata,并知道recursion_limit数的到底是什么 - 用
get_state_history看完整执行轨迹,用update_state打补丁,回到任意检查点重跑甚至分叉 - 知道节点抛异常时状态停在哪,以及怎么修完数据从断点继续
- 产出:一张可运行、可观测、可恢复的状态图
参考文档:
1.1. 为什么「编译」值得单独一章 #
有一个疑问是:图都搭好了,compile() 不就是个形式吗?直接跑不行吗?
不行,而且原因不止一个:
| 编译带来的东西 | 如果没有编译这一步 |
|---|---|
| 一次结构校验(悬空边、没入口) | 这些错要等跑到那一步才炸,而且报错位置很难看懂 |
| 把 builder 冻结成快照 | 边跑边改图,行为无法复现 |
| 挂上 checkpointer / 中断点等运行期设施 | 没有地方安放这些「不属于图结构」的配置 |
| 产出一个 Runnable | 图没法嵌进 LCEL 链,也没法当别的图的节点 |
其中最容易被低估的是第三条:同一张图纸,你会想编译出好几个版本:测试版不要持久化,生产版要,调试版还要在某个节点前停下来。这些差异全都由 compile() 的参数表达,而不是改图。§2.4 和 §8 会把这条用起来。
1.2. 贯穿全章的实验用图 #
本章除 §4.4 外都用同一张两节点小图做实验。它只有两个节点:一个把字符串两端空白去掉,一个数字符数。
# operator 提供 operator.add,第 21 章用它做 list 累加型 reducer
import operator
# Annotated 用来给字段挂 reducer;TypedDict 是最常用的状态类型
from typing import Annotated, TypedDict
# END / START 是两个特殊节点常量,StateGraph 是图纸类
from langgraph.graph import END, START, StateGraph
# 状态定义:两个字段,一个默认覆盖,一个累加
class S(TypedDict):
# text 没挂 reducer,后写的会覆盖先写的
text: str
# log 挂了 operator.add,每个节点的写入会追加到列表尾部
log: Annotated[list[str], operator.add]
# 第一个节点:把 text 两端的空白去掉,并记一笔日志
def clean(state: S) -> dict:
# 返回的 dict 只需包含要改的字段:text 被覆盖,log 被追加
return {"text": state["text"].strip(), "log": ["clean"]}
# 第二个节点:数一下 text 有多少个字符,只写日志
def count(state: S) -> dict:
# 注意这里读到的是 clean 处理过的 text(同一个超步内的顺序执行,见第 21 章)
return {"log": [f"count={len(state['text'])}"]}
# 注意这个函数返回的是 builder(没编译),方便每个实验单独 compile
def build():
"""返回 builder(没编译),方便每个实验单独 compile。"""
# 用状态类型初始化图纸
g = StateGraph(S)
# 注册第一个节点,节点名是字符串 "clean"
g.add_node("clean", clean)
# 注册第二个节点,节点名是字符串 "count"
g.add_node("count", count)
# 入口:START -> clean
g.add_edge(START, "clean")
# 串行:clean -> count
g.add_edge("clean", "count")
# 出口:count -> END
g.add_edge("count", END)
# 返回图纸本身,编译动作交给调用方
return g
# 默认版本:不带 checkpointer,§2~§5 大部分实验用它
graph = build().compile()为什么
build()返回的是没编译的图纸? 因为本章要反复演示「同一张图纸编译成不同版本」:有的带 checkpointer,有的带中断点,有的什么都不带。把编译动作留给调用方,是 LangGraph 项目里很常见的写法,§8 的实战会把它固化成build_graph()。
2. compile():从图纸到可执行对象 #
这一节回答三个问题:编译出来是什么、会替你挡掉哪些错、编译时能拧哪些旋钮。
2.1. 两个类型,两种身份 #
先用最直接的方式看清「编译前」和「编译后」是两个不同的类:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# build() 返回的是图纸
builder = build()
# 打印类名,确认它是 StateGraph
print("builder:", type(builder).__name__)
# 编译,得到可执行对象
graph = builder.compile()
# 打印类名,确认它变成了 CompiledStateGraph
print("compiled:", type(graph).__name__)
# 检查三个 Runnable 标准方法是否都在
print("是 Runnable 吗:", hasattr(graph, "invoke"), hasattr(graph, "stream"), hasattr(graph, "batch"))运行输出:
builder: StateGraph
compiled: CompiledStateGraph
是 Runnable 吗: True True Truehasattr 只能说明「有这几个方法」,不足以说明「它就是 Runnable」。用 isinstance 和继承链把话说实:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# Runnable 是 LangChain 的统一执行接口基类(第 7 章)
from langchain_core.runnables import Runnable
# 图纸和编译产物各留一份,下面要分别判断
builder = build()
graph = builder.compile()
# 编译后的图是不是 Runnable
print("compiled 是 Runnable:", isinstance(graph, Runnable))
# 图纸是不是 Runnable
print("builder 是 Runnable:", isinstance(builder, Runnable))
# 打印完整继承链,看清它是怎么变成 Runnable 的
print("继承链:", [c.__name__ for c in type(graph).__mro__])compiled 是 Runnable: True
builder 是 Runnable: False
继承链: ['CompiledStateGraph', 'Pregel', 'PregelProtocol', 'Runnable', 'ABC', 'Generic', 'object']这条继承链在第 19 章分析 create_agent 的返回值时见过一次,现在它对上了:create_agent 返回的和你自己 compile() 出来的,是同一个类型的东西。第 19 章那句「Agent 只是一张预置的图」,在类型层面就是这条链。
分工总结:
StateGraph |
CompiledStateGraph |
|
|---|---|---|
| 身份 | 图纸(builder) | 可执行对象 |
| 能干什么 | add_node / add_edge / add_conditional_edges |
invoke / stream / batch / get_state |
| 能跑吗 | 不能 | 能 |
| 是 Runnable 吗 | 不是 | 是(所以能进 LCEL 链,见第 7 章) |
| 能改结构吗 | 能 | 不能(改了也不生效,见 §2.4) |
CompiledStateGraph 是 Runnable 这件事有三个直接后果:
- 能进 LCEL 链:
prompt | graph | parser是合法的。 - 能当别的图的节点:
add_node("sub", 子图),因为节点接受任何 Runnable(第 25 章会用)。 - 免费获得 Runnable 全家桶:
batch/ainvoke/astream/with_retry/with_config都不是 LangGraph 单独实现的,而是继承来的(§3.3、§3.4 就是这么来的)。
2.2. compile() 检查什么 #
编译时会做一遍结构检查。注意它检查的是「图的拓扑」,不是「业务逻辑」。一共两类会当场拦下来。
2.2.1 第一种:边指向了不存在的节点。 #
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# 新建一张图纸
bad = StateGraph(S)
# 注册一个什么都不做的节点(lambda 返回空 dict 表示不改状态)
bad.add_node("a", lambda s: {})
# 入口边,正常
bad.add_edge(START, "a")
# 出边指向一个没注册过的名字:这是打错字最常见的形态
bad.add_edge("a", "不存在的节点")
try:
# 编译时才会发现这个问题(add_edge 时不会报)
bad.compile()
except Exception as e:
# 打印异常类型和完整消息
print("报错:", type(e).__name__, str(e))报错: ValueError Found edge ending at unknown node `不存在的节点`注意报错时机:add_edge 那一行是不报错的,因为那时候你完全有可能还没注册目标节点(先连边后加节点是允许的)。校验只能推迟到编译。
2.2.2 第二种:压根没有入口。 #
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# 新建一张图纸(这两段的图纸都故意搭坏,所以不叫 builder)
bad = StateGraph(S)
# 注册一个节点
bad.add_node("a", lambda s: {})
# 只连了出口边,忘了从 START 连进来
bad.add_edge("a", END) # 只有出口,没从 START 连进来
try:
# 编译时检查入口
bad.compile()
except Exception as e:
# 打印异常类型和完整消息
print("报错:", type(e).__name__, str(e))报错: ValueError Graph must have an entrypoint: add at least one edge from START to another node这两条报错信息都写得很清楚,属于「大声失败」,不需要费劲排查。真正麻烦的是下一节那些静默问题。
2.3. compile() 不检查什么 #
这一点第 20 章 §7.1 和 §10 已经强调过,这里再列一遍,因为它决定了你必须自己做哪些验证:
compile() 会拦 |
compile() 不会拦 |
|---|---|
| 边指向不存在的节点 | 节点返回了状态里没有的键(静默丢弃) |
没有从 START 出发的边 |
某个节点没有出边(图提前结束) |
| — | 有节点连不到(永远不执行) |
| — | 输入缺字段(跑到节点里才 KeyError) |
| — | 节点里的业务逻辑对不对 |
两类失败的区别值得再点一次:「边指向了不存在的节点」会当场报错,「忘了连某条边」却一声不响。前者你打错了一个名字,后者你少写了一行。少写的那行没有任何痕迹可查,编译器无从知晓你的意图。
一句话:
compile()只保证「这张图能跑」,不保证「这张图跑的是你想要的」。 后者要靠stream(§4)和draw_mermaid()自己看。
2.4. builder 和编译产物是两份东西 #
这是很多人踩过的坑:改完图发现行为没变。
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# 拿一张干净的图纸
builder = build()
# 先编译一次
graph = builder.compile()
# 编译之后再往图纸上加节点
builder.add_node("extra", lambda s: {}) # 编译之后再改 builder
# 检查已编译的图里有没有这个新节点
print("旧的 compiled 图受影响吗:", "extra" in graph.get_graph().nodes)运行输出(还会附带一条警告):
Adding a node to a graph that has already been compiled. This will not be reflected in the compiled graph.
旧的 compiled 图受影响吗: False编译是一次快照。 改 builder 不会影响已经编译出来的图,想生效必须重新 compile()。
那条提示只是警告,不是异常:程序照常往下跑,你要是没看控制台就完全无感。这也是「改了图但行为没变」这个问题难查的原因:线索是有的,只是很容易被淹没在其它输出里。
重新编译就能拿到新版本:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
builder = build()
# 用改过的图纸编译一次
graph_a = builder.compile()
builder.add_node("extra", lambda s: {})
# 再编译一次
graph_b = builder.compile()
# 两次编译产出的是不同对象,互不干扰
print("两次 compile 得到不同对象:", graph_a is not graph_b)
# 新编译出来的图包含刚加的节点
print("新编译的含 extra 吗:", "extra" in graph_b.get_graph().nodes)两次 compile 得到不同对象: True
新编译的含 extra 吗: True这个特性很实用:同一张图纸,编译出「带 checkpointer 的生产版」和「不带的测试版」,§8 的实战就是这么做的。反过来,如果你在 notebook 里反复执行 add_node,同一个 builder 会被越加越多。图纸是有状态的可变对象,而编译产物是不可变快照,这个区别记牢了能省下很多困惑。
2.5. compile() 的参数 #
想知道当前版本到底能传什么,最可靠的办法是直接问 Python:
from langgraph.graph import StateGraph
# inspect 可以在运行时读出函数签名,不用翻文档
import inspect
# 只打印参数名列表,比打印完整签名(带一堆类型标注)清爽得多
print(list(inspect.signature(StateGraph.compile).parameters))['self', 'checkpointer', 'cache', 'store', 'interrupt_before', 'interrupt_after', 'debug', 'name', 'transformers']想看完整签名(含类型标注和默认值)就去掉
list(...)直接print(inspect.signature(...)),输出会长得多,但能看到checkpointer之后是*,说明除checkpointer外全都是关键字参数,必须写参数名。
现阶段需要关心的是这几个:
| 参数 | 作用 | 本文哪里讲 |
|---|---|---|
checkpointer |
持久化状态,get_state 的前提 |
§6、第 11 章 |
name |
给图起名,出现在 trace 和嵌套图里 | §2.6 |
interrupt_before / interrupt_after |
在指定节点前 / 后暂停 | §7.4 预览,第 26 章详讲 |
store |
长期记忆(跨 thread) | 后续章节 |
cache |
节点级缓存 | 后续章节 |
debug |
把执行过程打到控制台 | 等价于 stream(print_mode=...) |
这张表里最该记住的一件事:这些参数没有一个是关于「图长什么样」的。节点、边、状态类型全在图纸里定;compile() 的参数清一色是「运行期设施」。这条界线划清了,你就知道该改哪儿:改流程去动 builder,改运行方式来动 compile 参数。
2.6. 顺手起个名字 #
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# 编译时传 name,给这张图起个人类可读的名字
graph = build().compile(name="我的图")
# name 会挂在编译产物上
print("name:", graph.name)
# 顺手看一眼图里有哪些节点(含两个特殊节点)
print("nodes:", list(graph.get_graph().nodes))
# 对比:不传 name 时的默认值
print("不传 name 时:", repr(build().compile().name))name: 我的图
nodes: ['__start__', 'clean', 'count', '__end__']
不传 name 时: 'LangGraph'图多了以后一定要起名,否则 LangSmith trace 里全是 LangGraph。注意这不是比喻,默认值就是这个字符串,十张图在 trace 列表里会显示成十个一模一样的 LangGraph。起名的成本是一个参数,收益是排障时能一眼定位。
3. invoke:只要最终结果 #
invoke 是最简单的入口:给一份输入,跑完,拿结果。简单到容易忽略两个细节:它返回的到底是什么,以及输入允许省略什么。
3.1. 返回完整的最终状态 #
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# 传入初始状态(注意 text 两端有空格,用来观察 clean 有没有生效)
out = graph.invoke({"text": " 你好世界 "})
# 返回值是普通 dict,不是 StateSnapshot 之类的包装对象
print("类型:", type(out).__name__)
# 打印全部内容
print("内容:", out)类型: dict
内容: {'text': '你好世界', 'log': ['clean', 'count=4']}注意返回的是整个状态字典,不是某个节点的返回值。想要哪个字段自己取:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
out = graph.invoke({"text": " 你好世界 "})
# log 的最后一条来自 count 节点
print(out["log"][-1]) # count=4这个设计有个很实际的好处:调用方不需要知道图里最后一个节点是谁。图内部怎么重构、加几个节点、换执行顺序,只要状态字段名不变,out["notice"] 这类读法就一直有效。反过来,如果 invoke 返回的是「最后一个节点的返回值」,那每次调整图结构,所有调用方都得跟着改。
对照一下 count 节点:它只返回了 {"log": [...]},但最终结果里 text 也在,因为返回值是状态,而不是某一步的输出。
3.2. 哪些字段可以不给,哪些必须给 #
第 20 章提过这个坑,这里再确认一次:用 TypedDict 时缺字段不会在入口拦下:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
try:
# 传一个完全空的输入
print(graph.invoke({}))
except Exception as e:
# 打印异常类型和消息
print("报错:", type(e).__name__, str(e))报错: KeyError 'text'报错发生在 clean 节点内部访问 state["text"] 的那一刻,而不是 invoke 的入口。
但这里有个第 20 章没说透的区别:挂了 reducer 的字段可以省略:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# 只给 text,不给 log
print("只缺 log:", graph.invoke({"text": " a "}))只缺 log: {'text': 'a', 'log': ['clean', 'count=1']}一点问题都没有。原因在第 21 章 §3.3 讲过:带 reducer 的字段,LangGraph 会按类型标注推出一个「零值」当初始值(list → []、int → 0、str → '')。所以:
| 字段 | 有 reducer 吗 | 输入里能省略吗 | 省略后的初值 |
|---|---|---|---|
text: str |
没有 | 不能,节点里会 KeyError |
无 |
log: Annotated[list[str], operator.add] |
有 | 能 | [](按标注推出的零值) |
这条规则解释了一个常见现象:为什么图只传 {"messages": [...]} 就能跑? 因为 MessagesState 的唯一字段挂着 add_messages reducer,其它自定义字段要么也挂了 reducer,要么就必须在输入里给全。
两种防法:
防法一:写个工厂函数,保证初始状态完整
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# 定义状态类型
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
# 防法一:写个工厂函数,保证初始状态完整
def new_input(text: str) -> S:
# 把所有字段都填上,包括可以省略的那些,让初始状态一目了然
return {"text": text, "log": []}
# 节点定义
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
# 图构建
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# 使用工厂函数创建输入,传给 graph
result = graph.invoke(new_input(" hello "))
print("用工厂函数补全初始状态:", result)防法二:状态改用 Pydantic,入口就会校验(第 20 章 §4.3)
# 防法二:状态改用 Pydantic,入口就会校验(第 20 章 §4.3)
import operator
from typing import Annotated
from pydantic import BaseModel
from langgraph.graph import END, START, StateGraph
# 防法二:状态改用 Pydantic,入口就会校验(第 20 章 §4.3)
# 定义状态类型
class S(BaseModel):
text: str
log: Annotated[list[str], operator.add]
# 节点定义
def clean(state: S) -> dict:
return {"text": state.text.strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state.text)}"]}
# 图构建
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# 尝试直接传不完整的初始状态:Pydantic 会校验,缺字段时报错
try:
result = graph.invoke({"text": " hello "})
print("Pydantic 校验通过:", result)
except Exception as e:
print("Pydantic 校验未通过,错误信息:", str(e))推荐第一种:它不改状态类型,成本最低,而且函数签名本身就是一份「必须给什么」的文档。§8 的实战里 new_order() 就是这个角色。
3.3. batch:一次跑多份输入 #
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# 传一个输入列表,返回一个结果列表,顺序与输入一一对应
outs = graph.batch([{"text": " a "}, {"text": " bbb "}])
# 逐个打印结果
for o in outs:
# 每份结果都是一份独立的完整最终状态
print(" ", o) {'text': 'a', 'log': ['clean', 'count=1']}
{'text': 'bbb', 'log': ['clean', 'count=3']}batch 来自 Runnable 接口(第 7 章),你没为它写过任何代码,它是 §2.1 那条继承链白送的。多份输入之间互相独立,不共享状态。
那如果传了带 thread_id 的 config 会怎样?
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# InMemorySaver 是最简单的 checkpointer:状态存在内存里,进程退出就没了
# 编译一个带 checkpointer 的版本,才能事后查状态
graph = build().compile(checkpointer=InMemorySaver())
# 故意让三份输入共用同一个 thread_id
shared = {"configurable": {"thread_id": "shared"}}
# 三份输入一起跑
res = graph.batch([{"text": " a "}, {"text": " bbb "}, {"text": " cc "}], shared)
# 先看返回结果对不对
for o in res:
# 三份结果各算各的,没有互相污染
print(" 结果:", o)
# 再看这个 thread 里最终存下来的是哪一份
print(" thread 里最终存的:", graph.get_state(shared).values) 结果: {'text': 'a', 'log': ['clean', 'count=1']}
结果: {'text': 'bbb', 'log': ['clean', 'count=3']}
结果: {'text': 'cc', 'log': ['clean', 'count=2']}
thread 里最终存的: {'text': 'a', 'log': ['clean', 'count=1']} ← 每次运行都可能不一样返回结果是对的:三份各算各的,log 也没串味,前三行每次运行都一样。出问题的是落盘的那份状态,只剩一份,而且是哪一份完全看谁最后写完,没有任何规律。
这句话不是推测。把上面这段连跑 20 次,统计最后落盘的是谁:
'bbb': 8 次
'cc' : 8 次
'a' : 4 次三份输入都可能赢,包括第一份。 所以如果你跑出来的最后一行和上面不一样,不是你写错了。这正是这个坑的本质:它不会报错,也不稳定复现,因此测试环境很可能一次都碰不到,上线之后才随机丢状态。 结论很简单:batch 不要和 thread_id 一起用,要么每份输入给各自的 thread_id,要么就别传 config。
翻一眼历史就知道有多乱:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
shared = {"configurable": {"thread_id": "shared"}}
graph.batch([{"text": " a "}, {"text": " bbb "}, {"text": " cc "}], shared)
# 把这个 thread 的检查点历史全打出来
for h in graph.get_state_history(shared):
# 重点看 step 有没有重复、values 属于哪一份输入
print(f" step={h.metadata.get('step'):>2} next={h.next} values={h.values}") step= 2 next=() values={'text': 'a', 'log': ['clean', 'count=1']}
step= 2 next=() values={'text': 'bbb', 'log': ['clean', 'count=3']}
step= 1 next=('count',) values={'text': 'bbb', 'log': ['clean']}
step= 0 next=('clean',) values={'text': ' bbb ', 'log': []}
step= 1 next=('count',) values={'text': 'a', 'log': ['clean']}
step= 2 next=() values={'text': 'cc', 'log': ['clean', 'count=2']}
step= 1 next=('count',) values={'text': 'cc', 'log': ['clean']}
step= 0 next=('clean',) values={'text': ' cc ', 'log': []}
step= 0 next=('clean',) values={'text': ' a ', 'log': []}
step=-1 next=('__start__',) values={'log': []}
step=-1 next=('__start__',) values={'log': []}
step=-1 next=('__start__',) values={'log': []}每个 step 都出现了三次,还乱序交错。 条数是稳定的(三份输入各留下一套),但交错的具体次序每次都不一样,你跑出来的行序和上面对不上属于正常。这就是「污染」的真实形态:
batch共用thread_id,不会算错结果,但会把检查点历史搅成一锅粥。get_state只能捞到最后落盘的那一份,事后审计和续跑全都指望不上。
正确做法是给每份输入配一份独立 config:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
# config 也传列表,与输入列表一一对应
res = graph.batch(
[{"text": " a "}, {"text": " bbb "}],
[{"configurable": {"thread_id": "b1"}}, {"configurable": {"thread_id": "b2"}}],
)
# 第一个 thread 只存着第一份输入的结果
print(" b1:", graph.get_state({"configurable": {"thread_id": "b1"}}).values)
# 第二个 thread 只存着第二份输入的结果,两者互不干扰
print(" b2:", graph.get_state({"configurable": {"thread_id": "b2"}}).values) b1: {'text': 'a', 'log': ['clean', 'count=1']}
b2: {'text': 'bbb', 'log': ['clean', 'count=3']}3.4. 异步:ainvoke #
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# asyncio 用来跑协程
import asyncio
# 定义一个协程函数
async def demo():
# ainvoke 是 invoke 的异步版本,用 await 等它完成
out = await graph.ainvoke({"text": " 异步 "})
# 返回值形态与 invoke 完全一致
print("ainvoke:", out)
# 启动事件循环并执行协程
asyncio.run(demo())ainvoke: {'text': '异步', 'log': ['clean', 'count=2']}节点是普通同步函数也能用 ainvoke,LangGraph 会在线程池里跑它们。但节点里若有真正的 I/O(调模型、查库),把节点本身写成 async def 才有并发收益。
一条实用建议:同一张图里的节点,同步异步可以混着写,LangGraph 会各自用合适的方式调度。但如果你打算在 Web 服务(FastAPI 等)里用这张图,从一开始就统一写 async def 会省掉不少线程池相关的麻烦,尤其是节点里用了只支持异步的客户端时。
4. stream:要执行过程 #
invoke 只给你终点,stream 给你沿途。两种场景下这是刚需:调试(哪一步出的问题)和 UI(让用户看到进度,而不是干等)。
有一件事要先说清,否则后面的例子会看不懂:stream() 返回的是生成器,不消费就一个节点都不会执行。
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# 只拿迭代器,不进 for 循环
it = graph.stream({"text": " hi "})
# 此刻图还一步都没跑
print("拿到的是:", type(it).__name__)拿到的是: generator所以本章凡是不写 for 的地方,都会用 list(...) 把它消费掉。看到 list(graph.stream(...)) 不要以为是在收集结果,那只是「让它真的跑起来」。
4.1. updates:每个节点改了什么(默认) #
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# 不传 stream_mode 时默认就是 "updates"
for chunk in graph.stream({"text": " hi "}):
# 每条是一个只有一个键的 dict:键是节点名,值是该节点的返回值
print(" ", chunk) {'clean': {'text': 'hi', 'log': ['clean']}}
{'count': {'log': ['count=2']}}格式是 {节点名: 该节点的返回值}。注意 count 那条只有 log,因为节点本来就只返回了这一个键。
这个「只有增量」的特性正是它适合调试的原因:你看到的就是节点 return 的那个 dict 本身,一个字都没多。所以当某个节点的键拼错了(第 20 章讲的静默丢弃),updates 会诚实地把错误的键名打出来,而 values 里只会「看起来少了点什么」。
这是调试时的首选模式。 它直接告诉你「谁写了什么」,能立刻发现「某个节点返回的键拼错了」这类第 20 章讲的静默失败。
4.2. values:每一步之后的完整状态 #
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# stream_mode="values" 要的是每一步之后的完整状态
for chunk in graph.stream({"text": " hi "}, stream_mode="values"):
# 每条都是完整状态 dict,不带节点名
print(" ", chunk) {'text': ' hi ', 'log': []}
{'text': 'hi', 'log': ['clean']}
{'text': 'hi', 'log': ['clean', 'count=2']}两个节点却有三条输出,第一条是还没执行任何节点时的初始状态。这个细节常让人困惑,记住就好:values 的条数 = 节点数 + 1。
顺便留意第一条里的 'log': []:输入里我们只给了 text,log 这个 [] 是 §3.2 讲的「按类型标注推出的零值」,在这里能直接看到它。
updates 和 values 的取舍:
updates |
values |
|
|---|---|---|
| 内容 | 只有增量 | 每步的完整状态 |
| 带节点名 | 带 | 不带 |
| 数据量 | 小 | 大(状态大时很占带宽) |
| 适合 | 调试、日志 | UI 实时渲染当前全貌 |
一个实际的选型例子:状态里存着一份检索到的 20 篇文档,values 模式每一步都会把这 20 篇原样推一遍;updates 只在检索节点那一步推一次。状态越大,两者的带宽差距越夸张。
4.3. 同时要好几种 #
stream_mode 可以传列表:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# 传列表表示两种模式都要
for chunk in graph.stream({"text": " hi "}, stream_mode=["updates", "values"]):
# 注意此时每条的形态变了:是 (模式名, 内容) 的二元组
print(" ", chunk) ('values', {'text': ' hi ', 'log': []})
('updates', {'clean': {'text': 'hi', 'log': ['clean']}})
('values', {'text': 'hi', 'log': ['clean']})
('updates', {'count': {'log': ['count=2']}})
('values', {'text': 'hi', 'log': ['clean', 'count=2']})传列表时每条会变成 (模式名, 内容) 的元组,单模式时则直接是内容。写消费代码时要注意这个形态差异。
这个差异很容易踩坑:本来写着 for chunk in ...,为了多看点信息给 stream_mode 加了个列表,结果消费代码全线崩。稳妥的写法是从一开始就按元组解,哪怕只用一种模式:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# 即使只要一种模式,也传列表,让消费代码的形态固定下来
for mode, payload in graph.stream({"text": " hi "}, stream_mode=["updates"]):
# mode 一定是 "updates",payload 是 {节点名: 增量}
# 以后想再加一种模式,这段代码一个字都不用改
print(f" [{mode}] {payload}") [updates] {'clean': {'text': 'hi', 'log': ['clean']}}
[updates] {'count': {'log': ['count=2']}}4.4. messages:模型 token 级流式 #
前面几种都是节点级的:节点跑完才有输出。但如果节点里在调模型,用户要的是一个字一个字往外蹦。这就是 messages 模式:
# 从 .env 读取 API Key
from dotenv import load_dotenv
# init_chat_model 统一各家模型的初始化(第 2 章)
from langchain.chat_models import init_chat_model
# MessagesState 是内置的、只有一个 messages 字段的状态(第 21 章 §7)
from langgraph.graph import END, START, MessagesState, StateGraph
load_dotenv(override=True)
# temperature=0 让输出尽量稳定,方便对照
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)
# 节点:把历史消息丢给模型,把回复追加进 messages
def say(state: MessagesState) -> dict:
# 注意这里用的是 invoke(不是 stream),下面会解释为什么照样能流式
return {"messages": [model.invoke(state["messages"])]}
# 搭一张只有一个节点的图
builder = StateGraph(MessagesState)
# 注册这个调模型的节点
builder.add_node("say", say)
# 入口
builder.add_edge(START, "say")
# 出口
builder.add_edge("say", END)
# 编译
graph = builder.compile()
# 计数器,用来统计一共收到多少片段
n = 0
# stream_mode="messages" 时,每条是 (消息片段, 元数据) 二元组
for chunk, meta in graph.stream(
{"messages": [{"role": "user", "content": "用一句话介绍杭州"}]},
stream_mode="messages",
):
# 每收到一个片段就加一
n += 1
# 只打印前三条,看清片段长什么样
if n <= 3:
# langgraph_node 说明这个 token 来自哪个节点
print(f" type={type(chunk).__name__} content={chunk.content!r} 节点={meta['langgraph_node']}")
# 打印总片段数
print(f" ... 共 {n} 个片段") type=AIMessageChunk content='' 节点=say
type=AIMessageChunk content='' 节点=say
type=AIMessageChunk content='' 节点=say
... 共 59 个片段(片段总数每次跑都不一样,取决于模型这次想了多久、答了多长,几十上下都正常。下面几段里出现的具体数字同理,看的是比例关系,不是那个数本身。)
三个要点:
- 每条是
(消息片段, 元数据)的元组,片段类型是AIMessageChunk(不是AIMessage),元数据里的langgraph_node告诉你这个 token 来自哪个节点。多个节点都调模型时,靠它区分。 - 开头一大串片段的
content是空的,下面专门讲这是什么。 - 节点里写的是
model.invoke而不是model.stream,token 流式照样生效。LangGraph 在底层接管了模型的流式回调,所以这个模式不要求你改节点代码。
空片段到底是什么
「开头几个片段是空的」这个说法太轻描淡写了。数一下就知道问题有多大:
# ↓ §4.4 的单节点模型图,为了本段能独立运行而重复一遍
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)
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)
def say(state: MessagesState) -> dict:
return {"messages": [model.invoke(state["messages"])]}
builder = StateGraph(MessagesState)
builder.add_node("say", say)
builder.add_edge(START, "say")
builder.add_edge("say", END)
graph = builder.compile()
# 片段总数
n = 0
# 开头连续空片段的个数
empty_head = 0
# 第一个非空片段的序号,None 表示还没遇到
first_nonempty = None
# 照常用 messages 模式遍历
for chunk, meta in graph.stream(
{"messages": [{"role": "user", "content": "用一句话介绍杭州"}]},
stream_mode="messages",
):
# 计数
n += 1
# 还没遇到过非空片段,且当前片段为空,就计入「开头空片段」
if chunk.content == "" and first_nonempty is None:
# 累加开头空片段数
empty_head += 1
# 第一次遇到非空片段,记下它的序号
elif first_nonempty is None:
# 之后这个分支不会再进来
first_nonempty = n
# 汇总三个数字
print(f" 共 {n} 个片段;开头连续空片段 {empty_head} 个;第 {first_nonempty} 条才开始有内容") 共 39 个片段;开头连续空片段 16 个;第 17 条才开始有内容39 个片段里,开头 16 条一个字都没有。 如果 UI 直接把 chunk.content 往界面上追加,用户会先盯着空白框等上一会儿,然后文字才突然开始出现,体验比不流式还差。(这三个数字每次跑都不同,同一个问题另一次实测是 62 个片段、开头空了 42 条,思维链越长空片段越多。)
这些空片段里装的不是「元信息」,而是模型的思考过程:
# ↓ §4.4 的单节点模型图,为了本段能独立运行而重复一遍
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)
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)
def say(state: MessagesState) -> dict:
return {"messages": [model.invoke(state["messages"])]}
builder = StateGraph(MessagesState)
builder.add_node("say", say)
builder.add_edge(START, "say")
builder.add_edge("say", END)
graph = builder.compile()
# 打印前四条片段的完整内容,看看空 content 的片段里还有什么
n = 0
# 问一个极短的问题,减少输出量
for chunk, meta in graph.stream(
{"messages": [{"role": "user", "content": "1+1=?"}]}, stream_mode="messages"
):
# 计数
n += 1
# additional_kwargs 是各家模型放非标准字段的地方
print(f" [{n}] content={chunk.content!r} addl={chunk.additional_kwargs}")
# 看前四条就够了
if n >= 4:
break [1] content='' addl={'reasoning_content': ''}
[2] content='' addl={'reasoning_content': 'We'}
[3] content='' addl={'reasoning_content': ' need'}
[4] content='' addl={'reasoning_content': ' answer'}答案清楚了:deepseek-v4-flash 是推理模型,它先流式输出思维链(reasoning_content),再流式输出正式回答(content)。思维链阶段每个 token 都会产生一个片段,只是 content 为空。
这带来两条实用结论:
| 需求 | 做法 |
|---|---|
| 只要正式回答(最常见) | 过滤 if not chunk.content: continue |
| 想像某些产品那样展示「思考中…」 | 读 chunk.additional_kwargs.get("reasoning_content") |
| 换成非推理模型 | 空片段会少很多,但仍需过滤 |
最后一条是实测过的。同一个问题分别问两个模型:
| 模型 | 总片段 | 空片段 | 开头连续空片段 |
|---|---|---|---|
deepseek-chat(非推理) |
24 | 3 | 1 |
deepseek-v4-flash(推理) |
46 | 18 | 16 |
非推理模型的空片段从「几十个」降到「个位数」,但没有降到零,所以过滤那一行不能省。
正确的消费写法:
# ↓ §4.4 的单节点模型图,为了本段能独立运行而重复一遍
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)
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)
def say(state: MessagesState) -> dict:
return {"messages": [model.invoke(state["messages"])]}
builder = StateGraph(MessagesState)
builder.add_node("say", say)
builder.add_edge(START, "say")
builder.add_edge("say", END)
graph = builder.compile()
# 收集所有非空片段,拼成完整回答
parts = []
# 照常用 messages 模式遍历
for chunk, meta in graph.stream(
{"messages": [{"role": "user", "content": "用一句话介绍杭州"}]},
stream_mode="messages",
):
# 空片段直接跳过:这一行就是 UI 不闪的关键
if not chunk.content:
continue
# 真正要往界面上追加的内容
parts.append(chunk.content)
# 真实 UI 里这里是「往前端推一个 token」,示例中用 print 模拟
print(" 拼接结果:", "".join(parts))
# 非空片段数远小于总片段数
print(" 非空片段数:", len(parts)) 拼接结果: 杭州是中国浙江省的省会,一座以西湖美景、千年历史和数字经济闻名的“人间天堂”城市。
非空片段数: 22片段总数每次都不一样,因为它取决于模型这次想了多久、答了多长。任何依赖「片段数量」的代码都是错的,要依赖的是内容本身。
流式输出
# ↓ §4.4 的单节点模型图,为了本段能独立运行而重复一遍
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)
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)
def say(state: MessagesState) -> dict:
return {"messages": [model.invoke(state["messages"])]}
builder = StateGraph(MessagesState)
builder.add_node("say", say)
builder.add_edge(START, "say")
builder.add_edge("say", END)
graph = builder.compile()
# 收集所有非空片段,拼成完整回答
parts = []
# 记录当前打印到哪一段,用来在思考和回答之间只插一次分隔标题
section = None
# 照常用 messages 模式遍历
for chunk, meta in graph.stream(
{"messages": [{"role": "user", "content": "用一句话介绍杭州"}]},
stream_mode="messages",
):
# content_blocks 把 deepseek 的 reasoning_content 和正文统一成带 type 的块
for block in chunk.content_blocks:
# 思考块的正文在 reasoning 字段,正式回答在 text 字段
kind = block["type"]
piece = block.get("reasoning") if kind == "reasoning" else block.get("text")
# 空片段(比如只带 usage 的收尾 chunk)直接跳过
if not piece:
continue
# 段落切换时打一个标题,思考和回答就不会糊在一行里
if kind != section:
print(f"\n【{'思考' if kind == 'reasoning' else '回答'}】", end="")
section = kind
# 流式输出思考过程,再流式输出最终回答,都靠 flush 即时刷屏
print(piece, end="", flush=True)
# 只把正式回答攒进 parts,思考过程不算最终答案
if kind == "text":
parts.append(piece)4.5. custom:自己往外推进度 #
节点内部跑一个耗时循环时,节点级流式帮不上忙(节点没跑完就没输出)。这时可以自己写:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
# get_stream_writer 只能在节点函数内部调用
from langgraph.config import get_stream_writer
# 一个假装很耗时的节点
def worker(state: S) -> dict:
# 拿到「往外推数据」的写入器
writer = get_stream_writer()
# 推第一条进度,内容格式完全由你决定
writer({"进度": "开始"})
# 真实场景里这里是一段耗时循环,每轮推一次进度
writer({"进度": "完成"})
# 节点该返回的状态更新照常返回,两件事互不影响
return {"log": ["worker"]}
# 搭一张只有这个节点的图
builder = StateGraph(S)
# 注册这个会推进度的节点
builder.add_node("worker", worker)
# 入口
builder.add_edge(START, "worker")
# 出口
builder.add_edge("worker", END)
# 编译
graph = builder.compile()
# 用 custom 模式消费,收到的就是 writer() 里传的东西
for c in graph.stream({"text": "x"}, stream_mode="custom"):
# 内容形态由你在 writer() 里决定,这里是个 dict
print(" ", c) {'进度': '开始'}
{'进度': '完成'}writer() 里传什么就收到什么,格式完全由你定。这是「检索中…」「正在生成第 3 页…」这类进度提示的标准做法。
两个必须知道的配套细节。第一,custom 模式下看不到节点的返回值,上面输出里没有 {'worker': {...}}。要两样都要就传列表:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
from langgraph.config import get_stream_writer
def worker(state: S) -> dict:
writer = get_stream_writer()
writer({"进度": "开始"})
writer({"进度": "完成"})
return {"log": ["worker"]}
builder = StateGraph(S)
builder.add_node("worker", worker)
builder.add_edge(START, "worker")
builder.add_edge("worker", END)
graph = builder.compile()
# 同时要自定义进度和节点增量
for mode, payload in graph.stream({"text": "x"}, stream_mode=["custom", "updates"]):
# 多模式下每条是 (模式名, 内容),用模式名区分来源
print(f" [{mode}] {payload}") [custom] {'进度': '开始'}
[custom] {'进度': '完成'}
[updates] {'worker': {'log': ['worker']}}第二,没请求 custom 模式时,writer() 推的东西被静默丢弃:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
from langgraph.config import get_stream_writer
def worker(state: S) -> dict:
writer = get_stream_writer()
writer({"进度": "开始"})
writer({"进度": "完成"})
return {"log": ["worker"]}
builder = StateGraph(S)
builder.add_node("worker", worker)
builder.add_edge(START, "worker")
builder.add_edge("worker", END)
graph = builder.compile()
# 节点里照样调了 writer,但这次只要 updates
for c in graph.stream({"text": "x"}, stream_mode="updates"):
# 只剩节点增量,进度不见了
print(" ", c) {'worker': {'log': ['worker']}}没有报错,也没有警告,进度就是没了。所以「进度不显示」的第一个排查点不是节点代码,而是调用方有没有请求 custom 模式。
4.6. debug 和 tasks:排障用的细粒度事件 #
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# tasks 模式给出节点级的「开始 / 结束」事件
for chunk in graph.stream({"text": " hi "}, stream_mode="tasks"):
# 每个节点两条:一条带 input/triggers,一条带 result/error
print(" ", chunk) {'id': '6e5af5c3-...', 'name': 'clean', 'input': {'text': ' hi ', 'log': []}, 'triggers': ('branch:to:clean',)}
{'id': '6e5af5c3-...', 'name': 'clean', 'error': None, 'result': {'text': 'hi', 'log': ['clean']}, 'interrupts': []}
{'id': '74e00ea5-...', 'name': 'count', 'input': {'text': 'hi', 'log': ['clean']}, 'triggers': ('branch:to:count',)}
{'id': '74e00ea5-...', 'name': 'count', 'error': None, 'result': {'log': ['count=2']}, 'interrupts': []}每个节点产生两条:开始(有 input 和 triggers)和结束(有 result 和 error)。同一个节点的两条 id 相同,这是配对的依据。节点并行执行时,事件会交错到达,靠 id 才能把「开始」和「结束」对上。
triggers 告诉你这个节点是被哪条边触发的(branch:to:clean 表示「有条边指向 clean」),排查「为什么这个节点跑了/没跑」时很有用;input 则记下了节点当时看到的状态,这比自己在节点里加 print 方便得多,尤其是并行分支各自看到的快照不同时(第 21 章 §5.3)。
debug 模式内容类似,多包了一层,并且带上了超步编号和时间戳:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# debug 模式:在 tasks 的基础上加了 step 和 timestamp
for chunk in graph.stream({"text": " hi "}, stream_mode="debug"):
# 真实事件很长,这里原样打印,重点看 step 字段
print(" ", chunk) {'step': 1, 'timestamp': '2026-08-13T14:59:15.407244+00:00', 'type': 'task', 'payload': {'id': '26e245b6-...', 'name': 'clean', ...}}
{'step': 1, 'timestamp': '2026-08-13T14:59:15.407331+00:00', 'type': 'task_result', 'payload': {'id': '26e245b6-...', 'name': 'clean', 'error': None, ...}}
{'step': 2, 'timestamp': '2026-08-13T14:59:15.407481+00:00', 'type': 'task', 'payload': {'id': '74fbaa4d-...', 'name': 'count', ...}}
{'step': 2, 'timestamp': '2026-08-13T14:59:15.407633+00:00', 'type': 'task_result', 'payload': {'id': '74fbaa4d-...', 'name': 'count', 'error': None, ...}}留意 step:clean 是第 1 步,count 是第 2 步。step 数的是超步,不是事件。串行图里一步一个节点;并行图里同一步会出现多个节点,这时 step 就是判断「谁和谁是并行的」的直接证据(第 21 章 §5)。时间戳则能算出每个节点的真实耗时。
还有一个 checkpoints 模式,只在编译时传了 checkpointer 才有输出:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from rich import print
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# 创建一个不带检查点保存器(checkpointer)的图对象
graph_plain = build().compile()
# 创建一个带内存检查点保存器的图对象
graph_cp = build().compile(checkpointer=InMemorySaver())
# 使用不带检查点的图对象进行推理流式输出(stream 模式为 checkpoints)
for item in graph_plain.stream({"text": " hi "}, stream_mode="checkpoints"):
# 打印没有检查点功能的流式输出结果
print("without-checkpoints", item)
# 使用带检查点的图对象进行推理流式输出,添加 thread_id 标识(stream 模式为 checkpoints)
for item in graph_cp.stream(
{"text": " hi "},
{"configurable": {"thread_id": "cp1"}},
stream_mode="checkpoints",
):
# 打印有检查点功能的流式输出结果
print("checkpoints", item)
checkpoints
{
'config': {
'configurable': {
'checkpoint_ns': '',
'thread_id': 'cp1',
'checkpoint_id': '1f1a85ef-884b-6c0c-bfff-7484cfe80993'
}
},
'parent_config': None,
'values': {'log': []},
'metadata': {'source': 'input', 'step': -1, 'parents': {}},
'next': ['__start__'],
'tasks': [
{
'id': '417e70c0-f596-1d38-90fb-9799ba41dca1',
'name': '__start__',
'interrupts': (),
'state': None
}
]
}
checkpoints
{
'config': {
'configurable': {
'checkpoint_ns': '',
'thread_id': 'cp1',
'checkpoint_id': '1f1a85ef-88c3-65c7-8000-02c60003b5a3'
}
},
'parent_config': {
'configurable': {
'checkpoint_ns': '',
'thread_id': 'cp1',
'checkpoint_id': '1f1a85ef-884b-6c0c-bfff-7484cfe80993'
}
},
'values': {'text': ' hi ', 'log': []},
'metadata': {'source': 'loop', 'step': 0, 'parents': {}},
'next': ['clean'],
'tasks': [
{
'id': '1b7369d9-4f42-b2d4-538e-632205f56fe8',
'name': 'clean',
'interrupts': (),
'state': None
}
]
}
checkpoints
{
'config': {
'configurable': {
'checkpoint_ns': '',
'thread_id': 'cp1',
'checkpoint_id': '1f1a85ef-88cb-646b-8001-c9b998b58e9e'
}
},
'parent_config': {
'configurable': {
'checkpoint_ns': '',
'thread_id': 'cp1',
'checkpoint_id': '1f1a85ef-88c3-65c7-8000-02c60003b5a3'
}
},
'values': {'text': 'hi', 'log': ['clean']},
'metadata': {'source': 'loop', 'step': 1, 'parents': {}},
'next': ['count'],
'tasks': [
{
'id': 'bfa83989-b31d-51f8-57dc-1abfa2c1c547',
'name': 'count',
'interrupts': (),
'state': None
}
]
}
checkpoints
{
'config': {
'configurable': {
'checkpoint_ns': '',
'thread_id': 'cp1',
'checkpoint_id': '1f1a85ef-88d1-6d90-8002-105a62b1abff'
}
},
'parent_config': {
'configurable': {
'checkpoint_ns': '',
'thread_id': 'cp1',
'checkpoint_id': '1f1a85ef-88cb-646b-8001-c9b998b58e9e'
}
},
'values': {'text': 'hi', 'log': ['clean', 'count=2']},
'metadata': {'source': 'loop', 'step': 2, 'parents': {}},
'next': [],
'tasks': []没有 checkpointer 时一条都没有,也不报错,又一个静默行为。有 checkpointer 时,两节点的图产生 4 个检查点(输入、每个节点后各一个、最终各算一次),和 §6.3 的历史条数正好对得上。
4.7. 七种模式选型表 #
模式名不用背,问代码就行:
# typing.get_args 能取出 Literal 类型里的所有取值
import typing
# StreamMode 是 LangGraph 定义的 Literal 类型
from langgraph.types import StreamMode
# 打印当前版本支持的全部模式名
print(typing.get_args(StreamMode))('values', 'updates', 'checkpoints', 'tasks', 'debug', 'messages', 'custom')| 模式 | 每条是什么 | 什么时候用 |
|---|---|---|
updates |
{节点名: 增量} |
调试首选、写日志 |
values |
每步后的完整状态 | UI 渲染当前全貌 |
messages |
(AIMessageChunk, 元数据) |
聊天界面打字机效果(记得过滤空片段) |
custom |
你 writer() 里传的任意内容 |
长任务进度提示 |
tasks |
节点开始 / 结束事件(带 input、triggers) |
查「谁触发了谁」 |
debug |
tasks 再加 step 和时间戳 |
深度排障、算耗时、看并行 |
checkpoints |
检查点写入事件 | 观察持久化(需 checkpointer) |
真实项目里最常见的组合是 ["messages", "updates"]:前者喂给聊天气泡,后者驱动「正在检索 / 正在调用工具」这类状态提示。
4.8. 两个省事的小参数 #
output_keys 只要某几个字段:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# output_keys 传单个字段名字符串
for c in graph.stream({"text": " hi "}, stream_mode="values", output_keys="log"):
# 注意输出直接就是 log 的值,没有外层字典
print(" ", c) ['clean']
['clean', 'count=2']传单个字段名时,输出直接就是那个字段的值(不再包一层字典)。传列表则保留字典形态:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# 同一个字段,但用列表形式传
for c in graph.stream({"text": " hi "}, stream_mode="values", output_keys=["log"]):
# 这次保留了 {字段名: 值} 的字典形态
print(" ", c) {'log': ['clean']}
{'log': ['clean', 'count=2']}这个形态差异和 §4.3 的元组差异是同一类坑:为了少写一层解包传了字符串,后来想多要一个字段改成列表,消费代码就崩了。统一传列表,形态就永远稳定。
顺便一提,output_keys 对 updates 模式也有效,裁剪的是每个节点增量里的字段:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# 在 updates 模式下裁剪字段
for c in graph.stream({"text": " hi "}, stream_mode="updates", output_keys="log"):
# 外层的节点名还在,只是每个节点的增量被裁成了 log
print(" ", c) {'clean': ['clean']}
{'count': ['count=2']}print_mode 不用自己写循环,直接打到控制台:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# print_mode 让 LangGraph 自己打印,list() 只是为了消费掉迭代器
list(graph.stream({"text": " hi "}, print_mode="updates"))[updates] {'clean': {'text': 'hi', 'log': ['clean']}}
[updates] {'count': {'log': ['count=2']}}调试时很顺手(实际输出带 ANSI 颜色高亮,模式名是加粗的)。它也能传列表,一次看两种:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
# print_mode 同样支持列表
list(graph.stream({"text": " hi "}, print_mode=["updates", "values"]))[values] {'text': ' hi ', 'log': []}
[updates] {'clean': {'text': 'hi', 'log': ['clean']}}
[values] {'text': 'hi', 'log': ['clean']}
[updates] {'count': {'log': ['count=2']}}
[values] {'text': 'hi', 'log': ['clean', 'count=2']}注意 print_mode 和 stream_mode 是两个独立参数:前者管「LangGraph 帮你打印什么」,后者管「迭代器交给你什么」。只想看一眼就用 print_mode + list();要拿数据做处理就用 stream_mode + for。
4.9. astream 与 astream_events #
- LangGraph 除了同步流式 (
stream) 方法,还提供了异步版本astream,接口参数完全一样,用于异步场景(例如在异步 Web 框架、Jupyter Notebook、GUI 等环境)。 astream_events进一步提供带进度事件的异步流式,适合「看进度、查节点事件」等高级用法。它返回的是更细粒度的操作事件(如start_node,end_node),适用于深度排障或自定义 UI 进度条。
astream:异步流式,常用,和同步用法几乎一致,推荐优先用这个astream_events:进阶用法,监控每一个关键事件(节点、状态变更等),复杂进度 UI 、调试等场景用
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile()
import asyncio
# astream 是 stream 的异步版,参数完全一致
async def demo():
# 用 async for 消费异步生成器
async for c in graph.astream({"text": " 异步 "}):
# 默认也是 updates 模式,形态和同步版一致
print(" ", c)
# 启动事件循环
asyncio.run(demo()) {'clean': {'text': '异步', 'log': ['clean']}}
{'count': {'log': ['count=2']}}astream_events 则给出更细的、贯穿整个调用栈的事件流:
# ↓ §4.4 的单节点模型图,为了本段能独立运行而重复一遍
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)
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)
def say(state: MessagesState) -> dict:
return {"messages": [model.invoke(state["messages"])]}
builder = StateGraph(MessagesState)
builder.add_node("say", say)
builder.add_edge(START, "say")
builder.add_edge("say", END)
graph = builder.compile()
import asyncio
# 统计各类事件出现的次数,而不是把事件全打出来(量太大)
async def demo_events():
# 用字典累计每种 event 的条数
kinds = {}
# astream_events 只有异步版本,没有同步的 stream_events
async for ev in graph.astream_events(
{"messages": [{"role": "user", "content": "1+1=?"}]}
):
# ev["event"] 是事件类型名,如 on_chat_model_stream
kinds[ev["event"]] = kinds.get(ev["event"], 0) + 1
# 打印统计结果
print("事件类型统计:", kinds)
# 启动事件循环
asyncio.run(demo_events())事件类型统计: {'on_chain_start': 2, 'on_chat_model_start': 1, 'on_chat_model_stream': 32,
'on_chat_model_end': 1, 'on_chain_stream': 2, 'on_chain_end': 2}| 事件类型 | 含义说明 | 本例次数 | 说明 |
|---|---|---|---|
on_chain_start |
Graph/链条开始执行事件,每一层 chain/图各产生一次 | 2 | 本例有图和单节点各计一次 chain |
on_chat_model_start |
Chat model 开始响应,进入模型处理阶段 | 1 | 仅一次,对应调用模型时刻 |
on_chat_model_stream |
Chat model 每产出一个 stream 响应(如每个 token 或 chunk)都来一条 | 18 | 分块/token 数,每次运行通常不同 |
on_chat_model_end |
Chat model 输出/推理完成 | 1 | 仅一次,模型输出完全结束 |
on_chain_stream |
节点(chain)每次有增量流式输出时产生,可用于 UI 等的逐步渲染 | 2 | 对应本例中流式输出的两步 |
on_chain_end |
节点(chain)执行结束,每层 chain/图各产生一次 | 2 | 结束标志,有助于测量 chain 周期 |
注:token 数(
on_chat_model_stream计数)和模型内容有关,每次运行会略有变化。
这些事件类型贯穿调用全过程,用于区分不同层级的开始、增量流内容、结束等关键节点,方便开发者跟踪推理细节和状态流转。
它连模型的开始、每个 token、结束都拆开了:on_chat_model_start 一条、on_chat_model_stream 每个片段一条(所以这个数字每次都不同)、on_chat_model_end 一条。on_chain_* 各两条是因为「图本身」和「节点」各算一层链。这种粒度能做到很精细的前端控制,比如收到 on_chat_model_start 就把「思考中」的转圈换成光标。代价是事件量大、处理起来啰嗦(on_chat_model_stream 的条数会随回答长度线性增长)。第 27 章会专门讲它。
选型建议:先用
stream_mode,它能满足九成需求;只有需要区分「模型开始输出」和「工具开始执行」这种细节时,才上astream_events。
5. config:运行时的旋钮 #
invoke / stream 的第二个位置参数是 config,它不改变图的结构,只影响这一次运行。
config 是个普通嵌套字典,有一处结构很容易搞错:有的键在顶层,有的键在 configurable 里面:
# config 就是这样一个普通嵌套字典
config = {
# 顶层键:LangGraph / LangChain 框架自己认的运行参数
"recursion_limit": 25,
"tags": ["demo"],
"metadata": {"user": "u1"},
# configurable 里面:会被透传给 checkpointer、以及第 23 章之后的可配置字段
"configurable": {"thread_id": "t1"},
}
# 顶层有哪些键
print("顶层键:", list(config))
# configurable 里有哪些键
print("configurable 里的键:", list(config["configurable"]))顶层键: ['recursion_limit', 'tags', 'metadata', 'configurable']
configurable 里的键: ['thread_id']记混的代价是静默失效(§5.2 会实测),所以先记住这条界线:thread_id 在里面,其它三个在外面。
5.1. configurable.thread_id #
这是最常用的一项,第 11 章已经见过,它决定状态存到哪个会话里:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# InMemorySaver:状态存内存,进程退出就没了
# 编译一个带持久化的版本
graph = build().compile(checkpointer=InMemorySaver())
# 第二个位置参数就是 config,thread_id 指定这次运行属于哪个会话
out = graph.invoke({"text": " 你好 "}, {"configurable": {"thread_id": "t1"}})
# 结果和不带 checkpointer 时一样,区别在于这次的状态被存下来了
print(out){'text': '你好', 'log': ['clean', 'count=2']}编译时传了 checkpointer 却不给 thread_id,会直接报错:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
try:
# 完全不传 config
graph.invoke({"text": " a "}) # 没给 config
except Exception as e:
# 打印异常类型和消息
print("报错:", type(e).__name__, str(e))报错: ValueError Checkpointer requires one or more of the following 'configurable' keys: thread_id, checkpoint_ns, checkpoint_id传了空的 configurable 也一样报这个错。它要的不是「有 config」,而是「能定位到一个会话」:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
try:
# 传了 config,但 configurable 是空的
graph.invoke({"text": " a "}, {"configurable": {}})
except Exception as e:
# 报的是同一个错:它要的是「能定位到会话」,不是「有 config」
print("报错:", type(e).__name__, str(e))报错: ValueError Checkpointer requires one or more of the following 'configurable' keys: thread_id, checkpoint_ns, checkpoint_id这条报错是「大声失败」,反而省心:忘了传 thread_id 会当场炸,而不是悄悄把所有用户的会话混在一起。
同一个 thread_id 多次 invoke,状态会按 reducer 规则累积:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
# 固定一个 thread,连着跑三次
cfg3 = {"configurable": {"thread_id": "t3"}}
# 第一次:log 里两条
print("第一次:", graph.invoke({"text": " a "}, cfg3))
# 第二次:log 累积到四条,text 被新值覆盖
print("第二次:", graph.invoke({"text": " bb "}, cfg3))
# 第三次:log 累积到六条
print("第三次:", graph.invoke({"text": " ccc "}, cfg3))第一次: {'text': 'a', 'log': ['clean', 'count=1']}
第二次: {'text': 'bb', 'log': ['clean', 'count=1', 'clean', 'count=2']}
第三次: {'text': 'ccc', 'log': ['clean', 'count=1', 'clean', 'count=2', 'clean', 'count=3']}text 被覆盖(没 reducer),log 累积(有 reducer),完全符合第 21 章的规则。
这个「累积」是双刃剑:多轮对话正是靠它把历史消息攒起来的,但如果你的字段本来只想存「这一次的结果」,同一个 thread 跑多次就会越滚越长。判断标准很简单:这个字段的语义是「历史」还是「当前」?是历史就挂 reducer,是当前就别挂。
5.2. recursion_limit:防死循环 #
图里有环时(第 24 章会大量用到),需要一个上限兜底。故意写一个死循环:
# 这一节的实验只用到一个计数字段,所以不复用 §1.2 那张图
from typing import TypedDict
from langgraph.graph import START, StateGraph
# 只有一个计数字段的状态
class LoopS(TypedDict):
# 计数器,没挂 reducer,每次被覆盖
n: int
# 搭一张自己连自己的图
builder = StateGraph(LoopS)
# 节点每次把 n 加一,永远不会停
builder.add_node("tick", lambda s: {"n": s["n"] + 1})
# 入口
builder.add_edge(START, "tick")
# 出边指回自己:这就是一个无条件死循环
builder.add_edge("tick", "tick") # 自己连自己
# 编译(注意 compile 不会检查有没有环,它管不了这个)
graph = builder.compile()
try:
# 不设上限,看默认值是多少
graph.invoke({"n": 0})
except Exception as e:
# 报错信息里会带上实际生效的上限值
print("报错:", type(e).__name__, str(e))报错: GraphRecursionError Recursion limit of 10007 reached without hitting a stop condition.
You can increase the limit by setting the `recursion_limit` config key.
For troubleshooting, visit: https://docs.langchain.com/oss/python/langgraph/errors/GRAPH_RECURSION_LIMIT默认值是 10007,不是很多资料里写的 25。 这个数字来自 LangGraph 源码 langgraph/_internal/_config.py(下面这行是源码摘录,不用自己跑):
DEFAULT_RECURSION_LIMIT = int(getenv("LANGGRAPH_DEFAULT_RECURSION_LIMIT", "10007"))而且它还能用环境变量 LANGGRAPH_DEFAULT_RECURSION_LIMIT 全局改掉。这里有三个实际影响:
- 别指望默认值救你。 一万步足够把 API 额度烧穿了,有环的图必须自己设一个小的上限。
- 这是版本相关的。 老版本 LangGraph 默认是 25,遇到行为不一致先确认版本。
- 它在私有模块里。 别写
from langgraph._internal._config import DEFAULT_RECURSION_LIMIT,下个版本可能就搬走了;要读当前值就os.getenv那个环境变量。
单次调用改成 5:
# ↓ §5.2 的死循环图,为了本段能独立运行而重复一遍
from typing import TypedDict
from langgraph.graph import START, StateGraph
class LoopS(TypedDict):
n: int
builder = StateGraph(LoopS)
builder.add_node("tick", lambda s: {"n": s["n"] + 1})
builder.add_edge(START, "tick")
builder.add_edge("tick", "tick")
graph = builder.compile()
try:
# recursion_limit 是 config 的顶层键
graph.invoke({"n": 0}, {"recursion_limit": 5})
except Exception as e:
# 只打印前 70 个字符,看清生效的数字就行
print(type(e).__name__, str(e)[:70])GraphRecursionError Recursion limit of 5 reached without hitting a stop condition. You can注意 recursion_limit 是 config 的顶层键,不在 configurable 里面。放错了会怎样?
# ↓ §5.2 的死循环图,为了本段能独立运行而重复一遍
from typing import TypedDict
from langgraph.graph import START, StateGraph
class LoopS(TypedDict):
n: int
builder = StateGraph(LoopS)
builder.add_node("tick", lambda s: {"n": s["n"] + 1})
builder.add_edge(START, "tick")
builder.add_edge("tick", "tick")
graph = builder.compile()
try:
# 故意把它塞进 configurable
graph.invoke({"n": 0}, {"configurable": {"recursion_limit": 5}})
except Exception as e:
# 看报错里的数字是 5 还是 10007,就知道到底生效没
print(type(e).__name__, str(e)[:70])GraphRecursionError Recursion limit of 10007 reached without hitting a stop condition. You数字还是 10007,你设的 5 被完全无视了,没有任何报错或警告。 这个坑很隐蔽:代码看起来「设了上限」,实际上跑了一万步。报错信息里的那个数字就是唯一的验收凭据,写完限制一定要故意触发一次,确认打出来的是你设的值。
它数的到底是什么
recursion_limit 数的是超步(superstep,第 20 章 §6.4 提过),不是节点数、也不是循环圈数。串行图里一步一个节点,所以可以直接实测出对应关系:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# 拿本章那张两节点图(clean -> count)
graph = build().compile()
# 从 1 开始试,找出最小的可用上限
for lim in (1, 2, 3):
try:
# 每次只改 recursion_limit
graph.invoke({"text": " hi "}, {"recursion_limit": lim})
# 没抛异常说明这个上限够用
print(f" limit={lim} 通过")
except Exception as e:
# 抛了就说明上限不够
print(f" limit={lim} -> {type(e).__name__}") limit=1 -> GraphRecursionError
limit=2 -> GraphRecursionError
limit=3 通过两个节点,要 3 才够。再验一张五节点串行图:
from typing import TypedDict
from langgraph.graph import END, START, StateGraph
# 五个节点串成一条线
class NS(TypedDict):
# 只有一个计数字段
n: int
# 新建图纸
builder = StateGraph(NS)
# 批量注册 n0~n4,每个都把 n 加一
for i in range(5):
# 节点名用 f-string 拼出来,函数体都一样
builder.add_node(f"n{i}", lambda s: {"n": s["n"] + 1})
# 入口连到第一个节点
builder.add_edge(START, "n0")
# 依次串起来:n0->n1->n2->n3->n4
for i in range(4):
# 第 i 个连到第 i+1 个
builder.add_edge(f"n{i}", f"n{i+1}")
# 最后一个连到 END
builder.add_edge("n4", END)
# 编译
graph = builder.compile()
# 在临界值附近试
for lim in (4, 5, 6):
try:
# 同一张图,只改上限
graph.invoke({"n": 0}, {"recursion_limit": lim})
# 通过说明上限够
print(f" limit={lim} 通过")
except Exception:
# 不够就会抛 GraphRecursionError
print(f" limit={lim} -> GraphRecursionError") limit=4 -> GraphRecursionError
limit=5 -> GraphRecursionError
limit=6 通过规律清楚了:
n 个节点的串行图,
recursion_limit至少要 n + 1。 多出来的那 1 步是把状态推到END。
所以设上限时不能掐得太紧。一个实用的算法:「一轮循环用几个超步 × 你允许它转几轮 + 图的固定串行长度 + 2」,宁可多给两三步。上限的作用是防失控,不是精确计数。
在 LangGraph 中,一个超步被定义为对图节点的一次迭代。它的核心规则是:
- 并行节点 → 同一超步:所有在逻辑上可以并行运行的节点,都属于同一个超步。
- 串行节点 → 不同超步:必须等待前序节点完成后才能运行的节点,则属于不同的超步。
5.3. tags 和 metadata #
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
list( # stream 是迭代器,要消费掉才会真正执行
graph.stream(
# 输入
{"text": " y "},
# config:三种键一起用,注意 thread_id 在 configurable 里,另两个在顶层
{
"configurable": {"thread_id": "t5"},
# tags 是一串标签,主要给 LangSmith 筛选用
"tags": ["demo"],
# metadata 是任意键值对,会被写进检查点
"metadata": {"user": "u1"},
},
stream_mode="updates",
)
)
# 跑完之后把检查点的 metadata 读出来
print("metadata:", graph.get_state({"configurable": {"thread_id": "t5"}}).metadata)metadata: {'source': 'loop', 'step': 2, 'parents': {}, 'user': 'u1'}| 字段 | 含义 | 示例值 | 说明 |
|---|---|---|---|
| source | 本次计算的来源标记 | 'loop' | 哪个节点或系统过程作为这次检查点的“来源”。如循环、节点名等。 |
| step | 累计运行的步数(超步/superstep) | 2 | 当前 session/thread 已执行的步数。 |
| parents | 前一状态的来源(分支合并用) | {} | 上一步结果的父节点信息。串行场景一般为空字典。 |
「source」「step」「parents」都是系统自动生成的元数据字段,用于追踪流程来源、步骤计数和流转关系。
你传的 metadata 会合并进 checkpoint 的 metadata,跟系统自带的 source / step 放在一起(user: 'u1' 就是我们传进去的那条)。可以用来记「哪个用户、哪个请求 ID 触发了这次运行」,事后排查能对上号。
由于是合并而不是替换,有一个小风险:如果你传的键名正好和系统键(source / step / parents)撞了,会把系统字段盖掉,而 §7.3 的排障恰好依赖 source。给自己的键加个前缀(app_user、req_id)是省心的做法。
tags 则主要给 LangSmith 用,方便在 trace 列表里筛选(第 29 章会用到)。它不进检查点,所以上面的输出里看不到 demo。
5.4. 输入输出裁剪:input_schema / output_schema #
有时候图内部状态字段很多,但对外只想暴露一小部分。参数名同样可以问代码:
import inspect
from langgraph.graph import StateGraph
# 看 StateGraph 构造函数收哪些参数
print(list(inspect.signature(StateGraph.__init__).parameters))['self', 'state_schema', 'context_schema', 'input_schema', 'output_schema', 'kwargs']| 参数名 | 类型 | 说明 | 是否必填 |
|---|---|---|---|
| state_schema | 类型/类 | 内部状态结构,定义整个状态的字段和类型。所有节点之间交互的数据都存在于state里。 | 必填 |
| context_schema | 类型/类 | 扩展/上下文信息结构,通常用于传递不属于「主状态」的数据(如API密钥、全局资源等)。 | 选填 |
| input_schema | 类型/类 | 输入裁剪:限制调用(如invoke)入口参数格式。只允许这些字段作为输入,其它被忽略。 | 选填 |
| output_schema | 类型/类 | 输出裁剪:限制最终输出的字段。比如只想对外暴露结果的某几个字段。 | 选填 |
| kwargs | 其他关键字参数 | 其它附加参数,通常是内部扩展用。 | 选填 |
总结:
state_schema定义了图内部通用的完整状态,input_schema/output_schema用于裁剪接口,context_schema传递额外上下文,kwargs支持其它特殊用途。
state_schema 是内部状态(必填),input_schema / output_schema 是可选的对外收窄:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# 对外只收 text
class InSchema(TypedDict):
# 入口只认这一个字段,其它字段会被丢掉
text: str
# 对外只吐 log
class OutSchema(TypedDict):
# invoke 的返回值里只会有这一个字段
log: list[str]
# 内部状态仍是 S,只是限定了入口和出口能看到什么
builder = StateGraph(S, input_schema=InSchema, output_schema=OutSchema)
# 节点复用前面定义好的两个函数,它们读写的还是内部状态 S
builder.add_node("clean", clean)
# 第二个节点同样不用改
builder.add_node("count", count)
# 入口
builder.add_edge(START, "clean")
# 串行
builder.add_edge("clean", "count")
# 出口
builder.add_edge("count", END)
# 编译
graph = builder.compile()
# 只有 log 会出现在返回值里
print("invoke 结果:", graph.invoke({"text": " hi "}))invoke 结果: {'log': ['clean', 'count=2']}text 没出现在返回值里,被 output_schema 裁掉了。注意节点代码一行都没改:节点看到的仍是完整的内部状态 S,裁剪只发生在图的边界上。
input_schema 那一侧则是静默丢弃:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
class InSchema(TypedDict):
text: str
class OutSchema(TypedDict):
log: list[str]
builder = StateGraph(S, input_schema=InSchema, output_schema=OutSchema)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
graph = builder.compile()
# 故意多传一个 input_schema 里没有的 log 字段
print("多传 log:", graph.invoke({"text": " hi ", "log": ["外部传入"]}))多传 log: {'log': ['clean', 'count=2']}'外部传入' 消失了,没有报错。这和第 20 章讲的「输入里多余的键被静默丢弃」是同一个机制,只不过这里的「多余」是相对 input_schema 而言的。收窄了入口,原本合法的字段也会变成多余。
还要注意一个不一致:output_schema 只作用于 invoke 的返回值,不影响 stream。
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
class InSchema(TypedDict):
text: str
class OutSchema(TypedDict):
log: list[str]
builder = StateGraph(S, input_schema=InSchema, output_schema=OutSchema)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
graph = builder.compile()
# 同一张裁剪过的图,用 stream 看
for c in graph.stream({"text": " hi "}, stream_mode="values"):
# 被 output_schema 裁掉的 text 又出现了
print(" ", c) {'text': ' hi ', 'log': []}
{'text': 'hi', 'log': ['clean']}
{'text': 'hi', 'log': ['clean', 'count=2']}text 还在。updates 模式也一样漏:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
class InSchema(TypedDict):
text: str
class OutSchema(TypedDict):
log: list[str]
builder = StateGraph(S, input_schema=InSchema, output_schema=OutSchema)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
graph = builder.compile()
# 换成 updates 模式,被裁的字段照样出现
for c in graph.stream({"text": " hi "}, stream_mode="updates"):
# clean 那条里能看到 text
print(" ", c) {'clean': {'text': 'hi', 'log': ['clean']}}
{'count': {'log': ['count=2']}}想在 stream 里也裁剪,得用 §4.8 的 output_keys。如果 output_schema 是为了不泄露敏感字段而设的,这个差异要特别小心:把 output_schema 当安全边界,然后在接口里换成流式返回,敏感字段就跟着流出去了。
| 手段 | 影响 invoke 返回值 |
影响 stream |
定在哪 |
|---|---|---|---|
output_schema |
裁 | 不裁 | 建图时(结构) |
output_keys |
裁 | 裁 | 每次调用时(运行期) |
真要拦敏感字段,两个一起上,或者干脆别把它放进状态。
6. get_state:查看运行现场 #
前面两个接口都是「跑」,这一个是「查」。它是 LangGraph 相对普通函数调用多出来的核心能力之一。
先想想普通函数是什么处境:函数跑完,局部变量全部释放,中途状态无从查证;函数中途抛异常,已经算好的那半也一起没了。get_state 系列 API 解决的就是这个问题:把「执行过程」变成了可以事后查询的数据。这也是本章 §6、§7 被排在最后的原因:它们才是 LangGraph 真正拉开差距的地方。
6.1. 前提:必须有 checkpointer #
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# 编译时什么都不传
graph = build().compile() # 没传 checkpointer
try:
# 尝试查状态
graph.get_state({"configurable": {"thread_id": "t1"}})
except Exception as e:
# 没存过就没得查
print("报错:", type(e).__name__, str(e))报错: ValueError No checkpointer set同一族的另外两个方法也一样:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# 同样是不带 checkpointer 的版本
graph = build().compile()
try:
# 查历史,同样需要 checkpointer;list() 用来触发生成器
list(graph.get_state_history({"configurable": {"thread_id": "t1"}}))
except Exception as e:
# 同一个错
print("get_state_history:", type(e).__name__, str(e))
try:
# 改状态,同样需要 checkpointer
graph.update_state({"configurable": {"thread_id": "t1"}}, {"log": ["x"]})
except Exception as e:
# 还是同一个错
print("update_state:", type(e).__name__, str(e))get_state_history: ValueError No checkpointer set
update_state: ValueError No checkpointer set道理很直白:没存过就没得查。 加上就行:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# 挂上 checkpointer,§6 之后的实验都这么编译
from langgraph.checkpoint.memory import InMemorySaver
# checkpointer 是 compile 的第一个位置参数,可以不写参数名
graph = build().compile(checkpointer=InMemorySaver())
# 先跑一次,让这个 thread 里有东西可查
cfg = {"configurable": {"thread_id": "t1"}}
graph.invoke({"text": " 你好 "}, cfg)
# 这次不再报错,能读到状态了
print("现在查得到了:", graph.get_state(cfg).values)现在查得到了: {'text': '你好', 'log': ['clean', 'count=2']}
InMemorySaver的状态跟着进程走:脚本重启就全没了。要跨进程保留,换成SqliteSaver或PostgresSaver(第 11 章讲过),API 一模一样,只有构造方式不同。
6.2. StateSnapshot 里有什么 #
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
# 先正常跑一次,才有东西可查
cfg = {"configurable": {"thread_id": "t1"}}
# 返回值这里不用,只为了在 t1 里留下检查点
graph.invoke({"text": " 你好 "}, cfg)
# 取当前快照
snap = graph.get_state(cfg)
# 它是个 NamedTuple
print("类型:", type(snap).__name__)
# _fields 是 NamedTuple 自带的属性,能列出全部字段名
print("字段:", snap._fields)
# 当前完整状态
print("values:", snap.values)
# 下一步要跑的节点
print("next:", snap.next)
# 这个检查点的元数据
print("metadata:", snap.metadata)
# 注意 config 里比你传进去的多了两个键
print("config:", snap.config)
# 上一个检查点的 config,用来往回追溯
print("parent_config:", snap.parent_config)
# 待执行 / 失败的任务
print("tasks:", snap.tasks)
# 中断信息(第 26 章)
print("interrupts:", snap.interrupts)
# 检查点写入时间
print("created_at:", snap.created_at)类型: StateSnapshot
字段:
('values', 'next', 'config', 'metadata', 'created_at', 'parent_config', 'tasks', 'interrupts')
values:
{'text': '你好', 'log': ['clean', 'count=2']}
next:
()
metadata:
{'source': 'loop', 'step': 2, 'parents': {}}
config:
{'configurable': {'thread_id': 't1', 'checkpoint_ns': '', 'checkpoint_id': '1f1a868d-f8ef-6165-8002-df7d9aa3e31c'}}
parent_config:
{'configurable': {'thread_id': 't1', 'checkpoint_ns': '', 'checkpoint_id': '1f1a868d-f8ed-6c4d-8001-74b452fd7c80'}}
tasks:
()
interrupts:
()
created_at: 2026-09-04T13:59:40.594817+00:00最常用的三个:
| 字段 | 含义 | 怎么读 |
|---|---|---|
values |
当前完整状态 | 就是 invoke 会返回的东西 |
next |
下一步要跑哪些节点 | 空元组 () = 已经跑完 |
tasks |
待执行 / 失败的任务 | 里面有 error 字段,排障关键 |
next 是判断「这个 thread 处于什么状态」的核心信号:
next == ()→ 跑完了next == ('某节点',)→ 停在这个节点之前(可能是中断,也可能是失败)
另外注意 config 里多了两个你没传的键:checkpoint_ns(命名空间,子图会用到)和 checkpoint_id。后者是「时间旅行」的钥匙(§6.5)。parent_config 指向上一个检查点,一路顺着它就能手工回溯整条链,不过通常直接用 get_state_history 更省事。
一个静默陷阱:查一个从没跑过的 thread
next == () 到底是不是「跑完了」?看这个对照实验:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
# 一个真正跑完的 thread
done_cfg = {"configurable": {"thread_id": "done"}}
# 正常跑一次,让它有检查点
graph.invoke({"text": " zz "}, done_cfg)
# 一个压根没跑过的 thread 名
never_cfg = {"configurable": {"thread_id": "never"}}
# 两个都查一下,对比字段
for name, c in (("跑完的", done_cfg), ("没跑过的", never_cfg)):
# 查不存在的 thread 不会报错
s = graph.get_state(c)
# 重点对比 next 和 metadata 这两列
print(f" {name}: next={s.next} values={s.values} metadata={s.metadata} created_at={s.created_at}") 跑完的: next=() values={'text': 'zz', 'log': ['clean', 'count=2']} metadata={'source': 'loop', 'step': 2, 'parents': {}} created_at=2026-08-13T15:05:04.426304+00:00
没跑过的: next=() values={} metadata=None created_at=None查一个不存在的 thread 不报错,返回一个空快照,next 也是 ()。 于是:
next == ()有两个截然不同的含义:「跑完了」和「压根没跑过」。 想区分,看metadata,没跑过的是None(created_at也是None)。
这个坑在真实系统里很常见:前端传来一个拼错的 thread_id,后端 get_state 一看 next == () 就回「任务已完成」,而实际上什么都没发生过。判断的正确写法是:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
# 判断这个 thread 到底有没有跑过
snap = graph.get_state({"configurable": {"thread_id": "never"}})
# metadata 为 None 就说明这个 thread 没有任何检查点
if snap.metadata is None:
# 这一支必须放在最前面判断
print(" 这个 thread 从未运行过")
# 有 metadata 且 next 为空,才是真的跑完了
elif not snap.next:
# 到这里才能安全地回「已完成」
print(" 已跑完")
# 否则就是中途停下了(中断或失败)
else:
# 具体是中断还是失败,看 tasks[].error(§7.4)
print(" 停在:", snap.next) 这个 thread 从未运行过6.3. get_state_history:完整执行轨迹 #
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
cfg = {"configurable": {"thread_id": "t1"}}
graph.invoke({"text": " 你好 "}, cfg)
# 遍历这个 thread 的所有检查点
for i, s in enumerate(graph.get_state_history(cfg)):
# step 和 source 都在 metadata 里;用 .get 避免 metadata 为 None 时崩掉
print(f" [{i}] next={s.next} values={s.values} step={s.metadata.get('step')} source={s.metadata.get('source')}") [0] next=() values={'text': '你好', 'log': ['clean', 'count=2']} step=2 source=loop
[1] next=('count',) values={'text': '你好', 'log': ['clean']} step=1 source=loop
[2] next=('clean',) values={'text': ' 你好 ', 'log': []} step=0 source=loop
[3] next=('__start__',) values={'log': []} step=-1 source=input四个要点:
- 顺序是从新到旧(
[0]是最终状态)。想按时间正序看就reversed(list(...))。 step=-1是输入刚写进来、还没开始跑的时刻,此时values里只有 reducer 推出来的零值{'log': []},连text都还没写进去。source区分了这个检查点的来源:input是用户输入写进来的,loop是图正常执行产生的,还有第三种update(§7.3 会见到)。- 每一步都是一个完整的
StateSnapshot,都能拿来重跑(§6.5)。
对照 §4.6 的 checkpoints 模式:那边推的 4 条事件,就是这里的 4 个检查点。 区别在于时机:
stream只能在运行时看,历史是跑完之后(甚至几天以后)还能查。 前者是监控,后者是审计。
6.4. update_state:手动改状态 #
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
cfg = {"configurable": {"thread_id": "t1"}}
graph.invoke({"text": " 你好 "}, cfg)
# 往 log 里插一条人工记录
graph.update_state(cfg, {"log": ["手动插入"]})
# 查看效果
print("改后:", graph.get_state(cfg).values)改后: {'text': '你好', 'log': ['clean', 'count=2', '手动插入']}注意它走的是 reducer:log 有 operator.add,所以是追加而不是覆盖。这一点和「节点返回值走 reducer」完全一致:在 LangGraph 眼里,update_state 就是「一次由人发起的状态更新」,没有特权。
想覆盖就得用没挂 reducer 的字段:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
cfg = {"configurable": {"thread_id": "t1"}}
graph.invoke({"text": " 你好 "}, cfg)
graph.update_state(cfg, {"log": ["手动插入"]})
# text 没挂 reducer,所以这次是真覆盖
graph.update_state(cfg, {"text": "覆盖了"})
# log 保持原样,只有 text 变了
print("改后:", graph.get_state(cfg).values)改后: {'text': '覆盖了', 'log': ['clean', 'count=2', '手动插入']}它还有个容易被忽略的返回值:新检查点的 config:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
cfg = {"configurable": {"thread_id": "t1"}}
graph.invoke({"text": " 你好 "}, cfg)
# update_state 会产生一个新检查点,并把它的 config 返回给你
ret = graph.update_state(cfg, {"log": ["再插一条"]})
# 返回值里的 checkpoint_id 就是「改完之后那个时刻」的坐标
print("update_state 返回:", ret)update_state 返回: {'configurable': {'thread_id': 't1', 'checkpoint_ns': '', 'checkpoint_id': '1f197282-dc32-63b7-8003-48799206844b'}}这个返回值在 §6.5 做「分叉」时是关键:它就是「改完之后那个时刻」的坐标。
还可以指定 as_node,效果是「假装某个节点刚跑完」:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
# 换一个干净的 thread
cfg2 = {"configurable": {"thread_id": "t2"}}
# 先完整跑一次
graph.invoke({"text": " x "}, cfg2)
# 关键:as_node="clean" 表示「就当 clean 刚跑完并写了这些」
graph.update_state(cfg2, {"text": "改过了"}, as_node="clean")
# 取快照,重点看 next
s = graph.get_state(cfg2)
# next 会变成 clean 的下游节点
print("改后:", s.values, "next=", s.next)改后: {'text': '改过了', 'log': ['clean', 'count=1']} next= ('count',)next 变回了 ('count',),图被「拨回」到 clean 刚跑完的位置,再 invoke(None, cfg2) 就会从 count 接着跑:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
cfg2 = {"configurable": {"thread_id": "t2"}}
graph.invoke({"text": " x "}, cfg2)
graph.update_state(cfg2, {"text": "改过了"}, as_node="clean")
# 第一个参数传 None:不给新输入,从检查点继续
print("续跑:", graph.invoke(None, cfg2))续跑: {'text': '改过了', 'log': ['clean', 'count=1', 'count=3']}留意 log 里现在有两条 count:count=1 是第一次跑留下的,count=3 是这次重跑 count 节点产生的('改过了' 三个字)。这是「拨回去重跑」的必然代价:已经落盘的历史不会被抹掉,重跑的写入是追加上去的。所以用 as_node 拨回时要想清楚:如果字段挂着累加型 reducer,重复记录是预期行为,别当成 bug;真不想要重复,就得把那个字段设计成覆盖型。
as_node 写了一个不存在的节点名会当场报错,不会静默:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
# 准备一个 thread
cfg2b = {"configurable": {"thread_id": "t2b"}}
# 先跑一次
graph.invoke({"text": " x "}, cfg2b)
try:
# 故意写一个没注册过的节点名
graph.update_state(cfg2b, {"text": "y"}, as_node="不存在")
except Exception as e:
# 这是「大声失败」,不会静默
print(" 报错:", type(e).__name__, str(e)) 报错: InvalidUpdateError Node 不存在 does not exist这是人工干预流程的底层机制(第 26 章的 HITL 用的就是它)。
6.5. 时间旅行:回到任意检查点重跑 #
get_state_history 里每个快照的 config 都带着 checkpoint_id。把它当 config 传给 invoke,就会从那个时刻重新开始:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from rich import print
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
# 用一个新 thread 做实验
tt = {"configurable": {"thread_id": "tt"}}
# 先完整跑一次,攒出一条轨迹
graph.invoke({"text": " 时间旅行 "}, tt)
# 把历史全取出来(生成器要转成 list 才能多次筛选)
history = list(graph.get_state_history(tt))
# 找到「clean 跑完、count 还没跑」的那个时刻
target = [h for h in history if h.next == ("count",)][0]
# 确认找对了:log 里只有 clean
print(
"目标检查点 next:",
target.next,
" values:",
target.values,
" config:",
target.config,
)
# 关键:把那个快照的 config 传进去,第一个参数给 None
replay = graph.invoke(None, target.config)
# count 又跑了一次,结果和原来一样
print("重放结果:", replay)目标检查点 next:
('count',)
values:
{'text': '时间旅行', 'log': ['clean']}
config:
{
'configurable': {
'thread_id': 'tt',
'checkpoint_ns': '',
'checkpoint_id': '1f1a86c4-a728-6d33-8001-df0679260075'
}
}
重放结果:
{'text': '时间旅行', 'log': ['clean', 'count=4']}注意第一个参数传的是 None,意思是「不给新输入,从检查点存的状态继续」。能这样,是因为 target.config 里带着 checkpoint_id,LangGraph 优先按这个 id 定位起点,而不是用「这个 thread 的最新状态」。
真正有用的地方:模型在第 3 步给了个烂结果,你不必从头重跑前两步(省钱省时间),直接回到第 3 步之前改改参数再来一次。
更常用的组合:回到过去 + 改状态 + 重放
光重放意义有限,同样的输入多半还是同样的结果。真正有用的是回到某个时刻,改点东西,再往下跑。做法是把 update_state 作用在旧检查点的 config 上:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
graph = build().compile(checkpointer=InMemorySaver())
tt = {"configurable": {"thread_id": "tt"}}
graph.invoke({"text": " 时间旅行 "}, tt)
history = list(graph.get_state_history(tt))
# 找到「什么都还没跑」的那个时刻
target0 = [h for h in history if h.next == ("clean",)][0]
# 在这个旧检查点上改状态,返回的是「改完之后」的新坐标
forked = graph.update_state(target0.config, {"text": " 换个输入 "})
# 这个 checkpoint_id 和原轨迹上的都不一样,是新分支的起点
print("fork 出来的新 config:", forked)
# 从这个新坐标往下跑
print("从 fork 点重放:", graph.invoke(None, forked))fork 出来的新 config: {'configurable': {'thread_id': 'tt', 'checkpoint_ns': '', 'checkpoint_id': '1f197282-dc48-61e2-8001-7d81bb360d34'}}
从 fork 点重放: {'text': '换个输入', 'log': ['clean', 'count=4']}回到过去、改掉输入、重新执行,而且原来那条轨迹还完整留在历史里。 这就是「时间旅行」这个说法的完整含义:不只是重放,还能分叉。
| 想干什么 | 怎么写 |
|---|---|
| 从头再跑一次 | invoke(新输入, {"configurable": {"thread_id": 新 id}}) |
| 从某检查点原样重放 | invoke(None, 那个快照的 config) |
| 从某检查点改点东西再跑 | 新 cfg = update_state(旧快照的 config, 补丁) → invoke(None, 新 cfg) |
| 从最新状态继续(失败 / 中断后) | invoke(None, {"configurable": {"thread_id": 原 id}}) |
7. 节点失败时会发生什么 #
这是本章最该记住的一节。换一张订单图来演示,它就是 §8 实战要用的那张:三个节点串行(validate → price → notify),第一个节点遇到非法金额会抛异常。
重点只在 validate 那几行:它是所有实验的触发点。注意编译时一定要带 checkpointer,不带的话炸完就查不了现场。
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
TAX_RATE = 0.06
class OrderState(TypedDict):
order_id: str
amount: float
tax: float
notice: str
log: Annotated[list[str], operator.add]
def validate(state: OrderState) -> dict:
if state["amount"] <= 0:
raise ValueError(f"金额非法: {state['amount']}")
return {"log": [f"validate ok: {state['amount']}"]}
def price(state: OrderState) -> dict:
tax = round(state["amount"] * TAX_RATE, 2)
return {"tax": tax, "log": [f"price tax={tax}"]}
def notify(state: OrderState) -> dict:
total = state["amount"] + state["tax"]
return {
"notice": f"订单 {state['order_id']} 已确认,含税合计 {total:.2f} 元。",
"log": ["notify sent"],
}
def build_graph(checkpointer=None):
builder = StateGraph(OrderState)
builder.add_node("validate", validate)
builder.add_node("price", price)
builder.add_node("notify", notify)
builder.add_edge(START, "validate")
builder.add_edge("validate", "price")
builder.add_edge("price", "notify")
builder.add_edge("notify", END)
return builder.compile(checkpointer=checkpointer, name="order-graph")
def new_order(order_id: str, amount: float) -> OrderState:
return {"order_id": order_id, "amount": amount, "tax": 0.0, "notice": "", "log": []}
# 编译一个带 checkpointer 的版本,后面几段都这么编译
graph = build_graph(checkpointer=InMemorySaver())
# 三个节点、五个状态字段,确认图搭好了
print("节点:", [n for n in graph.get_graph().nodes if not n.startswith("__")])
print("状态字段:", list(OrderState.__annotations__))节点: ['validate', 'price', 'notify']
状态字段: ['order_id', 'amount', 'tax', 'notice', 'log']上面那段里,真正的主角是 validate 这四行:
def validate(state: OrderState) -> dict:
if state["amount"] <= 0: # 金额不合法就抛异常
raise ValueError(f"金额非法: {state['amount']}") # 普通 ValueError,没有任何 LangGraph 专用写法
return {"log": [f"validate ok: {state['amount']}"]} # 合法就记一笔日志为什么专门用「抛异常」来讲?因为真实系统里节点失败是常态:模型超时、接口 429、数据库连接断开、上游返回的 JSON 解析失败。这一节要回答的是:出事之后,你手上还剩什么。
7.1. 异常直接抛出,不会被吞 #
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
TAX_RATE = 0.06
class OrderState(TypedDict):
order_id: str
amount: float
tax: float
notice: str
log: Annotated[list[str], operator.add]
def validate(state: OrderState) -> dict:
if state["amount"] <= 0:
raise ValueError(f"金额非法: {state['amount']}")
return {"log": [f"validate ok: {state['amount']}"]}
def price(state: OrderState) -> dict:
tax = round(state["amount"] * TAX_RATE, 2)
return {"tax": tax, "log": [f"price tax={tax}"]}
def notify(state: OrderState) -> dict:
total = state["amount"] + state["tax"]
return {
"notice": f"订单 {state['order_id']} 已确认,含税合计 {total:.2f} 元。",
"log": ["notify sent"],
}
def build_graph(checkpointer=None):
builder = StateGraph(OrderState)
builder.add_node("validate", validate)
builder.add_node("price", price)
builder.add_node("notify", notify)
builder.add_edge(START, "validate")
builder.add_edge("validate", "price")
builder.add_edge("price", "notify")
builder.add_edge("notify", END)
return builder.compile(checkpointer=checkpointer, name="order-graph")
def new_order(order_id: str, amount: float) -> OrderState:
return {"order_id": order_id, "amount": amount, "tax": 0.0, "notice": "", "log": []}
graph = build_graph(checkpointer=InMemorySaver())
# 用一个专门的 thread 记录这次失败
bad_cfg = {"configurable": {"thread_id": "bad-1"}}
try:
# 金额是 -5,validate 一定会抛
graph.invoke(new_order("A002", -5.0), bad_cfg)
except Exception as e:
# 打印异常类型和消息,看有没有被包装
print("报错:", type(e).__name__, e)报错: ValueError 金额非法: -5.0LangGraph 不做任何异常包装,你在节点里抛什么,外面就接到什么:类型是原本的 ValueError,消息一个字没改,traceback 也直通节点内部那一行。
这和第 20 章那四种「静默失败」形成鲜明对比:那些是结构问题(图默默跑偏),这里是业务问题(大声报错)。这个对比可以当成一条排查路线:
| 症状 | 大概率属于 | 先看哪里 |
|---|---|---|
| 抛异常、有 traceback | 业务问题 | 节点内部的逻辑和输入数据 |
| 不报错但结果不对 / 少了字段 | 结构问题 | stream_mode="updates"、图的边(第 20 章) |
7.2. 状态停在失败之前 #
关键在这里:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
TAX_RATE = 0.06
class OrderState(TypedDict):
order_id: str
amount: float
tax: float
notice: str
log: Annotated[list[str], operator.add]
def validate(state: OrderState) -> dict:
if state["amount"] <= 0:
raise ValueError(f"金额非法: {state['amount']}")
return {"log": [f"validate ok: {state['amount']}"]}
def price(state: OrderState) -> dict:
tax = round(state["amount"] * TAX_RATE, 2)
return {"tax": tax, "log": [f"price tax={tax}"]}
def notify(state: OrderState) -> dict:
total = state["amount"] + state["tax"]
return {
"notice": f"订单 {state['order_id']} 已确认,含税合计 {total:.2f} 元。",
"log": ["notify sent"],
}
def build_graph(checkpointer=None):
builder = StateGraph(OrderState)
builder.add_node("validate", validate)
builder.add_node("price", price)
builder.add_node("notify", notify)
builder.add_edge(START, "validate")
builder.add_edge("validate", "price")
builder.add_edge("price", "notify")
builder.add_edge("notify", END)
return builder.compile(checkpointer=checkpointer, name="order-graph")
def new_order(order_id: str, amount: float) -> OrderState:
return {"order_id": order_id, "amount": amount, "tax": 0.0, "notice": "", "log": []}
graph = build_graph(checkpointer=InMemorySaver())
# 先制造一次失败,才有现场可查
bad_cfg = {"configurable": {"thread_id": "bad-1"}}
try:
graph.invoke(new_order("A002", -5.0), bad_cfg)
except ValueError:
pass
# 图已经炸了,但状态还在,查一下
s = graph.get_state(bad_cfg)
# 失败时刻的完整状态
print("values:", s.values)
# 下一步该跑谁
print("next:", s.next)
# tasks 里挂着这次没跑成的任务,带异常信息
print("失败任务:", [(t.name, t.error) for t in s.tasks])values: {'order_id': 'A002', 'amount': -5.0, 'tax': 0.0, 'notice': '', 'log': []}
next: ('validate',)
失败任务: [('validate', "ValueError('金额非法: -5.0')")]三件事同时成立:
values是失败节点执行前的状态:log是空的,失败节点的写入没有落盘,不会留下半成品。next指向失败的那个节点,说明「它还欠着没跑」。tasks[0].error保存了异常信息(字符串形式),事后能查到底怎么挂的。
tasks 里的每个任务还有别的字段可用:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
TAX_RATE = 0.06
class OrderState(TypedDict):
order_id: str
amount: float
tax: float
notice: str
log: Annotated[list[str], operator.add]
def validate(state: OrderState) -> dict:
if state["amount"] <= 0:
raise ValueError(f"金额非法: {state['amount']}")
return {"log": [f"validate ok: {state['amount']}"]}
def price(state: OrderState) -> dict:
tax = round(state["amount"] * TAX_RATE, 2)
return {"tax": tax, "log": [f"price tax={tax}"]}
def notify(state: OrderState) -> dict:
total = state["amount"] + state["tax"]
return {
"notice": f"订单 {state['order_id']} 已确认,含税合计 {total:.2f} 元。",
"log": ["notify sent"],
}
def build_graph(checkpointer=None):
builder = StateGraph(OrderState)
builder.add_node("validate", validate)
builder.add_node("price", price)
builder.add_node("notify", notify)
builder.add_edge(START, "validate")
builder.add_edge("validate", "price")
builder.add_edge("price", "notify")
builder.add_edge("notify", END)
return builder.compile(checkpointer=checkpointer, name="order-graph")
def new_order(order_id: str, amount: float) -> OrderState:
return {"order_id": order_id, "amount": amount, "tax": 0.0, "notice": "", "log": []}
graph = build_graph(checkpointer=InMemorySaver())
# 先制造一次失败,才有现场可查
bad_cfg = {"configurable": {"thread_id": "bad-1"}}
try:
graph.invoke(new_order("A002", -5.0), bad_cfg)
except ValueError:
pass
s = graph.get_state(bad_cfg)
# 看一眼任务对象都有哪些字段
print("task 全字段:", s.tasks[0]._fields)task 全字段: ('id', 'name', 'path', 'error', 'interrupts', 'state', 'result')排障时最常用的就是 name(哪个节点)和 error(为什么)。
「停在失败之前」不等于「什么都没保存」
上面 log 是空的,容易让人误会成「失败会丢掉整轮的所有写入」。换成第二个节点失败就清楚了:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
TAX_RATE = 0.06
class OrderState(TypedDict):
order_id: str
amount: float
tax: float
notice: str
log: Annotated[list[str], operator.add]
def validate(state: OrderState) -> dict:
if state["amount"] <= 0:
raise ValueError(f"金额非法: {state['amount']}")
return {"log": [f"validate ok: {state['amount']}"]}
def price(state: OrderState) -> dict:
tax = round(state["amount"] * TAX_RATE, 2)
return {"tax": tax, "log": [f"price tax={tax}"]}
def notify(state: OrderState) -> dict:
total = state["amount"] + state["tax"]
return {
"notice": f"订单 {state['order_id']} 已确认,含税合计 {total:.2f} 元。",
"log": ["notify sent"],
}
def build_graph(checkpointer=None):
builder = StateGraph(OrderState)
builder.add_node("validate", validate)
builder.add_node("price", price)
builder.add_node("notify", notify)
builder.add_edge(START, "validate")
builder.add_edge("validate", "price")
builder.add_edge("price", "notify")
builder.add_edge("notify", END)
return builder.compile(checkpointer=checkpointer, name="order-graph")
def new_order(order_id: str, amount: float) -> OrderState:
return {"order_id": order_id, "amount": amount, "tax": 0.0, "notice": "", "log": []}
# 一个必炸的 price 节点
def boom(state: OrderState) -> dict:
# 无条件抛异常
raise RuntimeError("price 挂了")
# 重搭一张图:validate 正常,price 必炸
builder = StateGraph(OrderState)
# 第一个节点用原来那个正常的
builder.add_node("validate", validate)
# 第二个节点换成必炸的 boom
builder.add_node("price", boom)
# 入口
builder.add_edge(START, "validate")
# 串行
builder.add_edge("validate", "price")
# 出口
builder.add_edge("price", END)
# 必须带 checkpointer,否则炸完就查不了现场
graph = builder.compile(checkpointer=InMemorySaver())
# 这次金额是合法的,所以 validate 会成功,price 才炸
mid = {"configurable": {"thread_id": "mid"}}
try:
# 跑到第二个节点才失败
graph.invoke(new_order("A009", 10.0), mid)
except Exception as e:
# 异常仍然原样抛出
print(" 报错:", type(e).__name__, e)
# 关键:查 validate 的写入还在不在
sm = graph.get_state(mid)
# log 里应该有 validate 那条
print(" values:", sm.values)
# next 指向真正失败的 price
print(" next:", sm.next) 报错: RuntimeError price 挂了
values: {'order_id': 'A009', 'amount': 10.0, 'tax': 0.0, 'notice': '', 'log': ['validate ok: 10.0']}
next: ('price',)validate 的写入完整保留,next 指向真正失败的 price。 所以准确的说法是:
成功的节点,写入照常落盘;失败的节点,一个字都不留。 状态精确地停在「失败节点开跑之前」那一刻。这正是断点续跑能安全工作的前提。
这就是 checkpointer 的价值:崩了不等于全丢,现场还在。 也正因为「失败节点不留半成品」,续跑时它是干净地重新执行一次,而不是接着一个不完整的中间态往下走。
7.3. 修完数据从断点续跑 #
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
TAX_RATE = 0.06
class OrderState(TypedDict):
order_id: str
amount: float
tax: float
notice: str
log: Annotated[list[str], operator.add]
def validate(state: OrderState) -> dict:
if state["amount"] <= 0:
raise ValueError(f"金额非法: {state['amount']}")
return {"log": [f"validate ok: {state['amount']}"]}
def price(state: OrderState) -> dict:
tax = round(state["amount"] * TAX_RATE, 2)
return {"tax": tax, "log": [f"price tax={tax}"]}
def notify(state: OrderState) -> dict:
total = state["amount"] + state["tax"]
return {
"notice": f"订单 {state['order_id']} 已确认,含税合计 {total:.2f} 元。",
"log": ["notify sent"],
}
def build_graph(checkpointer=None):
builder = StateGraph(OrderState)
builder.add_node("validate", validate)
builder.add_node("price", price)
builder.add_node("notify", notify)
builder.add_edge(START, "validate")
builder.add_edge("validate", "price")
builder.add_edge("price", "notify")
builder.add_edge("notify", END)
return builder.compile(checkpointer=checkpointer, name="order-graph")
def new_order(order_id: str, amount: float) -> OrderState:
return {"order_id": order_id, "amount": amount, "tax": 0.0, "notice": "", "log": []}
graph = build_graph(checkpointer=InMemorySaver())
# 先制造一次失败,才有现场可查
bad_cfg = {"configurable": {"thread_id": "bad-1"}}
try:
graph.invoke(new_order("A002", -5.0), bad_cfg)
except ValueError:
pass
# 第一步:打补丁,把非法金额改成合法的
graph.update_state(bad_cfg, {"amount": 50.0}) # 打补丁
# 第二步:第一个参数传 None,表示「不给新输入,从断点继续」
out = graph.invoke(None, bad_cfg) # 从断点继续
# 看最终产物
print("续跑结果:", out["notice"])
# 看日志完整性
print("log:", out["log"])续跑结果: 订单 A002 已确认,含税合计 53.00 元。
log: ['validate ok: 50.0', 'price tax=3.0', 'notify sent']log 里没有任何重复,三条恰好对应三个节点。因为失败的 validate 从来没成功写入过,续跑时它是第一次真正执行。对比 §6.4 用 as_node 拨回去导致 count 记录两次,区别一目了然:从失败点续跑不会重复,从已完成的位置拨回去重跑才会。
也可以直接传部分输入重新 invoke,效果类似:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
TAX_RATE = 0.06
class OrderState(TypedDict):
order_id: str
amount: float
tax: float
notice: str
log: Annotated[list[str], operator.add]
def validate(state: OrderState) -> dict:
if state["amount"] <= 0:
raise ValueError(f"金额非法: {state['amount']}")
return {"log": [f"validate ok: {state['amount']}"]}
def price(state: OrderState) -> dict:
tax = round(state["amount"] * TAX_RATE, 2)
return {"tax": tax, "log": [f"price tax={tax}"]}
def notify(state: OrderState) -> dict:
total = state["amount"] + state["tax"]
return {
"notice": f"订单 {state['order_id']} 已确认,含税合计 {total:.2f} 元。",
"log": ["notify sent"],
}
def build_graph(checkpointer=None):
builder = StateGraph(OrderState)
builder.add_node("validate", validate)
builder.add_node("price", price)
builder.add_node("notify", notify)
builder.add_edge(START, "validate")
builder.add_edge("validate", "price")
builder.add_edge("price", "notify")
builder.add_edge("notify", END)
return builder.compile(checkpointer=checkpointer, name="order-graph")
def new_order(order_id: str, amount: float) -> OrderState:
return {"order_id": order_id, "amount": amount, "tax": 0.0, "notice": "", "log": []}
graph = build_graph(checkpointer=InMemorySaver())
# 另开一个 thread,同样先失败
another = {"configurable": {"thread_id": "bad-2"}}
try:
# 金额 -1,照样炸在 validate
graph.invoke(new_order("A003", -1.0), another)
except ValueError as e:
# 确认失败了,现场已经存下来
print("第一次失败:", e)
# 只传要改的字段,其余沿用已存下来的状态
print("续跑:", graph.invoke({"amount": 80.0}, another)["notice"])第一次失败: 金额非法: -1.0
续跑: 订单 A003 已确认,含税合计 84.80 元。注意 order_id 完全没传,却出现在结果里(84.80 = 80 × 1.06)。这份输入是「补丁」,不是「新的完整输入」,没提到的字段沿用检查点里存着的值。
两种方式的区别:
update_state + invoke(None) |
invoke({部分字段}) |
|
|---|---|---|
| 语义 | 改状态,然后继续 | 补一次新输入,然后继续 |
能否指定 as_node |
能 | 不能 |
历史里的 source |
多一条 update 记录 |
只有 input 和 loop |
| 适合 | 精确控制从哪继续、要留审计痕迹 | 简单场景 |
如果这次修改需要事后说得清(谁改的、改了什么),选 update_state,因为它会在历史里留下一条 source=update。
看一眼历史,update_state 那一步会被清楚地标出来:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
TAX_RATE = 0.06
class OrderState(TypedDict):
order_id: str
amount: float
tax: float
notice: str
log: Annotated[list[str], operator.add]
def validate(state: OrderState) -> dict:
if state["amount"] <= 0:
raise ValueError(f"金额非法: {state['amount']}")
return {"log": [f"validate ok: {state['amount']}"]}
def price(state: OrderState) -> dict:
tax = round(state["amount"] * TAX_RATE, 2)
return {"tax": tax, "log": [f"price tax={tax}"]}
def notify(state: OrderState) -> dict:
total = state["amount"] + state["tax"]
return {
"notice": f"订单 {state['order_id']} 已确认,含税合计 {total:.2f} 元。",
"log": ["notify sent"],
}
def build_graph(checkpointer=None):
builder = StateGraph(OrderState)
builder.add_node("validate", validate)
builder.add_node("price", price)
builder.add_node("notify", notify)
builder.add_edge(START, "validate")
builder.add_edge("validate", "price")
builder.add_edge("price", "notify")
builder.add_edge("notify", END)
return builder.compile(checkpointer=checkpointer, name="order-graph")
def new_order(order_id: str, amount: float) -> OrderState:
return {"order_id": order_id, "amount": amount, "tax": 0.0, "notice": "", "log": []}
graph = build_graph(checkpointer=InMemorySaver())
# 先制造一次失败,才有现场可查
bad_cfg = {"configurable": {"thread_id": "bad-1"}}
try:
graph.invoke(new_order("A002", -5.0), bad_cfg)
except ValueError:
pass
graph.update_state(bad_cfg, {"amount": 50.0})
graph.invoke(None, bad_cfg)
# 遍历失败又续跑成功的那个 thread
for snap in graph.get_state_history(bad_cfg):
# 用 :>2 右对齐 step,方便竖着看
print(f" step={snap.metadata.get('step'):>2} next={snap.next} source={snap.metadata.get('source')}") step= 4 next=() source=loop
step= 3 next=('notify',) source=loop
step= 2 next=('price',) source=loop
step= 1 next=('validate',) source=update
step= 0 next=('validate',) source=loop
step=-1 next=('__start__',) source=input从下往上读,整个故事都在里面:输入进来(step=-1)→ 图开始跑并停在 validate(step=0,就是失败那次)→ 人工打了补丁(step=1,source=update)→ 三个节点依次跑完(step=2/3/4)。
source 有三种值,正好对应三种状态来源:
source |
含义 |
|---|---|
input |
用户传进来的输入 |
loop |
图正常执行产生的 |
update |
人工 update_state 改的 |
审计的时候,这一列能直接回答「这个值是系统算的还是人改的」。金额、审批结论这类敏感字段一旦支持人工修改,这条记录就是唯一能追溯的凭据。
7.4. 顺便一提:interrupt_before #
失败是被动停下,interrupt_before 则是主动停下:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
# InMemorySaver:中断要能恢复,就必须有地方存状态
from langgraph.checkpoint.memory import InMemorySaver
# 编译时声明「在 count 之前停一下」
graph = build().compile(checkpointer=InMemorySaver(), interrupt_before=["count"])
# 中断要能恢复,所以必须有 thread_id
cfg4 = {"configurable": {"thread_id": "t4"}}
# 第一次 invoke 会在 count 之前返回,注意 log 里只有 clean
print("第一次 invoke 返回:", graph.invoke({"text": " 停一下 "}, cfg4))
# 取快照,看停在哪、任务有没有报错
st = graph.get_state(cfg4)
# 关键区别:这里的 error 是 None,说明是暂停而不是失败
print("next:", st.next, "| tasks:", [(t.name, t.error) for t in st.tasks])
# 恢复:和失败续跑一模一样的姿势
print("恢复执行:", graph.invoke(None, cfg4))第一次 invoke 返回: {'text': '停一下', 'log': ['clean']}
next: ('count',) | tasks: [('count', None)]
恢复执行: {'text': '停一下', 'log': ['clean', 'count=3']}恢复的姿势和失败续跑完全一样:invoke(None, config)。 这不是巧合:在 LangGraph 眼里,「暂停」和「失败」都只是「有节点还没跑完」的一种表现,恢复机制是同一套。
那怎么区分「它是暂停了还是挂了」?看 tasks[].error:
| 情况 | next |
tasks[].error |
|---|---|---|
| 跑完了 | () |
无任务 |
| 主动中断 | ('count',) |
None |
| 节点失败 | ('validate',) |
异常字符串 |
| 从没跑过 | () |
无任务(且 metadata is None,见 §6.2) |
interrupt_after 是同一件事的另一种写法:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
text: str
log: Annotated[list[str], operator.add]
def clean(state: S) -> dict:
return {"text": state["text"].strip(), "log": ["clean"]}
def count(state: S) -> dict:
return {"log": [f"count={len(state['text'])}"]}
def build():
builder = StateGraph(S)
builder.add_node("clean", clean)
builder.add_node("count", count)
builder.add_edge(START, "clean")
builder.add_edge("clean", "count")
builder.add_edge("count", END)
return builder
from langgraph.checkpoint.memory import InMemorySaver
# 换成「clean 跑完之后停」
graph = build().compile(checkpointer=InMemorySaver(), interrupt_after=["clean"])
# 换一个 thread
cfg5 = {"configurable": {"thread_id": "t5"}}
# 返回结果和 interrupt_before=["count"] 那次一模一样
print("第一次:", graph.invoke({"text": " 停一下 "}, cfg5))
# next 和 interrupt_before=["count"] 的结果完全一样
print("next:", graph.get_state(cfg5).next)
# 恢复姿势也一样
print("恢复:", graph.invoke(None, cfg5))第一次: {'text': '停一下', 'log': ['clean']}
next: ('count',)
恢复: {'text': '停一下', 'log': ['clean', 'count=3']}在这张两节点串行图上,interrupt_after=["clean"] 和 interrupt_before=["count"] 的效果完全相同,因为「clean 之后」和「count 之前」是同一个时刻。分支图上两者就不同了:clean 之后可能有好几个候选下游,interrupt_after 关心的是「谁刚跑完」,interrupt_before 关心的是「谁将要跑」。第 26 章会在这个基础上讲人机在环。
8. 实战:一张可运行、可观测、可恢复的状态图 #
8.1. 设计 #
前面的能力都是零散验证的,这里把它们收成一个像样的模块。核心思路是把「怎么搭图」和「怎么用图」分开:
build_graph(checkpointer) ← 图纸 + 编译,可编译出不同版本
↓
run(...) 要结果 → invoke
watch(...) 要过程 → stream
inspect(...) 要现场 → get_state + get_state_history
recover(...) 要续跑 → update_state + invoke(None)这四个函数不是随便凑的,它们和本章四节一一对应,也和你在真实项目里的四种诉求一一对应:
| 函数 | 对应本章 | 真实场景里谁会用 |
|---|---|---|
run |
§3 | 业务代码,只关心结果 |
watch |
§4 | 前端 / 调试时看进度 |
inspect |
§6 | 排障、客服查单、审计 |
recover |
§7 | 出事之后的人工兜底 |
业务是三步串行的订单处理:validate → price → notify,其中 validate 会对非法金额抛异常,正好用来演示恢复。选订单场景,是因为它刚好满足两个条件:每一步都是确定性计算(不需要模型,零成本可复现),而且「金额算错了要能查、能修」这个需求非常真实。
8.2. runnable_graph.py #
"""实战:一张可运行、可观测、可恢复的状态图"""
# 让类型标注延迟求值,函数签名里可以直接写 OrderState 这类前向引用
from __future__ import annotations
# operator.add 用作 log 字段的累加 reducer
import operator
# Annotated 挂 reducer,TypedDict 定义状态
from typing import Annotated, TypedDict
# 内存版 checkpointer;换 SqliteSaver / PostgresSaver 只需改这一行
from langgraph.checkpoint.memory import InMemorySaver
# 图的三件套
from langgraph.graph import END, START, StateGraph
# 税率常量,抽出来避免散落在节点里
TAX_RATE = 0.06
# 订单状态:四个普通字段 + 一个累加型日志
class OrderState(TypedDict):
# 订单号,全程只读
order_id: str
# 金额,可能被人工修正(§7.3 就是改它)
amount: float
# 税额,由 price 节点写
tax: float
# 通知文案,由 notify 节点写
notice: str
# 执行日志,挂 operator.add,每个节点追加一条
log: Annotated[list[str], operator.add]
def validate(state: OrderState) -> dict:
"""校验金额。非法就抛异常,这正是后面演示断点续跑的入口。"""
# 金额必须为正,否则抛业务异常
if state["amount"] <= 0:
# 异常会原样传到 invoke 的调用方(§7.1)
raise ValueError(f"金额非法: {state['amount']}")
# 校验通过,只写日志,不改业务字段
return {"log": [f"validate ok: {state['amount']}"]}
def price(state: OrderState) -> dict:
"""算税。"""
# 按税率算税并保留两位小数
tax = round(state["amount"] * TAX_RATE, 2)
# 同时写业务字段和日志
return {"tax": tax, "log": [f"price tax={tax}"]}
def notify(state: OrderState) -> dict:
"""生成通知文案。"""
# 含税合计 = 金额 + 税额,读的是 price 刚写进状态的 tax
total = state["amount"] + state["tax"]
# 返回文案和日志
return {
# :.2f 保证金额显示成两位小数
"notice": f"订单 {state['order_id']} 已确认,含税合计 {total:.2f} 元。",
"log": ["notify sent"],
}
def build_graph(checkpointer=None):
"""把「怎么搭」和「怎么编译」分开:同一张图能编译出不同用途的版本。"""
# 图纸
builder = StateGraph(OrderState)
# 注册校验节点
builder.add_node("validate", validate)
# 注册算税节点
builder.add_node("price", price)
# 注册通知节点
builder.add_node("notify", notify)
# 入口
builder.add_edge(START, "validate")
# 串行连接:校验完算税
builder.add_edge("validate", "price")
# 算完税发通知
builder.add_edge("price", "notify")
# 出口:显式连到 END,别依赖默认行为(第 20 章 §6.2)
builder.add_edge("notify", END)
# checkpointer 由调用方决定传不传;name 让 trace 里能认出这张图(§2.6)
return builder.compile(checkpointer=checkpointer, name="order-graph")
def new_order(order_id: str, amount: float) -> OrderState:
"""构造一份完整初始状态,避免节点里 KeyError(§3.2)。"""
# 把每个字段都显式填上,包括本可省略的 log,让初始状态一目了然
return {
"order_id": order_id,
"amount": amount,
# tax 和 notice 由下游节点填,这里给零值占位
"tax": 0.0,
"notice": "",
"log": [],
}
def run(graph, order: OrderState, thread: str) -> OrderState:
"""要结果:invoke。"""
# 一行搞定:传输入 + 指定 thread,返回完整最终状态
return graph.invoke(order, {"configurable": {"thread_id": thread}})
def watch(graph, order: OrderState, thread: str) -> None:
"""要过程:stream。"""
# updates 模式:每条是 {节点名: 该节点的增量}
for step in graph.stream(
order, {"configurable": {"thread_id": thread}}, stream_mode="updates"
):
# 一条里通常只有一个节点,但并行时会有多个,所以用 items() 遍历
for node, update in step.items():
# 打印「谁写了什么」
print(f" [{node}] {update}")
def inspect(graph, thread: str) -> None:
"""要现场:get_state + get_state_history。"""
# 复用同一个 config
cfg = {"configurable": {"thread_id": thread}}
# 当前快照
snap = graph.get_state(cfg)
# next 为空元组时打印「已结束」,比直接打印 () 易读
print(f" 当前 next: {snap.next or '(已结束)'}")
# 直接读日志字段,看跑到哪一步了
print(f" 当前 log : {snap.values['log']}")
# 失败的任务会挂在 tasks 上,带着异常信息
for task in snap.tasks:
# error 为 None 表示只是待执行(中断),不是失败(§7.4)
if task.error:
# 打印是哪个节点、什么异常
print(f" 失败节点: {task.name} -> {task.error}")
# 分隔一下,下面是历史
print(" 历史(新 -> 旧):")
# 遍历全部检查点
for h in graph.get_state_history(cfg):
# 每行打三列:超步号、来源、下一步
print(
# :>2 右对齐 step,负数和正数能竖着对齐
f" step={h.metadata.get('step'):>2} "
# :<6 左对齐 source,方便扫 input / loop / update
f"source={h.metadata.get('source'):<6} next={h.next}"
)
def recover(graph, thread: str, fix: dict) -> OrderState:
"""要续跑:update_state 打补丁,再 invoke(None) 从断点继续。"""
# 复用同一个 config
cfg = {"configurable": {"thread_id": thread}}
# 打补丁:这一步会在历史里留下 source=update 的痕迹(§7.3)
graph.update_state(cfg, fix)
# 第一个参数 None:不给新输入,从断点继续
return graph.invoke(None, cfg)
if __name__ == "__main__":
# 生产版:带持久化,才能 inspect 和 recover
graph = build_graph(checkpointer=InMemorySaver())
# 演示一:最普通的用法
print("=== 1. invoke:只要最终结果 ===")
# 正常订单,一次跑完
out = run(graph, new_order("A001", 100.0), "t-ok")
# 只取关心的字段
print(" ", out["notice"])
# 演示二:要过程
print("\n=== 2. stream:看每个节点做了什么 ===")
# 换一个 thread,避免和上面的状态混在一起
watch(graph, new_order("A002", 200.0), "t-watch")
# 演示三:制造一次失败
print("\n=== 3. 节点失败 ===")
try:
# 负金额,validate 会抛 ValueError
run(graph, new_order("A003", -5.0), "t-bad")
except ValueError as e:
# 异常原样抛出,没有被 LangGraph 包装(§7.1)
print(" 抛出:", e)
# 演示四:查失败现场
print("\n=== 4. get_state:查看失败现场 ===")
# 图已经炸了,但现场还在
inspect(graph, "t-bad")
# 演示五:修完继续
print("\n=== 5. 修数据后从断点续跑 ===")
# 只改 amount,其余字段沿用失败那次存下来的
fixed = recover(graph, "t-bad", {"amount": 50.0})
# 订单号还是失败那次的 A003,说明是接着现场往下走
print(" ", fixed["notice"])
# log 完整且无重复,证明 validate 是第一次真正执行
print(" log:", fixed["log"])
# 演示六:同一份图纸的另一个版本
print("\n=== 6. 无 checkpointer 的版本:能跑,但查不了 ===")
# 同一份 build_graph,这次不传 checkpointer;同一段里有两张图,所以加个后缀区分
graph_plain = build_graph()
# 业务照样能跑
print(" ", run(graph_plain, new_order("A004", 10.0), "t-none")["notice"])
try:
# 但查不了状态
graph_plain.get_state({"configurable": {"thread_id": "t-none"}})
except ValueError as e:
# 报的就是 §6.1 那个 No checkpointer set
print(" get_state ->", e)8.3. 运行结果 #
=== 1. invoke:只要最终结果 ===
订单 A001 已确认,含税合计 106.00 元。
=== 2. stream:看每个节点做了什么 ===
[validate] {'log': ['validate ok: 200.0']}
[price] {'tax': 12.0, 'log': ['price tax=12.0']}
[notify] {'notice': '订单 A002 已确认,含税合计 212.00 元。', 'log': ['notify sent']}
=== 3. 节点失败 ===
抛出: 金额非法: -5.0
=== 4. get_state:查看失败现场 ===
当前 next: ('validate',)
当前 log : []
失败节点: validate -> ValueError('金额非法: -5.0')
历史(新 -> 旧):
step= 0 source=loop next=('validate',)
step=-1 source=input next=('__start__',)
=== 5. 修数据后从断点续跑 ===
订单 A003 已确认,含税合计 53.00 元。
log: ['validate ok: 50.0', 'price tax=3.0', 'notify sent']
=== 6. 无 checkpointer 的版本:能跑,但查不了 ===
订单 A004 已确认,含税合计 10.60 元。
get_state -> No checkpointer set这段输出里有三个地方值得细看一下。
第 4 步的历史只有两行。 step=-1(输入写进来)和 step=0(图跑起来并停在 validate)。因为 validate 是第一个节点,它一失败,后面什么都没发生。对比 §7.3 那份续跑成功后的六行历史,就能直观看出「历史长度」本身就是一个进度指标。
第 5 步的订单号是 A003。 recover 只改了 amount,order_id 用的还是失败那次存下来的状态。这恰恰说明续跑不是重新开始,而是接着那个现场往下走。
第 6 步证明了两个版本共用一份图纸。 build_graph() 不传参数就得到测试版:业务逻辑一字不差,只是没有持久化,get_state 会直接报错。真实项目里这条特别实用:单元测试跑测试版(快、无副作用),线上跑生产版(可观测、可恢复),而被测的图结构是同一份。
8.4. 验收清单 #
invoke拿到完整最终状态,notice字段有值stream逐节点输出增量,能看出每个节点各写了什么- 异常原样抛出,类型和消息都没被包装
- 失败后
next指向失败节点,log为空(失败节点没留半成品) tasks里能读到异常信息- 续跑后
log完整且无重复 - 历史里
source=update标出了人工干预的那一步 - 没 checkpointer 的图能跑但
get_state报错,两个版本由同一份build_graph产出
9. 实用约定与坑 #
9.1. 值得固化成习惯的约定 #
| 约定 | 说明 |
|---|---|
把 build_graph() 单独写成函数 |
同一张图纸能编译出带 / 不带 checkpointer 的版本(§8.2) |
编译时一定传 name |
否则 trace 里全是默认的 LangGraph,分不清哪张图(§2.6) |
写个 new_xxx() 工厂构造初始状态 |
避免 TypedDict 缺字段导致节点内 KeyError(§3.2) |
调试用 stream_mode="updates" |
直接看到「谁写了什么」(§4.1) |
stream_mode / output_keys / print_mode 一律传列表 |
传字符串和传列表的返回形态不同,统一用列表就不会崩(§4.3、§4.8) |
聊天 UI 用 stream_mode="messages" |
并且必须过滤空片段(§4.4) |
有环的图必须自己设 recursion_limit |
默认 10007 起不到保护作用,且要留够 n+1 步(§5.2) |
设完 recursion_limit 故意触发一次 |
报错信息里的数字是唯一的验收凭据(§5.2) |
判断 thread 状态先看 metadata is None |
光看 next == () 分不清「跑完了」和「没跑过」(§6.2) |
排障先看 snap.next 和 snap.tasks |
一眼定位卡在哪、为什么(§7.2) |
需要留痕的修改用 update_state |
历史里会有 source=update,可审计(§7.3) |
9.2. 会当场报错的坑 #
这些都是「大声失败」,看报错信息就能定位:
| 报错 | 原因 | 处理 |
|---|---|---|
ValueError: Found edge ending at unknown node |
add_edge 的目标节点名打错或没注册 |
核对节点名(§2.2) |
ValueError: Graph must have an entrypoint |
忘了 add_edge(START, ...) |
补入口边(§2.2) |
ValueError: No checkpointer set |
编译时没传 checkpointer 却调 get_state / get_state_history / update_state |
compile(checkpointer=...)(§6.1) |
ValueError: Checkpointer requires ... thread_id |
有 checkpointer 但没给 thread_id(空 configurable 也算) |
config 里补上(§5.1) |
KeyError: 'text' 出现在节点里而不是入口 |
TypedDict 不做运行时校验;但挂了 reducer 的字段可以省略 |
工厂函数或改用 Pydantic(§3.2) |
GraphRecursionError |
图里有环且没有终止条件,或上限设得比串行长度还小 | 加终止条件;上限至少 n+1(§5.2) |
InvalidUpdateError: Node xxx does not exist |
update_state(as_node=...) 写了不存在的节点名 |
核对节点名(§6.4) |
9.3. 不报错但结果不对的坑 #
这些是真正难查的,因为程序看起来跑得很正常:
| 现象 | 原因 | 处理 |
|---|---|---|
| 改了 builder 但行为没变 | 编译是快照,控制台那条提示只是警告 | 重新 compile()(§2.4) |
recursion_limit 设了不生效 |
放进了 configurable,被静默忽略(报错里仍是 10007) |
放 config 顶层(§5.2) |
output_schema 设了但 stream 还能看到敏感字段 |
它只作用于 invoke 返回值 |
同时用 output_keys,或干脆别放进状态(§5.4) |
| 输入里的某个字段「没生效」 | 用了 input_schema 收窄入口,多余键被静默丢弃 |
检查 input_schema 覆盖了哪些字段(§5.4) |
节点里的 writer() 进度不显示 |
调用方没请求 custom 模式 |
stream_mode=["custom", ...](§4.5) |
custom 模式下看不到节点返回值 |
单一模式只推自己那一类 | 传列表把 updates 一起要上(§4.5) |
查到 next == () 就当成「已完成」 |
不存在的 thread 也返回空快照 | 先判 metadata is None(§6.2) |
batch 之后 get_state 只剩一份状态 |
几份输入共用了同一个 thread_id,检查点互相交错 |
每份配独立 config(§3.3) |
同 thread_id 反复 invoke 导致 log 越滚越长 |
状态在按 reducer 累积 | 换 thread,或该字段别用累积型 reducer(§5.1) |
自己传的 metadata 把 source / step 冲掉了 |
用户 metadata 是合并进系统 metadata 的 | 给自己的键加前缀(§5.3) |
9.4. 容易误解为 bug 的正常行为 #
| 现象 | 其实是 |
|---|---|
stream 的 values 比节点数多一条 |
第一条是初始状态,条数 = 节点数 + 1(§4.2) |
传了列表 stream_mode 后解包失败 |
多模式时每条是 (模式名, 内容) 元组(§4.3) |
messages 模式开头一大串片段是空的 |
推理模型在流式输出思维链,content 为空、reasoning_content 有值(§4.4) |
| 两次运行的 token 片段数不一样 | 取决于模型这次想多久、答多长,不要依赖片段数量(§4.4) |
update_state 想覆盖却变成了追加 |
该字段挂了 reducer,人工更新和节点更新走同一套规则(§6.4) |
as_node 拨回去重跑后日志出现重复 |
已落盘的历史不会被抹掉,重跑的写入是追加的(§6.4) |
| 从失败点续跑后日志没有重复 | 失败节点从未成功写入,续跑是它第一次真正执行(§7.3) |
没 checkpointer 时 checkpoints 模式一条都没有 |
预期行为,不报错(§4.6) |
interrupt_after=["a"] 和 interrupt_before=["b"] 效果一样 |
串行图上「a 之后」就是「b 之前」;分支图上两者不同(§7.4) |
口诀:
invoke要结果,stream要过程,get_state要现场。 前两个没 checkpointer 也能用,第三个必须有。
10. 练习 #
10.1. 机制验证(不写业务,只确认行为) #
- 看清楚 values 多的那一条:对同一张图分别用
updates和values跑一遍,数一数条数差,解释原因。 - 多模式解包:用
stream_mode=["updates", "custom"]跑一次,写出正确的消费代码;再改成只传"updates",观察这段代码为什么崩了。 - 空片段有多少:用
messages模式跑一次真实模型调用,统计空片段占比,并把additional_kwargs打出来确认它们是不是思维链。 - 递归上限数的是什么:搭一张三节点串行图,从
recursion_limit=1试到4,找出最小可用值,验证「n + 1」这个规律。 - 放错位置的上限:把
recursion_limit塞进configurable,跑死循环图,从报错信息里的数字说明它没生效。 - 区分「跑完」和「没跑过」:对一个没用过的
thread_id调get_state,打印next/values/metadata,写出一段能正确区分两种情况的判断代码。 output_schema的陷阱:给OrderState加一个secret字段,用output_schema排除它,然后证明stream的values和updates里仍然能看到它,最后用output_keys把它挡住。batch共用 thread 的后果:让三份输入共用一个thread_id,打印get_state_history,指出哪些地方说明历史被搅乱了。
10.2. 能力构建(往本章实战上加东西) #
- 自定义进度:给
price节点加get_stream_writer(),在一个假装耗时的循环里推三次进度,并用stream_mode=["custom", "updates"]同时看到进度和增量。 - 时间旅行:跑通 §8 的图后,从
get_state_history里找到next=('notify',)的检查点,改掉tax再重放,看notice是否跟着变。 - 分叉重跑:在同一个 thread 上,回到「什么都还没跑」的检查点,把
amount改成另一个值再invoke(None, ...),然后打印完整历史,指出两条轨迹分别在哪几步。 - 审计来源:在 §8 第 4、5 步的历史输出里,逐行说出
source为什么是input/loop/update。 - 两个版本:用同一个
build_graph()编译出带 / 不带 checkpointer 的两张图,验证前者能get_state而后者不能;再给测试版加上interrupt_before=["price"],确认三个版本共用一份图纸。 - (扩展)区分暂停与失败:写一个
status(graph, thread)函数,用metadata、next、tasks[].error三个信号,返回「未运行 / 已完成 / 已暂停 / 已失败」四种状态之一,并用四个 thread 分别验证。
11. 本章小结 #
StateGraph是图纸,CompiledStateGraph是可执行对象,继承链是CompiledStateGraph → Pregel → Runnable,所以能进 LCEL 链、能当别的图的节点,batch/ainvoke这些方法都是白送的。compile()只做结构检查(悬空边、没入口),不检查「跑出来对不对」;它的参数清一色是运行期设施,没有一个和图的结构有关。- 编译是快照:改了 builder 要重新编译(那条提示只是警告,很容易漏看);同一份 builder 可以编译出多个独立版本,这是「生产版 / 测试版 / 调试版」的实现方式。
invoke返回完整最终状态,不是某个节点的返回值,这让调用方不必知道最后一个节点是谁。- 输入不用给全:挂了 reducer 的字段会用类型标注推出的零值兜底;没挂 reducer 的字段缺了就会在节点内部
KeyError。 stream有七种模式:updates调试、values渲染、messages打字机、custom进度、tasks/debug排障、checkpoints观察持久化。stream()是生成器,不消费就不执行。values模式的条数 = 节点数 + 1,第一条是初始状态;传列表模式时每条变成(模式名, 内容)元组。output_keys传字符串 / 列表也有同样的形态差异,统一传列表最省事。messages模式的空片段是推理模型的思维链(content空、additional_kwargs["reasoning_content"]有值),某次实测 39 个片段里开头 16 条都是空的,换成非推理模型也仍有个位数的空片段,UI 里必须过滤;片段总数每次都不同,任何依赖它的代码都是错的。recursion_limit默认是 10007(在langgraph/_internal/_config.py,可用环境变量改),不是常见资料写的 25;它数的是超步,n 个节点的串行图至少要 n+1;放进configurable会被静默忽略。output_schema只裁剪invoke的返回值,不影响stream,想在流里裁剪要用output_keys;input_schema则会静默丢弃收窄之外的字段。get_state必须有 checkpointer。snap.next告诉你停在哪,snap.tasks[].error告诉你为什么。但next == ()有两个含义(跑完了 / 没跑过),要靠metadata is None区分。get_state_history里的source区分input/loop/update,能回答「这个值是系统算的还是人改的」;历史的长度本身就是进度指标。- 节点抛异常时,异常原样抛出、状态精确停在失败节点之前:成功的节点写入照常落盘,失败的节点一个字不留。
update_state打完补丁用invoke(None, config)就能从断点续跑,而且日志不会重复。 - 时间旅行不只是重放,还能分叉:
update_state(旧检查点的 config, 补丁)会返回一个新坐标,从它invoke(None, ...)就能在保留原轨迹的前提下走另一条路。 - 「暂停」和「失败」的恢复姿势是同一个:
invoke(None, config);区分二者看tasks[].error是不是None。这为第 26 章的人机在环打下了基础。 - 本章产出:一张可运行、可观测、可恢复的状态图。
run/watch/inspect/recover四个入口分别对应invoke/stream/get_state/update_state。
到这里,图的「搭」和「跑」都齐了,但所有例子还都是一条直线。下一章开始让图分叉:条件边与分支,根据状态决定下一步走哪条路。