时间旅行:翻看每一份存档、修改过去、从岔路重开
既然每个超步都存了 checkpoint(D33-40),那"回到过去"就是自然而然的能力:get_state_history() 把整条存档链读出来,update_state() 让你修改某个历史点的状态、甚至假装"是某个节点写的"。这就是时间旅行(time travel)——人在环里"人不仅能答,还能改、能倒带"的底座。
痛点:AI 走错一步,能不能倒回去重来?
invoke 重跑(浪费前 4 步、而且模型可能又随机选错)?我想要的是:倒回到第 5 步之前,手动把状态改对,从那里继续——就像游戏读一个早期存档、改改装备再往下打。①
get_state(cfg)——看当前状态(读最新档);②
get_state_history(cfg)——看全部历史(把链倒序读出来);③
update_state(cfg, values)——在某个历史点写一份新档,可作为分叉起点。然后把某个历史 checkpoint 的 config(带它的 checkpoint_id)传给
invoke,图就从那个点重新往下跑。类比:checkpoint 链 = 游戏的存档列表;get_state_history = 打开存档列表;update_state = 读档后改角色属性再另存;带着那个档 invoke = 从那个存档继续游戏。
StateSnapshot:一份存档读出来的样子
时间旅行的 API 都返回 StateSnapshot,定义在 types.py:643-659:
# types.py:643
class StateSnapshot(NamedTuple):
values: dict[str, Any] | Any # 这一步的状态值
next: tuple[str, ...] # 下一步该跑哪些节点
config: RunnableConfig # 这份快照的坐标(含 checkpoint_id)
metadata: CheckpointMetadata | None # 元数据(step/source/...)
created_at: str | None # 创建时间戳
parent_config: RunnableConfig | None# 父快照坐标(往回走用)
tasks: tuple[PregelTask, ...] # 本步的任务(含中断/错误信息)
values那一刻的状态(各通道当前值)。你最常看的。next下一步要跑的节点名。空元组 () = 图已结束;非空(如 ('n',))= 停在这、下一步是 n。判断"是否停住"就看它。config这份快照的"门牌号"——里面的 checkpoint_id 就是它在时间线上的定位。把它传给 invoke 就能从这里重开。parent_config父快照的门牌。一路顺着 parent 往回,就是整条历史(D35 讲的父子链在这兑现)。tasks本步的任务列表,含 interrupts 字段——所以你能从快照里看出"这步是被哪个中断卡住的"。_prepare_state_snapshot() 从 CheckpointTuple 组装而来。get_state:读"最新那一份"
main.py:1392 起,先看它怎么拿 checkpointer、怎么处理子图:
# main.py:1392
def get_state(self, config, *, subgraphs=False) -> StateSnapshot:
"""Get the current state of the graph."""
checkpointer = ensure_config(config)[CONF].get(
CONFIG_KEY_CHECKPOINTER, self.checkpointer)
if isinstance(checkpointer, BaseCheckpointSaver):
checkpointer = self._apply_checkpointer_allowlist(checkpointer)
if not checkpointer:
raise ValueError("No checkpointer set") # ① 没存档 = 没历史可看
if (checkpoint_ns := config[CONF].get(CONFIG_KEY_CHECKPOINT_NS, "")) \
and CONFIG_KEY_CHECKPOINTER not in config[CONF]:
recast = recast_checkpoint_ns(checkpoint_ns) # ② 子图:剥掉 task_id
for _, pregel in self.get_subgraphs(namespace=recast, recurse=True):
return pregel.get_state(...) # 转发给对应子图
else:
raise ValueError(f"Subgraph {recast} not found")
... # 否则读本图最新 checkpoint,组装成 StateSnapshot 返回
if not checkpointer: raise时间旅行的前提和中断一样:必须有 checkpointer。没存档,无历史可谈。checkpoint_ns 非空 → 转发子图如果 config 指向某子图命名空间,get_state 会递归转发给那个子图的 Pregel 实例——所以你能查看子图内部的状态,不只是顶层。recast_checkpoint_ns把命名空间里的 task_id 后缀剥掉("sub:abc123"→"sub"),才能按逻辑名找到子图。返回 StateSnapshot最终读最新 checkpoint、拼成 L02 的快照返回。get_state = "现在什么状态、下一步跑啥"。它读的是最新档,是时间旅行里的"现在时刻"。get_state_history:把整条时间线读出来
main.py:1480-1531,核心是调 checkpointer.list():
# main.py:1480
def get_state_history(self, config, *, filter=None, before=None, limit=None):
"""Get the history of the state of the graph."""
config = ensure_config(config)
checkpointer = config[CONF].get(CONFIG_KEY_CHECKPOINTER, self.checkpointer)
...
if not checkpointer:
raise ValueError("No checkpointer set")
# (子图转发逻辑,同 get_state) ...
config = merge_configs(self.config, config, {...})
# eagerly consume list() to avoid holding up the db cursor
for checkpoint_tuple in list( # ① 一次性取出,别占着游标
checkpointer.list(config, before=before, limit=limit, filter=filter)
):
yield self._prepare_state_snapshot( # ② 每份 → StateSnapshot
checkpoint_tuple.config, checkpoint_tuple)
checkpointer.list(...)直接用 D34 学的 list 接口——它按时间倒序(新→旧)返回该 thread 的所有 checkpoint。时间旅行不过是 list 的一个上层封装。before / limit / filter翻页与筛选:before 只看某个 checkpoint 之前的、limit 限条数、filter 按 metadata 过滤。适合历史很长时分页看。list(checkpointer.list(...))注释点明:先一次性 list() 消费完再逐个 yield,避免生成器边迭代边占着数据库游标(长时间持有 DB cursor 是坑)。_prepare_state_snapshot把每个 CheckpointTuple 翻译成 StateSnapshot(L02)。于是你 for s in app.get_state_history(cfg) 拿到一串快照。update_state:像"某个节点"那样写一份新档
update_state(main.py:2515,异步版 aupdate_state main.py:2528)不是"覆盖旧档",而是造一份新 checkpoint,模拟"某个节点执行并写入"。看它执行更新的核心 main.py:2413-2467:
# main.py:2413
for as_node, values, provided_task_id in valid_updates:
writers = self.nodes[as_node].flat_writers # ① 拿"这个节点"的写入器
if not writers:
raise InvalidUpdateError(f"Node {as_node} has no writers")
writes: deque[tuple[str, Any]] = deque()
task = PregelTaskWrites((), as_node, writes, [INTERRUPT])
task_id = provided_task_id or (... uuid5(UUID(checkpoint["id"]), INTERRUPT))
run = RunnableSequence(*writers) if len(writers) > 1 else writers[0]
await run.ainvoke(values, patch_config(config, ... # ② 用你的 values 跑写入器
configurable={ CONFIG_KEY_SEND: writes.extend, ... }))
# main.py:2468
apply_writes(checkpoint, channels, run_tasks, # ③ 把写 apply 到通道
checkpointer.get_next_version, self.trigger_to_nodes)
# main.py:2485
checkpoint = create_checkpoint(checkpoint, channels, step + 1, ...) # ④ 造新档
next_config = await checkpointer.aput(checkpoint_config, checkpoint, ...) # ⑤ 存
writers = nodes[as_node].flat_writers关键:update_state 不是硬塞值,而是走 as_node 这个节点的写入器——于是你的 values 会像该节点的真实输出一样,经过 reducer(如 add_messages 的合并)处理。run.ainvoke(values, ...)把你给的 values 喂给写入器执行,产出的写进 writes。等于"假装 as_node 刚跑完、输出了 values"。apply_writes(...)D21 学的:把这批写合并进通道、涨版本号。此刻状态真的变了。create_checkpoint + aput造 step+1 的新 checkpoint 并存盘。所以 update_state 会在时间线上新增一份档,而不是改旧的。返回 next_config返回新档的坐标——拿它 invoke,就从这个"人改过的点"继续。as_node 的推断与"从哪叉出去"
你可以显式指定 as_node="某节点";不指定时,框架会猜(main.py:2381-2397):
# main.py:2381
elif as_node is None:
last_seen_by_node = sorted(
(v, n)
for n, seen in checkpoint["versions_seen"].items()
if n in self.nodes
for v in seen.values())
# if two nodes updated the state at the same time, it's ambiguous
if last_seen_by_node:
if len(last_seen_by_node) == 1:
as_node = last_seen_by_node[0][1] # 只有一个候选,就是它
elif last_seen_by_node[-1][0] != last_seen_by_node[-2][0]:
as_node = last_seen_by_node[-1][1] # 最后动的那个节点
if as_node is None:
raise InvalidUpdateError("Ambiguous update, specify as_node") # 猜不出→报错
versions_seen 找"最后动手的节点"借 D43 见过的 versions_seen 账本,找"看到的版本最高"的节点——通常就是最近执行的那个,默认让新更新"接在它后面"。len==1 → 唯一只有一个节点动过,无歧义,直接用它。[-1][0]!=[-2][0] → 有明确最后者最后两个节点版本不同 → 有明确的"最新节点",用它。否则 raise Ambiguous并列第一时拒绝猜:两个节点同时更新、分不清谁最后,宁可报错让你显式指定 as_node,也不乱猜写错。versions_seen 版本并列最高,框架故意抛 InvalidUpdateError("Ambiguous update, specify as_node")。这不是缺陷,是安全设计:update_state 要"署名某个节点",署错名会让状态更新经过错误的 reducer、或从错误的节点位置分叉,后果隐蔽难查。与其猜错,不如逼你写清 as_node="a"。反模式:在并行图里偷懒不写 as_node。子图的 replay:ReplayState + 今日小结
回到过去重跑时,一个棘手问题:图里如果有子图,子图也有自己的 checkpoint。重放到过去某点时,子图该读哪份存档?这归 ReplayState 管(_replay.py:14-73):
# _replay.py:14
class ReplayState:
"""Tracks which subgraphs have already loaded their pre-replay checkpoint."""
__slots__ = ("checkpoint_id", "_visited_ns")
def __init__(self, checkpoint_id: str) -> None:
self.checkpoint_id = checkpoint_id # 重放的"目标时间点"
self._visited_ns: set[str] = set() # 记哪些子图已加载过
# _replay.py:52
def get_checkpoint(self, checkpoint_ns, checkpointer, checkpoint_config):
if self._is_first_visit(checkpoint_ns): # ① 该子图第一次被访问
for saved in checkpointer.list(
checkpoint_config,
before={"configurable": {"checkpoint_id": self.checkpoint_id}},
limit=1):
return saved # 取"重放点之前"的档
return None
return checkpointer.get_tuple(checkpoint_config) # ② 之后按正常最新档
_is_first_visit(ns)判断某子图是不是重放中第一次跑到(_replay.py:34,会剥掉 task_id 后缀识别"同一个逻辑子图")。第一次 → list(before=目标点)第一次访问时,取"重放目标时间点之前"的那份子图存档——保证子图也回到过去的正确状态,而不是用它现在的最新态。之后 → get_tuple 正常同一子图在循环里被再次访问(如 for 循环里的子图),就用正常最新档——因为那些是重放过程中新产生的,该用新的。共享一个实例类注释说明:一次父执行里所有派生 config 共享同一个 ReplayState(按引用传),才能全局记住"哪些子图已回到过去了"。_visited_ns 用一个集合记住"这个子图已经完成了它的一次性回溯",之后切回正常加载。用一个 set 的极小状态,精确处理了"回溯只发生一次、后续走正常"的微妙语义。👶 小白:从过去某点重跑,前面已经花过的 token / 调过的 API 会重复消耗吗?
👨🏫 从你选的那个 checkpoint 往后的节点会真实重跑(重新调模型、花 token);但那个 checkpoint 之前的节点不会——它们的结果已经在你选中的存档里了,直接读回,不重算。所以时间旅行省的正是"重跑前半段"的钱。这也呼应明天 D46 的幂等:已完成的写会被复用而非重做。反过来注意:你选的点往后如果有副作用节点,重跑会再执行,这点和 D42 的幂等提醒一致。
🧠 今天你应该能回答
- 时间旅行的三个 API 各干什么?(get_state 看当前 / get_state_history 翻历史 / update_state 改并分叉)
- StateSnapshot 里判断"是否停住"看哪个字段?(next,空=结束、非空=停在此)
- get_state_history 底层是什么?(checkpointer.list 的薄封装,倒序读同一批 checkpoint)
- update_state 是覆盖旧档还是新增?(新增一份署名 as_node 的新档,append-only)
- 为什么 update_state 要走节点的 writers?(让值经过 reducer 合并,像该节点真实输出)
- 并行图不指定 as_node 为何报错?(版本并列、署名歧义,故意拒绝乱猜)
- ReplayState 解决什么?(子图重放时"第一次回旧档、之后用新档")
✋ 10 分钟动手
# 1. 读三个 API 源码
sed -n '1480,1531p' libs/langgraph/langgraph/pregel/main.py # get_state_history
sed -n '2413,2506p' libs/langgraph/langgraph/pregel/main.py # update_state 核心
sed -n '14,73p' libs/langgraph/langgraph/_internal/_replay.py
# 2. 亲手时间旅行
python - <<'PY'
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import InMemorySaver
from typing_extensions import TypedDict
class S(TypedDict):
x: int
g=StateGraph(S)
g.add_node("a", lambda s:{"x":s["x"]+1})
g.add_node("b", lambda s:{"x":s["x"]+10})
g.add_edge(START,"a"); g.add_edge("a","b"); g.add_edge("b",END)
app=g.compile(checkpointer=InMemorySaver())
cfg={"configurable":{"thread_id":"t"}}
app.invoke({"x":0}, cfg)
hist=list(app.get_state_history(cfg))
print("历史档数:", len(hist)) # 若干份(倒序)
for h in hist: print(" step", h.metadata["step"], "x=", h.values.get("x"), "next=", h.next)
# 挑一个 a 刚跑完的历史点,改 x,再从那分叉
target=[h for h in hist if h.next==("b",)][0]
app.update_state(target.config, {"x": 100}, as_node="a")
print("从改后的点继续:", app.invoke(None, {"configurable":{"thread_id":"t"}}))
PY
durability 三档 "sync"/"async"/"exit" 决定"什么时候写、写得多勤"——性能与安全的权衡。看 _loop.py 里 do_checkpoint 和 _checkpointer_put_after_previous 怎么按档位决策。