1. 本章目标 #
第 10 章用过 HumanInTheLoopMiddleware,第 11 章用过 checkpointer 存对话历史,当时都是「照着用」,至于 Command(resume=...) 为什么能让 Agent 从中间继续,并没有解释。这一章把底层拆开。
核心是一个能力:
让一张图在半路停下来,等外面的人拍板,之后接着往下跑。
「停下来」不是 input() 卡在那儿等人。进程可以直接退出,图的现场存在数据库里;几分钟甚至几天后,另一个进程读出现场、喂进决定,流程从断点继续。审批流、风控确认、敏感操作二次校验,靠的就是这套机制。
学完你应能:
- 用
interrupt()在任意节点暂停,用Command(resume=...)恢复 - 避开本章头号坑:恢复时被中断的节点会从头重跑一遍,副作用会执行两次
- 分清动态
interrupt()和静态interrupt_before/interrupt_after,知道各自该用在哪 - 实现 HITL 的三种决策:批准、修改后批准、驳回
- 处理并行的多个 interrupt 和子图里的 interrupt
- 用
SqliteSaver做到跨进程暂停与恢复 - 产出:一条可暂停的退款审批流水线,带 CLI,小额免审、大额挂人工
参考文档:
1.1 为什么「暂停」值得单独讲一章 #
前面五章讲的都是「图自己怎么跑完」:顺序、分支、循环、子图,输入进去、输出出来,中间不需要人。真实业务里却总有一类动作不该由程序单方面决定:打款、删库、给客户发正式函件、调用需要授权的外部接口。
这类场景常见三种错误实现,本章的机制正好把它们一一换掉:
| 错误实现 | 具体问题 | 本章的做法 |
|---|---|---|
在节点里 input() 等人输入 |
进程必须一直活着;Web 服务里根本没有终端可读 | interrupt() 存现场后返回,进程可退出 |
| 拆成两个接口,中间状态自己塞 Redis | 状态格式要自己设计、自己序列化,图跑到哪一步要自己记 | checkpointer 把整张图的现场原样存下 |
| 在 prompt 里写「打款前请先征得同意」 | 那是请求不是保证,模型可以不听(第 25 章 §2) | 边和节点是结构约束,模型绕不过去 |
一句话概括本章的立场:审批不是提示词问题,是图结构问题。
2. checkpointer 不只是记忆 #
第 11 章介绍 checkpointer 时,用途是「让 Agent 记住上一轮说了什么」。那只是它的副产品。
它真正做的是:每个超步结束后,把整张图的完整状态存成一个 checkpoint。「完整」包含三部分:
values:所有状态字段此刻的值next:下一步该执行哪些节点(也就是「跑到哪了」)pending writes:本超步里已经产生、但还没并入状态的写入(包括人已经给出的 resume 值)
存下来之后能干四件事,第 22 章已经用过前三件:
| 能力 | 用什么 | 出现在 |
|---|---|---|
| 读当前现场 | get_state |
第 22 章 |
| 读历史轨迹 | get_state_history |
第 22 章 |
| 改状态后续跑 | update_state + invoke(None) |
第 22 章 |
| 停在半路等外部输入 | interrupt + Command(resume=) |
本章 |
第四件是前三件的自然延伸。状态既能完整存下、又能从任意一点继续,「停下来等人」不过是「暂时不继续」。暂停不需要进程一直活着,这是本章最反直觉、也最有用的一点。
反过来说,本章也没引入新机制:interrupt() 只是在「保存现场 + 停止推进」这两个已有能力上加了一层约定:让节点能说「我需要外部输入才能往下走」,并把要问的问题捎给调用方。
3. interrupt():在节点里暂停 #
3.1 最小例子 #
interrupt() 是个函数:在节点里调用它,图就停在这里。
# operator.add 用作 log 字段的 reducer,让每个节点都是追加而不是覆盖
import operator
# Annotated 用来给字段挂 reducer;TypedDict 用来声明状态结构
from typing import Annotated, TypedDict
# 内存版 checkpointer:本节只在同一个进程里演示,§8 换成文件版
from langgraph.checkpoint.memory import InMemorySaver
# 图的三件套:起点、终点、图构建器
from langgraph.graph import END, START, StateGraph
# Command 用于恢复,interrupt 用于暂停
from langgraph.types import Command, interrupt
# 状态:一笔待审批的支出
class S(TypedDict):
# 金额,只读不改
amount: int
# 人工决定,暂停期间是空的
approved: str
# 审计流水,挂 operator.add 所以是追加
log: Annotated[list[str], operator.add]
# 第一个节点:向人提问
def ask(state: S) -> dict:
# 这行 print 是本章的关键观测点,请记住它出现了几次
print(" [node] ask 开始执行")
# 括号里的东西会被原样送到外面给人看;这一行会让整张图停下来
decision = interrupt({"question": "批准这笔支出吗?", "amount": state["amount"]})
# 恢复之后,interrupt() 会直接返回人给的值,代码从这里继续
print(f" [node] ask 拿到人的答复: {decision}")
# 把人的决定写回状态
return {"approved": decision, "log": [f"人工决定: {decision}"]}
# 第二个节点:拍板之后的后续步骤
def done(state: S) -> dict:
# 打印一行,确认恢复之后流程真的继续往下走了
print(" [node] done")
# 记一行流水
return {"log": ["流程结束"]}
# 开始装图
builder = StateGraph(S)
# 注册提问节点
builder.add_node("ask", ask)
# 注册后续节点
builder.add_node("done", done)
# 入口连到提问
builder.add_edge(START, "ask")
# 提问之后走后续
builder.add_edge("ask", "done")
# 后续之后结束
builder.add_edge("done", END)
# 必须带 checkpointer 编译,否则恢复时无处读现场(下一段会实测不带的后果)
graph = builder.compile(checkpointer=InMemorySaver())
# 一次暂停恢复要跨多次 invoke,所以必须指定 thread_id 把它们串成一条线程
cfg = {"configurable": {"thread_id": "t1"}}
# 第一次调用:跑到 interrupt 就停
out = graph.invoke({"amount": 5000, "approved": "", "log": []}, cfg)
# 看看「停下来」是以什么形式返回的
print("第一次 invoke 返回:", out)
# 单独看一下多出来的那个键
print("返回值的键:", list(out.keys())) [node] ask 开始执行
第一次 invoke 返回: {'amount': 5000, 'approved': '', 'log': [],
'__interrupt__': [Interrupt(value={'question': '批准这笔支出吗?', 'amount': 5000},
id='8f2f30a30af2023a2a5eda4e8f1eff8b')]}
返回值的键: ['amount', 'approved', 'log', '__interrupt__']三个观察:
invoke正常返回了,没有抛异常。这一点很关键,暂停是正常返回状态,不是错误。Web 接口不必用try/except去捕它,而是检查返回值里有没有__interrupt__。- 返回值里多了一个
__interrupt__键,装着Interrupt对象列表。它是列表而不是单个对象,因为可能同时挂起多个(§7.1)。 Interrupt只有两个数据字段:value(你传给interrupt()的东西,原样带出来)和id(内部标识,用于精确恢复)。- 状态没有任何变化:
approved还是空,log还是空。ask节点的返回值根本没执行到,interrupt()那一行之后的代码一行都没跑。
容易误判的一点:没有 checkpointer 时,第一次调用照样「成功」。 很多人以为不配 checkpointer 会立刻报错,实测不是:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class S(TypedDict):
amount: int
approved: str
log: Annotated[list[str], operator.add]
def ask(state: S) -> dict:
print(" [node] ask 开始执行")
decision = interrupt({"question": "批准这笔支出吗?", "amount": state["amount"]})
print(f" [node] ask 拿到人的答复: {decision}")
return {"approved": decision, "log": [f"人工决定: {decision}"]}
def done(state: S) -> dict:
print(" [node] done")
return {"log": ["流程结束"]}
builder = StateGraph(S)
builder.add_node("ask", ask)
builder.add_node("done", done)
builder.add_edge(START, "ask")
builder.add_edge("ask", "done")
builder.add_edge("done", END)
# 同一张图,这次故意不带 checkpointer
graph_no_cp = builder.compile()
# 第一次调用:能跑,也能拿到 __interrupt__,看起来一切正常
out_no_cp = graph_no_cp.invoke({"amount": 100, "approved": "", "log": []})
# 打印一下确认它真的返回了中断信息
print("没有 checkpointer 也能挂起:", "__interrupt__" in out_no_cp)
# 但只要试着恢复就会翻车
try:
# 没有地方存现场,也就没有现场可读
graph_no_cp.invoke(Command(resume="同意"))
# 捕获运行时错误
except RuntimeError as e:
# 打印错误原文,注意它说的是 checkpointer 而不是 interrupt
print("恢复时报错:", e)
# 顺便看看读现场会怎样
try:
# 读状态同样依赖 checkpointer
graph_no_cp.get_state({"configurable": {"thread_id": "x"}})
# 这里抛的是 ValueError,不是 RuntimeError
except ValueError as e:
# 打印错误原文
print("读现场报错:", e) [node] ask 开始执行
没有 checkpointer 也能挂起: True
恢复时报错: Cannot use Command(resume=...) without checkpointer
读现场报错: No checkpointer set忘配 checkpointer 的症状是「提交能用、审批点不动」,而不是启动即失败。 联调若只跑通了「提交」路径,这类错误很容易蒙过去。写审批流的第一件事,就是确认 compile(checkpointer=...)。
3.2 暂停时的现场 #
用第 22 章的 get_state 看看图此刻认为自己在哪:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class S(TypedDict):
amount: int
approved: str
log: Annotated[list[str], operator.add]
def ask(state: S) -> dict:
print(" [node] ask 开始执行")
decision = interrupt({"question": "批准这笔支出吗?", "amount": state["amount"]})
print(f" [node] ask 拿到人的答复: {decision}")
return {"approved": decision, "log": [f"人工决定: {decision}"]}
def done(state: S) -> dict:
print(" [node] done")
return {"log": ["流程结束"]}
builder = StateGraph(S)
builder.add_node("ask", ask)
builder.add_node("done", done)
builder.add_edge(START, "ask")
builder.add_edge("ask", "done")
builder.add_edge("done", END)
graph = builder.compile(checkpointer=InMemorySaver())
cfg = {"configurable": {"thread_id": "t1"}}
graph.invoke({"amount": 5000, "approved": "", "log": []}, cfg)
# 读当前 checkpoint 的快照
snap = graph.get_state(cfg)
# next 是「下一步要执行谁」
print("next:", snap.next)
# values 是此刻的状态值
print("values:", snap.values)
# interrupts 是「正在等外部输入的那些提问」
print("interrupts:", snap.interrupts)
# tasks 是这一刻的待执行任务列表,interrupt 挂在任务上
print("tasks[0].name:", snap.tasks[0].name)
# snap.interrupts 其实就是把各个任务的 interrupts 拼起来
print("和 tasks[0].interrupts 相同吗:", snap.tasks[0].interrupts == snap.interrupts)
# created_at 是这个 checkpoint 的落库时间,§10.3 会用它做幂等判据
print("created_at:", snap.created_at)
# 到这里已经写了两个 checkpoint:输入落库一个、超步结束落库一个
print("历史条数:", len(list(graph.get_state_history(cfg)))) [node] ask 开始执行
next: ('ask',)
values: {'amount': 5000, 'approved': '', 'log': []}
interrupts: (Interrupt(value={'question': '批准这笔支出吗?', 'amount': 5000}, id='538b646b...'),)
tasks[0].name: ask
和 tasks[0].interrupts 相同吗: True
created_at: 2026-09-04T03:42:14.748294+00:00
历史条数: 2(第一行是前置代码跑到暂停时留下的,从 next: 开始才是本段的输出。)
next 指着 ask:图认为 ask 还没跑完。这和第 22 章节点抛异常时的现场一模一样:状态存的是节点执行之前的样子,节点要么整个成功,要么就当没发生过。这个细节直接埋下了 §4 的坑。
三个字段的分工,做审批界面时会反复用到:
| 字段 | 回答什么问题 | 典型用法 |
|---|---|---|
next |
图卡在哪个节点 | 判断「是否还需要人处理」 |
interrupts[i].value |
要问人什么 | 渲染审批页面 |
values |
业务数据现在什么样 | 渲染工单详情 |
3.3 用 Command(resume=...) 恢复 #
把人的决定包进 Command(resume=...),当作 invoke 的输入传进去。注意传的不是状态字典,而是 Command 对象,LangGraph 靠这个类型区分「新一轮输入」和「恢复上一次的暂停」:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class S(TypedDict):
amount: int
approved: str
log: Annotated[list[str], operator.add]
def ask(state: S) -> dict:
print(" [node] ask 开始执行")
decision = interrupt({"question": "批准这笔支出吗?", "amount": state["amount"]})
print(f" [node] ask 拿到人的答复: {decision}")
return {"approved": decision, "log": [f"人工决定: {decision}"]}
def done(state: S) -> dict:
print(" [node] done")
return {"log": ["流程结束"]}
builder = StateGraph(S)
builder.add_node("ask", ask)
builder.add_node("done", done)
builder.add_edge(START, "ask")
builder.add_edge("ask", "done")
builder.add_edge("done", END)
graph = builder.compile(checkpointer=InMemorySaver())
cfg = {"configurable": {"thread_id": "t1"}}
graph.invoke({"amount": 5000, "approved": "", "log": []}, cfg)
# resume 的值会成为 interrupt() 的返回值
out = graph.invoke(Command(resume="同意"), cfg)
# 看流程有没有走完
print("恢复后:", out)
# next 为空元组表示已经跑到 END
print("恢复后 next:", graph.get_state(cfg).next) [node] ask 开始执行
[node] ask 拿到人的答复: 同意
[node] done
恢复后: {'amount': 5000, 'approved': '同意', 'log': ['人工决定: 同意', '流程结束']}
恢复后 next: ()(前置代码首次跑到暂停时也会打印一行 [node] ask 开始执行,上面这段是恢复这一步自己的输出。)
流程走完了,approved 拿到了人的决定。这一次 interrupt() 返回了 resume 传进来的值:同一个调用,第一次把图挂起,第二次直接返回答复。
但请仔细看第一行:ask 节点的开场 print 又打印了一遍。第一次 invoke 时它打印过一次,这次恢复它又打印一次。这不是日志重复,而是下一节要讲的全部内容。
4. 头号坑:恢复时节点从头重跑 #
4.1 副作用会执行两次 #
ask 节点的第一行 print 在恢复后又执行了一次,说明整个函数体是从头重新执行的,不是从 interrupt() 那一行往下续。
正确的心智模型是:
每次 resume,被中断的节点函数都从第一行重新执行。已经答过的
interrupt()不再挂起,直接返回记录下来的答案;没答过的那个才会重新把图挂起。
Python 没法把普通函数「冻结在某一行」再解冻,那需要协程或生成器。LangGraph 选的方案是重放:重跑函数,靠已记录的 resume 值把之前的 interrupt() 一个个「跳过」。
函数里只有 print 时无所谓,但真实节点里 interrupt() 之前往往有正事:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class S(TypedDict):
amount: int
approved: str
log: Annotated[list[str], operator.add]
# 用一个列表记录扣款发生了几次,模拟真实世界里不可撤销的副作用
charged = []
# 错误示范:把副作用写在 interrupt 之前
def bad(state: S) -> dict:
# 假装这里真的调了支付网关
charged.append(1)
# 打印第几次,方便数
print(f" [!] 扣款执行第 {len(charged)} 次")
# 扣完钱才去问人同不同意,顺序本身就有问题,但先看技术后果
d = interrupt("确认吗")
# 人答完之后记一行流水
return {"log": [f"确认 {d}"]}
# 单独装一张只有这个节点的图
builder = StateGraph(S)
# 注册这个有副作用的节点
builder.add_node("bad", bad)
# 入口直接连它
builder.add_edge(START, "bad")
# 它之后就结束
builder.add_edge("bad", END)
# 带 checkpointer 编译
graph = builder.compile(checkpointer=InMemorySaver())
# 用一条新线程
cfg_bad = {"configurable": {"thread_id": "bad"}}
# 第一次:扣款 + 挂起
graph.invoke({"amount": 5000, "approved": "", "log": []}, cfg_bad)
# 恢复:函数从头重跑,于是又扣一次
graph.invoke(Command(resume="ok"), cfg_bad)
# 数一数总共扣了几次
print(f"扣款一共执行 {len(charged)} 次") [!] 扣款执行第 1 次
[!] 扣款执行第 2 次
扣款一共执行 2 次钱扣了两次。 而且这是一笔「批准之后才该扣款」的流程:第一次扣款发生在人还没看到请求时,第二次发生在人点了同意之后。就算审批人点的是「驳回」,第一笔钱也已经出去了。
原因就是 §3.2 那个 next: ('ask',):LangGraph 存的是节点执行前的状态,恢复时它认为这个节点还没做过,于是整函数重跑。interrupt() 内部靠已记录的 resume 值直接返回,它前面的普通代码没有这种记忆。
4.2 正确写法:让副作用独占一个节点 #
解法不是加标志位(那要自己管幂等,容易漏,也容易在并发下失效),而是顺着图的结构拆开:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class S(TypedDict):
amount: int
approved: str
log: Annotated[list[str], operator.add]
# 同样用列表计数,这次验证只执行一次
charged2 = []
# 提问节点:只读状态、只提问、只把决定写回状态
def gate(state: S) -> dict:
# 唯一的动作就是问;重跑一百次也没有任何外部影响
d = interrupt("确认吗")
# 把决定写进流水
return {"log": [f"确认 {d}"]}
# 副作用节点:它只在恢复之后才会被调度,所以只跑一次
def do_charge(state: S) -> dict:
# 真实系统里这里是支付网关
charged2.append(1)
# 打印次数
print(f" [ok] 扣款执行第 {len(charged2)} 次")
# 记一行流水
return {"log": ["已扣款"]}
# 装图:把「问」和「做」拆成前后两个节点
builder = StateGraph(S)
# 注册提问节点
builder.add_node("gate", gate)
# 注册扣款节点
builder.add_node("charge", do_charge)
# 入口先问
builder.add_edge(START, "gate")
# 拍板通过之后才走到扣款
builder.add_edge("gate", "charge")
# 扣完结束
builder.add_edge("charge", END)
# 带 checkpointer 编译
graph = builder.compile(checkpointer=InMemorySaver())
# 换一条新线程
cfg_good = {"configurable": {"thread_id": "good"}}
# 第一次:只挂起,什么钱都没动
graph.invoke({"amount": 5000, "approved": "", "log": []}, cfg_good)
# 恢复:gate 重跑(无害),然后才第一次执行 charge
out_good = graph.invoke(Command(resume="ok"), cfg_good)
# 数一数
print(f"扣款一共执行 {len(charged2)} 次")
# 顺便看流水顺序
print("log:", out_good["log"]) [ok] 扣款执行第 1 次
扣款一共执行 1 次
log: ['确认 ok', '已扣款']规则:带 interrupt() 的节点必须是纯粹的:只读状态、只提问、只把决定写回状态。 任何有外部影响的操作(扣款、发消息、写库、调第三方 API)都放到下一个节点。
这条规则值得写进团队规范。还有个额外好处:gate 节点重跑无害,你可以放心反复 get_state 查看现场,甚至让人改主意再 resume 一遍。
顺便一提,这也解释了第 24 章为什么强调「路由函数不能有副作用」:图会在重放时反复执行这些代码。重放是 LangGraph 的常态,不是异常。
4.3 副作用写在 interrupt() 之后就安全吗 #
有人会想:既然是从头重跑,把副作用挪到 interrupt() 之后不就行了?第一次跑到 interrupt() 就停,副作用根本没机会执行;恢复时 interrupt() 立即返回,副作用才第一次执行。
实测确实只执行一次。但这个结论有个前提:节点里只有一个 interrupt()。一旦有两个,中间的代码就会被跑两次:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class S(TypedDict):
amount: int
approved: str
log: Annotated[list[str], operator.add]
# 记录两次提问之间的代码被执行了几次
between = []
# 一个节点里问两次:先问要不要办,再问用哪种方式办
def two_asks(state: S) -> dict:
# 第一问:第一次 invoke 停在这里
d1 = interrupt("第一问")
# 这段代码在「答完第一问」和「答完第二问」时各跑一次
between.append(1)
# 打印次数,这就是证据
print(f" [!] 两问之间的代码执行第 {len(between)} 次")
# 第二问:第一次 resume 停在这里
d2 = interrupt("第二问")
# 两个答案都拿到了才写状态
return {"log": [f"{d1}/{d2}"]}
# 装一张只有这个节点的图
builder = StateGraph(S)
# 注册双问节点
builder.add_node("two", two_asks)
# 入口连它
builder.add_edge(START, "two")
# 之后结束
builder.add_edge("two", END)
# 带 checkpointer 编译
graph = builder.compile(checkpointer=InMemorySaver())
# 新线程
cfg2 = {"configurable": {"thread_id": "two"}}
# 第一次:停在第一问
result = graph.invoke({"amount": 1, "approved": "", "log": []}, cfg2)
# 确认停在哪一问
print("第一次挂起在:", result["__interrupt__"][0].value)
# 答第一问:函数重跑,第一问直接返回,中间代码执行,停在第二问
result = graph.invoke(Command(resume="答1"), cfg2)
# 确认现在停在第二问
print("第二次挂起在:", result["__interrupt__"][0].value)
# 答第二问:函数再重跑,两个 interrupt 都直接返回,中间代码又执行一次
result = graph.invoke(Command(resume="答2"), cfg2)
# 最终状态
print("最终 log:", result["log"])
# 中间那段代码总共执行了几次
print("中间代码执行次数:", len(between))第一次挂起在: 第一问
[!] 两问之间的代码执行第 1 次
第二次挂起在: 第二问
[!] 两问之间的代码执行第 2 次
最终 log: ['答1/答2']
中间代码执行次数: 2规律很清楚:函数体里第 k 个 interrupt() 之前的代码,会被执行 k 次(每次启动都要走到那儿)。所以「放在 interrupt() 之后」不是安全保证,只是「这个节点恰好只有一个 interrupt()」时的巧合。
还有一个反直觉的场景:update_state 也会清掉待答的提问。实测在暂停期间调 update_state,next 仍指着那个节点,但 tasks 里的 interrupt 消失了;下一次 Command(resume=...) 会让节点函数从头再跑一遍(总共执行 2 次),并把新的 resume 值交给 interrupt()。
# operator.add 给 log 字段当「追加」reducer
import operator
# Annotated 挂 reducer,TypedDict 声明状态
from typing import Annotated, TypedDict
# 内存 checkpointer:暂停现场必须落盘,否则无法恢复
from langgraph.checkpoint.memory import InMemorySaver
# 图的三件套
from langgraph.graph import END, START, StateGraph
# Command 用来恢复,interrupt 用来暂停
from langgraph.types import Command, interrupt
# 状态:一笔待审批的支出
class S(TypedDict):
# 金额,只读
amount: int
# 人工决定,暂停期间是空的
approved: str
# 审计流水,挂 operator.add 所以是追加
log: Annotated[list[str], operator.add]
# 计数器:验证节点函数到底执行了几次
runs = []
# 提问节点:跑到 interrupt 就停
def ask(state: S) -> dict:
# 每次进入函数体都记一笔,用来证明「从头再跑」
runs.append(1)
# 打印第几次,这就是「总共执行 2 次」的证据
print(f" [node] ask 开始执行,第 {len(runs)} 次")
# 第一次走到这里会把图停住;若已有 resume 值,则直接返回那个值
decision = interrupt({"question": "批准这笔支出吗?", "amount": state["amount"]})
# 恢复之后才会打印,resume 进来的值就在这里
print(f" [node] ask 拿到人的答复: {decision!r}")
# 把人的决定写回状态
return {"approved": decision, "log": [f"人工决定: {decision}"]}
# 后续节点:用来确认恢复后流程继续
def done(state: S) -> dict:
# 打印一行
print(" [node] done")
# 记流水
return {"log": ["流程结束"]}
# 装图
builder = StateGraph(S)
# 注册提问节点
builder.add_node("ask", ask)
# 注册后续节点
builder.add_node("done", done)
# 入口连到提问
builder.add_edge(START, "ask")
# 提问之后走后续
builder.add_edge("ask", "done")
# 后续之后结束
builder.add_edge("done", END)
# 必须带 checkpointer,暂停现场才有地方存
graph = builder.compile(checkpointer=InMemorySaver())
# 固定线程,多次 invoke / update_state 才能看到同一条现场
cfg = {"configurable": {"thread_id": "upd"}}
# 第一次:跑到 interrupt 就停
out = graph.invoke({"amount": 5000, "approved": "", "log": []}, cfg)
# 确认确实挂起了
print("第一次挂起在:", out["__interrupt__"][0].value)
# 暂停期间的快照
snap = graph.get_state(cfg)
print("=== update_state 之前 ===")
# next 指着还没跑完的那个节点
print("next:", snap.next)
# 待答的提问挂在 tasks 上
print("tasks:", [(t.name, [i.value for i in t.interrupts]) for t in snap.tasks])
# 关键操作:暂停期间改状态。这会清掉待答的提问
graph.update_state(cfg, {"log": ["人工改了一笔"]})
# 再读快照,对照两个变化
snap2 = graph.get_state(cfg)
print("=== update_state 之后 ===")
# next 仍指着 ask:图认为这个节点还没跑完
print("next:", snap2.next)
# 状态里能看到刚打的补丁
print("values.log:", snap2.values["log"])
# 但 interrupts 空了
print("interrupts:", snap2.interrupts)
# tasks 里那个 ask 还在,挂着的 interrupt 却消失了
print("tasks:", [(t.name, [i.value for i in t.interrupts]) for t in snap2.tasks])
# 带着 resume 再启动:节点函数从头再跑一遍,interrupt() 直接拿到这个值
print("=== Command(resume='同意') ===")
out2 = graph.invoke(Command(resume="同意"), cfg)
# 最终状态:补丁还在,批准结果是这次 resume 给的
print("最终:", out2)
# 第一次 invoke 一次 + resume 再跑一次 = 2
print("ask 总共执行次数:", len(runs))
结论仍是那一条:副作用独占一个节点。 这条规则不依赖「有几个 interrupt()」「有没有人调过 update_state」这类你控制不住的细节。
4.4 一张表记住重跑规则 #
| 代码位置 | 执行次数 | 能否放副作用 |
|---|---|---|
interrupt() 之前 |
每次启动都跑,1 次暂停 = 2 次 | 绝对不行 |
两个 interrupt() 之间 |
后面还剩几次 resume 就多跑几次 | 不行 |
唯一的 interrupt() 之后 |
1 次 | 能,但依赖「只有一个 interrupt」这个前提,不推荐 |
| 下游独立节点里 | 1 次 | 推荐 |
5. 静态 interrupt:不改节点代码就能暂停 #
interrupt() 要写进节点里。若只想在某个节点前后卡一下(调试,或给现成的图临时加观察点),可以在 compile() 时声明,节点代码一行都不用动:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class S(TypedDict):
amount: int
approved: str
log: Annotated[list[str], operator.add]
# 计数器:验证被卡住的节点到底执行了几次
counter = []
# 一个有副作用的节点,故意不写 interrupt
def counting(state: S) -> dict:
# 每次执行都记一笔
counter.append(1)
# 打印第几次
print(f" [node] counting 第 {len(counter)} 次")
# 记流水
return {"log": [f"计数 {len(counter)}"]}
# 后续节点,用来确认恢复后流程继续
def notify(state: S) -> dict:
# 打印一行
print(" [node] notify")
# 记流水
return {"log": ["已通知"]}
# 装图:两个普通节点,没有任何 interrupt 代码
builder = StateGraph(S)
# 注册第一个节点,名字叫 charge
builder.add_node("charge", counting)
# 注册第二个节点
builder.add_node("notify", notify)
# 入口连到第一个节点
builder.add_edge(START, "charge")
# 第一个节点之后走通知
builder.add_edge("charge", "notify")
# 通知之后结束
builder.add_edge("notify", END)
# 关键在这一行:编译时声明「进 charge 之前先停一下」
graph = builder.compile(checkpointer=InMemorySaver(), interrupt_before=["charge"])
# 新线程
cfg_b = {"configurable": {"thread_id": "before"}}
# 第一次调用:直接停在 charge 之前
o = graph.invoke({"amount": 100, "approved": "", "log": []}, cfg_b)
# 返回值里没有中断信息,只有状态
print("返回:", o)
# 只能靠 next 知道停在哪
print("next:", graph.get_state(cfg_b).next)
# 静态中断不产生 Interrupt 对象
print("有 __interrupt__ 吗:", "__interrupt__" in o)
# 快照里的 interrupts 也是空的
print("snap.interrupts:", graph.get_state(cfg_b).interrupts)
# 静态中断用 None 恢复:没有值要传进去
print("--- 用 invoke(None) 恢复 ---")
# 恢复并跑到底
print(graph.invoke(None, cfg_b))
# 关键结论:节点只执行了一次
print(f"一共执行 {len(counter)} 次")返回: {'amount': 100, 'approved': '', 'log': []}
next: ('charge',)
有 __interrupt__ 吗: False
snap.interrupts: ()
--- 用 invoke(None) 恢复 ---
[node] counting 第 1 次
[node] notify
{'amount': 100, 'approved': '', 'log': ['计数 1', '已通知']}
一共执行 1 次和动态 interrupt() 比,有三处重要区别:
一、返回值里没有 __interrupt__,快照里的 interrupts 也是空的。 静态中断不携带任何「要问什么」的信息,外面只能靠 get_state().next 知道图停在哪。它适合调试,不适合做审批界面,没法告诉前端该展示什么。
二、恢复用 invoke(None, cfg)。 None 的含义是「不给新输入,接着上次的现场往下跑」,第 22 章 update_state 之后续跑用的是同一个写法。
三、节点不会重跑。 实测计数器只有 1 次。interrupt_before 是在节点开始之前卡住的,节点压根没启动过,谈不上重跑。这是它和 §4 那个坑的本质差别,也意味着静态中断对副作用是安全的。
interrupt_after=["charge"] 则是节点跑完、状态存好之后再卡:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class S(TypedDict):
amount: int
approved: str
log: Annotated[list[str], operator.add]
counter = []
def counting(state: S) -> dict:
counter.append(1)
print(f" [node] counting 第 {len(counter)} 次")
return {"log": [f"计数 {len(counter)}"]}
def notify(state: S) -> dict:
print(" [node] notify")
return {"log": ["已通知"]}
builder = StateGraph(S)
builder.add_node("charge", counting)
builder.add_node("notify", notify)
builder.add_edge(START, "charge")
builder.add_edge("charge", "notify")
builder.add_edge("notify", END)
# 清空计数器重新数
counter.clear()
# 同一张图,这次改成「跑完 charge 再停」
graph = builder.compile(checkpointer=InMemorySaver(), interrupt_after=["charge"])
# 新线程
cfg_a = {"configurable": {"thread_id": "after"}}
# 第一次调用
graph.invoke({"amount": 200, "approved": "", "log": []}, cfg_a)
# 读现场
snap_a = graph.get_state(cfg_a)
# charge 的结果已经并进状态了,next 指向的是它的下游
print("暂停时 log:", snap_a.values["log"], " next:", snap_a.next)
# 同样用 None 恢复
print("恢复:", graph.invoke(None, cfg_a))
# 依然只执行一次
print(f"一共执行 {len(counter)} 次") [node] counting 第 1 次
暂停时 log: ['计数 1'] next: ('notify',)
[node] notify
恢复: {'amount': 200, 'approved': '', 'log': ['计数 1', '已通知']}
一共执行 1 次怎么选很直观:想在某个动作发生前审查输入,用 interrupt_before;想在它做完之后检查输出再决定要不要继续,用 interrupt_after。
调试时还有个省事写法:interrupt_before="*" 表示每个节点前都停一下,等价于单步执行。接受的是节点名列表,或字符串 "*"。
5.1 一个容易踩的组合错误 #
动态 interrupt() 不能用 invoke(None, cfg) 恢复。试一下:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class S(TypedDict):
amount: int
approved: str
log: Annotated[list[str], operator.add]
charged2 = []
def gate(state: S) -> dict:
d = interrupt("确认吗")
return {"log": [f"确认 {d}"]}
def do_charge(state: S) -> dict:
charged2.append(1)
print(f" [ok] 扣款执行第 {len(charged2)} 次")
return {"log": ["已扣款"]}
builder = StateGraph(S)
builder.add_node("gate", gate)
builder.add_node("charge", do_charge)
builder.add_edge(START, "gate")
builder.add_edge("gate", "charge")
builder.add_edge("charge", END)
graph = builder.compile(checkpointer=InMemorySaver())
# 重建 §4.2 那张「gate → charge」的图,它用的是动态 interrupt
cfg_x = {"configurable": {"thread_id": "wrong-resume"}}
# 第一次调用:正常挂起
graph.invoke({"amount": 1, "approved": "", "log": []}, cfg_x)
# 用静态中断那套写法去恢复动态中断
out_x = graph.invoke(None, cfg_x)
# 看看发生了什么
print("结果:", out_x)
# 图有没有往前走
print("还在暂停吗:", graph.get_state(cfg_x).next)
# 换成正确的写法就好了
print("改用 Command(resume=) 后:", graph.invoke(Command(resume="ok"), cfg_x))结果: {'amount': 1, 'approved': '', 'log': [], '__interrupt__': [Interrupt(value='确认吗', id='22d0469e...')]}
还在暂停吗: ('gate',)
[ok] 扣款执行第 1 次
改用 Command(resume=) 后: {'amount': 1, 'approved': '', 'log': ['确认 ok', '已扣款']}(最后那次扣款是这笔工单被正确恢复后的第一次扣款,不是重复扣款:charge 是独立节点,invoke(None) 那次根本没走到它。)
图原地不动,又把同一个 interrupt 抛了出来。 不报错,就是不往前走。Web 服务里写错这一处,表现往往是「点了确认没反应」,日志里什么异常都没有。
更糟的是:那次失败的 invoke(None) 并不是空转:被中断的节点函数会被完整重跑一遍(interrupt() 之前的代码又执行一次),只是跑到 interrupt() 时没有 resume 值可用,又一次挂起。写错恢复方式不仅不前进,还会白白触发一次 §4 的副作用。
反方向倒是宽容的:静态中断用 Command(resume=...) 恢复不会报错,那个 resume 值没人接收,会被安静忽略,图照常往下跑。所以只有一个方向会出问题,记住「动态必须用 Command」就够了。
两者的完整对照:
动态 interrupt() |
静态 interrupt_before/after |
|
|---|---|---|
| 在哪声明 | 节点函数内部 | compile() 参数 |
| 能否携带数据 | 能,interrupt(value) |
不能 |
返回 __interrupt__ |
是 | 否 |
snap.interrupts |
有内容 | 空元组 |
| 恢复方式 | 只能 Command(resume=值) |
invoke(None, cfg),也接受 Command |
| 节点是否重跑 | 是(§4) | 否 |
| 对副作用安全吗 | 不安全,必须拆节点 | 安全 |
| 适合 | 生产审批流 | 调试、临时卡点 |
6. HITL 的三种决策 #
审批场景通常要三种结果:同意、改一改再同意、驳回。一个节点就能全覆盖,interrupt() 的返回值是任意 Python 对象,可以让前端回传结构化字典,节点内部再分派。
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
# Literal 用来声明路由函数的返回值范围(第 23 章)
from typing import Literal
# 一笔待审批的请求
class ReqState(TypedDict):
# 动作名,比如「退款」
action: str
# 金额,审批人可以改
amount: int
# 流转状态
status: str
# 审计流水
log: Annotated[list[str], operator.add]
# 第一个节点:登记申请
def submit(state: ReqState) -> dict:
# 只记一行,代表前置的确定性步骤
return {"log": [f"申请:{state['action']} {state['amount']} 元"]}
# 审批节点:一个 interrupt 覆盖三种决策
def human_gate(state: ReqState) -> dict:
# 把审批界面需要的所有信息一次性装齐
d = interrupt({
# 做什么
"action": state["action"],
# 多少钱
"amount": state["amount"],
# 告诉前端可以回传哪几种决策
"options": ["approve", "edit", "reject"],
})
# 用结构化字典而不是裸字符串,后面加字段不用改协议
kind = d["type"]
# 第一种:原样批准
if kind == "approve":
# 只改状态,动钱交给下游节点
return {"status": "approved", "log": ["人工:批准原样执行"]}
# 第二种:改金额后批准
if kind == "edit":
# 一次返回三个字段:状态、新金额、流水
return {
# 状态同样是已批准
"status": "approved",
# 人改了金额,直接写回状态;下游读到的就是新值
"amount": d["amount"],
# 流水里留痕
"log": [f"人工:改成 {d['amount']} 元后批准"],
}
# 第三种:驳回,理由也记进流水
return {"status": "rejected", "log": [f"人工:驳回({d.get('reason', '')})"]}
# 批准后的执行节点:副作用独占一个节点(§4 的规则)
def execute(state: ReqState) -> dict:
# 这里才是真正动钱的地方
return {"log": [f"执行:{state['action']} {state['amount']} 元"]}
# 驳回后的终止节点
def abort(state: ReqState) -> dict:
# 明确记录「什么都没做」,事后可查
return {"log": ["终止:未执行任何操作"]}
# 条件边的路由函数:看 status 决定去哪(第 23 章)
def route(state: ReqState) -> Literal["execute", "abort"]:
# 只有明确批准才执行
return "execute" if state["status"] == "approved" else "abort"
# 装图
builder = StateGraph(ReqState)
# 注册登记节点
builder.add_node("submit", submit)
# 注册审批节点
builder.add_node("human_gate", human_gate)
# 注册执行节点
builder.add_node("execute", execute)
# 注册终止节点
builder.add_node("abort", abort)
# 入口先登记
builder.add_edge(START, "submit")
# 登记后送审
builder.add_edge("submit", "human_gate")
# 审批后按决策分流
builder.add_conditional_edges("human_gate", route, ["execute", "abort"])
# 执行完结束
builder.add_edge("execute", END)
# 终止也是结束
builder.add_edge("abort", END)
# 带 checkpointer 编译
graph = builder.compile(checkpointer=InMemorySaver())
# 三种决策各跑一遍,每种用一条独立线程
for name, d in (
# 原样批准,不带额外字段
("approve", {"type": "approve"}),
# 改额批准,带上新金额
("edit", {"type": "edit", "amount": 500}),
# 驳回,带上理由
("reject", {"type": "reject", "reason": "凭证不全"}),
):
# 每种决策一条线程,互不干扰
cfg_d = {"configurable": {"thread_id": f"d-{name}"}}
# 先跑到审批点
graph.invoke({"action": "退款", "amount": 1280, "status": "", "log": []}, cfg_d)
# 再喂决策
out_d = graph.invoke(Command(resume=d), cfg_d)
# 打印最终状态
print(f" 决策 {name} -> status={out_d['status']}")
# 打印完整流水,看清走了哪条路
for line in out_d["log"]:
# 缩进一点,和上面的标题区分开
print(f" {line}")
# 空行分隔
print() 决策 approve -> status=approved
申请:退款 1280 元
人工:批准原样执行
执行:退款 1280 元
决策 edit -> status=approved
申请:退款 1280 元
人工:改成 500 元后批准
执行:退款 500 元
决策 reject -> status=rejected
申请:退款 1280 元
人工:驳回(凭证不全)
终止:未执行任何操作edit 那条路值得多看一眼:申请的是 1280,审批人改成 500,最终执行的是 500。改动经 human_gate 写回状态,execute 读到的就是新值,不需要额外的传参机制。
6.1 Command 的四个参数 #
「改一改」有两种写法。上面是在节点里根据决定改状态;另一种是用 Command 同时做两件事:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class S(TypedDict):
amount: int
approved: str
log: Annotated[list[str], operator.add]
charged2 = []
def gate(state: S) -> dict:
d = interrupt("确认吗")
return {"log": [f"确认 {d}"]}
def do_charge(state: S) -> dict:
charged2.append(1)
print(f" [ok] 扣款执行第 {len(charged2)} 次")
return {"log": ["已扣款"]}
builder = StateGraph(S)
builder.add_node("gate", gate)
builder.add_node("charge", do_charge)
builder.add_edge(START, "gate")
builder.add_edge("gate", "charge")
builder.add_edge("charge", END)
graph = builder.compile(checkpointer=InMemorySaver())
# 复用 §4.2 的 gate → charge 图,这次连带改状态
cfg_u = {"configurable": {"thread_id": "with-update"}}
# 先挂起
graph.invoke({"amount": 1280, "approved": "", "log": []}, cfg_u)
# resume 喂答复,update 直接改状态字段
out_u = graph.invoke(Command(resume="同意", update={"amount": 500}), cfg_u)
# 金额已经被外部改掉了
print("resume + update:", out_u) [ok] 扣款执行第 1 次
resume + update: {'amount': 500, 'approved': '', 'log': ['确认 同意', '已扣款']}Command 一共有四个参数,可以组合使用:
| 参数 | 作用 | 相关章节 |
|---|---|---|
resume |
喂给 interrupt() 的值 |
本章 |
update |
直接改状态,会走 reducer | 第 21 章 |
goto |
指定下一个节点 | 第 23 章 |
graph |
指定作用于哪一层图(父图 / 子图) | 第 25 章 |
update「会走 reducer」这一点要特别注意:它不是覆盖赋值。对挂了 operator.add 的字段做 update,是追加:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class S(TypedDict):
amount: int
approved: str
log: Annotated[list[str], operator.add]
charged2 = []
def gate(state: S) -> dict:
d = interrupt("确认吗")
return {"log": [f"确认 {d}"]}
def do_charge(state: S) -> dict:
charged2.append(1)
print(f" [ok] 扣款执行第 {len(charged2)} 次")
return {"log": ["已扣款"]}
builder = StateGraph(S)
builder.add_node("gate", gate)
builder.add_node("charge", do_charge)
builder.add_edge(START, "gate")
builder.add_edge("gate", "charge")
builder.add_edge("charge", END)
graph = builder.compile(checkpointer=InMemorySaver())
# 新线程,初始 log 里先放一项
cfg_u2 = {"configurable": {"thread_id": "update-reducer"}}
# 跑到挂起
graph.invoke({"amount": 1280, "approved": "", "log": ["原有"]}, cfg_u2)
# 对带 reducer 的字段做 update
out_u2 = graph.invoke(Command(resume="同意", update={"log": ["被 update 加的"]}), cfg_u2)
# 结果是追加,不是覆盖
print("log:", out_u2["log"]) [ok] 扣款执行第 1 次
log: ['原有', '被 update 加的', '确认 同意', '已扣款']选型建议: 优先在节点里处理,改状态的逻辑和校验(金额不能为负、不能超过原单)能写在一起,改动也会自然留在 log 里。Command(update=...) 适合调用方明确知道要覆盖什么字段的场合,比如管理后台的强制改单。
最后一个提醒:恢复时不要顺手带上 goto。实测 Command(resume="同意", goto="charge") 会让 charge 执行两次:一次来自被中断节点恢复后正常走的边,一次来自 goto 额外派发的任务。想跳步就用条件边表达,别在恢复时临时改路线。
7. 多个 interrupt 与子图 #
7.1 并行的多个 interrupt #
两个并行节点都调 interrupt(),会同时挂起。对应「一笔支出要财务和法务分别签字」这类真实需求:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
# 一个只有流水字段的简单状态
class P(TypedDict):
# 两个并行节点都往里追加
log: Annotated[list[str], operator.add]
# 第一个审批人
def a_node(state: P) -> dict:
# A 自己的提问
d = interrupt("A 要确认")
# 把 A 的答复写进流水
return {"log": [f"A={d}"]}
# 第二个审批人
def b_node(state: P) -> dict:
# B 自己的提问
d = interrupt("B 要确认")
# 把 B 的答复写进流水
return {"log": [f"B={d}"]}
# 装图:从 START 同时扇出到两个节点
builder = StateGraph(P)
# 注册 A
builder.add_node("a", a_node)
# 注册 B
builder.add_node("b", b_node)
# 从 START 连到 A
builder.add_edge(START, "a")
# 也从 START 连到 B,两条边让 a 和 b 落在同一个超步里并行
builder.add_edge(START, "b")
# A 之后结束
builder.add_edge("a", END)
# B 之后也结束
builder.add_edge("b", END)
# 带 checkpointer 编译
graph = builder.compile(checkpointer=InMemorySaver())
# 新线程
cfg_p = {"configurable": {"thread_id": "par"}}
# 一次 invoke,两个节点各挂起一次
out_p = graph.invoke({"log": []}, cfg_p)
# __interrupt__ 是列表,这就是它为什么是列表
print("__interrupt__ 条数:", len(out_p["__interrupt__"]))
# 逐条看:每个挂起的任务有自己的 id
for i in out_p["__interrupt__"]:
# id 用于精确恢复,value 是给人看的内容
print(f" id={i.id} value={i.value}")
# next 里两个节点都在
print("next:", graph.get_state(cfg_p).next)__interrupt__ 条数: 2
id=0b3417ea3069ad89babdda0c80df199f value=A 要确认
id=8fb9fe241d7e28b885be83eebe4aedba value=B 要确认
next: ('a', 'b')这时用单个值 resume 会直接报错,LangGraph 不会猜你想答哪一个:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class P(TypedDict):
log: Annotated[list[str], operator.add]
def a_node(state: P) -> dict:
d = interrupt("A 要确认")
return {"log": [f"A={d}"]}
def b_node(state: P) -> dict:
d = interrupt("B 要确认")
return {"log": [f"B={d}"]}
builder = StateGraph(P)
builder.add_node("a", a_node)
builder.add_node("b", b_node)
builder.add_edge(START, "a")
builder.add_edge(START, "b")
builder.add_edge("a", END)
builder.add_edge("b", END)
graph = builder.compile(checkpointer=InMemorySaver())
cfg_p = {"configurable": {"thread_id": "par"}}
out_p = graph.invoke({"log": []}, cfg_p)
# 故意用单值恢复,看它怎么拒绝
try:
# 有两个待答提问时这样写是歧义的
graph.invoke(Command(resume="批准"), cfg_p)
# 捕获运行时错误
except RuntimeError as e:
# 错误信息里直接告诉你「必须指定 id」
print("报错:", e)报错: When there are multiple pending interrupts, you must specify the interrupt id when resuming. Docs: https://docs.langchain.com/oss/python/langgraph/add-human-in-the-loop#resume-multiple-interrupts-with-one-invocation.(上面这条 URL 是 LangGraph 1.2.11 错误信息里硬编码的原文,官方文档已经改版,该地址现在是 404,正确的页面是 interrupts#handling-multiple-interrupts。
正确做法是传一个 {id: 值} 的字典:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class P(TypedDict):
log: Annotated[list[str], operator.add]
def a_node(state: P) -> dict:
d = interrupt("A 要确认")
return {"log": [f"A={d}"]}
def b_node(state: P) -> dict:
d = interrupt("B 要确认")
return {"log": [f"B={d}"]}
builder = StateGraph(P)
builder.add_node("a", a_node)
builder.add_node("b", b_node)
builder.add_edge(START, "a")
builder.add_edge(START, "b")
builder.add_edge("a", END)
builder.add_edge("b", END)
graph = builder.compile(checkpointer=InMemorySaver())
cfg_p = {"configurable": {"thread_id": "par"}}
out_p = graph.invoke({"log": []}, cfg_p)
# 用返回值里的 id 作键,构造「哪个提问 → 什么答复」的映射
resume_map = {i.id: f"批准-{i.value[0]}" for i in out_p["__interrupt__"]}
# 一次把两个都答掉
print("结果:", graph.invoke(Command(resume=resume_map), cfg_p))结果: {'log': ['A=批准-A', 'B=批准-B']}也可以只批一个,另一个继续挂着,这才是「两个审批人各自处理自己那条」的真实形态:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class P(TypedDict):
log: Annotated[list[str], operator.add]
def a_node(state: P) -> dict:
d = interrupt("A 要确认")
return {"log": [f"A={d}"]}
def b_node(state: P) -> dict:
d = interrupt("B 要确认")
return {"log": [f"B={d}"]}
builder = StateGraph(P)
builder.add_node("a", a_node)
builder.add_node("b", b_node)
builder.add_edge(START, "a")
builder.add_edge(START, "b")
builder.add_edge("a", END)
builder.add_edge("b", END)
graph = builder.compile(checkpointer=InMemorySaver())
cfg_p = {"configurable": {"thread_id": "par"}}
out_p = graph.invoke({"log": []}, cfg_p)
# 换一条新线程重来
cfg_p2 = {"configurable": {"thread_id": "par2"}}
# 再次挂起两个
out_p2 = graph.invoke({"log": []}, cfg_p2)
# 挑出 A 的那一条
a_itr = [i for i in out_p2["__interrupt__"] if i.value.startswith("A")][0]
# 只答 A
out_one = graph.invoke(Command(resume={a_itr.id: "只批A"}), cfg_p2)
# A 的结果已经并入状态,返回值里还剩 B 的提问
print("只 resume 一个后:", out_one)
# next 只剩 b
print("next:", graph.get_state(cfg_p2).next)只 resume 一个后: {'log': ['A=只批A'], '__interrupt__': [Interrupt(value='B 要确认', id='63d1019e...')]}
next: ('b',)这里有个必须知道的陷阱:部分恢复之后,snap.interrupts 会把已经批过的那条也列出来。 因为重放时 a 会被再次调度、再次走到 interrupt()(resume 值虽已记录,但快照的展现方式仍是「这个任务有一个 interrupt」)。所以:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class P(TypedDict):
log: Annotated[list[str], operator.add]
def a_node(state: P) -> dict:
d = interrupt("A 要确认")
return {"log": [f"A={d}"]}
def b_node(state: P) -> dict:
d = interrupt("B 要确认")
return {"log": [f"B={d}"]}
builder = StateGraph(P)
builder.add_node("a", a_node)
builder.add_node("b", b_node)
builder.add_edge(START, "a")
builder.add_edge(START, "b")
builder.add_edge("a", END)
builder.add_edge("b", END)
graph = builder.compile(checkpointer=InMemorySaver())
cfg_p = {"configurable": {"thread_id": "par"}}
out_p = graph.invoke({"log": []}, cfg_p)
cfg_p2 = {"configurable": {"thread_id": "par2"}}
out_p2 = graph.invoke({"log": []}, cfg_p2)
a_itr = [i for i in out_p2["__interrupt__"] if i.value.startswith("A")][0]
out_one = graph.invoke(Command(resume={a_itr.id: "只批A"}), cfg_p2)
print("真正等人的:", out_one["__interrupt__"])
# 读现场
snap_p = graph.get_state(cfg_p2)
# 直接看 snap.interrupts:A 也在里面,具有误导性
print("snap.interrupts:", [i.value for i in snap_p.interrupts])
# tasks 里同样两个都在
print("tasks:", [(t.name, [i.value for i in t.interrupts]) for t in snap_p.tasks])
# 用 next 过滤才是可靠判据:只认此刻真正待执行的节点
pending = [i for t in snap_p.tasks if t.name in snap_p.next for i in t.interrupts]
# 这才是「真正还在等人」的那一条
print("真正等人的:", [i.value for i in pending])snap.interrupts: ['A 要确认', 'B 要确认']
tasks: [('a', ['A 要确认']), ('b', ['B 要确认'])]
真正等人的: ['B 要确认']判断「谁还在等人」有两个可靠来源: 上一次 invoke 返回值里的 __interrupt__;或从快照里用 next 过滤 tasks。直接用 snap.interrupts,在「部分恢复」时会多列出已办项,§10 的实战里专门封了一个 pending_interrupts() 做这件事。
顺便提一个静默失败:拿一个已经答过的旧 id 去 resume,什么都不会发生,不报错,图也不前进。所以待办列表要现查现用,不要把 id 缓存在前端。
把剩下的那条也批掉,流程才走完:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class P(TypedDict):
log: Annotated[list[str], operator.add]
def a_node(state: P) -> dict:
d = interrupt("A 要确认")
return {"log": [f"A={d}"]}
def b_node(state: P) -> dict:
d = interrupt("B 要确认")
return {"log": [f"B={d}"]}
builder = StateGraph(P)
builder.add_node("a", a_node)
builder.add_node("b", b_node)
builder.add_edge(START, "a")
builder.add_edge(START, "b")
builder.add_edge("a", END)
builder.add_edge("b", END)
graph = builder.compile(checkpointer=InMemorySaver())
cfg_p = {"configurable": {"thread_id": "par"}}
out_p = graph.invoke({"log": []}, cfg_p)
cfg_p2 = {"configurable": {"thread_id": "par2"}}
out_p2 = graph.invoke({"log": []}, cfg_p2)
a_itr = [i for i in out_p2["__interrupt__"] if i.value.startswith("A")][0]
out_one = graph.invoke(Command(resume={a_itr.id: "只批A"}), cfg_p2)
# 从上一次的返回值里取最新的 B 提问
b_itr = out_one["__interrupt__"][0]
# 答掉它
print("最终:", graph.invoke(Command(resume={b_itr.id: "再批B"}), cfg_p2))
# 现在才真的跑完
print("next:", graph.get_state(cfg_p2).next)最终: {'log': ['A=只批A', 'B=再批B']}
next: ()好消息是:a 虽在重放里被再次调度,状态更新只会并入一次(log 里只有一个 A=只批A),写在 interrupt() 之后的副作用实测也只执行一次。但这依赖「a 只有一个 interrupt」,所以 §4 的规则照旧。
7.2 子图里的 interrupt #
第 25 章把 create_agent 当子图用过。子图内部的 interrupt() 会一路冒泡到最外层,调用方不必知道它来自第几层:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class P(TypedDict):
log: Annotated[list[str], operator.add]
# 子图的提问节点
def inner_ask(state: P) -> dict:
# 子图内部的 interrupt
d = interrupt("子图要确认")
# 写进共享的 log 字段(这里故意共享,§7.3 会看到后果)
return {"log": [f"inner={d}"]}
# 装子图
inner_builder = StateGraph(P)
# 注册子图节点
inner_builder.add_node("ask", inner_ask)
# 子图入口连到提问节点
inner_builder.add_edge(START, "ask")
# 提问之后子图就结束
inner_builder.add_edge("ask", END)
# 子图不需要自己的 checkpointer,它会用父图的
subgraph = inner_builder.compile()
# 父图的前置节点
def pre(state: P) -> dict:
# 记一行,代表审批前的确定性步骤
return {"log": ["pre"]}
# 父图的后置节点
def post(state: P) -> dict:
# 记一行,代表审批后的确定性步骤
return {"log": ["post"]}
# 装父图
builder = StateGraph(P)
# 前置
builder.add_node("pre", pre)
# 直接把编译好的子图当节点用(第 25 章的写法)
builder.add_node("sub", subgraph)
# 后置
builder.add_node("post", post)
# 入口先跑前置
builder.add_edge(START, "pre")
# 前置之后进子图
builder.add_edge("pre", "sub")
# 子图出来跑后置
builder.add_edge("sub", "post")
# 后置之后结束
builder.add_edge("post", END)
# 只有父图挂 checkpointer
graph = builder.compile(checkpointer=InMemorySaver())
# 新线程
cfg_s = {"configurable": {"thread_id": "sub"}}
# 跑到子图里的 interrupt 就停
out_s = graph.invoke({"log": []}, cfg_s)
# 中断信息冒泡到了最外层的返回值里
print("暂停:", out_s)
# 读父图现场
snap_s = graph.get_state(cfg_s)
# 父图只知道自己停在 sub 这个节点上
print("父图 next:", snap_s.next)
# 默认情况下 tasks[0].state 只是一个 config 字典,不是快照
print("tasks[0].state 的类型:", type(snap_s.tasks[0].state).__name__)
# 想看子图内部停在哪,必须显式要求展开子图
snap_deep = graph.get_state(cfg_s, subgraphs=True)
# 这时候拿到的才是子图的 StateSnapshot
print("展开后的类型:", type(snap_deep.tasks[0].state).__name__)
# 子图内部具体停在哪个节点
print("子图 next:", snap_deep.tasks[0].state.next)
# 子图看到的状态值
print("子图 values:", snap_deep.tasks[0].state.values)暂停: {'log': ['pre'], '__interrupt__': [Interrupt(value='子图要确认', id='3d467db0...')]}
父图 next: ('sub',)
tasks[0].state 的类型: dict
展开后的类型: StateSnapshot
子图 next: ('ask',)
子图 values: {'log': ['pre']}注意 tasks[0].state 那两行:不加 subgraphs=True 时它只是一个 config 字典(装着子图的 checkpoint_ns),访问 .next 会抛 AttributeError。这和第 25 章 §6.3 是同一件事:想看嵌套层的现场,必须显式要求展开。
stream 模式下 interrupt 会作为一个单独的 chunk 出现,键就是 __interrupt__:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class P(TypedDict):
log: Annotated[list[str], operator.add]
def inner_ask(state: P) -> dict:
d = interrupt("子图要确认")
return {"log": [f"inner={d}"]}
inner_builder = StateGraph(P)
inner_builder.add_node("ask", inner_ask)
inner_builder.add_edge(START, "ask")
inner_builder.add_edge("ask", END)
subgraph = inner_builder.compile()
builder = StateGraph(P)
builder.add_node("pre", lambda s: {"log": ["pre"]})
builder.add_node("sub", subgraph)
builder.add_node("post", lambda s: {"log": ["post"]})
builder.add_edge(START, "pre")
builder.add_edge("pre", "sub")
builder.add_edge("sub", "post")
builder.add_edge("post", END)
graph = builder.compile(checkpointer=InMemorySaver())
# 新线程,用流式跑一遍
cfg_st = {"configurable": {"thread_id": "sub-stream"}}
# 默认 updates 模式:每个节点产出一个 chunk
for chunk in graph.stream({"log": []}, cfg_st):
# 暂停会以一个特殊的 chunk 出现在流里
print(" ", chunk) {'pre': {'log': ['pre']}}
{'__interrupt__': (Interrupt(value='子图要确认', id='c65de5bf...'),)}做前端时这点很实用:不必等 stream 结束再查状态,收到 __interrupt__ 这个 chunk 就可以立刻把审批卡片推给用户。恢复照常用 Command(resume=...),不用指定层级,Command 的 graph 参数留给「要明确操作某一层」的高级场景,日常不需要。
7.3 一个隐蔽的重复累积 #
把上面那个子图例子恢复,看 log:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
class P(TypedDict):
log: Annotated[list[str], operator.add]
def inner_ask(state: P) -> dict:
d = interrupt("子图要确认")
return {"log": [f"inner={d}"]}
inner_builder = StateGraph(P)
inner_builder.add_node("ask", inner_ask)
inner_builder.add_edge(START, "ask")
inner_builder.add_edge("ask", END)
subgraph = inner_builder.compile()
builder = StateGraph(P)
builder.add_node("pre", lambda s: {"log": ["pre"]})
builder.add_node("sub", subgraph)
builder.add_node("post", lambda s: {"log": ["post"]})
builder.add_edge(START, "pre")
builder.add_edge("pre", "sub")
builder.add_edge("sub", "post")
builder.add_edge("post", END)
graph = builder.compile(checkpointer=InMemorySaver())
cfg_s = {"configurable": {"thread_id": "sub"}}
graph.invoke({"log": []}, cfg_s)
# 恢复子图里的提问
out_r = graph.invoke(Command(resume="子图批准"), cfg_s)
# 注意 pre 出现了两次
print("恢复后:", out_r)恢复后: {'log': ['pre', 'pre', 'inner=子图批准', 'post']}pre 出现了两次。原因和第 25 章 §4.3 相同:子图返回的是完整状态,父图再用 reducer 合并一次。 父图传给子图 log=['pre'],子图算完返回 log=['pre', 'inner=...'],父图的 operator.add 把它整个追加上去,于是 ['pre'] + ['pre', 'inner=...']。
第 25 章的 messages 场景没出现这个问题,因为 add_messages 会按 id 去重;operator.add 不会。
最省事的解法是给子图定义独立的状态字段,不和父图共享带 reducer 的键:
import operator
from typing import Annotated, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
# 子图的状态:只有它自己的字段
class Inner2(TypedDict):
# 子图专用流水
inner_log: Annotated[list[str], operator.add]
# 父图的状态:包含子图字段,这样才能接收子图的回传
class Outer2(TypedDict):
# 父图自己的流水
log: Annotated[list[str], operator.add]
# 承接子图回传的字段
inner_log: Annotated[list[str], operator.add]
# 子图节点:只写自己的字段
def inner_ask2(state: Inner2) -> dict:
# 同样是 interrupt
d = interrupt("子图要确认")
# 写进子图专用字段,父图的 log 就不会被重复追加
return {"inner_log": [f"inner={d}"]}
# 装子图
inner_builder = StateGraph(Inner2)
# 注册节点
inner_builder.add_node("ask", inner_ask2)
# 子图入口连到提问节点
inner_builder.add_edge(START, "ask")
# 提问之后子图结束
inner_builder.add_edge("ask", END)
# 装父图
builder = StateGraph(Outer2)
# 前置节点用 lambda 简写,只写父图字段
builder.add_node("pre", lambda s: {"log": ["pre"]})
# 嵌入子图
builder.add_node("sub", inner_builder.compile())
# 后置节点
builder.add_node("post", lambda s: {"log": ["post"]})
# 入口先跑前置
builder.add_edge(START, "pre")
# 前置之后进子图
builder.add_edge("pre", "sub")
# 子图出来跑后置
builder.add_edge("sub", "post")
# 后置之后结束
builder.add_edge("post", END)
# 带 checkpointer 编译
graph = builder.compile(checkpointer=InMemorySaver())
# 新线程
cfg_sep = {"configurable": {"thread_id": "sep"}}
# 跑到挂起
graph.invoke({"log": [], "inner_log": []}, cfg_sep)
# 恢复并看结果
print("独立字段后:", graph.invoke(Command(resume="子图批准"), cfg_sep))独立字段后: {'log': ['pre', 'post'], 'inner_log': ['inner=子图批准']}干净了。另一个解法是第 25 章 §4.3 的包装写法:把子图裹进普通节点,自己控制回传什么。没法改子图状态定义时(比如子图是 create_agent 造的),只能走这条路。
8. 真正的持久化:跨进程暂停 #
前面都用 InMemorySaver,进程一退状态就没了,「暂停」只是同一次运行内的事。换成文件或数据库后端,才能体现这套机制的真正价值。
下面两个脚本是两个独立进程,中间没有任何东西在内存里等着。
进程 A:提交申请,跑到审批点就退出。
"""p1_start.py:提交申请,跑到审批点就退出。"""
# operator.add 用作 log 的 reducer
import operator
# 声明状态结构
from typing import Annotated, TypedDict
# 文件型 checkpointer,状态落在 sqlite 文件里
from langgraph.checkpoint.sqlite import SqliteSaver
# 图的三件套
from langgraph.graph import END, START, StateGraph
# 只有暂停要用 interrupt,恢复在另一个进程里做
from langgraph.types import interrupt
# 状态文件路径:两个进程靠这个文件对话
DB = "hitl_demo.sqlite"
# 一笔退款申请
class S(TypedDict):
# 金额
amount: int
# 审批人署名,暂停期间是空的
approver: str
# 审计流水
log: Annotated[list[str], operator.add]
# 第一步:起草,确定性步骤
def draft(state: S) -> dict:
# 打印一行,方便对照两个进程各跑了哪些节点
print(" [node] draft 起草退款单")
# 记流水
return {"log": [f"起草退款 {state['amount']} 元"]}
# 第二步:请人审批。这个节点只提问,不动钱(§4 的规则)
def review(state: S) -> dict:
# 这一行在进程 A 和进程 B 里各会打印一次,正是 §4 说的重跑
print(" [node] review 请求人工审批")
# 把审批界面需要的信息装齐
decision = interrupt({"action": "退款", "amount": state["amount"]})
# 恢复后才会执行到这里,把决定写回状态
return {"approver": decision["by"], "log": [f"审批: {decision['ok']} by {decision['by']}"]}
# 第三步:真正打款,副作用独占一个节点
def execute(state: S) -> dict:
# 只有恢复之后才会走到这里,只执行一次
print(" [node] execute 真的打款了")
# 记流水
return {"log": ["打款完成"]}
# 图的定义写成函数,方便两个进程复用同一份结构
def build(checkpointer):
# 用状态类型初始化构建器
builder = StateGraph(S)
# 注册起草节点
builder.add_node("draft", draft)
# 注册审批节点
builder.add_node("review", review)
# 注册打款节点
builder.add_node("execute", execute)
# 入口先起草
builder.add_edge(START, "draft")
# 起草后送审
builder.add_edge("draft", "review")
# 审批通过后打款
builder.add_edge("review", "execute")
# 打款之后结束
builder.add_edge("execute", END)
# checkpointer 从外面传进来,换后端不用改这里
return builder.compile(checkpointer=checkpointer)
# 只有直接运行才执行,被 import 时只提供 DB 和 build
if __name__ == "__main__":
# with 管理 sqlite 连接,退出时正常关闭
with SqliteSaver.from_conn_string(DB) as cp:
# 装图
graph = build(cp)
# thread_id 用业务主键,这里是退款单号
cfg = {"configurable": {"thread_id": "refund-001"}}
# 跑到审批点就返回
out = graph.invoke({"amount": 1280, "approver": "", "log": []}, cfg)
# 确认停在哪
print("进程 A 停在:", graph.get_state(cfg).next)
# 把待办内容打出来,真实系统里这块是写进待办表或推给前端
print("待办:", [i.value for i in out["__interrupt__"]])
# 强调一下:这个进程接下来就退出了
print("进程 A 退出,图的状态留在", DB) [node] draft 起草退款单
[node] review 请求人工审批
进程 A 停在: ('review',)
待办: [{'action': '退款', 'amount': 1280}]
进程 A 退出,图的状态留在 hitl_demo.sqlite进程 B:一个完全独立的进程,读出现场并批准。
"""p2_resume.py:另一个进程读出现场并批准。"""
# 同一个文件型后端
from langgraph.checkpoint.sqlite import SqliteSaver
# 恢复要用 Command
from langgraph.types import Command
# 只复用图的定义和文件路径,不复用任何运行时状态
from p1_start import DB, build
# 直接运行时执行
if __name__ == "__main__":
# 打开同一个 sqlite 文件
with SqliteSaver.from_conn_string(DB) as cp:
# 重新装一遍图:图的结构是代码,状态才在文件里
graph = build(cp)
# 用同一个 thread_id 定位那笔退款
cfg = {"configurable": {"thread_id": "refund-001"}}
# 先读现场,这一步完全不执行任何节点
snap = graph.get_state(cfg)
# 打印表头
print("进程 B 读到的现场:")
# 卡在哪
print(" next:", snap.next)
# 业务数据长什么样
print(" values:", snap.values)
# 要问人什么
print(" 待办:", [i.value for i in snap.interrupts])
# 喂进决定,从断点继续
print("\n进程 B 批准:")
# resume 的值就是 review 里 interrupt() 的返回值
out = graph.invoke(Command(resume={"ok": True, "by": "王经理"}), cfg)
# 完整结果
print(" 最终:", out)
# 确认跑完了
print(" next:", graph.get_state(cfg).next)进程 B 读到的现场:
next: ('review',)
values: {'amount': 1280, 'approver': '', 'log': ['起草退款 1280 元']}
待办: [{'action': '退款', 'amount': 1280}]
进程 B 批准:
[node] review 请求人工审批
[node] execute 真的打款了
最终: {'amount': 1280, 'approver': '王经理', 'log': ['起草退款 1280 元', '审批: True by 王经理', '打款完成']}
next: ()这就是审批流的完整实现。 进程 A 可以是接用户请求的 Web 接口,把待办写进数据库就返回;进程 B 可以是几天后管理员点「批准」时触发的另一个请求。中间没有任何东西在内存里等着。
三个值得留意的细节:
draft只在进程 A 打印,review在两个进程里各打印一次。 已经跑完的节点不会重放,被中断的那个会,§4 的规则跨进程同样成立,而且这次「两次」发生在两个不同的进程里,光看单个进程的日志根本发现不了。- 进程 B 里
get_state是纯读操作,不执行任何节点,所以做「待办列表」页面很安全。 - 进程 B 只 import 了
DB和build。真实项目里这两个进程通常是同一个代码库的不同入口(Web handler 和后台命令),图的定义放在共享模块里。
后端的选择:
| 后端 | 适用 | 说明 |
|---|---|---|
InMemorySaver |
单元测试、Demo | 进程退出即失效 |
SqliteSaver |
本地开发、单机小应用 | 一个文件,零配置 |
PostgresSaver |
生产 | 见 md/PostgreSQL.md |
换后端不用改图的代码,只换 compile(checkpointer=...) 那一个参数。异步场景对应 AsyncSqliteSaver / AsyncPostgresSaver,配合 ainvoke 使用。
9. 和第 10 章的 HITL 中间件是什么关系 #
第 10 章用 HumanInTheLoopMiddleware 给 Agent 的高危工具加审批,当时只说「用 Command(resume=...) 恢复」,没讲原理。翻一下源码(langchain/agents/middleware/human_in_the_loop.py),核心就一行:
# 中间件的 after_model 钩子里,把要审批的工具调用打包问人
decisions = interrupt(hitl_request)["decisions"]它就是本章的 interrupt()。 中间件替你做的四件封装是:
- 从最后一条
AIMessage里挑出命中interrupt_on的tool_calls - 把它们组装成结构化的
hitl_request(含action_requests和review_configs) - 一次
interrupt()问所有工具调用,要求回传的decisions条数与之匹配,不匹配就抛ValueError - 按决策改写
tool_calls,被驳回的那些补上「人造」ToolMessage
实际跑一遍,看看它长什么样:
# 从 .env 读 DEEPSEEK_API_KEY
from dotenv import load_dotenv
# override=True 让 .env 覆盖同名的系统环境变量
load_dotenv(override=True)
# 第 9 章的 Agent 工厂
from langchain.agents import create_agent
# 第 10 章的 HITL 中间件
from langchain.agents.middleware import HumanInTheLoopMiddleware
# 工具装饰器
from langchain_core.tools import tool
# 中间件底层是 interrupt(),所以 Agent 也必须挂 checkpointer
from langgraph.checkpoint.memory import InMemorySaver
# 一个高危工具:真实系统里它会动钱
@tool
def refund(order_id: str, amount: int) -> str:
"""给指定订单退款。"""
# 打印一行,用来确认它到底有没有被真的执行
print(f" [$] 真的退款了 {order_id} {amount}")
# 返回给模型的观察结果
return f"已退款 {amount} 元"
# 造一个带审批的 Agent
agent = create_agent(
# 模型标识
model="deepseek:deepseek-v4-flash",
# 只给它这一个工具
tools=[refund],
# interrupt_on 声明「调 refund 之前必须问人」
middleware=[HumanInTheLoopMiddleware(interrupt_on={"refund": True})],
# 中间件底层是 interrupt(),所以必须有 checkpointer
checkpointer=InMemorySaver(),
)
# 看看中间件在图里加了什么
print("nodes:", list(agent.get_graph().nodes))
# 新线程
cfg_h = {"configurable": {"thread_id": "hitl"}}
# 提出一个会触发工具调用的请求
out_h = agent.invoke({"messages": [{"role": "user", "content": "给订单 A-9 退款 1280 元"}]}, cfg_h)
# 中断信息的结构:这就是要交给审批界面的东西
print("interrupt value:", out_h["__interrupt__"][0].value)
# 关键:图停在中间件自己的节点上,而不是 model 节点
print("next:", agent.get_state(cfg_h).next)nodes: ['__start__', 'model', 'tools', 'HumanInTheLoopMiddleware.after_model', '__end__']
interrupt value: {'action_requests': [{'name': 'refund', 'args': {'order_id': 'A-9', 'amount': 1280},
'description': "Tool execution requires approval\n\nTool: refund\nArgs: {'order_id': 'A-9', 'amount': 1280}"}],
'review_configs': [{'action_name': 'refund',
'allowed_decisions': ['approve', 'edit', 'reject', 'respond']}]}
next: ('HumanInTheLoopMiddleware.after_model',)next 那一行很有意思:中间件的钩子是一个独立的节点(HumanInTheLoopMiddleware.after_model),interrupt() 在它里面。所以恢复时重跑的是这个纯粹的中间件节点,不会重新调用一次模型,它恰好遵守了 §4 的规则,这也是为什么用中间件不会莫名多烧一次 token。
恢复时传的是 {"decisions": [...]},每个决策对应一个被拦下的工具调用:
from dotenv import load_dotenv
load_dotenv(override=True)
from langchain.agents import create_agent
from langchain.agents.middleware import HumanInTheLoopMiddleware
from langchain_core.tools import tool
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.types import Command
@tool
def refund(order_id: str, amount: int) -> str:
"""给指定订单退款。"""
print(f" [$] 真的退款了 {order_id} {amount}")
return f"已退款 {amount} 元"
agent = create_agent(
model="deepseek:deepseek-v4-flash",
tools=[refund],
middleware=[HumanInTheLoopMiddleware(interrupt_on={"refund": True})],
checkpointer=InMemorySaver(),
)
cfg_h = {"configurable": {"thread_id": "hitl"}}
agent.invoke({"messages": [{"role": "user", "content": "给订单 A-9 退款 1280 元"}]}, cfg_h)
# 批准:decisions 的条数必须和被拦下的工具调用数一致
out_h2 = agent.invoke(Command(resume={"decisions": [{"type": "approve"}]}), cfg_h)
# 打印完整消息序列,看它怎么把审批结果接回对话
for m in out_h2["messages"]:
# 类型 + 内容前 40 字
print(f" {type(m).__name__}: {str(m.content)[:40]!r}") [$] 真的退款了 A-9 1280
HumanMessage: '给订单 A-9 退款 1280 元'
AIMessage: ''
ToolMessage: '已退款 1280 元'
AIMessage: '已为订单 A-9 退款 1280 元。'换成驳回,中间件会替你补一条 ToolMessage,工具本身根本不执行:
from dotenv import load_dotenv
load_dotenv(override=True)
from langchain.agents import create_agent
from langchain.agents.middleware import HumanInTheLoopMiddleware
from langchain_core.tools import tool
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.types import Command
@tool
def refund(order_id: str, amount: int) -> str:
"""给指定订单退款。"""
print(f" [$] 真的退款了 {order_id} {amount}")
return f"已退款 {amount} 元"
agent = create_agent(
model="deepseek:deepseek-v4-flash",
tools=[refund],
middleware=[HumanInTheLoopMiddleware(interrupt_on={"refund": True})],
checkpointer=InMemorySaver(),
)
# 新线程重来一遍
cfg_h2 = {"configurable": {"thread_id": "hitl-reject"}}
# 先跑到审批点
agent.invoke({"messages": [{"role": "user", "content": "给订单 B-7 退款 999 元"}]}, cfg_h2)
# 驳回,并带上给模型看的说明
out_h3 = agent.invoke(
# type 是 reject,message 是给模型看的驳回理由
Command(resume={"decisions": [{"type": "reject", "message": "凭证不全,不予退款"}]}),
cfg_h2,
)
# 观察消息序列:有 ToolMessage,但没有 [$] 打印
for m in out_h3["messages"]:
# 那条 ToolMessage 是中间件伪造的,工具本身没执行
print(f" {type(m).__name__}: {str(m.content)[:70]!r}") HumanMessage: '给订单 B-7 退款 999 元'
AIMessage: ''
ToolMessage: 'User rejected the tool call for `refund` with reason: 凭证不全,不予退款'
AIMessage: '退款未成功。系统提示“凭证不全,不予退款”,因此订单 **B-7** 的 999 元退款未能办理。\n\n建议您先补充相关的退款凭证,然后我'注意那条 ToolMessage 的内容:中间件用的是自己的英文模板 User rejected the tool call for ...,你传的 message 被接在 with reason: 后面,而不是整条替换。想完全掌控这句话的措辞,就得走 §9.1 或自己搭图的路子。
这条人造 ToolMessage 也不是可选的礼貌行为,而是协议要求:第 24 章 §7.2 讲过,AIMessage 里每个 tool_call_id 都必须有对应的 ToolMessage,否则模型 API 会拒绝整段消息序列(400)。中间件替你补了;自己写就得自己补。
9.1 第三条路:直接在工具函数里 interrupt() #
除了「中间件」和「自己搭图」,还有一种官方写法:把 interrupt() 写进工具函数内部。适合「这个工具无论谁调、在哪调都要审批」的场景:
# 从 .env 读 DEEPSEEK_API_KEY
from dotenv import load_dotenv
# override=True 让 .env 覆盖同名的系统环境变量
load_dotenv(override=True)
# 第 9 章的 Agent 工厂
from langchain.agents import create_agent
# 工具装饰器
from langchain_core.tools import tool
# 工具里调 interrupt 同样要求 Agent 挂 checkpointer
from langgraph.checkpoint.memory import InMemorySaver
# Command 用于恢复,interrupt 用于在工具内部暂停
from langgraph.types import Command, interrupt
# 记录真实打款次数
PAID = []
# 工具内部自己问人
@tool
def refund2(order_id: str, amount: int) -> str:
"""给指定订单退款,执行前必须人工确认。"""
# 工具函数里一样可以调 interrupt,图会在 tools 节点上挂起
ok = interrupt({"tool": "refund", "order_id": order_id, "amount": amount})
# 人不同意就直接返回,钱不动
if not ok.get("approve"):
# 这句话会成为 ToolMessage 的内容,模型据此组织对客话术
return f"人工驳回:{ok.get('reason', '未说明')}"
# 同意才动钱;这行在 interrupt 之后,只会执行一次
PAID.append(amount)
# 打印次数
print(f" [$] 真的打款 {amount}(第 {len(PAID)} 次)")
# 返回观察结果
return f"已退款 {amount} 元"
# 这次不挂中间件,审批逻辑在工具里
agent = create_agent(
# 同一个模型
model="deepseek:deepseek-v4-flash",
# 换成自带审批的工具
tools=[refund2],
# 明确要求它直接调工具,别在回复里反问用户
system_prompt="你是退款机器人。用户提出退款请求时直接调用工具,不要在回复里反问用户确认。",
# 依然需要 checkpointer
checkpointer=InMemorySaver(),
)
# 新线程
cfg_t = {"configurable": {"thread_id": "tool-hitl"}}
# 提出请求,模型发出工具调用,工具内部挂起
out_t = agent.invoke({"messages": [{"role": "user", "content": "给订单 A-9 退款 1280 元"}]}, cfg_t)
# 提问内容是工具自己定义的,比中间件的格式更贴业务
print("挂起内容:", out_t["__interrupt__"][0].value)
# 这次停在 tools 节点上
print("next:", agent.get_state(cfg_t).next)
# 此刻一分钱都没动
print("打款次数:", len(PAID))
# 批准
agent.invoke(Command(resume={"approve": True}), cfg_t)
# 恢复之后才动钱,且只动一次
print("恢复后打款次数:", len(PAID))挂起内容: {'tool': 'refund', 'order_id': 'A-9', 'amount': 1280}
next: ('tools',)
打款次数: 0
[$] 真的打款 1280(第 1 次)
恢复后打款次数: 1但这条路有个必须知道的坑:同一条 AIMessage 里的兄弟工具调用会被重复执行。 tools 节点是一个整体,恢复时整节点重跑,那些不需要审批的工具就跑了第二次:
from dotenv import load_dotenv
load_dotenv(override=True)
from langchain.agents import create_agent
from langchain_core.tools import tool
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
PAID = []
@tool
def refund2(order_id: str, amount: int) -> str:
"""给指定订单退款,执行前必须人工确认。"""
ok = interrupt({"tool": "refund", "order_id": order_id, "amount": amount})
if not ok.get("approve"):
return f"人工驳回:{ok.get('reason', '未说明')}"
PAID.append(amount)
print(f" [$] 真的打款 {amount}(第 {len(PAID)} 次)")
return f"已退款 {amount} 元"
# AIMessage 用来手工构造带两个工具调用的模型输出
from langchain_core.messages import AIMessage
# MessagesState 自带 messages 字段和 add_messages reducer
from langgraph.graph import MessagesState
# ToolNode 是第 24 章的工具执行节点
from langgraph.prebuilt import ToolNode
# 记录短信发送次数
SMS = []
# 一个不需要审批的普通工具
@tool
def send_sms(to: str) -> str:
"""给客户发短信。"""
# 记一次
SMS.append(to)
# 打印第几次,这就是证据
print(f" [sms] 发短信给 {to}(第 {len(SMS)} 次)")
# 返回观察结果
return "已发送"
# 用伪模型代替真模型,保证每次都发出同样的两个工具调用
def fake_model(state: MessagesState) -> dict:
# 打印一行,用来确认模型节点没有重跑
print(" [model] 伪模型发出两个工具调用")
# 一条 AIMessage 里带两个 tool_calls:一个要审批,一个不要
return {"messages": [AIMessage(
# 有 tool_calls 时 content 可以为空
content="",
# 两个调用会在同一个 tools 节点里执行
tool_calls=[
# 不需要审批的那个
{"name": "send_sms", "args": {"to": "13800000000"}, "id": "c1", "type": "tool_call"},
# 需要审批的那个
{"name": "refund2", "args": {"order_id": "A-9", "amount": 1280}, "id": "c2",
"type": "tool_call"},
],
)]}
# 手搭一张 model → tools 的图
builder = StateGraph(MessagesState)
# 伪模型节点
builder.add_node("model", fake_model)
# 工具节点,同时挂两个工具
builder.add_node("tools", ToolNode([send_sms, refund2]))
# 入口连到伪模型
builder.add_edge(START, "model")
# 模型之后执行工具
builder.add_edge("model", "tools")
# 工具执行完就结束,不回模型(本例不需要循环)
builder.add_edge("tools", END)
# 带 checkpointer 编译
graph = builder.compile(checkpointer=InMemorySaver())
# 新线程
cfg_sib = {"configurable": {"thread_id": "sibling"}}
# 第一次:短信发出去了,退款卡在审批
graph.invoke({"messages": [{"role": "user", "content": "退款"}]}, cfg_sib)
# 此刻短信 1 次
print(f"挂起时:短信 {len(SMS)} 次 / 打款 {len(PAID)} 次")
# 批准退款
graph.invoke(Command(resume={"approve": True}), cfg_sib)
# 短信变成 2 次,兄弟工具被重复执行了
print(f"恢复后:短信 {len(SMS)} 次 / 打款 {len(PAID)} 次") [model] 伪模型发出两个工具调用
[sms] 发短信给 13800000000(第 1 次)
挂起时:短信 1 次 / 打款 0 次
[sms] 发短信给 13800000000(第 2 次)
[$] 真的打款 1280(第 1 次)
恢复后:短信 2 次 / 打款 1 次打款只有一次,短信却发了两次。 model 节点只打印了一次,它上个超步就跑完了,不会重放;但 tools 节点整节点重跑,send_sms 跟着又跑一遍。结论:工具内 interrupt() 只在「同一批工具调用里没有其他有副作用的工具」时才安全。 拿不准就用中间件,或把审批点提到图层面(§10 的做法)。
9.2 三种写法怎么选 #
| 第 10 章中间件 | 工具内 interrupt() |
本章自己搭图 | |
|---|---|---|---|
| 粒度 | 工具调用级 | 单个工具 | 任意节点、任意位置 |
| 配置 | interrupt_on={"refund": True} |
写在工具函数里 | 自己写在节点里 |
| 提问内容 | 中间件规定的结构 | 自己定 | 自己定 |
驳回时补 ToolMessage |
自动 | 靠 return 一句话 |
不涉及 |
| 兄弟工具重复执行 | 不会 | 会(§9.1) | 不涉及 |
| 代码量 | 一行 | 几行 | 要自己搭图 |
选型: 审批点正好是「某个工具调用要不要放行」,用中间件,一行搞定。审批点落在业务流程里(比如「起草完退款单之后、打款之前」,中间还有别的确定性步骤),就用本章的写法。工具内 interrupt() 适合「审批规则属于这个工具本身」的小场景。§10 的实战属于第二类。
10. 实战:可暂停的退款审批流水线 #
10.1 设计 #
一个带 CLI 的退款审批系统:小额自动通过,大额挂人工。
START → validate ──┬─ auto_approve(小额)──→ payout ──→ notify → END
├─ human_gate(大额,interrupt)─┬─ payout ─↗
│ └─ abort ──↗
└─ abort(校验失败)──────────────────────↗几个刻意设计,每一条都对应前面某一节的教训:
| 设计 | 防的是什么 | 出处 |
|---|---|---|
human_gate 只提问,payout 才动钱 |
恢复时重复打款 | §4 |
payout 和 abort 都汇入 notify |
驳回的单子没人通知客户 | 第 25 章 |
小额走 auto_approve,完全不触发 interrupt |
让人去点几十块钱的退款 | 本节 |
thread_id 用工单号 |
无法追查、无法幂等 | §10.3 |
提交入口检查 created_at |
重复提交导致二次打款 | §10.3 |
待办列表用 next 过滤 tasks |
已办项混进待办列表 | §7.1 |
10.2 refund_flow.py #
"""实战:可暂停的多步退款审批图。每条 CLI 命令都是一个独立进程。"""
# 让类型标注延迟求值,Annotated[...] 的写法不受定义顺序限制
from __future__ import annotations
# operator.add 用作 log 字段的 reducer,实现「追加而不是覆盖」
import operator
# 读命令行参数
import sys
# Annotated 给字段挂 reducer;Literal 声明路由函数的返回值范围
from typing import Annotated, Literal, TypedDict
# 文件型 checkpointer:状态落在 sqlite 文件里,进程退出也不丢
from langgraph.checkpoint.sqlite import SqliteSaver
# 图的三件套:起点、终点、图构建器
from langgraph.graph import END, START, StateGraph
# Command 用来恢复,interrupt 用来暂停
from langgraph.types import Command, interrupt
# 状态文件路径:换成 PostgresSaver 时只有这一行和 main() 里的连接方式要改
DB = "refund_flow.sqlite"
# 小于这个金额自动通过,不打扰人
AUTO_LIMIT = 200
# 一张退款工单的全部状态;这些字段会被完整存进 checkpoint
class RefundState(TypedDict):
# 工单号,同时用作 thread_id,天然可查、可幂等
ticket_id: str
# 退款金额,审批人可以在 human_gate 里改小
amount: int
# 退款事由,会原样送到审批界面
reason: str
# 流转状态:"" 表示还没定,approved / rejected 是终态
status: str
# 审批人署名,自动通过时写「系统自动」
approver: str
# 全流程审计流水,挂 operator.add 所以每个节点都是追加
log: Annotated[list[str], operator.add]
def validate(state: RefundState) -> dict:
"""确定性校验:金额必须为正。"""
# 金额非正数直接判死,后面的路由会把它送到 abort
if state["amount"] <= 0:
# 只写状态,不做任何外部动作;approver 也填上,notify 才不会渲染出空括号
return {"status": "rejected", "approver": "系统校验", "log": ["校验失败:金额必须大于 0"]}
# 校验通过,记一行受理流水
return {"log": [f"受理工单 {state['ticket_id']}:退款 {state['amount']} 元({state['reason']})"]}
def auto_approve(state: RefundState) -> dict:
"""小额免审,确定性放行。"""
# 直接把状态写成已批准,审批人记成系统
return {
# 后面的固定边会把它送到 payout
"status": "approved",
# 审批人字段不留空,方便 notify 统一渲染
"approver": "系统自动",
# 把「为什么没走人工」写进流水,事后可解释
"log": [f"金额 {state['amount']} < {AUTO_LIMIT},自动通过"],
}
def human_gate(state: RefundState) -> dict:
"""大额必须人工拍板。这个节点只负责问,不做任何副作用,恢复时它会整个重跑。"""
# interrupt 的参数是给审批界面看的,要一次装齐渲染所需的全部信息
decision = interrupt({
# 给人看的业务主键,不要让人看 interrupt 的 id
"ticket_id": state["ticket_id"],
# 待审金额
"amount": state["amount"],
# 退款事由
"reason": state["reason"],
# 告诉前端可以回传哪几种决策
"allowed": ["approve", "edit", "reject"],
})
# 决策用结构化字典而不是裸字符串,后面加字段不用改协议
kind = decision["type"]
# 署名缺失时给个兜底值,避免 notify 渲染出 None
who = decision.get("by", "未署名")
# 第一种决策:原样批准
if kind == "approve":
# 只改状态,钱由下一个节点去动
return {"status": "approved", "approver": who, "log": [f"{who} 批准原样退款"]}
# 第二种决策:改金额后批准
if kind == "edit":
# 前端传来的数字可能是字符串,统一转成 int
new_amount = int(decision["amount"])
# 一次返回四个字段:状态、审批人、改后金额、流水
return {
# 状态同样是已批准
"status": "approved",
# 记下是谁改的
"approver": who,
# 把改后的金额写回状态,payout 读到的就是新值
"amount": new_amount,
# 流水里留下改动痕迹
"log": [f"{who} 将金额改为 {new_amount} 元后批准"],
}
# 第三种决策:驳回
return {
# 终态是 rejected,路由会送到 abort
"status": "rejected",
# 驳回也要留下是谁驳的
"approver": who,
# 驳回理由要落进流水,这是对客解释的依据
"log": [f"{who} 驳回:{decision.get('reason', '未说明')}"],
}
def payout(state: RefundState) -> dict:
"""副作用独占一个节点:恢复后才执行,只会执行一次。"""
# 真实系统里这里是调支付网关;打印一行方便在输出里数它执行了几次
print(f" [$] 实际打款 {state['amount']} 元 -> 工单 {state['ticket_id']}")
# 打款金额取自状态,所以审批人改过的值会自动生效
return {"log": [f"打款成功 {state['amount']} 元"]}
def abort(state: RefundState) -> dict:
"""终止路径:明确记录「没有发生资金变动」。"""
# 什么都不做,但一定要留一行流水,否则事后看不出这单为什么没钱动
return {"log": ["流程终止,未发生资金变动"]}
def notify(state: RefundState) -> dict:
"""确定性后处理:所有路径都汇到这里。"""
# 无论批准还是驳回都要通知客户,所以这个节点必须是必经之路
return {"log": [f"已通知客户,最终状态 {state['status']}(审批人 {state['approver']})"]}
def need_human(state: RefundState) -> Literal["auto_approve", "human_gate", "abort"]:
"""校验之后的三分岔。"""
# 校验就没过,不用惊动任何人
if state["status"] == "rejected":
# 直接去终止节点,连审批点都不进
return "abort"
# 小额免审,大额挂人工
return "auto_approve" if state["amount"] < AUTO_LIMIT else "human_gate"
def after_decision(state: RefundState) -> Literal["payout", "abort"]:
"""人工拍板之后的两分岔。"""
# 只有明确批准才走到动钱的节点
return "payout" if state["status"] == "approved" else "abort"
def build(checkpointer):
"""装图。checkpointer 从外面传进来,换后端不用改这里。"""
# 用状态类型初始化构建器
builder = StateGraph(RefundState)
# 注册校验节点,名字就是后面连边和 next 里看到的名字
builder.add_node("validate", validate)
# 注册小额免审节点
builder.add_node("auto_approve", auto_approve)
# 注册人工审批节点,唯一带 interrupt 的节点
builder.add_node("human_gate", human_gate)
# 注册打款节点,唯一有副作用的节点
builder.add_node("payout", payout)
# 注册终止节点
builder.add_node("abort", abort)
# 注册通知节点,所有路径的必经出口
builder.add_node("notify", notify)
# 入口固定走校验
builder.add_edge(START, "validate")
# 校验后三分岔:小额自动、大额人工、校验失败终止
builder.add_conditional_edges("validate", need_human, ["auto_approve", "human_gate", "abort"])
# 人工决策后两分岔:批准去打款、驳回去终止
builder.add_conditional_edges("human_gate", after_decision, ["payout", "abort"])
# 自动通过也必须经过同一个打款节点,副作用只有一处入口
builder.add_edge("auto_approve", "payout")
# 打款后一定通知
builder.add_edge("payout", "notify")
# 终止也一定通知
builder.add_edge("abort", "notify")
# 通知完才结束
builder.add_edge("notify", END)
# 带 checkpointer 编译,interrupt 才有地方存现场
return builder.compile(checkpointer=checkpointer)
def pending_interrupts(snap):
"""从快照里取出「真正还在等人」的 interrupt。
不能直接用 snap.interrupts:部分恢复之后,已经拍板过的任务也会
在重放里再抛一次 interrupt。用 next 过滤才准。
"""
# next 里是这一刻真正待执行的节点名,只认它们名下的 interrupt
return [i for t in snap.tasks if t.name in snap.next for i in t.interrupts]
def show(graph, cfg, title):
"""把一个工单的现场打印成人能读的样子。"""
# 读当前现场(第 22 章的 get_state)
snap = graph.get_state(cfg)
# 标题行,前面留一个空行便于分段
print(f"\n--- {title} ---")
# 按顺序打印审计流水
for line in snap.values["log"]:
# 缩进对齐,便于扫读
print(f" {line}")
# next 非空说明图还没跑完,此刻只可能是卡在审批点
if snap.next:
# 打印卡在哪个节点,这是排查的第一现场
print(f" [暂停中] next={snap.next}")
# 把待审内容原样打出来,真实系统里这块是给前端渲染的
for i in pending_interrupts(snap):
# value 就是 human_gate 里传给 interrupt() 的那个字典
print(f" [待审批] {i.value}")
else:
# next 为空且状态非空,说明已经跑到 END
print(f" [已完结] status={snap.values['status']}")
def cmd_submit(graph, ticket_id: str, amount: int, reason: str):
"""提交一张新工单。"""
# thread_id 用工单号,同一张单永远落在同一条线程上
cfg = {"configurable": {"thread_id": ticket_id}}
# 幂等保护:重复提交会在旧状态上继续跑、二次打款,必须先挡住。
# created_at 是唯一可靠的「这个 thread 存在过吗」信号(见 §10.3)
if graph.get_state(cfg).created_at is not None:
# 已存在就只展示现场,绝不再 invoke
show(graph, cfg, f"{ticket_id} 已存在,拒绝重复提交")
# 提前返回,这是防止二次打款的那一行
return
# 全新工单,喂完整初始状态启动
graph.invoke(
{"ticket_id": ticket_id, "amount": amount, "reason": reason,
# 三个字段给显式初值,避免节点里读到 KeyError
"status": "", "approver": "", "log": []},
cfg,
)
# 跑到 END 或者卡在审批点,都用同一个函数展示
show(graph, cfg, f"提交 {ticket_id}")
def cmd_pending(graph, checkpointer):
"""列出所有卡在人工审批的工单。"""
# 必须先把 thread_id 收集完再查 get_state:
# 迭代 list() 的游标期间用同一个连接发起查询,SqliteSaver 会静默死锁
thread_ids = sorted({
ckpt.config["configurable"]["thread_id"] for ckpt in checkpointer.list(None)
})
# 清单标题
print("\n--- 待审批清单 ---")
# 计数器,用来判断要不要打「无待办」
found = 0
# 游标已经关闭,现在逐个查现场是安全的
for tid in thread_ids:
# 逐个读快照
snap = graph.get_state({"configurable": {"thread_id": tid}})
# 取出真正等人的 interrupt
waiting = pending_interrupts(snap)
# 有等人的才算待办
if waiting:
# interrupt 的 value 就是当初装给前端的那个字典
v = waiting[0].value
# 只展示业务字段,不展示 interrupt 的 id
print(f" {tid}: {v['amount']} 元({v['reason']})")
# 计数加一
found += 1
# 一条都没有时给一句明确的空态提示
if not found:
# 空态也要有输出,否则分不清「没待办」和「命令没跑」
print(" (无待办)")
def cmd_decide(graph, ticket_id: str, decision: dict):
"""喂进一个人工决策,让工单从断点继续。"""
# 同样按工单号定位线程
cfg = {"configurable": {"thread_id": ticket_id}}
# 动态 interrupt 只能用 Command(resume=...) 恢复,用 invoke(None) 会原地不动
graph.invoke(Command(resume=decision), cfg)
# 恢复后再展示一次现场
show(graph, cfg, f"处理 {ticket_id}")
def main(argv: list[str]):
"""CLI 入口:每次运行都是一个独立进程,状态全靠 sqlite 文件传递。"""
# 用 with 管理连接,退出时正常关闭
with SqliteSaver.from_conn_string(DB) as cp:
# 每个进程都要重新装一次图,图的定义是代码、状态才在文件里
graph = build(cp)
# 第一个参数是子命令,没给就画图
action = argv[0] if argv else "demo"
# 提交:工单号、金额、事由
if action == "submit":
# 事由是可选参数,没给就写「未说明」
cmd_submit(graph, argv[1], int(argv[2]), argv[3] if len(argv) > 3 else "未说明")
# 查待办
elif action == "pending":
# 需要把 checkpointer 也传进去,因为要列出所有 thread
cmd_pending(graph, cp)
# 原样批准:工单号、审批人
elif action == "approve":
# 决策字典的 type 决定 human_gate 走哪个分支
cmd_decide(graph, argv[1], {"type": "approve", "by": argv[2]})
# 改额批准:工单号、审批人、新金额
elif action == "edit":
# 多带一个 amount 字段,human_gate 会把它写回状态
cmd_decide(graph, argv[1], {"type": "edit", "by": argv[2], "amount": int(argv[3])})
# 驳回:工单号、审批人、理由
elif action == "reject":
# 理由会落进审计流水,是对客解释的依据
cmd_decide(graph, argv[1], {"type": "reject", "by": argv[2], "reason": argv[3]})
# 没给子命令就打印结构图,方便贴到文档里
else:
# draw_mermaid 输出可以直接粘到 Markdown 的 mermaid 代码块里
print(graph.get_graph().draw_mermaid())
# 只有直接运行时才执行,被 import 时只提供 build/DB 供复用
if __name__ == "__main__":
# 去掉脚本名,把剩下的参数交给 main
main(sys.argv[1:])10.3 跑起来 #
每一条命令都是一个独立进程,这正是要演示的重点。Windows 下建议先 chcp 65001,并设置 set PYTHONIOENCODING=utf-8,否则中文参数会乱码。
python refund_flow.py submit T-001 99 "商品有瑕疵" [$] 实际打款 99 元 -> 工单 T-001
--- 提交 T-001 ---
受理工单 T-001:退款 99 元(商品有瑕疵)
金额 99 < 200,自动通过
打款成功 99 元
已通知客户,最终状态 approved(审批人 系统自动)
[已完结] status=approved小额一次跑完,没有暂停。再用同一个工单号提交一次:
--- T-001 已存在,拒绝重复提交 ---
受理工单 T-001:退款 99 元(商品有瑕疵)
金额 99 < 200,自动通过
打款成功 99 元
已通知客户,最终状态 approved(审批人 系统自动)
[已完结] status=approved判据是 created_at:全新的 thread_id 拿到的是一个空快照:
'从来没用过的-id': created_at = None values = {} next = ()
'T-001': created_at = '2026-09-04T03:46:52.103102+00:00' values = {...} next = ()注意不能用 snap.values 或 snap.next 判断:新线程的 values 是 {}、next 是 (),已跑完的线程 next 同样是 (),两者区分不开。created_at 是唯一可靠的「这个 thread 存在过吗」信号。
这个幂等保护不是可选的。 去掉那三行再提交一次,validate 会在旧状态上重跑一遍,payout 会再打一次款,log 里出现整套重复记录(实测):
受理工单 T-900:退款 99 元(试试)
金额 99 < 200,自动通过
打款成功 99 元
已通知客户,最终状态 approved(审批人 系统自动)
受理工单 T-900:退款 99 元(试试) <- 第二次提交从这里开始
金额 99 < 200,自动通过
打款成功 99 元 <- 钱又出去了一次
已通知客户,最终状态 approved(审批人 系统自动)提交三笔大额和一笔金额非法的:
--- 提交 T-002 ---
受理工单 T-002:退款 1280 元(七天无理由)
[暂停中] next=('human_gate',)
[待审批] {'ticket_id': 'T-002', 'amount': 1280, 'reason': '七天无理由', 'allowed': ['approve', 'edit', 'reject']}
--- 提交 T-003 ---
受理工单 T-003:退款 5000 元(客户投诉)
[暂停中] next=('human_gate',)
--- 提交 T-004 ---
受理工单 T-004:退款 8800 元(要求全额退)
[暂停中] next=('human_gate',)
--- 提交 T-005 ---
校验失败:金额必须大于 0
流程终止,未发生资金变动
已通知客户,最终状态 rejected(审批人 系统校验)
[已完结] status=rejectedT-005 值得留意:校验没过的单子一次都没惊动人,直接走 abort → notify 跑完,而且照样有通知记录。
换个进程查待办:
--- 待审批清单 ---
T-002: 1280 元(七天无理由)
T-003: 5000 元(客户投诉)
T-004: 8800 元(要求全额退)三种决策,各走一个独立进程:
python refund_flow.py approve T-002 "王经理"
python refund_flow.py edit T-003 "李总监" 2000
python refund_flow.py reject T-004 "张主管" "凭证不全" [$] 实际打款 1280 元 -> 工单 T-002
--- 处理 T-002 ---
受理工单 T-002:退款 1280 元(七天无理由)
王经理 批准原样退款
打款成功 1280 元
已通知客户,最终状态 approved(审批人 王经理)
[已完结] status=approved
[$] 实际打款 2000 元 -> 工单 T-003
--- 处理 T-003 ---
受理工单 T-003:退款 5000 元(客户投诉)
李总监 将金额改为 2000 元后批准
打款成功 2000 元
已通知客户,最终状态 approved(审批人 李总监)
[已完结] status=approved
--- 处理 T-004 ---
受理工单 T-004:退款 8800 元(要求全额退)
张主管 驳回:凭证不全
流程终止,未发生资金变动
已通知客户,最终状态 rejected(审批人 张主管)
[已完结] status=rejected三处关键结果:
- T-002 只打印了一次
[$] 实际打款。human_gate在恢复时重跑了(输出上看不出来,因为它没有 print),但payout是独立节点,只执行一次。这就是 §4 的规则在本章实战里的体现。 - T-003 申请 5000,审批人改成 2000,实际打款 2000。 改动经
human_gate写回状态,payout读到的就是新值。 - T-004 驳回后走了
abort,log里明确记着「未发生资金变动」,但照样进了notify,这正是第 25 章那条「所有分支汇入必经节点」的规则。
全部处理完再查一次,待办清空:
--- 待审批清单 ---
(无待办)不带参数运行会打印图结构,可以贴到文档或 PR 里:
---
config:
flowchart:
curve: linear
---
graph TD;
__start__([<p>__start__</p>]):::first
validate(validate)
auto_approve(auto_approve)
human_gate(human_gate)
payout(payout)
abort(abort)
notify(notify)
__end__([<p>__end__</p>]):::last
__start__ --> validate;
abort --> notify;
auto_approve --> payout;
human_gate -.-> abort;
human_gate -.-> payout;
payout --> notify;
validate -.-> abort;
validate -.-> auto_approve;
validate -.-> human_gate;
notify --> __end__;虚线是条件边,实线是固定边(第 23 章)。从图上一眼能看出两件事:payout 是动钱的唯一入口,notify 是所有路径的唯一出口。
10.4 验收清单 #
- 小额工单一次跑完,全程没有
[暂停中] - 大额工单挂起后,关掉终端再开一个,
pending仍能查到它 - approve / edit / reject 三条路都跑通,
edit的打款金额是改后的值 - 驳回的工单没有
[$] 实际打款输出,但有notify记录 - 金额填 0 的工单一次都没惊动人,但也走完了
notify - 重复提交被幂等保护挡住;去掉那三行,确认
payout会执行两次 - 把
human_gate里的interrupt()挪到函数最后一行,观察前面的代码怎么被执行两次 - 用
Command(resume=..., update={"amount": 1})恢复一笔,确认打款金额被外部改掉了
11. 实用约定与坑 #
约定
- 带
interrupt()的节点必须无副作用。 只读状态、只提问、只写决定。这是本章第一规则。 - 副作用独占一个节点,放在审批之后。 打款、发消息、写库都适用。别指望「写在
interrupt()之后」能救你(§4.3)。 - 生产用
SqliteSaver或PostgresSaver,InMemorySaver只配 Demo。 thread_id用业务主键(工单号、订单号),别用随机 UUID,业务主键可查、可幂等。- 对外暴露的提交入口一定加幂等保护。 检查
get_state(cfg).created_at是否为None。 interrupt()的参数要装齐前端渲染所需的一切,它是你和审批界面之间唯一的通道。- 审批结果用结构化字典(
{"type": "approve", "by": ...}),别用裸字符串,后面加字段不用改协议。 - 判断「谁还在等人」用返回值的
__interrupt__,或用next过滤tasks,别直接用snap.interrupts(§7.1)。 - 所有分支汇入必经的后处理节点(第 25 章),审批流尤其如此。
- 待办 id 现查现用。 interrupt 的
id跟 checkpoint 绑定,不要缓存在前端、更不要写进代码。
会大声报错的
- 多个 interrupt 并行时用单值 resume →
RuntimeError: When there are multiple pending interrupts...,必须传{id: 值},见 §7.1。 - 没有
checkpointer就Command(resume=)→RuntimeError: Cannot use Command(resume=...) without checkpointer;get_state则抛ValueError: No checkpointer set。 - 不加
subgraphs=True就访问tasks[0].state.next→AttributeError: 'dict' object has no attribute 'next',见 §7.2。 - HITL 中间件回传的
decisions条数不对 →ValueError,条数必须和被拦下的工具调用数一致,见 §9。
会静默出错的(更危险)
- 恢复时被中断的节点从头重跑 →
interrupt()之前的副作用执行两次,钱扣两次。本章头号坑,见 §4。 - 动态
interrupt()用invoke(None, cfg)恢复 → 图原地不动,又抛出同一个 interrupt,不报错;而且节点白跑一遍,副作用再来一次,见 §5.1。 - 忘了配
checkpointer→ 提交能用、审批点永远动不了,启动阶段完全没有提示,见 §3.1。 - 同一个
thread_id重复invoke初始输入 → 在旧状态上继续累加,副作用重复执行,log出现整套重复记录,见 §10.3。 - 用陈旧的 interrupt id 去 resume → 什么都不发生,不报错也不前进,见 §7.1。
- 部分恢复后直接读
snap.interrupts→ 已经批过的那条也在里面,待办列表会显示假待办,见 §7.1。 - 子图与父图共享带 reducer 的字段 → 子图返回完整状态被父图再追加一次,出现重复项。
add_messages按 id 去重所以看不出来,operator.add会暴露,见 §7.3。 - 工具内
interrupt()时同批还有别的有副作用的工具 →tools节点整体重跑,兄弟工具执行两次,见 §9.1。 - 恢复时顺手带
goto→ 下游节点执行两次(原边一次、goto一次),见 §6.1。 - 迭代
checkpointer.list()的同时调get_state()→SqliteSaver死锁,没有报错、直接卡住,而且同一连接后续所有查询都跟着卡死。先把thread_id收集成集合再查,见 §10.2。 - 审批界面直接展示 interrupt 的
id→ 它是内部标识,换条线程、换次运行都不一样。给用户看业务主键。
正常行为,别当 bug
invoke返回时带__interrupt__而不抛异常 → 暂停是正常返回状态,检查键、不要try/except。- 恢复后
[node] xxx 开始执行又打印一遍 → 这就是重放,符合预期,见 §4。 - 静态
interrupt_before不返回__interrupt__,snap.interrupts是空的 → 它本来就不携带数据,只能做调试,见 §5。 - 静态中断用
Command(resume=...)恢复不报错 → resume 值被安静忽略,图正常前进。只有反方向(动态用None)会出问题。 - 暂停期间进程可以退出 → 这正是持久化的意义,见 §8。
12. 练习 #
- 复现头号坑。 把 §10 的打款逻辑搬进
human_gate,分别试两个位置:放在interrupt()之前(数一数[$]打印几次),再放到interrupt()之后(这次是几次?)。想清楚为什么两个位置的结果不同,以及为什么第二种依然不该写。 - 复现幂等漏洞。 注释掉
cmd_submit里的三行幂等检查,对同一个工单号提交两次,观察log和打款次数。 - 写错恢复方式。 对一笔挂起的大额工单执行
graph.invoke(None, cfg),确认它既不报错也不前进;再在human_gate的interrupt()之前加一行 print,看这次「无效恢复」是不是也让它执行了一遍。 - 加二级审批。 金额超过 5000 需要两个人分别签字。提示:用 §7.1 的并行 interrupt,或者串两个
human_gate节点,想清楚两种做法在「谁能先批」上的差别。 - 加超时自动驳回。 在状态里记下
submitted_at,pending命令顺便把超过 24 小时的工单自动 reject。为什么这件事适合放在一个独立的定时进程里做? - 换成 PostgresSaver。 参照
md/PostgreSQL.md把后端换掉,确认图的代码一行都不用改。 - 接一个 Agent。 把第 25 章的
create_agent塞进validate和human_gate之间,让它先根据退款原因起草一段审批说明,写进interrupt()的内容里供审批人参考。注意别把 Agent 的messages和父图的log搅在一起(§7.3)。
13. 本章小结 #
checkpointer的第四种用法是「停在半路等外部输入」,前三种(读现场、读历史、改状态续跑)第 22 章已经讲过。既然状态能完整存下、能从任意点继续,暂停就只是「暂时不继续」。interrupt()让图在节点内暂停,invoke正常返回并带上__interrupt__;Command(resume=值)恢复,interrupt()这一次会直接返回那个值。没配checkpointer时第一次调用照样成功,只有恢复才报错,症状是「提交能用、审批点不动」。- 本章头号坑:恢复时被中断的节点从头重跑。 因为存的是节点执行之前的状态。准确的规则是:每次 resume 都从函数第一行重跑,已答过的
interrupt()直接返回答案,所以第 k 个interrupt()之前的代码会执行 k 次。实测钱扣了两次。规则:带interrupt()的节点必须无副作用,副作用独占下一个节点。 - 重放是 LangGraph 的常态,不是异常。 这也解释了第 24 章为什么要求路由函数无副作用。
- 动态和静态 interrupt 是两套东西:静态
interrupt_before不携带数据、节点不重跑、对副作用安全,适合调试;动态interrupt()能携带数据、节点会重跑、必须用Command(resume=),适合生产。动态用invoke(None)恢复会原地不动且不报错,还白白触发一次副作用;反过来静态用Command恢复则是安全的。 - HITL 三种决策(批准 / 改后批准 / 驳回)一个节点就能覆盖,后面接条件边分流。「改一改」优先在节点里处理,这样校验逻辑能和改动写在一起;
Command(update=...)会走 reducer,对operator.add字段是追加而不是覆盖。 - 并行的多个 interrupt 要用
{id: 值}恢复,可以只批一个。部分恢复之后snap.interrupts会把已办项也列出来,判断「谁还在等人」要用返回值的__interrupt__或用next过滤tasks;拿旧 id 去 resume 是静默无效。 - 跨进程暂停恢复是这个机制的真正价值:进程 A 提交后即可退出,几天后进程 B 读出现场并批准。换后端只改
compile(checkpointer=...)一个参数。 - 第 10 章的
HumanInTheLoopMiddleware底层就是interrupt()(源码里是decisions = interrupt(hitl_request)["decisions"]),它把提问放在独立的after_model节点里,所以恢复时不会重复烧模型调用。第三条路是把interrupt()写进工具函数,但要小心同批的兄弟工具会被重复执行。 - 两个只在真实系统里才会撞到的坑:同一个
thread_id重复提交会在旧状态上继续跑并重复打款,所以对外入口必须用created_at做幂等;迭代checkpointer.list()时调get_state()会让SqliteSaver静默死锁,要先把 id 收集完。
到这里,图的核心能力就全部拿下了:顺序、分支、循环、子图、暂停。下一章处理最后一块,高级流式:怎么让这些多步骤、带嵌套、可暂停的图,在前端给出流畅的实时反馈。