LangGraph 源码深潜:Interrupt 与 Human-in-the-Loop 机制解剖
核心问题: 人工确认如何从“业务上需要人审”落到 LangGraph 的暂停、持久化与恢复执行机制?
源码主线:
interrupt()→GraphInterrupt→checkpointer→StateSnapshot.interrupts→Command(resume=...)→ 节点重放 → 后续节点继续执行前置文章: 第 6 篇
StateGraph、第 7 篇Conditional Edge、第 8 篇Reducer、第 11 篇Checkpoint / Thread / Durable Execution依赖基线:
langgraph==1.2.7源码基线:
https://github.com/langchain-ai/langgraph/tree/1.2.7,发布 tag1.2.7阅读边界: 本文只分析 LangGraph 中
interrupt()、Command(resume=...)、checkpoint 与 HITL 的运行机制;不展开 LangChain Agent 的 HITL middleware 业务封装,不展开前端审批 UI,也不展开生产 checkpointer 的数据库部署细节。
0. 本篇在源码学习主线中的位置
前面的文章已经解释了:
StateGraph 负责把节点、边、状态 schema 编译成可执行图。
Conditional Edge 负责让图根据 state 决定下一个节点。
Reducer 负责让多个节点对同一 state key 的更新可合并。
Checkpoint / Thread 负责按 thread_id 保存状态快照,支持恢复、回放和 time travel。本篇继续进入 LangGraph 的生产级能力:
Workflow / Router / Reducer / Checkpoint ↓Interrupt 与 Human-in-the-loop ↓可暂停、可审批、可恢复的生产级 Agent Workflow本篇只解决:
interrupt()如何让图在节点内部暂停。- 暂停时 state、next task、interrupt payload 保存在哪里。
Command(resume=...)如何把人工输入传回interrupt()。- 为什么 interrupt 必须依赖 checkpointer 与
thread_id。 - 多个 interrupt 同时存在时如何恢复。
- 业务上如何把审批、编辑、补充信息、安全确认映射成 LangGraph 机制。
本篇不展开:
- LangChain Agent HITL middleware 的高级封装。
- LangGraph Platform 的部署式审批队列。
- 前端审批页面、通知系统和权限系统实现。
- 数据库型 checkpointer 的表结构和运维部署。
1. 本篇问题、学习目标与能力边界
1.1 核心问题
LangGraph 如何把“等待人工输入”变成一个可持久化、可恢复、可审计的图运行时机制?
1.2 学习目标
完成本篇后,读者必须能够:
- 从构建期解释为什么
interrupt()必须和checkpointer、thread_id配合使用。 - 从运行时解释
interrupt(value)如何暂停图,并把 payload 暴露给调用方。 - 从恢复路径解释
Command(resume=...)如何变成interrupt()的返回值。 - 从源码边界解释为什么恢复时节点会从头重新执行。
- 从工程角度判断哪些节点需要 HITL,哪些节点不应该用 interrupt 解决。
- 解释多 interrupt、并行 interrupt、复杂 payload、try/except 包裹 interrupt 的风险。
- 说明 Human-in-the-loop 不能替代权限、策略、幂等和业务 verifier。
1.3 能力边界
| 能力 | 本篇是否覆盖 | 说明 |
|---|---|---|
interrupt() 动态暂停 | 是 | 覆盖节点内部调用、payload 暴露与暂停语义 |
Command(resume=...) 恢复 | 是 | 覆盖单值恢复、按 interrupt id 恢复、多 interrupt 恢复 |
| checkpoint 保存位置 | 是 | 从机制层说明 state snapshot 与 pending interrupt 保存 |
thread_id 会话隔离 | 是 | 说明恢复必须使用同一 thread |
| 静态 breakpoint | 部分 | 只作为对比,不展开 interrupt_before/after |
| LangChain HITL middleware | 否 | 属于 Agent 高层封装,后续单独分析 |
| 生产审批 UI | 否 | 属于应用层,不属于 LangGraph 核心源码 |
| 权限 / Policy Engine | 否 | 本篇只说明 HITL 不能替代它们 |
2. 核心概念与最小心智模型
2.1 一句话定义
interrupt()是 LangGraph 在节点运行时触发的动态暂停机制,负责把“当前执行需要外部输入”转换成可恢复的 graph interrupt;它不负责审批业务本身,也不负责权限校验和工具幂等。
2.2 最小心智模型
节点执行到 interrupt(payload) ↓抛出可恢复中断 ↓图暂停,checkpoint 保存当前执行状态 ↓调用方拿到 payload,展示给人 ↓人输入 resume value ↓graph.invoke(Command(resume=value), 同一个 thread_id) ↓节点从头重放 ↓interrupt() 返回 resume value ↓节点继续执行并写回 state2.3 核心术语
| 术语 | 源码对象 | 语义 | 不要误解为 |
|---|---|---|---|
| 动态中断 | interrupt() | 节点内部主动暂停并向调用方暴露 payload | 普通 Python input() |
| 恢复命令 | Command(resume=...) | 再次调用图时传入恢复值 | 节点返回值 |
| 中断对象 | Interrupt | 保存 interrupt 的 value 与 id | 业务审批记录 |
| 图状态快照 | StateSnapshot | 当前 state、next、tasks、interrupts 的只读快照 | 数据库完整业务状态 |
| 线程 | thread_id | 一条图执行线 / 会话 / 长任务 ID | 操作系统线程 |
| 检查点 | checkpoint | 某个 super-step 边界保存的图状态 | 长期记忆 |
| 恢复列表 | task resume list | 同一 task 中多个 interrupt 的 resume 值列表 | 全局会话变量 |
| 静态断点 | interrupt_before/after | 节点前后暂停 | 节点内部动态判断 |
2.4 与相邻抽象的边界
| 对象 | 负责什么 | 不负责什么 | 与本篇对象的关系 |
|---|---|---|---|
StateGraph | 构造节点、边、状态 schema | 不负责人工审批 UI | HITL 图仍然由 StateGraph 构建 |
Checkpointer | 保存与读取 checkpoint | 不做业务权限判断 | interrupt 依赖它保存暂停位置 |
Command | 表达恢复值、状态更新或跳转 | 不自动审批业务 | Command(resume=...) 是恢复 interrupt 的入口 |
StateSnapshot | 暴露 state、next、tasks、interrupts | 不修改 state | 调用方可通过它观察 pending interrupts |
ToolNode | 执行工具调用 | 不天然保证人审 | 工具内部也可以调用 interrupt |
| Policy / Verifier | 做业务规则与安全校验 | 不负责图暂停恢复 | 高风险任务中必须与 HITL 配合 |
3. 完整执行链路
3.1 高层链路
用户输入旅行计划草案 ↓graph.invoke(input, config={"thread_id": ...}) ↓approval_node 执行 ↓interrupt(payload) ↓图暂停并返回 Interrupt ↓外部系统展示 payload ↓用户点击确认 / 拒绝 ↓graph.invoke(Command(resume=True/False), 同一个 config) ↓approval_node 从头重新执行 ↓interrupt() 返回 True/False ↓approval_node 写入 {"approved": ...} ↓final_node 执行 ↓END3.2 最小代码骨架
观察目标:这段代码只展示最小 HITL 链路,不展示前端、数据库审批单和权限系统。
from typing_extensions import TypedDictfrom langgraph.graph import StateGraph, START, ENDfrom langgraph.types import interrupt, Commandfrom langgraph.checkpoint.memory import InMemorySaver
class TravelState(TypedDict): plan: str approved: bool | None
def approval_node(state: TravelState): approved = interrupt({ "question": "是否确认使用该旅行计划?", "plan": state["plan"], }) return {"approved": approved}
def final_node(state: TravelState): if state["approved"]: return {"plan": state["plan"] + "\n用户已确认。"} return {"plan": "用户未确认,需要重新规划。"}
builder = StateGraph(TravelState)
builder.add_node("approval_node", approval_node)builder.add_node("final_node", final_node)
builder.add_edge(START, "approval_node")builder.add_edge("approval_node", "final_node")builder.add_edge("final_node", END)
graph = builder.compile(checkpointer=InMemorySaver())
config = {"configurable": {"thread_id": "travel-approval-001"}}
graph.invoke( {"plan": "东京 5 天行程草案", "approved": None}, config=config,)
graph.invoke( Command(resume=True), config=config,)第一轮调用不会正常走到 final_node,而是在 approval_node 中暂停。
第二轮调用不会重新传原始 state,而是传 Command(resume=True),并且必须使用同一个 thread_id。
3.3 对象流转
| 阶段 | 输入类型 | 核心函数 | 输出类型 | 状态变化 |
|---|---|---|---|---|
| 构建图 | StateGraph(TravelState) | add_node() / add_edge() | StateGraph | 注册节点与边 |
| 编译图 | checkpointer=InMemorySaver() | compile() | CompiledStateGraph | 图获得持久化能力 |
| 初次执行 | dict + config.thread_id | invoke() | interrupt output / partial result | 执行到 approval_node |
| 触发暂停 | payload | interrupt() | 抛出可恢复中断 | checkpoint 保存 pending interrupt |
| 外部审批 | Interrupt.value | 调用方展示 | 人工输入 | 图不运行 |
| 恢复执行 | Command(resume=True) | invoke() | final state | approved=True 写入 state |
| 结束 | state | final_node | dict update | 图到达 END |
3.4 时序链路
Caller │ │ graph.invoke(initial_state, config(thread_id)) ▼CompiledStateGraph / Pregel │ │ 执行 approval_node ▼approval_node │ │ interrupt(payload) ▼GraphInterrupt │ │ 保存 checkpoint,返回 Interrupt(value, id) ▼Caller / UI │ │ 用户确认 ▼Caller │ │ graph.invoke(Command(resume=True), same config) ▼CompiledStateGraph / Pregel │ │ 重放 approval_node,interrupt() 返回 True ▼approval_node │ │ return {"approved": True} ▼final_node │ │ return {"plan": "...用户已确认。"} ▼Caller3.5 正常结束条件
一次 HITL 图正常完成的条件是:
1. 节点触发 interrupt 后,调用方拿到 pending interrupt。2. 外部系统用同一个 thread_id 传入 Command(resume=...)。3. interrupt() 在恢复执行中返回 resume value。4. 后续节点完成执行。5. 图到达 END。最终结果保存在:
graph.invoke(Command(resume=...), config) 的返回 state如果使用 streaming API,则可以通过:
stream.interruptsstream.output观察暂停与最终结果。
4. 源码地图、关键文件与阅读顺序
4.1 核心目录
langgraph/├── types.py├── graph/│ └── state.py├── pregel/│ ├── main.py│ ├── loop.py│ ├── runner.py│ └── io.py└── checkpoint/ ├── base/ └── memory/4.2 关键文件
| 优先级 | 文件 | 核心对象 | 阅读目的 |
|---|---|---|---|
| 1 | langgraph/types.py | interrupt / Command / Interrupt / StateSnapshot | 理解 HITL 的公开类型契约 |
| 2 | langgraph/graph/state.py | StateGraph.compile() | 理解 checkpointer 如何进入编译后的图 |
| 3 | langgraph/pregel/main.py | Pregel.invoke() / stream() / get_state() | 理解图运行、恢复和状态读取入口 |
| 4 | langgraph/pregel/loop.py | Pregel loop | 理解 super-step、checkpoint 和 interrupt 处理 |
| 5 | langgraph/pregel/runner.py | task 执行 | 理解节点任务如何捕获 interrupt |
| 6 | langgraph/checkpoint/base | BaseCheckpointSaver | 理解 checkpoint 保存读取协议 |
| 7 | langgraph/checkpoint/memory | InMemorySaver | 理解调试型 checkpointer 的存储方式 |
4.3 推荐阅读顺序
1. `types.py::interrupt`2. `types.py::Command`3. `types.py::Interrupt`4. `types.py::StateSnapshot`5. `state.py::StateGraph.compile`6. `pregel/main.py::Pregel.invoke / stream`7. `pregel/main.py::get_state`8. checkpoint saver 的 get_tuple / put / put_writes9. interrupt 文档中的规则与多 interrupt 示例4.4 不建议的阅读顺序
不建议一开始阅读 Pregel loop 的全部源码。原因是 Pregel 运行时同时处理:
streamcheckpointretrytimeoutcachesubgraphdebuginterruptwriteschannels如果没有先理解 interrupt() 与 Command 的公开契约,很容易把 HITL 机制误解为普通异常处理或普通状态更新。
5. 对象模型、继承关系与协议边界
5.1 核心对象关系
StateGraph ↓ compile(checkpointer)CompiledStateGraph / Pregel ↓ invoke / streamPregel task ↓ node functioninterrupt() ↓ GraphInterruptCheckpointSaver ↓ persisted checkpointCommand(resume=...) ↓ resume valuenode replay5.2 对象职责
| 对象 | 生命周期 | 输入 | 输出 | 核心职责 |
|---|---|---|---|---|
interrupt() | 运行时节点内部 | JSON-serializable payload | resume value 或 GraphInterrupt | 暂停并暴露外部输入请求 |
Interrupt | 暂停结果 | payload + id | 可展示对象 | 让调用方知道需要人输入什么 |
Command | 恢复调用或节点返回 | resume/update/goto | 命令对象 | 表达恢复值、状态更新或跳转 |
Checkpointer | 图编译后运行期间 | checkpoint tuple / writes | 持久化记录 | 保存暂停点和状态快照 |
StateSnapshot | 状态读取时 | checkpoint | snapshot | 暴露当前 values、next、tasks、interrupts |
Pregel | 编译后运行时 | input + config | state / stream parts | 执行图、调度任务、处理 checkpoint |
thread_id | 每次调用 config | string | checkpoint namespace | 区分不同会话或任务执行线 |
5.3 协议边界
interrupt 协议 负责:暂停当前节点,暴露 payload,恢复后返回 resume value。 不负责:保存业务审批单、校验审批人权限、保证外部动作幂等。
checkpoint 协议 负责:保存图状态、next task、writes 和 pending interrupt。 不负责:跨用户长期记忆、审批业务表、业务审计报表。
Command 协议 负责:把 resume value 传回中断点,或作为节点返回值表达 update/goto。 不负责:判断 resume value 是否符合业务规则。
thread_id 协议 负责:定位同一条图执行线。 不负责:认证用户身份。5.4 稳定接口与内部实现
| 类型 | 对象 | 文章中的使用原则 |
|---|---|---|
| 公共 API | interrupt() | 可用于工程示例 |
| 公共 API | Command(resume=...) | 可用于恢复中断 |
| 公共 API | graph.invoke() / graph.stream_events() | 可用于运行与恢复 |
| 公共 API | graph.get_state() | 可用于读取 pending interrupts |
| 扩展接口 | BaseCheckpointSaver | 可实现自定义持久化 |
| 内部实现 | Pregel loop 细节 | 只用于理解,不建议业务代码依赖 |
| 内部实现 | resume list / task id 细节 | 只解释行为,不直接依赖 |
6. 源码阅读策略与证据标准
6.1 本篇阅读策略
先读公开契约:interrupt / Command / Interrupt ↓再读图编译:checkpointer 如何进入 compiled graph ↓再读运行时:invoke / stream 如何驱动节点 ↓再读暂停路径:GraphInterrupt 如何被捕获并写入 snapshot ↓再读恢复路径:Command(resume=...) 如何进入同一个 thread ↓最后读多 interrupt 与异常规则6.2 证据等级
| 标记 | 含义 | 写作要求 |
|---|---|---|
| 源码事实 | 当前正式版源码可直接证明 | 使用 1.2.7 tag 源码链接 |
| 官方契约 | 官方文档或 API Reference 明确说明 | 使用官方文档链接 |
| 简化伪代码 | 压缩真实控制流 | 明确标注不是源码逐字复制 |
| 作者推断 | 根据调用链得出的设计理解 | 使用“从调用关系可以推断” |
| 工程建议 | 面向项目实践 | 说明适用条件 |
6.3 本篇证据清单
| 结论 | 证据类型 | 文件或文档 | 定位 |
|---|---|---|---|
interrupt() 会暂停图执行并暴露 payload | 官方契约 / 源码事实 | types.py / Interrupt docs | interrupt |
Command.resume 可为单值或 interrupt id 映射 | 官方契约 / 源码事实 | types.py::Command | resume 字段 |
| interrupt 依赖 checkpointer | 官方契约 | Interrupt docs / API Reference | 规则说明 |
| 恢复时节点从头重新执行 | 官方契约 | Interrupt docs | resume key points |
| 多 interrupt 推荐按 id 映射恢复 | 官方契约 | Interrupt docs | Handling multiple interrupts |
| StateSnapshot 包含 interrupts | API Reference / 源码事实 | types.py::StateSnapshot | interrupts 字段 |
7. 构建期源码解剖
本章回答:
为了让图具备暂停和恢复能力,构建期需要准备哪些对象?
7.1 构建期职责
| 输入 | 归一化动作 | 构建结果 |
|---|---|---|
TravelState | 解析 state schema | 图知道有哪些 state key |
approval_node | 包装为可运行节点 | 节点可被 Pregel 调度 |
interrupt() | 不在构建期执行 | 只作为节点函数内的运行时调用 |
checkpointer | 校验并挂到 compiled graph | 图获得持久化能力 |
thread_id | 构建期不需要 | 运行时 config 提供 |
Command(resume=...) | 构建期不需要 | 恢复时作为 input 提供 |
7.2 构建期总链路
StateGraph(TravelState) ↓add_node("approval_node", approval_node) ↓add_node("final_node", final_node) ↓add_edge(START, "approval_node") ↓add_edge("approval_node", "final_node") ↓compile(checkpointer=InMemorySaver()) ↓CompiledStateGraph / Pregel7.3 compile(checkpointer=...) 源码解剖
职责与所处阶段
compile() 处于构建期末尾,负责把 builder 变成可执行图,并把 checkpointer 注入 Pregel runtime。
真实源码签名
以下签名来自当前正式版公开 API,具体泛型细节以源码为准:
def compile( self, checkpointer: Checkpointer = None, *, cache: BaseCache | None = None, store: BaseStore | None = None, interrupt_before: All | list[str] | None = None, interrupt_after: All | list[str] | None = None, debug: bool = False, name: str | None = None,) -> CompiledStateGraph: ...调用方与被调用方
业务代码 ↓builder.compile(checkpointer=InMemorySaver()) ↓StateGraph.compile ↓CompiledStateGraph / Pregel 构造输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | Checkpointer | None、False、True 或 BaseCheckpointSaver |
| 输出 | CompiledStateGraph | 可执行图 |
| 状态变化 | compiled graph 持有 checkpointer | 运行时可保存 / 读取 checkpoint |
| 副作用 | 无外部业务副作用 | 只是构造对象 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def compile(self, checkpointer=None, **options): # 1. 校验 graph 结构 self.validate( interrupt_before=options.get("interrupt_before"), interrupt_after=options.get("interrupt_after"), )
# 2. 规范化 checkpointer checkpointer = ensure_valid_checkpointer(checkpointer)
# 3. 构造 Pregel 节点、channel、边、branch compiled = CompiledStateGraph( builder=self, channels=self.channels, nodes={}, checkpointer=checkpointer, store=options.get("store"), cache=options.get("cache"), interrupt_before=options.get("interrupt_before"), interrupt_after=options.get("interrupt_after"), debug=options.get("debug"), name=options.get("name"), )
# 4. 把 builder 中的 nodes / edges / branches 附加到 compiled graph for node_name, node_spec in self.nodes.items(): compiled.attach_node(node_name, node_spec)
for start, end in self.edges: compiled.attach_edge(start, end)
for source, branches in self.branches.items(): for branch_name, branch_spec in branches.items(): compiled.attach_branch(source, branch_name, branch_spec)
# 5. 校验 compiled graph return compiled.validate()逐段解释
第一段校验图结构,保证 START、END、节点边界和静态 interrupt 配置合法。
第二段校验 checkpointer 类型。如果传入的不是合法 saver,构建期直接失败,而不是等 interrupt 发生时再失败。
第三段构造 compiled graph。此时 checkpointer 成为运行时的一部分,后续 invoke()、stream()、get_state() 都可以通过它保存和读取状态。
第四段把 builder 阶段注册的节点、边和 branch 附加到 compiled graph。interrupt() 本身不需要被注册,它是节点内部运行时调用。
正常路径
checkpointer 合法 ↓compiled graph 构造完成 ↓运行时可保存 checkpoint关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
checkpointer=None | 不启用持久化 | interrupt 无法正常用于可恢复暂停 |
checkpointer=InMemorySaver() | 启用内存保存 | 适合学习和测试 |
checkpointer 类型非法 | 抛出 TypeError | 构建期失败 |
interrupt_before/after 节点名非法 | 校验失败 | 构建期失败 |
设计原因与工程影响
HITL 的核心不是“暂停一下程序”,而是“暂停后还能在稍后、甚至不同进程中恢复”。因此 checkpointer 必须在构建期进入 compiled graph,否则运行时没有地方保存暂停点。
源码证据
langgraph/graph/state.py::StateGraph.compilelanggraph/types.py::ensure_valid_checkpointer- https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/types.py
7.4 Command 数据结构源码解剖
职责与所处阶段
Command 是恢复阶段和节点返回阶段共用的命令对象。本文重点关注作为图输入的 Command(resume=...)。
真实源码签名
@dataclassclass Command(Generic[N], ToolOutputMixin): graph: str | None = None update: Any | None = None resume: dict[str, Any] | Any | None = None goto: Send | Sequence[Send | N] | N = ()调用方与被调用方
外部调用方 ↓Command(resume=True) ↓graph.invoke(Command(resume=True), config) ↓Pregel runtime ↓恢复 pending interrupt输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | resume | 单个值或 {interrupt_id: value} 映射 |
| 输出 | Command | 传入图运行时的输入对象 |
| 状态变化 | 无直接 state 修改 | resume 值进入 interrupt 恢复机制 |
| 副作用 | 无 | 只是数据对象 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
@dataclassclass Command: graph: str | None = None update: Any | None = None resume: dict[str, Any] | Any | None = None goto: Send | Sequence[Send | str] | str = ()
def _update_as_tuples(self): if isinstance(self.update, dict): return list(self.update.items())
if is_list_of_key_value_tuples(self.update): return self.update
if is_annotated_dataclass_or_model(self.update): return extract_update_fields(self.update)
if self.update is not None: return [("__root__", self.update)]
return []逐段解释
resume 是本文核心,它和 interrupt() 配套使用。它可以是单个值,也可以是 interrupt id 到恢复值的映射。
update 与 goto 是节点返回 Command 时常用的能力,表示“更新状态”和“跳转节点”。但官方文档明确提醒:作为 invoke() 输入恢复 interrupt 时,应该使用 Command(resume=...),不要把 Command(update=...) 当作多轮对话输入。
正常路径
Command(resume=True) ↓同一个 thread_id 调用 graph ↓resume value 进入 pending interrupt ↓interrupt() 返回 True关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
resume 是单值 | 恢复下一个 pending interrupt | 简单审批场景 |
resume 是 {id: value} | 按 interrupt id 恢复 | 多 interrupt 场景 |
resume=None | 不提供恢复值 | pending interrupt 不会被解决 |
update/goto 作为 invoke 输入 | 不符合推荐用法 | 容易误解恢复语义 |
设计原因与工程影响
Command 把“继续运行”从普通输入 dict 中分离出来,避免框架无法区分“这是新的用户输入”还是“这是恢复暂停点的值”。
源码证据
langgraph/types.py::Command- https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/types.py
7.5 Interrupt 与 StateSnapshot 类型源码解剖
职责与所处阶段
Interrupt 是暂停暴露给调用方的信息对象;StateSnapshot 是读取当前图状态时的快照对象。
真实源码签名
@dataclassclass Interrupt: value: Any id: str
class StateSnapshot(NamedTuple): values: dict[str, Any] | Any next: tuple[str, ...] config: RunnableConfig metadata: CheckpointMetadata | None created_at: str | None parent_config: RunnableConfig | None tasks: tuple[PregelTask, ...] interrupts: tuple[Interrupt, ...]调用方与被调用方
interrupt(value) ↓Interrupt(value, id) ↓result["__interrupt__"] 或 stream.interrupts ↓调用方展示 payload
graph.get_state(config) ↓StateSnapshot ↓snapshot.interrupts输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | value | 要展示给人的 payload |
| 输出 | Interrupt | 带 value 和 id |
| 状态变化 | StateSnapshot.interrupts | 保存 pending interrupt |
| 副作用 | 无业务副作用 | 只是运行时状态记录 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
class Interrupt: def __init__(self, value, id=None, ns=None): self.value = value
if id is not None: self.id = id elif ns is not None: self.id = hash_namespace(ns) else: self.id = placeholder_or_generated_id
class StateSnapshot(NamedTuple): values: dict next: tuple[str, ...] config: RunnableConfig metadata: dict | None created_at: str | None parent_config: RunnableConfig | None tasks: tuple[PregelTask, ...] interrupts: tuple[Interrupt, ...]逐段解释
Interrupt.value 是业务层需要展示给人的数据,例如审批问题、计划草案、工具参数或待编辑文本。
Interrupt.id 用来区分多个 pending interrupt。并行节点同时 interrupt 时,不能只传一个裸值,应该用 {interrupt_id: resume_value} 映射。
StateSnapshot.interrupts 让外部系统可以通过 get_state() 或 streaming 方式看到当前还有哪些中断待处理。
正常路径
节点 interrupt ↓runtime 生成 Interrupt ↓snapshot.interrupts 包含 pending interrupt ↓外部系统读取并展示关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
| 单个 interrupt | 可以裸值 resume | 最简单 |
| 多个 interrupt | 需要 id 映射 resume | 避免错配 |
| value 不可序列化 | checkpointer 可能失败 | 生产中禁止 |
| id 丢失 | 难以精确恢复 | 多 interrupt 风险 |
设计原因与工程影响
Interrupt 必须携带 id,否则并行审批或多节点中断无法稳定恢复。生产审批系统应该把 id、thread_id、payload、审批人、审批结果一起落库审计。
源码证据
langgraph/types.py::Interruptlanggraph/types.py::StateSnapshot- https://reference.langchain.com/python/langgraph/types/
7.6 构建期产物
| 产物 | 保存的信息 | 运行时用途 |
|---|---|---|
CompiledStateGraph | 节点、边、state schema、checkpointer | 执行和恢复 |
Checkpointer | checkpoint 保存读取协议 | 保存暂停点 |
Command 类型 | resume/update/goto 字段 | 恢复或控制跳转 |
Interrupt 类型 | value/id | 暴露 pending interrupt |
StateSnapshot 类型 | values/next/tasks/interrupts | 读取当前暂停状态 |
8. 运行时主链源码解剖
本章回答:
一次
invoke()如何运行到 interrupt,暂停后又如何通过Command(resume=...)继续运行?
8.1 运行时入口
| 调用方式 | 公开入口 | 核心内部入口 | 返回类型 |
|---|---|---|---|
| 同步调用 | graph.invoke(input, config) | Pregel.stream() / loop | final state 或 interrupt result |
| 异步调用 | graph.ainvoke(input, config) | async Pregel loop | final state 或 interrupt result |
| 流式调用 | graph.stream(...) | Pregel streaming | chunk / interrupt |
| 事件流 | graph.stream_events(..., version="v3") | v3 event stream | stream.interrupts / stream.output |
| 批量调用 | graph.batch([...]) | 多次 invoke | list[result] |
8.2 运行时总链路
graph.invoke(initial_state, config(thread_id)) ↓读取或创建 thread checkpoint ↓启动 Pregel loop ↓执行 approval_node task ↓node 内部调用 interrupt(payload) ↓抛出 GraphInterrupt ↓task 记录 interrupt ↓loop 在 super-step 边界保存 checkpoint ↓向调用方返回 Interrupt ↓graph.invoke(Command(resume=value), same config) ↓从同一 thread 读取 checkpoint ↓注入 resume value ↓重新执行 approval_node ↓interrupt() 返回 resume value ↓节点继续返回 partial update ↓后续节点执行8.3 interrupt(value) 源码解剖
职责与所处阶段
interrupt() 处于运行时节点内部,负责在需要外部输入的位置暂停当前 task。
真实源码签名
def interrupt(value: Any) -> Any: ...调用方与被调用方
approval_node ↓interrupt(payload) ↓scratchpad / runtime interrupt state ↓GraphInterrupt 或 resume value输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | Any | 建议 JSON-serializable payload |
| 输出 | Any | 恢复后返回 Command.resume 的值 |
| 状态变化 | pending interrupt | 首次执行时生成 |
| 副作用 | 抛出可恢复中断 | 让 task 暂停 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def interrupt(value): # 1. 获取当前 task 的运行时 scratchpad scratchpad = get_current_task_scratchpad()
# 2. 计算当前节点内第几次 interrupt 调用 index = scratchpad.interrupt_counter()
# 3. 如果当前 task 已经有对应 index 的 resume 值 if index < len(scratchpad.resume_values): return scratchpad.resume_values[index]
# 4. 如果 resume 是按 interrupt_id 映射传入,也尝试用当前 interrupt id 匹配 interrupt_id = make_interrupt_id_from_task_path_and_index( scratchpad.task_path, index, )
if interrupt_id in scratchpad.resume_map: value = scratchpad.resume_map[interrupt_id] scratchpad.resume_values.append(value) return value
# 5. 没有 resume 值,说明这是首次执行到这个 interrupt interrupt_obj = Interrupt( value=value, id=interrupt_id, )
# 6. 抛出可恢复中断,让 Pregel task 停止 raise GraphInterrupt((interrupt_obj,))逐段解释
第一段获取当前 task 的 scratchpad。interrupt() 不是普通函数,它需要知道自己在哪个 graph task 中执行。
第二段计算 index。这个 index 是理解多个 interrupt 的关键:同一节点内多个 interrupt 是按调用顺序匹配 resume 值的。
第三段处理恢复路径。如果当前 task 已有 resume 值,说明这是恢复执行,interrupt() 不再暂停,而是直接返回值。
第四段处理多 interrupt id 映射恢复。并行节点同时中断时,推荐用 {interrupt_id: value} 精确恢复。
第五段创建 Interrupt 对象。它包含 payload 和 id,用于展示给调用方。
第六段通过特殊异常打断当前 task。这个异常不能被业务代码吞掉,否则 graph runtime 捕获不到中断。
正常路径
首次执行:interrupt(payload) ↓没有 resume 值 ↓raise GraphInterrupt ↓图暂停
恢复执行:interrupt(payload) ↓发现 resume 值 ↓return resume_value ↓节点继续执行关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
| 没有 resume 值 | 抛出 GraphInterrupt | 图暂停 |
| 有按 index 的 resume 值 | 返回值 | 节点继续 |
| 有按 id 的 resume 值 | 返回对应值 | 精确恢复 |
| 被 try/except 捕获 | runtime 接收不到 | 中断失效 |
| payload 不可序列化 | checkpoint 可能失败 | 生产风险 |
设计原因与工程影响
interrupt() 设计成抛出特殊异常,是为了从任意节点内部立即跳出当前执行,把控制权交还给 graph runtime。业务代码不能把它当成普通异常处理。
源码证据
langgraph/types.py::interrupt- https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/types.py
8.4 初次执行暂停链路源码解剖
职责与所处阶段
初次执行阶段负责运行到审批节点,并在 interrupt() 处暂停。
真实源码签名
公开入口:
graph.invoke(input: InputT, config: RunnableConfig | None = None, **kwargs) -> OutputT调用方与被调用方
graph.invoke(initial_state, config) ↓Pregel.stream / loop ↓PregelRunner.tick ↓approval_node ↓interrupt() ↓GraphInterrupt ↓checkpoint save输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | dict | 初始 state |
| Config | RunnableConfig | 必须包含 thread_id |
| 输出 | interrupt result / stream interruption | 告知调用方需要人工输入 |
| 状态变化 | checkpoint | 保存 values、next、tasks、interrupts |
| 副作用 | checkpointer 写入 | 保存暂停点 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def invoke(input_state, config): # 1. 规范化 config,读取 thread_id config = ensure_config(config) thread_id = config["configurable"]["thread_id"]
# 2. 从 checkpointer 读取当前 thread 的 checkpoint checkpoint = checkpointer.get_tuple(config)
# 3. 如果没有 checkpoint,创建初始 checkpoint if checkpoint is None: state = initialize_state(input_state) next_tasks = ("approval_node",) else: state = read_state_from_checkpoint(checkpoint) next_tasks = checkpoint.pending_tasks
# 4. 执行当前 super-step 中的 tasks task_results = [] task_interrupts = []
for task in next_tasks: try: result = run_node(task, state, config) task_results.append(result)
except GraphInterrupt as interrupt_exc: task_interrupts.extend(interrupt_exc.interrupts)
# 5. 如果出现 interrupt,不继续调度后续节点 if task_interrupts: snapshot = build_state_snapshot( values=state, next=next_tasks, tasks=mark_tasks_interrupted(next_tasks, task_interrupts), interrupts=tuple(task_interrupts), )
# 6. 保存 checkpoint checkpointer.put( config=config, checkpoint=snapshot_to_checkpoint(snapshot), metadata=build_metadata(source="loop", writes=None), )
# 7. 返回 interrupt 给调用方 return { "__interrupt__": tuple(task_interrupts) }
# 8. 没有 interrupt,应用节点写入并继续调度 state = apply_writes(state, task_results) next_tasks = compute_next_tasks(state)
return continue_until_end(state, next_tasks)逐段解释
第一段读取 thread_id。它是 checkpoint 的主索引,不同 thread_id 对应不同会话或任务执行线。
第二段尝试读取历史 checkpoint。如果这是第一次执行,就没有 checkpoint;如果是恢复或继续执行,就会读取到旧 checkpoint。
第四段执行 task。节点内部如果调用 interrupt(),不会正常返回,而是抛出 GraphInterrupt。
第五段发现 interrupt 后,不继续往后调度。因为此时业务需要外部输入。
第六段保存 checkpoint。此时保存的是“图暂停在这里”的状态,包括当前 values、next、tasks 和 pending interrupt。
第七段把 interrupt 暴露给调用方。默认 invoke() 可能通过 __interrupt__ 返回,事件流会通过 stream.interrupts 暴露。
正常路径
initial invoke ↓approval_node ↓interrupt ↓checkpoint saved ↓caller receives Interrupt关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
| 节点正常返回 | 应用 state update | 继续下一节点 |
| 节点 interrupt | 保存 checkpoint | 暂停 |
| 缺少 thread_id | 无法定位 checkpoint | 持久化/恢复失败 |
| 无 checkpointer | 无法可靠恢复 | interrupt 不可用 |
| 节点普通异常 | 按 retry/异常机制处理 | 不等于 HITL |
设计原因与工程影响
HITL 的暂停必须发生在 super-step 边界被保存下来。否则外部系统拿到审批请求后,图状态可能丢失,无法在用户确认后继续。
源码证据
langgraph/pregel/main.py::invokelanggraph/pregel/loop.pylanggraph/types.py::StateSnapshot- https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/types.py
8.5 Command(resume=...) 恢复链路源码解剖
职责与所处阶段
恢复链路负责把外部输入值传回暂停处,使 interrupt() 返回该值并继续执行。
真实源码签名
Command(resume: dict[str, Any] | Any | None = None)调用方与被调用方
graph.invoke(Command(resume=True), same config) ↓Pregel runtime ↓读取 checkpoint ↓把 resume value 注入 task scratchpad ↓重新执行 approval_node ↓interrupt() 返回 True输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | Command(resume=True) | 人工审批结果 |
| Config | 同一 thread_id | 定位暂停的 checkpoint |
| 输出 | final state | 图继续执行到 END |
| 状态变化 | approved=True | 节点返回 update 后写入 state |
| 副作用 | 读取并更新 checkpoint | 执行继续推进 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def invoke(input_value, config): config = ensure_config(config)
if isinstance(input_value, Command) and input_value.resume is not None: return resume_from_interrupt( command=input_value, config=config, )
return start_or_continue_normal_input(input_value, config)
def resume_from_interrupt(command, config): # 1. 用同一个 thread_id 读取暂停 checkpoint checkpoint = checkpointer.get_tuple(config)
if checkpoint is None: raise ValueError("No checkpoint found for resume")
# 2. 读取 pending interrupts 与 next tasks pending_tasks = checkpoint.pending_tasks pending_interrupts = extract_interrupts(checkpoint)
# 3. 把 resume 值转换成 task 可读取的 resume 数据 if isinstance(command.resume, dict): resume_map = command.resume else: resume_map = assign_to_next_interrupt( pending_interrupts, command.resume, )
# 4. 将 resume_map 注入 task scratchpad runtime = build_runtime_from_checkpoint( checkpoint=checkpoint, resume_map=resume_map, )
# 5. 重新执行被中断的节点 task_results = []
for task in pending_tasks: result = run_node_again( task=task, state=checkpoint.state, runtime=runtime, ) task_results.append(result)
# 6. 此次 interrupt() 不再抛中断,而是返回 resume value # approval_node 得到 approved=True 后正常返回 {"approved": True}
# 7. 应用写入,继续后续节点 state = apply_writes( checkpoint.state, task_results, )
next_tasks = compute_next_tasks(state)
return continue_until_end(state, next_tasks)逐段解释
第一段判断输入是否为 Command(resume=...)。这是恢复中断的特殊入口,不等同于普通 state dict 输入。
第二段用同一个 thread_id 读取 checkpoint。如果 thread_id 不一致,runtime 找不到之前暂停的位置。
第三段构造 resume map。单个 interrupt 可以传裸值;多个 interrupt 更安全的做法是传 {interrupt_id: value}。
第四段把 resume 值注入 task 的运行上下文。恢复时节点会从头重新执行,所以 interrupt() 必须能在同样位置读到对应 resume 值。
第五段重新执行节点。注意:不是从 interrupt() 下一行继续执行,而是从节点函数开头重跑。
第七段把节点返回的 partial update 合并进 state,然后继续调度后续节点。
正常路径
Command(resume=True) ↓读取 checkpoint ↓注入 resume value ↓重放 approval_node ↓interrupt() 返回 True ↓return {"approved": True} ↓final_node ↓END关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
| 同一 thread_id | 找到 checkpoint | 正常恢复 |
| thread_id 不同 | 找不到暂停点 | 恢复失败 |
| 单 interrupt + 裸值 | 恢复下一个 interrupt | 简单可用 |
| 多 interrupt + 裸值 | 可能不明确 | 不推荐 |
| 多 interrupt + id 映射 | 精确恢复 | 推荐 |
| 节点重放前有副作用 | 副作用重复执行 | 必须幂等或移到 interrupt 后 |
设计原因与工程影响
恢复时节点从头执行,是 LangGraph 能保持可恢复执行的一种设计:它不依赖 Python 调用栈持久化,而是依赖 checkpoint + deterministic replay。因此 interrupt 前的代码必须可重放。
源码证据
langgraph/types.py::Commandlanggraph/types.py::interruptlanggraph/pregel/main.py- https://reference.langchain.com/python/langgraph/types/
8.6 StateSnapshot 与 get_state() 读取暂停状态源码解剖
职责与所处阶段
get_state(config) 允许外部系统读取当前 thread 的最新状态,包括 pending interrupts。
真实源码签名
def get_state( self, config: RunnableConfig, *, subgraphs: bool = False,) -> StateSnapshot: ...调用方与被调用方
审批系统 ↓graph.get_state(config) ↓checkpointer.get_tuple(config) ↓StateSnapshot(values, next, tasks, interrupts)输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | RunnableConfig | 包含 thread_id |
| 输出 | StateSnapshot | 当前图状态快照 |
| 状态变化 | 无 | 只读 |
| 副作用 | 无 | 不修改 checkpoint |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def get_state(config, subgraphs=False): # 1. 校验 checkpointer if self.checkpointer is None: raise ValueError("No checkpointer set")
# 2. 用 thread_id / checkpoint_id 读取 checkpoint tuple checkpoint_tuple = self.checkpointer.get_tuple(config)
if checkpoint_tuple is None: return empty_state_snapshot(config)
# 3. 从 checkpoint 读取 channel values values = read_channels( checkpoint_tuple.checkpoint, self.output_channels, )
# 4. 根据 pending sends / tasks 推断 next next_nodes = infer_next_tasks(checkpoint_tuple)
# 5. 构造 tasks,包含 error / interrupts / subgraph state tasks = build_pregel_tasks( checkpoint_tuple, include_subgraphs=subgraphs, )
# 6. 从 tasks 中收集 interrupts interrupts = tuple( interrupt for task in tasks for interrupt in task.interrupts )
# 7. 返回 StateSnapshot return StateSnapshot( values=values, next=next_nodes, config=checkpoint_tuple.config, metadata=checkpoint_tuple.metadata, created_at=checkpoint_tuple.checkpoint["ts"], parent_config=checkpoint_tuple.parent_config, tasks=tasks, interrupts=interrupts, )逐段解释
get_state() 不执行节点,只读取 checkpointer 保存的最新 checkpoint。
values 是当前 state 的可见值。
next 表示下一步将执行哪些节点。如果图暂停在 interrupt,next 往往指向被中断或待恢复的 task。
tasks 包含当前 step 的任务信息,任务中可能携带 interrupt。
interrupts 是外部审批系统最关心的字段。
正常路径
graph.get_state(config) ↓读取最新 checkpoint ↓构造 StateSnapshot ↓snapshot.interrupts 展示给外部系统关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
| 有 checkpoint | 返回 snapshot | 可观察暂停状态 |
| 无 checkpoint | 返回空或抛错 | 取决于实现 |
subgraphs=True | 递归读取子图状态 | 用于复杂 HITL |
| 无 checkpointer | 无法读取 | 抛错 |
设计原因与工程影响
生产系统不一定只依赖 invoke() 返回结果来处理审批,也可以把审批任务队列化:后台读取 snapshot.interrupts,创建审批单,用户审批后再 Command(resume=...)。
源码证据
langgraph/pregel/main.py::get_statelanggraph/types.py::StateSnapshot- https://reference.langchain.com/python/langgraph/types/
8.7 运行时主链总结
初次执行:input state ↓approval_node ↓interrupt(payload) ↓GraphInterrupt ↓checkpoint 保存 pending interrupt ↓caller 获取 Interrupt
恢复执行:Command(resume=value) ↓same thread_id ↓读取 checkpoint ↓节点从头重放 ↓interrupt() 返回 value ↓节点继续 ↓END9. 关键分支、异常与边界
本章回答:
当 HITL 场景不是单个审批按钮时,LangGraph 如何分流、恢复、终止或失败?
9.1 分支矩阵
| 分支类型 | 触发条件 | 核心函数 | 结果 |
|---|---|---|---|
| 单个审批 | 一个节点调用一次 interrupt() | Command(resume=value) | 裸值恢复 |
| 多个顺序 interrupt | 同一节点多次调用 interrupt() | index-based resume | 按顺序恢复 |
| 并行 interrupt | 多个节点同一步暂停 | Command(resume={id: value}) | 按 id 恢复 |
| 审批拒绝 | resume=False | route / node logic | 进入 cancel / replan |
| 审批编辑 | resume=dict/str | node return update | 写回编辑内容 |
| 无 checkpoint | 未配置 saver | interrupt() | 无法可靠恢复 |
| 错误捕获 | try/except 包住 interrupt | GraphInterrupt 被吞 | 中断失效 |
| side effect 重放 | interrupt 前执行外部写操作 | node replay | 可能重复执行 |
9.2 同步与异步分支
| 维度 | 同步路径 | 异步路径 |
|---|---|---|
| 入口 | invoke() / stream() | ainvoke() / astream() |
| 调度方式 | 同步 Pregel loop | async Pregel loop |
| 中断暴露 | __interrupt__ / stream chunks | async stream parts |
| 恢复方式 | Command(resume=...) | Command(resume=...) |
| 资源语义 | 阻塞当前调用 | 适合 WebSocket / 异步服务 |
9.3 Batch、Stream、Parallel 或路由分支
当前对象支持 streaming 和并行 interrupt,但 HITL 更推荐使用事件流驱动:
graph.stream_events(..., version="v3") ↓stream.interruptedstream.interruptsstream.outputinvoke() 也能工作,但它更适合简单调试。交互式前端、审批系统或长任务建议使用 stream events,因为它能同时观察消息、状态和 interrupt。
9.4 异常分类
| 异常类别 | 抛出位置 | 是否可恢复 | 处理策略 | 是否反馈上层 |
|---|---|---|---|---|
GraphInterrupt | interrupt() | 是 | 由 graph runtime 捕获并暂停 | 是 |
| 普通节点异常 | 节点业务代码 | 视情况 | retry / fallback / fail | 是 |
thread_id 缺失 | config 校验 / checkpoint | 否 | 要求调用方补充 | 是 |
| checkpointer 缺失 | interrupt 使用时 | 否 | compile 时启用 checkpointer | 是 |
| resume 错配 | 恢复路径 | 视情况 | 用 id 映射恢复 | 是 |
| payload 序列化失败 | checkpoint saver | 否 | 使用 JSON-serializable payload | 是 |
9.5 异常路径源码解剖
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def approval_node(state): # 错误示例:不要这样包住 interrupt try: approved = interrupt({"plan": state["plan"]}) except Exception as exc: # 这会吞掉 GraphInterrupt,runtime 看不到中断 return {"approved": False}
return {"approved": approved}
def safe_approval_node(state): # 正确示例:interrupt 与易出错业务逻辑分离 approved = interrupt({ "question": "是否确认?", "plan": state["plan"], })
try: validate_approval(approved) except ValueError: return {"approved": False}
return {"approved": approved}逐段解释
第一段错误示例把 interrupt() 放进宽泛的 try/except。由于 interrupt() 依赖特殊异常暂停图,吞掉异常就意味着 runtime 无法识别暂停。
第二段先调用 interrupt(),恢复后再校验人工输入。这样中断异常不会被业务异常处理吞掉。
9.6 Retry、Fallback 与恢复边界
| 机制 | 适用条件 | 不适用条件 | 幂等要求 |
|---|---|---|---|
| Retry | 节点普通网络失败 | 已经产生外部副作用 | 工具必须幂等 |
| Fallback | 审批系统不可用 | 高风险动作自动放行 | 不能降低安全边界 |
| Repair | 用户输入格式不合法 | 业务规则冲突 | 不应自动伪造审批 |
| HITL | 需要人工判断或确认 | 纯自动规则可判断 | interrupt 前代码可重放 |
| Policy | 明确规则可判定 | 主观质量评价 | 不依赖人工 |
| Verifier | 可机器校验 | 需要责任人审批 | 与 HITL 配合 |
9.7 停止条件与保护上限
正常结束:resume value 被 interrupt() 返回,后续节点到达 END。提前结束:审批拒绝后路由到 cancel / final_response。人工中断:interrupt() 触发 pending interrupt。框架保护:recursion_limit 防止恢复后循环无限执行。异常失败:缺少 checkpoint、thread_id 错误、payload 不可序列化、普通节点异常。9.8 多个 interrupt 同时存在时如何恢复
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
# 初次执行后,外部系统拿到多个 pending interruptsnapshot = graph.get_state(config)pending = snapshot.interrupts
# pending 可能类似:# [# Interrupt(id="int_a", value={"question": "审批 A"}),# Interrupt(id="int_b", value={"question": "审批 B"}),# ]
# 外部系统分别收集审批结果resume_map = { pending[0].id: True, pending[1].id: False,}
# 一次性恢复多个 interruptresult = graph.invoke( Command(resume=resume_map), config=config,)逐段解释
并行节点同时触发 interrupt 时,裸值 Command(resume=True) 无法表达哪个值对应哪个 interrupt。使用 {interrupt_id: value} 可以明确匹配。
同一节点内多个 interrupt 则依赖调用顺序。恢复时节点从头执行,遇到第一个 interrupt 消耗第一个 resume 值,遇到第二个 interrupt 消耗第二个 resume 值。因此不能重排、条件跳过或非确定性循环调用 interrupt。
9.9 能力边界
| 容易误判的能力 | 实际提供者 | 本篇对象的真实职责 |
|---|---|---|
| 审批权限 | 业务权限系统 | interrupt 只暂停并收集输入 |
| 审批审计 | 业务数据库 / 日志系统 | checkpoint 不是审批表 |
| 动作安全 | Policy / Verifier / Tool Guard | HITL 只是其中一层 |
| 幂等控制 | 工具层 / 业务 API | interrupt 恢复可能重放节点 |
| 长期记忆 | Store / DB | checkpoint 是 thread-scoped state |
| 前端交互 | UI / API 服务 | LangGraph 只暴露 payload 和 resume 入口 |
10. 扩展机制与框架协作
本章回答:
HITL 如何与 checkpoint、Command、conditional edge、tool、subgraph、streaming 协作?
10.1 扩展点总览
| 扩展点 | 扩展方式 | 执行时机 | 可修改内容 | 约束 |
|---|---|---|---|---|
interrupt() | 节点内调用 | 需要人工输入时 | 暂停 payload | payload 应可序列化 |
Command(resume=...) | 外部恢复调用 | 用户输入后 | resume value | 必须同 thread_id |
Command(goto=...) | 节点返回 | 恢复后路由 | 下一节点 | 只跳合法节点 |
Command(update=...) | 节点返回 | 恢复后写 state | partial update | 遵守 reducer |
checkpointer | compile 参数 | 每个 super-step | 保存 checkpoint | 生产用持久化 saver |
stream_events | runtime API | 运行期间 | 观察 interrupts | 适合交互式应用 |
get_state() | runtime API | 暂停后 | 查看 snapshot | 只读 |
ToolNode / tool | 工具内部 interrupt | 高风险工具前 | 暂停工具执行 | 工具副作用必须在 interrupt 后 |
10.2 审批后路由 Command(goto=...) 源码解剖
职责与所处阶段
审批节点可以在恢复后直接返回 Command(goto=...),把人工决定映射成下一节点。
真实源码签名
Command( update: Any | None = None, resume: dict[str, Any] | Any | None = None, goto: Send | Sequence[Send | N] | N = (),)调用方与被调用方
approval_node ↓interrupt(payload) ↓resume value ↓return Command(goto="proceed" or "cancel") ↓Pregel routing输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | approved | interrupt 返回值 |
| 输出 | Command | 控制后续节点 |
| 状态变化 | 可选 update | 可同时写 state |
| 副作用 | 改变图路由 | 进入 proceed/cancel |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def approval_node(state): approved = interrupt({ "question": "是否执行高风险动作?", "details": state["pending_action"], })
if approved: return Command( update={"approved": True}, goto="proceed", )
return Command( update={"approved": False}, goto="cancel", )逐段解释
interrupt() 只负责拿到人工输入。人工输入如何改变流程,由恢复后的节点代码决定。
Command(update=..., goto=...) 让审批节点同时完成两件事:写状态和路由。
正常路径
resume=True ↓approved=True ↓Command(update, goto="proceed") ↓proceed node
resume=False ↓approved=False ↓Command(update, goto="cancel") ↓cancel node关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
| 审批通过 | goto proceed | 执行动作 |
| 审批拒绝 | goto cancel | 取消动作 |
| resume 值格式错误 | 进入校验失败分支 | 重新询问或拒绝 |
| goto 非法节点 | 编译或运行校验失败 | 需要修正图 |
设计原因与工程影响
HITL 不一定要配合 add_conditional_edges();如果审批节点自己返回 Command(goto=...),可以把“审批结果 → 路由”放在同一个节点内。但复杂业务仍建议使用显式 conditional edge,便于测试和可视化。
源码证据
langgraph/types.py::Command- https://reference.langchain.com/python/langgraph/types/
10.3 工具调用中的 interrupt 协作
职责与所处阶段
工具内部也可以调用 interrupt(),用于在真正执行高风险工具前暂停。
真实源码签名
def interrupt(value: Any) -> Any: ...调用方与被调用方
ToolNode ↓tool.invoke() ↓tool function ↓interrupt(tool_call_review_payload) ↓GraphInterrupt输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | 工具参数 | 待审批动作 |
| 输出 | 审批结果 / 修改后的参数 | 继续执行工具 |
| 状态变化 | 取决于工具返回 | 工具执行后写回 |
| 副作用 | 可能有外部动作 | 必须在 interrupt 后执行 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def send_email_tool(to, subject, body): decision = interrupt({ "action": "send_email", "to": to, "subject": subject, "body": body, "question": "是否发送这封邮件?", })
if decision is False: return "用户拒绝发送邮件。"
if isinstance(decision, dict): subject = decision.get("subject", subject) body = decision.get("body", body)
# 高风险副作用必须放在 interrupt 之后 return actually_send_email(to, subject, body)逐段解释
interrupt() 之前只准备 payload,不执行副作用。
恢复后如果用户拒绝,工具直接返回取消结果。
如果用户修改参数,工具使用修改后的参数。
真正发送邮件、退款、下单等副作用必须在 interrupt() 之后,否则恢复重放会重复执行。
正常路径
工具准备参数 ↓interrupt 审批 ↓用户同意 / 编辑 / 拒绝 ↓工具继续或取消关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
| approve | 执行工具 | 正常 |
| reject | 不执行工具 | 安全取消 |
| edit | 修改参数后执行 | 人工修正 |
| interrupt 前有副作用 | 恢复重放时重复 | 严重风险 |
设计原因与工程影响
工具 HITL 是生产 Agent 中最常见的高风险场景,但它必须和工具幂等、权限系统、审计日志配合使用。
源码证据
langgraph/types.py::interruptlanggraph/prebuilt/tool_node.py- LangGraph Interrupts 官方文档
10.4 与 checkpoint 的协作
interrupt ↓GraphInterrupt ↓Pregel task interrupted ↓checkpoint 保存: values next tasks interrupts metadata ↓Command(resume) ↓读取 checkpoint ↓恢复执行checkpoint 是 HITL 的地基。没有 checkpoint,图暂停后外部系统无法安全地在稍后恢复。
10.5 与 conditional edge 的协作
有两种方式:
方式 A:审批节点返回普通 state update ↓route_by_approval(state) ↓进入 proceed/cancel
方式 B:审批节点直接返回 Command(goto=...) ↓跳转 proceed/cancel建议:
简单审批:Command(goto=...)复杂业务:state update + add_conditional_edges10.6 与 subgraph 的协作
子图里也可以 interrupt,但生产系统要注意:
1. 父图和子图的 checkpoint 继承或隔离策略。2. get_state(subgraphs=True) 才能观察子图状态。3. Command.PARENT 可用于从子图向父图发送控制命令。10.7 选择 interrupt 还是静态 breakpoint
| 条件 | 选择 interrupt | 选择 interrupt_before/after |
|---|---|---|
| 节点内部根据业务条件暂停 | 是 | 否 |
| 每次进入某节点前都暂停 | 否 | 是 |
| 需要展示自定义 payload | 是 | 部分 |
| 审批点在工具参数生成后 | 是 | 视情况 |
| 调试图执行 | 否 | 是 |
| 生产人工审批 | 是 | 视场景 |
11. 工程决策与适用场景
11.1 适用场景
| 场景 | 是否推荐 | 原因 |
|---|---|---|
| 高风险工具调用前审批 | 是 | 需要人确认才能执行 |
| LLM 生成内容人工编辑 | 是 | resume value 可作为编辑后内容 |
| 用户补充缺失槽位 | 是 | 中断等待用户输入 |
| 金融转账 / 退款 / 删除 | 是,但必须配合 Policy | HITL 只是防线之一 |
| 普通 FAQ | 否 | 不需要暂停图 |
| 可由规则自动判断的低风险分支 | 否 | 用 conditional edge 即可 |
| 模型质量提升 | 不优先 | 应使用 Reflection / Evaluator |
| 长期用户偏好保存 | 否 | 应使用 Store / Memory |
11.2 工程决策表
| 决策点 | 推荐选择 | 前提 | 风险 |
|---|---|---|---|
| checkpointer | 生产用 Postgres/Redis 等持久化 saver | 需要跨进程恢复 | InMemorySaver 重启丢失 |
| thread_id | 使用业务会话/任务 ID | 能稳定关联用户任务 | 用随机临时 ID 会无法恢复 |
| interrupt payload | JSON-serializable dict | 前端可展示 | 复杂对象不可序列化 |
| resume 方式 | 多 interrupt 用 id 映射 | 并行审批 | 裸值可能错配 |
| interrupt 位置 | 副作用之前 | 避免重复执行 | 放在副作用后无意义 |
| 审批记录 | 业务库单独落表 | 合规审计 | checkpoint 不等于审批审计 |
| 恢复 API | Command(resume=...) | 同 thread_id | 重新传初始 state 会错误 |
11.3 性能、可靠性与安全边界
性能: interrupt 本身不是重计算瓶颈,但恢复时节点会从头执行,interrupt 前的重计算会重复发生。
可靠性: 依赖 checkpointer、thread_id 和节点可重放性。中断前副作用不幂等会导致严重问题。
安全: HITL 不能替代权限、Policy、Verifier、审计和工具幂等。审批人身份必须由业务系统保证。
可观测性: 必须记录 thread_id、interrupt_id、payload、审批结果、审批人、恢复时间、恢复后状态。11.4 旅行规划助手中的落地建议
| 节点 | 是否需要 HITL | 原因 |
|---|---|---|
generate_draft_plan | 否 | 普通生成,可自动执行 |
evaluate_plan | 否 | 用 evaluator 即可 |
approval_node | 是 | 用户确认是否采用计划 |
book_hotel_tool | 是 | 涉及真实预订 |
send_email_tool | 是 | 有外部副作用 |
final_response | 通常否 | 仅输出文本 |
save_user_preference | 视情况 | 涉及长期记忆写入 |
12. 常见误区与源码纠正
12.1 误区:interrupt 就是 Python input
错误原因:
表面上 interrupt() 看起来像“问用户一个问题,然后拿到回答”。
源码事实:
interrupt() 首次执行会抛出可恢复中断;图暂停并保存 checkpoint;恢复时通过 Command(resume=...) 返回值。工程影响:
把它当 input() 会忽略 checkpoint、thread_id、节点重放和副作用幂等。
12.2 误区:恢复时从 interrupt 下一行继续执行
错误原因:
许多语言中的协程或断点会让人以为运行栈被保留。
源码事实:
LangGraph 恢复时节点从头重新执行;interrupt() 根据 resume 值在相同调用位置返回。工程影响:
interrupt 前的代码会重复执行。所有副作用必须放在 interrupt 后,或保证幂等。
12.3 误区:只要有 interrupt 就不需要 Policy
错误原因:
人工确认看起来比规则更强。
源码事实:
interrupt 只负责暂停和恢复;不校验审批人权限,不判断动作是否合法。工程影响:
高风险动作必须先过 Policy / Verifier,再进入 HITL,最后执行工具。
12.4 误区:Command(resume=True) 可以换一个 thread_id 用
错误原因:
把 resume 当成全局输入值。
源码事实:
thread_id 是 checkpoint 的定位指针;恢复必须使用触发 interrupt 时的同一个 thread_id。工程影响:
换 thread_id 会找不到 pending interrupt,或者错误恢复到其他会话。
12.5 误区:多个 interrupt 可以随便写
错误原因:
示例里经常只有一个 interrupt。
源码事实:
同一节点内多个 interrupt 按调用顺序匹配;并行多个 interrupt 推荐用 interrupt id 映射恢复。工程影响:
条件跳过、重排、循环调用 interrupt 会导致 resume 值错配。
12.6 误区:checkpoint 就是审批记录
错误原因:
checkpoint 里确实保存了 interrupt payload 和 state。
源码事实:
checkpoint 是图运行恢复数据;审批记录需要业务系统单独保存。工程影响:
合规审计、审批人身份、审批意见、权限判断不能只依赖 checkpoint。
12.7 误区:InMemorySaver 可以生产使用
错误原因:
学习示例都用 InMemorySaver()。
源码事实:
InMemorySaver 是内存实现,进程重启会丢失。工程影响:
生产 HITL 需要持久化 checkpointer,否则用户审批后可能无法恢复。
13. 最终心智模型与掌握检查
13.1 构建期心智模型
StateGraph ↓add_node / add_edge ↓compile(checkpointer) ↓CompiledStateGraph with persistence13.2 运行时心智模型
initial input + thread_id ↓node runs ↓interrupt(payload) ↓GraphInterrupt ↓checkpoint saves pending interrupt ↓external human input ↓Command(resume=value) + same thread_id ↓node replay ↓interrupt returns value ↓state update ↓END13.3 分支与异常心智模型
正常路径 → 中断后 resume,节点继续,图到 END拒绝路径 → resume=False,进入 cancel/replan编辑路径 → resume=edited_value,写回 state多个中断 → resume={interrupt_id: value}可恢复异常 → interrupt 被 graph runtime 捕获不可恢复异常 → thread_id/checkpointer/payload/普通异常导致失败保护上限 → recursion_limit 防止恢复后循环不终止13.4 一句话总结
interrupt()通过可恢复中断把节点执行暂停在运行时,checkpointer保存暂停点和 state,Command(resume=...)在同一thread_id上注入人工输入,恢复时节点从头重放并让interrupt()返回该输入,从而把 Human-in-the-loop 从业务概念落成可持久化、可恢复、可审计的图执行机制。
13.5 掌握检查
- 能说清
interrupt()和静态 breakpoint 的区别。 - 能说明为什么 interrupt 必须依赖 checkpointer。
- 能画出初次执行暂停链路。
- 能画出
Command(resume=...)恢复链路。 - 能解释为什么恢复时节点从头执行。
- 能说明 interrupt 前为什么不能放非幂等副作用。
- 能解释
thread_id在恢复中的作用。 - 能说明单 interrupt、多个顺序 interrupt、并行 interrupt 的恢复差异。
- 能说明
StateSnapshot.interrupts的作用。 - 能判断一个业务节点是否应该使用 HITL。
- 能解释 HITL 与 Policy / Verifier / 审计的边界。
14. 参考资料与下一篇衔接
14.1 官方概念文档
-
LangGraph Interrupts
https://docs.langchain.com/oss/python/langgraph/interrupts -
LangGraph Persistence
https://docs.langchain.com/oss/python/langgraph/persistence -
LangChain Human-in-the-loop
https://docs.langchain.com/oss/python/langchain/human-in-the-loop -
LangGraph Graph API
https://docs.langchain.com/oss/python/langgraph/graph-api
14.2 官方 API Reference
-
interrupt()
https://reference.langchain.com/python/langgraph/types/#langgraph.types.interrupt -
Command
https://reference.langchain.com/python/langgraph/types/#langgraph.types.Command -
Interrupt
https://reference.langchain.com/python/langgraph/types/#langgraph.types.Interrupt -
StateSnapshot
https://reference.langchain.com/python/langgraph/types/#langgraph.types.StateSnapshot -
Checkpointing
https://reference.langchain.com/python/langgraph/checkpoints/
14.3 官方源码
-
langgraph/types.py::interrupt
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/types.py -
langgraph/types.py::Command
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/types.py -
langgraph/types.py::Interrupt
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/types.py -
langgraph/types.py::StateSnapshot
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/types.py -
langgraph/pregel/main.py
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py -
langgraph/checkpoint/base
https://github.com/langchain-ai/langgraph/tree/1.2.7/libs/checkpoint/langgraph/checkpoint/base -
langgraph/checkpoint/memory
https://github.com/langchain-ai/langgraph/tree/1.2.7/libs/checkpoint/langgraph/checkpoint/memory
14.4 下一篇衔接
下一篇进入:
第 13 篇:Memory / Store 与长期记忆源码解剖需要继续回答:
checkpoint 和 store 的边界是什么?thread-scoped state 与 cross-thread memory 有什么区别?memory write policy 如何工程化?长期记忆如何避免污染上下文?Store 如何与 Agent State、Tool、Retriever 协作?