Command(resume=):把答案喂回去,图从原地复活
昨天看透了"停":节点抛 GraphInterrupt,现场落盘。今天看"续"——你用 Command(resume="yes") 再 stream 一次,图怎么就知道"这个 yes 是回答刚才那个 interrupt 的"?答案藏在一个叫 便签本(PregelScratchpad)的小对象里,和一个把恢复值"发牌"到便签本的函数 _scratchpad()。
痛点:恢复不是"新一次调用",是"接着上次"
app.stream({"x":0}, cfg),那不就是重新跑一遍吗?框架怎么区分"我要重开"和"我要接着上次那个中断答 yes"?而且这个 "yes" 要精准送到那一个 interrupt 手里,不能塞错。app.stream(Command(resume="yes"), cfg)。cfg 里的 thread_id 没变 → 框架去 checkpointer 把昨天冻结的现场读回来;Command.resume="yes" → 这就是"要喂给中断的答案"。框架把 "yes" 存成一条特殊的 RESUME 写,然后重跑那个节点。节点里的 interrupt() 这次一查便签本"有答案",直接返回 "yes",代码继续往下——就像从没停过。类比:昨天的存档是"游戏退到桌面"。今天的
Command(resume=) 是"双击图标、点继续",而不是"新建游戏"。thread_id 就是那个存档槽位。Command 对象:一个信封装四种意图
Command 定义在 types.py:759-808,是个能装多种"命令"的数据类:
# types.py:759
@dataclass(...)
class Command(Generic[N], ToolOutputMixin):
graph: str | None = None # 发给哪张图(None=当前,PARENT=父图)
update: Any | None = None # 顺便更新一下状态
resume: dict[str, Any] | Any | None = None # ★ 今天的主角:恢复值
goto: Send | Sequence[Send | N] | N = () # 顺便跳到某节点(D15 讲过)
PARENT: ClassVar[Literal["__parent__"]] = "__parent__" # types.py:808
resume本日核心。文档(types.py:768-772)说它两种形态:① 单个值——喂给"下一个待恢复的中断";② {中断id: 值} 映射——精确指定回答哪个中断(用昨天的 Interrupt.id)。update恢复的同时还能改状态。比如人不仅点了"通过",还顺手改了金额。goto恢复的同时指定下一步跳哪去(Command 本来就是 D15 的跨节点跳转对象,resume 只是它的一个字段)。graph=PARENT子图里可以 Command(graph=Command.PARENT, ...) 把命令发给父图——人在环跨层时用得上。resume 只是往这个信封里多加了一格"给中断的答案"。所以你能一次性 Command(resume="yes", update={"note":"已复核"}, goto="audit") ——恢复 + 改状态 + 指定去向,一个信封全办了。PregelScratchpad:任务级的"便签本"
昨天 interrupt() 里那个 scratchpad 究竟是什么?定义在 _scratchpad.py:8-19,短到一屏:
# _scratchpad.py:8
@dataclasses.dataclass(**_DC_KWARGS)
class PregelScratchpad:
step: int
stop: int
# call
call_counter: Callable[[], int]
# interrupt
interrupt_counter: Callable[[], int] # ① 节点里 interrupt 的编号器
get_null_resume: Callable[[bool], Any] # ② 取"通用恢复值"的函数
resume: list[Any] # ③ 这个任务的恢复值列表
# subgraph
subgraph_counter: Callable[[], int]
一个任务一本关键理解:便签本是任务级的(一次节点执行一本),不是全图共享。所以昨天文档说"恢复值 scoped 到具体任务、不跨任务共享"。interrupt_counter昨天 L02 见过:给节点里的 interrupt 从 0 开始编号,让第 idx 个 interrupt 对上第 idx 个恢复值。resume: list已认领的恢复值列表。Command(resume=) 里针对这个任务的值会被放进来。interrupt 按 idx 从这里取。get_null_resume取"通用恢复值"(没指定 id 的那个单值)的函数。谁先叫谁拿走(消费型)。L05 细看。interrupt_counter 是个可调用的计数器。为什么?见 L07 的 LazyAtomicCounter:并发下 +=1 不安全,用函数封装原子操作。_scratchpad():从存档里"发牌"恢复值
便签本由谁填?由 _scratchpad()(_algo.py:1280)在每次准备任务时构造。它从 pending_writes(还记得吗,昨天 RESUME 就存这里)里找恢复值:
# _algo.py:1280
def _scratchpad(parent_scratchpad, pending_writes, task_id,
namespace_hash, resume_map, step, stop):
if len(pending_writes) > 0:
# find global resume value —— 找"通用恢复值"(谁都能领的那份)
for w in pending_writes:
if w[0] == NULL_TASK_ID and w[1] == RESUME: # ① 特殊 task_id + RESUME
null_resume_write = w
break
else:
null_resume_write = None
# find task-specific resume value —— 找"点名给本任务的"恢复值
for w in pending_writes:
if w[0] == task_id and w[1] == RESUME: # ② 就是本 task_id 的 RESUME
task_resume_write = w[2]
if not isinstance(task_resume_write, list):
task_resume_write = [task_resume_write]
break
else:
task_resume_write = []
# find namespace and task-specific resume value —— 按中断 id 精确映射
if resume_map and namespace_hash in resume_map: # ③ Command(resume={id:值})
mapped_resume_write = resume_map[namespace_hash]
task_resume_write.append(mapped_resume_write)
else:
null_resume_write = None
task_resume_write = []
NULL_TASK_ID + RESUME分支①:Command(resume="单个值") 会被存成 (NULL_TASK_ID, RESUME, 值)——不点名,谁的 interrupt 先要谁拿。这就是 get_null_resume 要取的那份。w[0] == task_id分支②:点名给这个任务的恢复值(多值场景),直接进 task_resume_write 列表,成为便签本的 resume。resume_map[namespace_hash]分支③:Command(resume={中断id: 值}) 时,按命名空间哈希(=昨天 Interrupt.id 的来源)精确匹配到本任务,追加进恢复列表。else: []没有 pending_writes(全新运行)→ 恢复列表空。于是 interrupt 走"抛异常"分支。首尾呼应昨天 L02。Command(resume="yes") 给个单值,框架不用你操心它该给谁——"谁在等谁拿"。点名值(task_id / 按 id 映射)服务复杂场景:并行分支同时中断了多个(比如 Send 扇出的多个任务各自 interrupt),这时一个笼统的单值无法区分,必须用 {中断id: 值} 精确投递。框架用同一套 pending_writes 存两类,靠 task_id 是不是 NULL_TASK_ID 区分——简单场景零负担,复杂场景有精度。这是"常见路径极简 + 高级路径可用"的经典分层。get_null_resume:通用恢复值"取一次就没了"
_scratchpad() 内部定义了闭包 get_null_resume(_algo.py:1320-1331),交给便签本:
# _algo.py:1320
def get_null_resume(consume: bool = False) -> Any:
if null_resume_write is None:
if parent_scratchpad is not None:
return parent_scratchpad.get_null_resume(consume) # ① 本层没有?问父层
return None
if consume: # ② 消费模式
try:
pending_writes.remove(null_resume_write) # 从 pending 里删掉!
return null_resume_write[2] # 返回值
except ValueError:
return None
return null_resume_write[2] # ③ 只看不消费
null_resume_write is None本任务没找到通用恢复值 → 向父便签本递归查(子图场景:恢复值可能是发给父图的)。父也没有则 None。if consume: pending_writes.remove(...)核心:consume=True 时,取值的同时把它从 pending_writes 删除。昨天 interrupt() 里正是 get_null_resume(True)——领了就消费掉。为什么要删保证"一个通用恢复值只被一个 interrupt 领走"。领完就不在池子里了,下一个 interrupt 再查就是 None → 该停还得停。避免一个 "yes" 被多个中断重复消费。return ..[2](不删)consume=False 时只窥视不拿走——某些地方要"看看有没有"但不真正认领。a=interrupt("Q1"); b=interrupt("Q2"),第一次跑停在 Q1。你 Command(resume="A1") 恢复 → 节点重跑:Q1 这次 get_null_resume(True) 拿到 "A1" 并消费删除,返回;执行到 Q2 时再查,通用恢复值已被消费掉了 → None → 再次抛 GraphInterrupt 停在 Q2。所以多中断节点需要多轮恢复(一次 resume 只喂一个通用值),或者一次用 resume={id1:A1, id2:A2} 映射式喂全。反模式:以为一个 Command(resume=单值) 能同时答完一个节点里的所有 interrupt。闭环:interrupt 重跑时如何"直接返回答案"
现在把便签本填好后,回到昨天的 interrupt() 前半段(types.py:912-925)——这次它走的是另一条路:
# types.py:912
scratchpad = conf[CONFIG_KEY_SCRATCHPAD] # 拿到 _scratchpad() 刚填好的便签本
idx = scratchpad.interrupt_counter() # 本 interrupt 编号
# ① 便签本里已有点名恢复值(task-specific)
if scratchpad.resume:
if idx < len(scratchpad.resume):
conf[CONFIG_KEY_SEND]([(RESUME, scratchpad.resume)])
return scratchpad.resume[idx] # 按编号取答案,直接返回!
# ② 或者有通用恢复值
v = scratchpad.get_null_resume(True) # 消费型取值
if v is not None:
assert len(scratchpad.resume) == idx, (scratchpad.resume, idx)
scratchpad.resume.append(v) # 记进便签本(下次同 idx 可复用)
conf[CONFIG_KEY_SEND]([(RESUME, scratchpad.resume)])
return v # 返回,节点继续往下跑
if scratchpad.resume: return resume[idx]点名恢复值路径:便签本里第 idx 项就是答案,取出返回。多个 interrupt 各按自己的 idx 拿各自的值。v = get_null_resume(True)通用恢复值路径:消费掉一个通用值。scratchpad.resume.append(v)妙笔:把领到的通用值追加进便签本的 resume 列表。这样如果本任务因为别的原因再重跑一次,同一个 idx 就能从 resume[idx] 稳定取到,不用再去消费池子——幂等(D46 主题)。CONFIG_KEY_SEND([(RESUME,...)])两条路都会把 RESUME 写回去——记录"这个恢复值已被用掉",随现场落盘,避免重复消费。raise。今天便签本里有值 → 在 if 里就 return 了,压根到不了 raise。你的节点函数一字没改,行为却从"停"变成"继续"——切换开关就是便签本里那几个恢复值,而它们由 _scratchpad() 从 Command.resume 转存的 RESUME 写填充。整个"停—存—喂—续"闭环到此合拢。LazyAtomicCounter + 今日小结
最后补一个细节:为什么 interrupt_counter 是个函数而不是 self.idx += 1?看 _algo.py:1426-1439:
# _algo.py:1426
class LazyAtomicCounter:
__slots__ = ("_counter",)
def __init__(self) -> None:
self._counter = None
def __call__(self) -> int:
if self._counter is None:
with LAZY_ATOMIC_COUNTER_LOCK: # 加锁,双重检查
if self._counter is None:
self._counter = itertools.count(0).__next__ # 惰性建计数器
return self._counter() # itertools.count 的 __next__ 是原子的
itertools.count(0).__next__用 C 实现的 itertools.count,它的 __next__ 在 CPython 里是原子操作,天然线程安全——比手写 x+=1(读-加-写三步,会撞车)稳。惰性(lazy)创建只有第一次真正调用计数时才建计数器。大多数节点根本没 interrupt,省掉无谓开销。_scratchpad 里赋值回看 _algo.py:1334-1344:interrupt_counter=LazyAtomicCounter()、call_counter、subgraph_counter 全用它——interrupt/call/subgraph 各一个独立原子计数器。self.idx += 1 是"读旧值→加一→写回"三步,两个线程交错会拿到相同编号(竞态),导致恢复值配错 interrupt。itertools.count().__next__ 把"取号+自增"压成一个 C 层原子操作,杜绝撞号。代价是多一层函数调用和一个锁(还用了惰性初始化把锁的开销摊薄)。这是"用极小的并发安全成本,换恢复值配对的绝对正确"。👶 小白:恢复时整个节点重跑,那节点前半段有副作用(比如已经写了数据库)怎么办?会写两遍吗?
👨🏫 好问题,这正是幂等问题(D46 专讲)。原则:interrupt 之前的代码,恢复时会再执行一遍。所以带副作用的操作(写库、发消息、扣款)要么放在 interrupt 之后,要么自己做幂等保护。LangGraph 的 put_writes 层对"通道写"做了去重(D35 L05 见过),但你节点里手写的外部副作用它管不着。这是人在环编程必须记住的铁律:把 interrupt 想象成一道"存档线",线之前的代码可能重放。
🧠 今天你应该能回答
- 怎么恢复一个中断的图?(stream/invoke 传 Command(resume=值),复用同一 thread_id)
- Command.resume 的两种形态?(单值给"下一个中断";{id:值} 精确映射)
- PregelScratchpad 是什么级别的?(任务级便签本,记 interrupt 编号和恢复值)
- _scratchpad() 从哪找恢复值?(pending_writes 里的 RESUME 写,分 NULL_TASK_ID 通用 / task_id 点名 / resume_map 映射)
- get_null_resume(True) 的"消费"是什么意思?(取值同时从 pending 删除,防重复领)
- 一个节点两个 interrupt 只给一个 resume 会怎样?(第二个再次中断,需多轮恢复)
✋ 10 分钟动手
# 1. 读便签本定义与发牌逻辑
sed -n '8,19p' libs/langgraph/langgraph/_internal/_scratchpad.py
sed -n '1280,1344p' libs/langgraph/langgraph/pregel/_algo.py
# 2. 完整跑一次"停→续"
python - <<'PY'
from langgraph.graph import StateGraph, START
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.types import interrupt, Command
from typing_extensions import TypedDict
class S(TypedDict):
x: int
def node(s):
ans = interrupt("确认?")
return {"x": 1 if ans=="yes" else 0}
g=StateGraph(S); g.add_node("n",node); g.add_edge(START,"n")
app=g.compile(checkpointer=InMemorySaver())
cfg={"configurable":{"thread_id":"t1"}}
print("第一次:", list(app.stream({"x":0}, cfg))) # 停,出 __interrupt__
print("恢复后:", list(app.stream(Command(resume="yes"), cfg))) # 续,节点跑完
print("最终:", app.get_state(cfg).values) # {'x': 1}
PY
interrupt()(动态断点),LangGraph 还能在 compile/stream 时用 interrupt_before=["n"] 设静态断点——不改节点代码就能"在某节点前/后停"。看 should_interrupt() 怎么靠 versions_seen 判断该不该停。