注:以下所提的“文档”,是指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

文档分几个阶段来说明:

  1. 运行图(invoke / stream)
    先用 invoke(...) 或 stream(...) 执行图 / 工作流,以产生一条执行记录和一系列 checkpoint。
  2. 识别 checkpoint

使用 get_state_history(config) 获取这个线程 (thread_id) 的历史检查点列表

每个历史记录包括 checkpoint_id 和 next(接下来要执行的节点)等元数据

选出你想回溯的那个 checkpoint(例如刚进入某节点之前的状态)

也可以在某个节点前设置 interrupt / breakpoint,使得执行停住,你知道对应 checkpoint 是哪个。

  1. 可选:修改状态
    在回溯 checkpoint 上,你可以用 update_state(...) 修改一些 state 字段(例如纠正某个变量、改变输入、注入新值)。这样新的执行路径就不是完完全全的 replay,而是带着人工或程序修改的“分支”。
  2. 从 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 细节;未来升级时请关注官方变更说明。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐