LangGraph-时间旅行(Time Travel)
注:以下所提的“文档”,是指LangGraph官方指南(https://langchain-ai.github.io/langgraph/guides/),参考版本是0.6.8。
以下代码的开发环境:
[project]
name = "my-langgraph"
version = "0.1.0"
requires-python = ">=3.12"
dependencies = [
"dotenv>=0.9.9",
"langchain>=0.3.27",
"langchain-openai>=0.3.33",
"langfuse>=3.5.0",
"langgraph>=0.6.7",
"langgraph-checkpoint-sqlite>=2.0.11",
"numpy>=2.3.3",
]
一、Time Travel 概念篇:Time Travel 概览
指南文档这部分主要介绍 LangGraph 的时间旅行(Time Travel)功能的动机、核心能力和用途。下面是其核心内容:
1. 为什么要 Time Travel / 回溯
在使用 LLM 驱动的系统中,因为很多决策 /输出是非确定性的(随机性、模型抽样、外部 API 结果等),有时我们希望:
- 了解推理过程:查看某次执行中每一步是如何得出的
- 调试错误:找出流程哪一步出错、为何出错
- 探索分支 / 试验不同路径:在某一点回溯并修改状态或输入,从那里继续执行不同策略,进行“分叉探索”
文档指出:Time Travel 是为了支持这些 use case。
2. 核心能力与机制
- 从历史 checkpoint 恢复执行
LangGraph 的时间旅行能力允许你选取一个先前执行的 checkpoint(检查点),从那个时间点(节点边界)恢复执行,既可以完全重放,也可以在该点修改状态后继续新的分支。 - 产生新的历史 / 分支
在回溯点继续执行会生成一个新的历史分支(fork),不会破坏原有历史;即多个可能路径可以共存。 - 与持久化 / 检查点密切耦合
时间旅行依赖于 Graph / StateGraph 在执行过程中不断做检查点 (checkpoint) 的机制。没有 persistence / checkpointer,就无法回溯。
文档里 “Tip” 部分还提示:要查看 “Use time travel” / “Time travel using Server API” 来获取更具体用法。
二、How-To 篇:Use Time Travel in 工作流 / Human-in-the-Loop 场景
指南文档这部分更具操作性,展示如何在 LangGraph 中实际用 time travel 恢复、修改状态、重演流程。下面是其步骤、示例与要点。
操作流程:Use Time Travel
文档分几个阶段来说明:
- 运行图(invoke / stream)
先用 invoke(...) 或 stream(...) 执行图 / 工作流,以产生一条执行记录和一系列 checkpoint。 - 识别 checkpoint
使用 get_state_history(config) 获取这个线程 (thread_id) 的历史检查点列表
每个历史记录包括 checkpoint_id 和 next(接下来要执行的节点)等元数据
选出你想回溯的那个 checkpoint(例如刚进入某节点之前的状态)
也可以在某个节点前设置 interrupt / breakpoint,使得执行停住,你知道对应 checkpoint 是哪个。
- 可选:修改状态
在回溯 checkpoint 上,你可以用 update_state(...) 修改一些 state 字段(例如纠正某个变量、改变输入、注入新值)。这样新的执行路径就不是完完全全的 replay,而是带着人工或程序修改的“分支”。 - 从 checkpoint 恢复执行
使用 invoke 或 stream,传入配置里带上相同 thread_id 和指定的 checkpoint_id,让执行从那个 checkpoint 恢复。输入可以是 None 或新的输入,视具体场景决定。
文档里还列了一个具体工作流示例:先生成一个笑话题目节点(generate_topic),再写笑话(write_joke),然后展示如何取得历史 checkpoint 并恢复重演或改写后继续。
示例摘录
- 获取历史检查点:
states = list(graph.get_state_history(config))
for st in states:
print(st.next, st.config["configurable"]["checkpoint_id"])
- 从 checkpoint 恢复:
graph.invoke(None, config={"configurable": {"thread_id": thread_id, "checkpoint_id": some_id}})
- 用 update_state(...) 注入修改:
graph.update_state(config, {"topic": "chickens"}, checkpoint_id=some_id)
- 然后再 resume:
graph.invoke(None, config=...)
文档也提醒:这部分 API 将在 LangGraph v1.0 中可能有重大变化 /废弃。
三、两篇文档对比、关系与补充视角
|
方面 |
概念篇 |
How-To / 操作篇 |
|
目的 |
介绍为什么要 time travel、是什么能力、设计意图 |
讲如何实操:跑、取历史、更新 state、恢复 |
|
重点 |
支持回溯、分支、debug / 分析模型决策 |
API 方法、checkpoint 识别、重演 /修改流程 |
|
依赖 |
强依赖持久化 /检查点机制 |
展示代码 /方法依赖 get_state_history、update_state、invoke 等 |
|
警示 |
概念层提到未来版本可能变动 |
明确说这些 docs 在 v1.0 之后可能废弃 /重构。 |
补充视角 /注意点:
- 回溯只能到节点边界:你不能从任意行号恢复,而是从某个节点 /检查点的状态继续。
- 原历史不被覆盖:回溯后的执行是一个新的历史分支,不破坏原有执行路径。
- 状态架构兼容性:如果你在回溯后修改 state schema(例如给 state 添加字段 / 删除字段),对旧 checkpoint 的兼容性可能需要 reducer /迁移机制。社区讨论中有人问这个问题。
- 与 Human-in-the-Loop 配合:time travel 是 Human-in-the-Loop 的有力补充——人工可以从某个 checkpoint 回溯、改写状态 /输入、再继续执行。How-To 那篇其实在 “human_in_loop/time_travel” 路径下,就是强调这一整合。
- 版本升级风险:这类功能在未来版本可能有 breaking change,尤其 checkpoint /state /resume API;在设计时要为版本升级预留接口兼容性。
四、例子代码
import argparse, os, time, json
import sqlite3
from typing import TypedDict, Annotated, Dict, Any, List
from operator import add
from langgraph.graph import StateGraph, START, END
from langgraph.config import get_stream_writer
from langgraph.checkpoint.sqlite import SqliteSaver
# === LLM ===
from langchain_openai import ChatOpenAI
from langchain_core.messages import SystemMessage, HumanMessage
"""
列出某个线程(thread_id)的历史检查点;
选定任一 checkpoint_id 回溯恢复执行(不覆盖原历史,会形成新分支);
在某个检查点上 修改状态(patch) 后再从该点恢复;
依然保留 SQLite 持久化、人审暂停、resume、流式观察。
"""
def get_llm():
llm = ChatOpenAI(model="trs-m5", api_key="empty", base_url="http://192.168.5.82:8118/v1", temperature=0.3)
return llm
# ---------- 1) 状态 ----------
class FlowState(TypedDict, total=False):
text: str
draft: str
approved: bool | None
feedback: str
needs_human: bool
result: str
log: Annotated[List[str], add]
# ---------- 2) 节点 ----------
# 生成一个草稿,放在状态的draft中
def node_draft(state: FlowState) -> Dict[str, Any]:
w = get_stream_writer()
topic = state.get("text", "No topic")
sys = ("You are a concise technical writer. Produce a crisp, actionable draft with:\n"
"- Goal / Why it matters\n- 3-5 step plan\n- Risks & mitigations (1-2 bullets)\n"
"Keep it skimmable. Avoid fluff.")
user = f"Draft a short plan about: {topic}"
w({"node": "draft", "event": "start", "topic": topic})
try:
resp = get_llm().invoke([SystemMessage(content=sys), HumanMessage(content=user)])
draft = resp.content.strip()
except Exception as e:
draft = f"- Goal: {topic}\n- Plan: step1, step2, step3\n- (fallback due to error: {e})"
w({"node": "draft", "event": "done"})
return {"draft": draft, "approved": None, "needs_human": False, "log": [f"draft ready for '{topic}'"]}
# 如果approved是None时,则需要人审阅
def node_review(state: FlowState) -> Dict[str, Any]:
w = get_stream_writer()
approved = state.get("approved", None)
if approved is None:
w({"node": "review", "event": "pause_for_human",
"hint": "run `approve` with --ok yes/no and optional --feedback"})
return {"needs_human": True, "log": ["awaiting human review (approved? yes/no, with feedback)"]}
else:
w({"node": "review", "event": "resume_with_human_input",
"approved": approved, "feedback": state.get("feedback", "")})
return {"needs_human": False, "log": [f"human says: {approved}"]}
def node_apply_feedback(state: FlowState) -> Dict[str, Any]:
w = get_stream_writer()
draft = state.get("draft", "")
approved = bool(state.get("approved", False))
fb = (state.get("feedback") or "").strip()
if approved:
sys = ("You are an editor. Lightly polish the draft for clarity and formatting. "
"If 'Incorporate' notes are provided, weave them in minimally. Keep content terse.")
user = f"Draft:\n{draft}\n\nIncorporate (optional): {fb or '(none)'}"
else:
sys = ("You are a reviser. Rewrite the draft to address the reviewer feedback precisely. "
"Keep it concise, structured, and actionable.")
user = f"Original draft:\n{draft}\n\nReviewer feedback (must address): {fb or '(no details)'}"
try:
resp = get_llm().invoke([SystemMessage(content=sys), HumanMessage(content=user)])
new_draft = resp.content.strip()
except Exception as e:
tag = "LIGHT_EDIT" if approved else "REWORKED"
new_draft = f"[{tag} - fallback]\n{draft}\n\n(LLM error: {e})"
w({"node": "apply_feedback", "event": "light_edit" if approved else "rework"})
return {"draft": new_draft, "log": ["applied feedback via LLM"]}
def node_finalize(state: FlowState) -> Dict[str, Any]:
w = get_stream_writer()
w({"node": "finalize", "event": "start"})
time.sleep(0.05)
result = f"✅ FINAL DOC\n{state.get('draft','')}\n---\nApproved: {state.get('approved')}"
w({"node": "finalize", "event": "done"})
return {"result": result, "log": ["finalized"]}
# ---------- 3) 条件路由 ----------
# 需要人审阅则退出流程
def route_after_review(state: FlowState) -> str:
return "END" if state.get("needs_human") else "apply_feedback"
# ---------- 4) 组装图 ----------
def build_app(db_path: str = "hitl.sqlite"):
conn = sqlite3.connect(db_path, check_same_thread=False)
saver = SqliteSaver(conn)
g = StateGraph(FlowState)
g.add_node("draft", node_draft)
g.add_node("review", node_review)
g.add_node("apply_feedback", node_apply_feedback)
g.add_node("finalize", node_finalize)
g.add_edge(START, "draft")
g.add_edge("draft", "review")
g.add_conditional_edges("review", route_after_review, path_map={"END": END, "apply_feedback": "apply_feedback"})
g.add_edge("apply_feedback", "finalize")
g.add_edge("finalize", END)
return g.compile(checkpointer=saver)
# ---------- 5) CLI 子命令 ----------
def _print_stream(app, init, config):
for mode, chunk in app.stream(
init, config=config, stream_mode=["updates", "custom", "values"], durability="sync"
):
if mode == "custom":
print("🟣 [custom ]", chunk)
elif mode == "updates":
print("🔸 [updates]", chunk)
elif mode == "values":
print("🟢 [values ]", chunk)
def cmd_run(args):
app = build_app(args.db)
cfg = {"configurable": {"thread_id": args.thread, }}
print(f"\n=== RUN (thread={args.thread}) ===")
_print_stream(app, {"text": args.text}, cfg)
st = app.get_state(cfg)
print("\n--- SNAPSHOT ---")
print("values:", st.values)
print("next :", st.next)
print("\nIf you see needs_human=True, run `approve` to resume.\n")
def cmd_approve(args):
app = build_app(args.db)
cfg = {"configurable": {"thread_id": args.thread}}
approved = (args.ok.lower() in ["yes", "y", "true", "1"])
fb = args.feedback or ""
print(f"\n=== APPROVE (thread={args.thread}, ok={approved}, feedback={fb!r}) ===")
app.update_state(cfg, values={"approved": approved, "feedback": fb, "needs_human": False}, as_node="review")
_print_stream(app, {}, cfg)
st = app.get_state(cfg)
hist = list(app.get_state_history(cfg))
print("\n--- FINAL ---")
print("result:\n", st.values.get("result", ""))
print("\nhistory length:", len(hist))
def cmd_history(args):
app = build_app(args.db)
cfg = {"configurable": {"thread_id": args.thread}}
hist = list(app.get_state_history(cfg))
print(f"\n=== HISTORY (thread={args.thread}) ===")
if not hist:
print("(empty)")
return
for i, st in enumerate(hist):
# 大多数实现可从 st.config["configurable"]["checkpoint_id"] 读到 checkpoint_id
cp = (st.config or {}).get("configurable", {}).get("checkpoint_id", None)
cns = (st.config or {}).get("configurable", {}).get("checkpoint_ns", "defaultxx")
print(f"[{i}] node_next={st.next!r} checkpoint_id={cp!r} checkpoint_ns={cns!r}")
# 也可视化关键信息
vals = st.values or {}
preview = {k: vals.get(k) for k in ["text", "draft", "approved", "needs_human", "result"] if k in vals}
print(" values:", json.dumps(preview, ensure_ascii=False)[:220])
def cmd_travel(args):
"""
从给定 checkpoint_id 恢复执行(不修改状态,纯 replay/继续)。
"""
app = build_app(args.db)
cfg = {"configurable": {"thread_id": args.thread, "checkpoint_id": args.checkpoint}}
print(f"\n=== TIME-TRAVEL RESUME (thread={args.thread}, checkpoint_id={args.checkpoint}) ===")
_print_stream(app, None, cfg) # 恢复不需要新的输入
st = app.get_state(cfg)
print("\n--- AFTER RESUME ---")
print("next:", st.next)
print("values:", {k: st.values.get(k) for k in ["approved", "needs_human", "result"]})
def cmd_patch(args):
"""
在某个 checkpoint_id 上“打补丁”(修改 state),然后从该点恢复。
例如:改 'approved'、'feedback'、甚至重设 'text' 再从回溯点继续。
"""
app = build_app(args.db)
# 对于按checkpoint_id更新状态,需要checkpoint_ns的值
cfg = {"configurable": {"thread_id": args.thread, "checkpoint_id": args.checkpoint, "checkpoint_ns": ""}}
try:
patch_vals = json.loads(args.values)
assert isinstance(patch_vals, dict)
except Exception as e:
raise SystemExit(f"--values 需要是 JSON 对象,例如: '{{\"approved\": true, \"feedback\": \"add A/B\"}}' // {e}")
print(f"\n=== PATCH @checkpoint (thread={args.thread}) ===")
print("checkpoint_id:", args.checkpoint)
print("patch values :", patch_vals)
# 把修改写入到指定 checkpoint 的状态里
# 要以新的config来执行
new_cfg = app.update_state(config=cfg, values=patch_vals)
# 再从该 checkpoint 恢复执行
_print_stream(app, None, new_cfg)
st = app.get_state(new_cfg)
print("\n--- AFTER PATCH+RESUME ---")
print("next:", st.next)
print("values:", {k: st.values.get(k) for k in ["text", "approved", "feedback", "needs_human", "result"]})
"""
怎么用(演练脚本)
# 1) 先跑一遍,使其在 review 节点暂停(等待人审)
python hitl_llm_timetravel_demo.py run --thread t-42 --text "Prepare incident response checklist"
# 2) 查看历史检查点,拿到 checkpoint_id
python hitl_llm_timetravel_demo.py history --thread t-42
# 3a) 直接从某个 checkpoint 回溯恢复(不改状态)
python hitl_llm_timetravel_demo.py travel --thread t-42 --checkpoint <ID_FROM_HISTORY>
# 3b) 或者在某个 checkpoint 打补丁(例如直接判定通过 + 附反馈),再恢复
python hitl_llm_timetravel_demo.py patch --thread t-42 \
--checkpoint <ID_FROM_HISTORY> \
--values '{"approved": true, "feedback": "Add on-call rotation and escalation path."}'
# 4) 也可以按原 HITL 流程注入人审意见继续
python hitl_llm_timetravel_demo.py approve --thread t-42 --ok yes --feedback "Looks good."
注:LangGraph 在 1.0 附近版本可能调整 time-travel / resume / checkpoint 的 API 细节;未来升级时请关注官方变更说明。
"""
class _Args:
# 一个简单的“命名空间”对象,用来模拟 argparse.Namespace
def __init__(self, **kwargs):
for k, v in kwargs.items():
setattr(self, k, v)
def run_via_code(thread: str="t-42", text: str="Prepare incident response checklist", db: str="hitl_llm_timetravel_example.sqlite"):
"""
等价于命令:
python hitl_llm_timetravel_demo.py run --thread t-42 --text "Prepare incident response checklist"
"""
args = _Args(thread=thread, text=text, db=db)
cmd_run(args)
def approve_via_code(thread: str, ok: str, feedback: str, db:str= "hitl_llm_timetravel_example.sqlite"):
"""
等价于命令:
python hitl_llm_timetravel_demo.py approve --thread t-42 --ok yes --feedback "Looks good."
"""
if ok not in ("yes", "no"):
raise ValueError("ok 必须是 'yes' 或 'no'")
args = _Args(thread=thread, ok=ok, feedback=feedback, db=db)
cmd_approve(args)
def history_via_code(thread: str="t-42", db:str= "hitl_llm_timetravel_example.sqlite"):
"""
等价于命令:
python hitl_llm_timetravel_demo.py history --thread t-42
"""
args = _Args(thread=thread, db=db)
cmd_history(args)
def travel_via_code(thread: str, checkpoint:str, db:str= "hitl_llm_timetravel_example.sqlite"):
"""
等价于命令:
python hitl_llm_timetravel_demo.py travel --thread t-42 --checkpoint <ID_FROM_HISTORY>
"""
args = _Args(thread=thread, checkpoint=checkpoint, db=db)
cmd_travel(args)
def patch_via_code(thread: str, checkpoint:str, values:str, db:str= "hitl_llm_timetravel_example.sqlite"):
"""
等价于命令:
python hitl_llm_timetravel_demo.py patch --thread t-42 \
--checkpoint <ID_FROM_HISTORY> \
--values '{"approved": true, "feedback": "Add on-call rotation and escalation path."}'
"""
args = _Args(thread=thread, checkpoint=checkpoint, values=values, db=db)
cmd_patch(args)
def main():
# run_via_code()
# approve_via_code(thread="t-42", ok="yes", feedback="Looks good.")
# history_via_code(thread="t-42") # checkpoint_id='1f09cddd-3a4e-6df6-8002-57c40a4bb118'
# 这里用的checkpoint的值是从history_via_code的输出获取
# travel_via_code(thread="t-42", checkpoint="1f09cdec-3d1f-68ec-8001-80dad25912f9")
patch_via_code(thread="t-42", checkpoint="1f09d007-8ecd-60ff-8001-abb1521ca548",
values='{"approved": true, "feedback": "Add on-call rotation and escalation path."}',
db="hitl_llm_timetravel_example.sqlite")
def run_console():
parser = argparse.ArgumentParser(description="HITL + Real LLM + Time-Travel Demo (SQLite)")
sub = parser.add_subparsers(dest="cmd", required=True)
# 运行到审阅点
p_run = sub.add_parser("run", help="start or continue until human review pause")
p_run.add_argument("--thread", default="t-1001")
p_run.add_argument("--text", default="Draft a rollout plan for feature flags")
p_run.add_argument("--db", default="hitl.sqlite")
p_run.set_defaults(func=cmd_run)
# 人工审批
p_ok = sub.add_parser("approve", help="inject human decision and resume")
p_ok.add_argument("--thread", default="t-1001")
p_ok.add_argument("--ok", required=True, choices=["yes", "no"])
p_ok.add_argument("--feedback", default="")
p_ok.add_argument("--db", default="hitl.sqlite")
p_ok.set_defaults(func=cmd_approve)
# 查看历史检查点
p_hist = sub.add_parser("history", help="list checkpoints for a thread")
p_hist.add_argument("--thread", default="t-1001")
p_hist.add_argument("--db", default="hitl.sqlite")
p_hist.set_defaults(func=cmd_history)
# 从某个检查点恢复执行(不修改 state)
p_tt = sub.add_parser("travel", help="resume from a checkpoint_id (time-travel)")
p_tt.add_argument("--thread", default="t-1001")
p_tt.add_argument("--checkpoint", required=True, help="checkpoint_id from 'history'")
p_tt.add_argument("--db", default="hitl.sqlite")
p_tt.set_defaults(func=cmd_travel)
# 在检查点上修改状态,然后从该点恢复
p_patch = sub.add_parser("patch", help="patch state at checkpoint_id and resume")
p_patch.add_argument("--thread", default="t-1001")
p_patch.add_argument("--checkpoint", required=True, help="checkpoint_id from 'history'")
p_patch.add_argument("--values", required=True,
help='JSON dict, e.g. \'{"approved": true, "feedback": "shorten plan"}\'')
p_patch.add_argument("--db", default="hitl.sqlite")
p_patch.set_defaults(func=cmd_patch)
args = parser.parse_args()
if not os.getenv("OPENAI_API_KEY"):
print("WARN: OPENAI_API_KEY not set. The demo will fallback to placeholder text.")
args.func(args)
def hitl_llm_timetravel_example():
main()
if __name__ == "__main__":
hitl_llm_timetravel_example()
注:LangGraph 在 1.0 附近版本可能调整 time-travel / resume / checkpoint 的 API 细节;未来升级时请关注官方变更说明。
更多推荐


所有评论(0)