目 录CONTENT

文章目录

LangGraph 持久化与人在回路:Checkpoint、中断与时间回溯

PySuper
2025-12-06 / 0 评论 / 0 点赞 / 0 阅读 / 0 字
温馨提示:
本文最后更新于2026-05-22,若内容或图片失效,请留言反馈。 所有牛逼的人都有一段苦逼的岁月。 但是你只要像SB一样去坚持,终将牛逼!!! ✊✊✊

作者:PySuper | 来源zhengxingtao.com | 更新日期:2026-11-01

关联阅读:第19篇 · LangGraph 状态机设计

你写了一个 Agent,跑得挺丝滑——直到进程挂了、对话丢了、审批流断了。生产环境不讲武德,没有持久化,Agent 就是沙滩上的城堡。

LangGraph 从 0.2 版本起就内置了 Checkpoint 持久化层,每个 superstep 自动存快照,配合 interrupt、Command、Time Travel 三板斧,让 Agent 真正「可控、可恢复、可审计」。

本文是实战导向,不讲废话。看完你就能在生产环境落地。

一、为什么需要持久化

先看没有持久化是什么体验:

plaintext

┌──────────────────────────────────────────────────────┐
│                   无持久化的灾难场景                    │
├──────────────────────────────────────────────────────┤
│                                                      │
│  进程崩溃 ──▶ 内存中的 State 瞬间蒸发                   │
│              用户:"我刚说的啥?"                        │
│              Agent:"您好,我是 AI 助手……"              │
│                                                      │
│  多轮对话 ──▶ 每次 invoke 都是全新开始                  │
│              用户:"接着上次的继续"                      │
│              Agent:"什么上次?"                        │
│                                                      │
│  审计追溯 ──▶ 无历史记录,出问题只能猜                   │
│              运维:"第三步到底发生了什么?"               │
│              系统:"不知道,反正现在坏了"                 │
│                                                      │
└──────────────────────────────────────────────────────┘

持久化的核心价值就三个词:

表格

能力

说明

关键 API

断点续跑

进程挂了重启后,从最后一个成功的 checkpoint 继续

checkpointer.put() / get_tuple()

多轮对话

同一个 thread_id 的多次 invoke 共享状态

config["configurable"]["thread_id"]

审计追溯

完整的状态历史,任意时间点可回放

get_state_history()

二、Checkpoint 机制详解

2.1 Checkpoint 是什么

LangGraph 把图的执行拆成 superstep(一个 tick 里所有并行节点执行完毕 = 一个 superstep)。每个 superstep 结束后,框架自动调用 checkpointer.put() 保存当前 State 的完整快照,这就是 Checkpoint

plaintext

  superstep 0        superstep 1        superstep 2        superstep 3
┌─────────────┐  ┌─────────────┐  ┌─────────────┐  ┌─────────────┐
│   INPUT     │  │  Node A     │  │  Node B     │  │   END       │
│  {初始状态}  │──│  {A的输出}   │──│  {B的输出}   │──│  {最终状态}  │
└─────────────┘  └─────────────┘  └─────────────┘  └─────────────┘
       │               │               │               │
       ▼               ▼               ▼               ▼
  checkpoint_0    checkpoint_1    checkpoint_2    checkpoint_3

核心概念:

  • StateSnapshot:一个 checkpoint 的数据结构,包含 values(当前状态)、next(接下来要执行的节点)、config(含 checkpoint_id)、metadata(步骤信息)、parent_config(父检查点)

  • 自动性:你不需要手动 save,框架在 superstep 边界自动存

2.2 Thread 概念

thread_id 是 Checkpoint 的核心索引——就像游戏存档的文件名。同一个 thread_id 的所有 checkpoint 串成一条链,不同 thread_id 互不干扰。

python

# 两个独立会话
config_alice = {"configurable": {"thread_id": "alice-001"}}
config_bob   = {"configurable": {"thread_id": "bob-002"}}

# Alice 的对话
graph.invoke({"messages": [{"role": "user", "content": "我叫Alice"}]}, config_alice)
graph.invoke({"messages": [{"role": "user", "content": "我叫什么?"}]}, config_alice)
# → "你叫Alice"

# Bob 的对话(完全隔离)
graph.invoke({"messages": [{"role": "user", "content": "我叫Bob"}]}, config_bob)
graph.invoke({"messages": [{"role": "user", "content": "我叫什么?"}]}, config_bob)
# → "你叫Bob"

2.3 四种 Checkpointer 对比

plaintext

┌───────────────┬──────────────────┬───────────────────┬──────────────────┐
│               │  MemorySaver     │  SqliteSaver      │  PostgresSaver   │  RedisSaver(社区)
├───────────────┼──────────────────┼───────────────────┼──────────────────┤
│ 适用场景       │ 开发/测试/调试    │ 单机轻量/个人项目   │ 生产环境/分布式   │ 高频缓存场景
│ 持久化能力     │ ❌ 进程重启丢失   │ ✅ 本地文件持久化   │ ✅ 分布式持久化   │ ✅ 可配置持久化
│ 性能          │ ⚡ 极高(内存)     │ 🔵 中等(本地I/O)   │ 🟢 强(连接池)    │ ⚡ 极高(内存级)
│ 分布式支持     │ ❌               │ ❌                 │ ✅               │ ✅
│ 异步版本       │ —               │ —                  │ AsyncPostgresSaver│ AsyncRedisSaver
│ 首次初始化     │ 无需             │ 无需               │ .setup()         │ .setup()
│ 生产推荐度     │ ⭐⭐             │ ⭐⭐⭐              │ ⭐⭐⭐⭐⭐          │ ⭐⭐⭐⭐
└───────────────┴──────────────────┴───────────────────┴──────────────────┘

MemorySaver —— 开发调试首选

零配置,开箱即用。数据全在内存,进程一挂啥都没了。

python

from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import StateGraph, START, END
from typing import TypedDict

class State(TypedDict):
    count: int

def increment(state: State) -> dict:
    return {"count": state.get("count", 0) + 1}

builder = StateGraph(State)
builder.add_node("increment", increment)
builder.add_edge(START, "increment")
builder.add_edge("increment", END)

# 零配置,一行搞定
checkpointer = InMemorySaver()
graph = builder.compile(checkpointer=checkpointer)

config = {"configurable": {"thread_id": "demo-001"}}

# 第一次调用
result1 = graph.invoke({"count": 0}, config)
print(result1)  # {'count': 1}

# 第二次调用(同一 thread_id,自动接续)
result2 = graph.invoke({"count": 0}, config)
print(result2)  # {'count': 2}

# 第三次调用(哪怕不传 input,也从 checkpoint 恢复)
result3 = graph.invoke(None, config)
print(result3)  # {'count': 3}

SqliteSaver —— 单机轻量方案

本地文件持久化,进程重启数据不丢。适合小项目和个人工具。

python

import sqlite3
from langgraph.checkpoint.sqlite import SqliteSaver
from langgraph.graph import StateGraph, START, END
from typing import TypedDict

class State(TypedDict):
    count: int
    message: str

def increment(state: State) -> dict:
    return {"count": state.get("count", 0) + 1}

# 方式一:上下文管理器(推荐)
with SqliteSaver.from_conn_string("./checkpoints.db") as checkpointer:
    builder = StateGraph(State)
    builder.add_node("increment", increment)
    builder.add_edge(START, "increment")
    builder.add_edge("increment", END)

    graph = builder.compile(checkpointer=checkpointer)

    config = {"configurable": {"thread_id": "sqlite-demo"}}

    result1 = graph.invoke({"count": 0, "message": "hello"}, config)
    print(result1)  # {'count': 1, 'message': 'hello'}

# 方式二:手动管理连接(生产级配置 + WAL 模式优化)
import os

DB_PATH = "./checkpoints/agent.db"
os.makedirs(os.path.dirname(DB_PATH), exist_ok=True)

conn = sqlite3.connect(DB_PATH, check_same_thread=False)
conn.execute("PRAGMA journal_mode=WAL")     # 提高并发性能
conn.execute("PRAGMA foreign_keys=ON")       # 数据完整性

checkpointer = SqliteSaver(conn)
graph = builder.compile(checkpointer=checkpointer)

# 用完记得关
# conn.close()

注意:SqliteSaver 当前是全量快照存储(不按 channel 版本去重),长流程 + 大状态会导致数据库膨胀。参见 langgraph#7843

PostgresSaver —— 生产环境首选

事务支持、连接池、channel 版本增量存储、分布式部署。LangGraph Cloud 官方使用此方案。

python

from langgraph.checkpoint.postgres import PostgresSaver
from langgraph.graph import StateGraph, START, END, MessagesState
from typing import TypedDict
import os

DB_URI = os.getenv(
    "POSTGRES_CONN_STR",
    "postgresql://postgres:postgres@localhost:5432/langgraph?sslmode=disable"
)

class State(TypedDict):
    count: int

def increment(state: State) -> dict:
    return {"count": state.get("count", 0) + 1}

# 同步版本
with PostgresSaver.from_conn_string(DB_URI) as checkpointer:
    # 首次使用时初始化表结构(后续调用会跳过)
    checkpointer.setup()

    builder = StateGraph(State)
    builder.add_node("increment", increment)
    builder.add_edge(START, "increment")
    builder.add_edge("increment", END)

    graph = builder.compile(checkpointer=checkpointer)

    config = {"configurable": {"thread_id": "pg-demo-001"}}

    # 首次运行
    result1 = graph.invoke({"count": 0}, config)
    print(result1)  # {'count': 1}

# ---- 模拟服务重启(进程重启后的真实场景) ----

with PostgresSaver.from_conn_string(DB_URI) as checkpointer:
    # 新连接,但同一个 thread_id,从数据库恢复状态
    graph_restarted = builder.compile(checkpointer=checkpointer)

    result2 = graph_restarted.invoke(None, config)
    print(result2)  # {'count': 2}  ← 从上次结果继续!

异步版本(FastAPI / asyncio 场景必备):

python

import asyncio
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver

async def main():
    DB_URI = os.getenv("POSTGRES_CONN_STR")

    async with AsyncPostgresSaver.from_conn_string(DB_URI) as checkpointer:
        await checkpointer.setup()  # 异步初始化

        graph = builder.compile(checkpointer=checkpointer)

        config = {"configurable": {"thread_id": "pg-async-demo"}}
        result = await graph.ainvoke({"count": 0}, config)
        print(result)

asyncio.run(main())

依赖安装:

bash

pip install -U "psycopg[binary,pool]" langgraph langgraph-checkpoint-postgres

RedisSaver(社区) —— 高性能缓存场景

适合高频状态更新、低延迟要求的场景。支持 TTL 过期策略和 Shallow 模式(只存最新 checkpoint)。

python

from langgraph.checkpoint.redis import RedisSaver
from langgraph.graph import StateGraph, START, END
from typing import TypedDict

class State(TypedDict):
    count: int

def increment(state: State) -> dict:
    return {"count": state.get("count", 0) + 1}

DB_URI = "redis://localhost:6379"

with RedisSaver.from_conn_string(DB_URI) as checkpointer:
    checkpointer.setup()  # 首次使用创建索引

    builder = StateGraph(State)
    builder.add_node("increment", increment)
    builder.add_edge(START, "increment")
    builder.add_edge("increment", END)

    graph = builder.compile(checkpointer=checkpointer)

    config = {"configurable": {"thread_id": "redis-demo"}}
    result = graph.invoke({"count": 0}, config)
    print(result)  # {'count': 1}

ShallowRedisSaver(只保留最新快照,常量级存储):

python

from langgraph.checkpoint.redis.shallow import ShallowRedisSaver

saver = ShallowRedisSaver(
    redis_url="redis://localhost:6379",
    ttl={"default_ttl": 60, "refresh_on_read": True},  # 60分钟过期
    key_cache_max_size=2000,
    channel_cache_max_size=200,
)

依赖安装:

bash

pip install -U langgraph langgraph-checkpoint-redis

2.4 Checkpoint 的五个设计原则

plaintext

┌──────────────────────────────────────────────────────────────────┐
│                 Checkpoint 五大设计原则                           │
├──────────────────────────────────────────────────────────────────┤
│                                                                  │
│  1️⃣  自动性    每个 superstep 边界自动存,无需手动 save            │
│  2️⃣  选择性    可按 channel 版本增量存储(PostgresSaver)          │
│  3️⃣  可压缩    大对象走 Store 而不是塞进 State                    │
│  4️⃣  可回溯    任意历史 checkpoint 可恢复(Time Travel)           │
│  5️⃣  可过期    Redis TTL / 手动清理,避免无限膨胀                  │
│                                                                  │
└──────────────────────────────────────────────────────────────────┘

三、Human-in-the-Loop(人在回路)

Agent 不可靠,关键操作必须有人把关。LangGraph 提供两套中断机制:静态中断动态中断

3.1 静态中断:interrupt_before / interrupt_after

在编译或调用时指定,无条件在节点前/后暂停。主要用于调试。

python

from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import StateGraph, START, END

# 编译时设置(永久生效)
graph = builder.compile(
    checkpointer=InMemorySaver(),
    interrupt_before=["risk_node"],   # 在 risk_node 执行前暂停
    interrupt_after=["review_node"],  # 在 review_node 执行后暂停
)

# 运行时设置(单次生效)
graph.invoke(
    inputs,
    config={"configurable": {"thread_id": "1"}},
    interrupt_before=["risk_node"],
)

# 恢复(传 None 继续)
graph.invoke(None, config={"configurable": {"thread_id": "1"}})

局限:静态中断是"只出不进"的暂停——你只能看状态,不能把人工决策传回节点内部。要真正实现审批,得上动态中断。

3.2 动态中断:interrupt() 函数

在节点内部根据运行时数据决定是否中断,这是实现人在回路的正道

python

from langgraph.types import interrupt, Command
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import StateGraph, START, END
from typing import TypedDict

class State(TypedDict):
    some_text: str

def human_node(state: State):
    # interrupt() 暂停执行,把 payload 展示给人工
    # 恢复时,interrupt() 的返回值就是人工传入的数据
    value = interrupt(
        {"text_to_revise": state["some_text"]}
    )
    return {"some_text": value}

builder = StateGraph(State)
builder.add_node("human_node", human_node)
builder.add_edge(START, "human_node")
builder.add_edge("human_node", END)

checkpointer = InMemorySaver()
graph = builder.compile(checkpointer=checkpointer)

config = {"configurable": {"thread_id": "hil-demo"}}

# 第一次运行 → 触发 interrupt,暂停
result = graph.invoke({"some_text": "原始文本"}, config)
print(result["__interrupt__"])
# [Interrupt(value={'text_to_revise': '原始文本'}, id='xxx')]

# 用 Command(resume=...) 恢复,传入人工编辑后的值
result2 = graph.invoke(Command(resume="编辑后的文本"), config)
print(result2)  # {'some_text': '编辑后的文本'}

关键细节:恢复时,节点从头重新执行(不是从 interrupt() 那一行继续)。interrupt() 的返回值变成 Command(resume=...) 传入的值,所以不会再暂停。这意味着 interrupt() 之前的代码会重跑——节点必须幂等

plaintext

  首次运行:                     恢复运行:
  ┌─────────────┐              ┌─────────────┐
  │ 节点开始     │              │ 节点开始     │  ← 从头重新执行
  │  ...code... │              │  ...code... │  ← 重跑(需幂等)
  │ interrupt() │──暂停──▶     │ interrupt() │  ← 返回 resume 值,不再暂停
  │  ...code... │              │  ...code... │  ← 正常继续
  │ return      │              │ return      │
  └─────────────┘              └─────────────┘

最佳实践:把 interrupt() 放在节点开头,或放在专用节点里,减少重跑的代码量。

3.3 Command 恢复机制

Command 是恢复执行的原语,支持三种用法:

python

from langgraph.types import Command

# 1. 基本恢复(传值)
graph.invoke(Command(resume="批准"), config)

# 2. 传复杂对象
graph.invoke(Command(resume={"approved": True, "modifier": "admin"}), config)

# 3. 多个并行 interrupt 同时恢复(mapping: interrupt_id → value)
interrupts = graph.get_state(config).interrupts
resume_map = {i.id: f"edited: {i.value}" for i in interrupts}
graph.invoke(Command(resume=resume_map), config)

3.4 三大设计模式

模式一:批准/拒绝(Approve/Reject)

金融交易审批场景:Agent 生成交易,人工审核后放行或拦截。

python

from typing import Literal
from langgraph.types import interrupt, Command
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import InMemorySaver
from typing import TypedDict

class TxnState(TypedDict):
    transaction: dict
    status: str

def approval_node(state: TxnState) -> Command[Literal["execute", "reject"]]:
    is_approved = interrupt(
        {
            "question": "是否批准此交易?",
            "transaction": state["transaction"],
        }
    )
    if is_approved:
        return Command(goto="execute")
    else:
        return Command(goto="reject")

def execute_node(state: TxnState) -> dict:
    print(f"执行交易: {state['transaction']}")
    return {"status": "executed"}

def reject_node(state: TxnState) -> dict:
    print(f"交易被拒绝: {state['transaction']}")
    return {"status": "rejected"}

builder = StateGraph(TxnState)
builder.add_node("approval", approval_node)
builder.add_node("execute", execute_node)
builder.add_node("reject", reject_node)
builder.add_edge(START, "approval")
builder.add_edge("execute", END)
builder.add_edge("reject", END)

graph = builder.compile(checkpointer=InMemorySaver())

config = {"configurable": {"thread_id": "txn-001"}}

# 运行到审批节点 → 暂停
graph.invoke({"transaction": {"amount": 50000, "to": "Bob"}, "status": "pending"}, config)

# 人工批准
result = graph.invoke(Command(resume=True), config)
print(result)  # {'transaction': {...}, 'status': 'executed'}

模式二:编辑状态(Edit)

修正 Agent 输出后再继续。

python

from langgraph.types import interrupt, Command
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import StateGraph, START, END
from typing import TypedDict

class EditState(TypedDict):
    draft: str
    final: str

def review_node(state: EditState) -> dict:
    # 展示草稿给人工,收集编辑后的版本
    edited = interrupt({"draft_to_review": state["draft"]})
    return {"final": edited}

def publish_node(state: EditState) -> dict:
    print(f"发布: {state['final']}")
    return {"final": state["final"]}

builder = StateGraph(EditState)
builder.add_node("review", review_node)
builder.add_node("publish", publish_node)
builder.add_edge(START, "review")
builder.add_edge("review", "publish")
builder.add_edge("publish", END)

graph = builder.compile(checkpointer=InMemorySaver())

config = {"configurable": {"thread_id": "edit-demo"}}

# 运行到 review → 暂停
graph.invoke({"draft": "初稿内容", "final": ""}, config)

# 人工编辑后恢复
result = graph.invoke(Command(resume="修正后的内容"), config)
print(result)  # {'draft': '初稿内容', 'final': '修正后的内容'}

模式三:获取输入(Input)

运行时收集用户补充信息。

python

from langgraph.types import interrupt, Command
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import StateGraph, START, END
from typing import TypedDict

class InputState(TypedDict):
    query: str
    user_info: dict

def collect_info_node(state: InputState) -> dict:
    info = interrupt({"question": "请提供您的账户ID和验证码"})
    return {"user_info": info}

def process_node(state: InputState) -> dict:
    print(f"处理 {state['query']},用户: {state['user_info']}")
    return state

builder = StateGraph(InputState)
builder.add_node("collect_info", collect_info_node)
builder.add_node("process", process_node)
builder.add_edge(START, "collect_info")
builder.add_edge("collect_info", "process")
builder.add_edge("process", END)

graph = builder.compile(checkpointer=InMemorySaver())

config = {"configurable": {"thread_id": "input-demo"}}

# 运行到收集信息 → 暂停
graph.invoke({"query": "查询余额", "user_info": {}}, config)

# 人工提供信息后恢复
result = graph.invoke(
    Command(resume={"account_id": "ACC-12345", "code": "888888"}),
    config,
)
print(result)  # {'query': '查询余额', 'user_info': {'account_id': 'ACC-12345', 'code': '888888'}}

3.5 完整实战:金融交易审批工作流

python

"""
金融交易审批工作流
流程: 提交交易 → 风控检查 → 人工审批 → 执行/拒绝
"""
from typing import Literal, TypedDict
from langgraph.types import interrupt, Command
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import InMemorySaver


class TransactionState(TypedDict):
    transaction: dict    # 交易详情
    risk_score: float    # 风控评分 0-100
    approval_status: str # pending / approved / rejected
    result: str          # 执行结果


def submit_transaction(state: TransactionState) -> dict:
    """提交交易(模拟)"""
    txn = state["transaction"]
    print(f"📩 收到交易请求: {txn['type']} ¥{txn['amount']} → {txn['to']}")
    return {}


def risk_check(state: TransactionState) -> dict:
    """风控评分(模拟)"""
    amount = state["transaction"]["amount"]
    # 简单规则:金额越大风险越高
    risk_score = min(amount / 1000, 100)
    print(f"🔍 风控评分: {risk_score:.1f}/100")
    return {"risk_score": risk_score}


def human_approval(state: TransactionState) -> Command[Literal["execute", "reject"]]:
    """人工审批:根据风控评分决定是否需要人工介入"""
    txn = state["transaction"]
    risk = state["risk_score"]

    # 低风险自动放行
    if risk < 30:
        print("✅ 低风险交易,自动放行")
        return Command(goto="execute")

    # 高风险需要人工审批
    decision = interrupt({
        "question": "此交易需要人工审批",
        "transaction": txn,
        "risk_score": risk,
        "recommendation": "拒绝" if risk > 70 else "谨慎批准",
    })

    if decision.get("approved", False):
        return Command(goto="execute")
    else:
        return Command(goto="reject")


def execute_transaction(state: TransactionState) -> dict:
    """执行交易"""
    txn = state["transaction"]
    print(f"💰 交易执行成功: {txn['type']} ¥{txn['amount']} → {txn['to']}")
    return {"approval_status": "approved", "result": "success"}


def reject_transaction(state: TransactionState) -> dict:
    """拒绝交易"""
    txn = state["transaction"]
    print(f"❌ 交易已拒绝: {txn['type']} ¥{txn['amount']} → {txn['to']}")
    return {"approval_status": "rejected", "result": "rejected_by_reviewer"}


# ---- 构建工作流 ----
builder = StateGraph(TransactionState)
builder.add_node("submit", submit_transaction)
builder.add_node("risk_check", risk_check)
builder.add_node("approval", human_approval)
builder.add_node("execute", execute_transaction)
builder.add_node("reject", reject_transaction)

builder.add_edge(START, "submit")
builder.add_edge("submit", "risk_check")
builder.add_edge("risk_check", "approval")
builder.add_edge("execute", END)
builder.add_edge("reject", END)

checkpointer = InMemorySaver()
graph = builder.compile(checkpointer=checkpointer)


# ---- 运行测试 ----

# 场景1: 低风险交易(自动放行)
print("=" * 50)
print("场景1: 低风险交易")
print("=" * 50)
config_low = {"configurable": {"thread_id": "low-risk"}}
result = graph.invoke(
    {
        "transaction": {"type": "转账", "amount": 5000, "to": "Alice"},
        "risk_score": 0.0,
        "approval_status": "pending",
        "result": "",
    },
    config_low,
)
print(f"结果: {result['approval_status']}\n")


# 场景2: 高风险交易(需要人工审批)
print("=" * 50)
print("场景2: 高风险交易")
print("=" * 50)
config_high = {"configurable": {"thread_id": "high-risk"}}
result = graph.invoke(
    {
        "transaction": {"type": "转账", "amount": 80000, "to": "Unknown"},
        "risk_score": 0.0,
        "approval_status": "pending",
        "result": "",
    },
    config_high,
)
# → 触发 interrupt,暂停等待人工审批

# 查看中断信息
state = graph.get_state(config_high)
print(f"中断信息: {state.interrupts}")

# 人工批准
result = graph.invoke(Command(resume={"approved": True}), config_high)
print(f"结果: {result['approval_status']}")

四、Time Travel(时间回溯)

Checkpoint 不仅是存档,更是完整的执行历史。你可以回到任意时间点,重放或分叉。

4.1 获取历史状态

python

from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import InMemorySaver
from typing import TypedDict, NotRequired

class State(TypedDict):
    topic: NotRequired[str]
    joke: NotRequired[str]

def generate_topic(state: State) -> dict:
    return {"topic": "socks in the dryer"}

def write_joke(state: State) -> dict:
    return {"joke": f"Why do {state['topic']} disappear? They elope!"}

checkpointer = InMemorySaver()
graph = (
    StateGraph(State)
    .add_node("generate_topic", generate_topic)
    .add_node("write_joke", write_joke)
    .add_edge(START, "generate_topic")
    .add_edge("generate_topic", "write_joke")
    .compile(checkpointer=checkpointer)
)

config = {"configurable": {"thread_id": "time-travel-demo"}}

# 运行图
result = graph.invoke({}, config)
print(result)  # {'topic': 'socks in the dryer', 'joke': 'Why do socks in the dryer disappear? They elope!'}

# 获取完整历史(按时间倒序,最新在前)
history = list(graph.get_state_history(config))
for snapshot in history:
    step = snapshot.metadata.get("step", -1)
    next_nodes = snapshot.next
    cp_id = snapshot.config["configurable"]["checkpoint_id"]
    print(f"Step {step}: next={next_nodes}, checkpoint_id={cp_id}")

输出类似:

plaintext

Step 2: next=(), checkpoint_id=1ef663ba-28fe-6528-8002-5a559208592c
Step 1: next=('write_joke',), checkpoint_id=1ef663ba-28f9-6ec4-8001-31981c2c39f8
Step 0: next=('generate_topic',), checkpoint_id=1ef663ba-28f8-6b49-8000-...
Step -1: next=(), checkpoint_id=1ef663ba-28f7-6334-ffff-...

4.2 从任意检查点恢复执行(Replay)

从某个 checkpoint 恢复,后续节点会重新执行(不是读缓存)。

python

# 找到 generate_topic 之后、write_joke 之前的 checkpoint
before_joke = next(s for s in history if s.next == ("write_joke",))

# 从这个 checkpoint 重新执行 write_joke
replay_result = graph.invoke(None, before_joke.config)
print(replay_result)  # write_joke 重新执行

4.3 "What-if" 调试:修改历史状态重新执行(Fork)

这是 Time Travel 最强大的能力——回到过去,改写状态,看看会发生什么。

python

# 找到 generate_topic 之后、write_joke 之前的 checkpoint
history = list(graph.get_state_history(config))
before_joke = next(s for s in history if s.next == ("write_joke",))

# Fork: 修改 topic,创建新分支
fork_config = graph.update_state(
    before_joke.config,
    values={"topic": "chickens"},  # 改成 "chickens"
)

# 从 fork 点继续执行 — write_joke 用新的 topic 重新执行
fork_result = graph.invoke(None, fork_config)
print(fork_result["joke"])
# "Why do chickens disappear? They elope!"  ← topic 变了!

plaintext

  原始执行:                              Fork 分支:
  ┌─────────────┐                      ┌─────────────┐
  │ topic=socks  │                      │ topic=chickens│  ← 修改后的状态
  └──────┬──────┘                      └──────┬──────┘
         │                                    │
         ▼                                    ▼
  ┌─────────────┐                      ┌─────────────┐
  │ joke about  │                      │ joke about  │
  │ socks       │                      │ chickens    │  ← 重新执行
  └─────────────┘                      └─────────────┘

update_state 不会回滚线程,而是在指定 checkpoint 处创建一个新分支。原始执行历史完整保留。

4.4 实际应用场景

  • 调试错误工具调用:回到出错的步骤,修改 tool_call 参数,看结果是否正确

  • 撤销操作:Agent 做了蠢事,回退到之前的状态重新来

  • A/B 测试:同一起点,不同参数,对比结果

  • 复现问题:用户报 bug,用 checkpoint 历史精确复现

4.5 完整代码示例

python

"""
Time Travel 完整示例:Replay + Fork
"""
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import InMemorySaver
from typing import TypedDict, NotRequired


class State(TypedDict):
    topic: NotRequired[str]
    joke: NotRequired[str]
    rating: NotRequired[str]


def generate_topic(state: State) -> dict:
    return {"topic": "AI agents"}


def write_joke(state: State) -> dict:
    return {"joke": f"Why did the {state['topic']} cross the road? To get to the other checkpoint!"}


def rate_joke(state: State) -> dict:
    return {"rating": "😂😂😂"}


# 构建图
checkpointer = InMemorySaver()
graph = (
    StateGraph(State)
    .add_node("generate_topic", generate_topic)
    .add_node("write_joke", write_joke)
    .add_node("rate_joke", rate_joke)
    .add_edge(START, "generate_topic")
    .add_edge("generate_topic", "write_joke")
    .add_edge("write_joke", "rate_joke")
    .compile(checkpointer=checkpointer)
)

config = {"configurable": {"thread_id": "time-travel-full"}}

# 1. 正常执行
result = graph.invoke({}, config)
print("正常执行结果:", result)
# {'topic': 'AI agents', 'joke': 'Why did the AI agents cross the road?...', 'rating': '😂😂😂'}

# 2. 查看历史
print("\n--- 历史检查点 ---")
history = list(graph.get_state_history(config))
for s in history:
    print(f"  Step {s.metadata.get('step')}: next={s.next}")

# 3. Replay: 从 write_joke 之前重放
before_joke = next(s for s in history if s.next == ("write_joke",))
print(f"\n--- Replay from step {before_joke.metadata['step']} ---")
replay_result = graph.invoke(None, before_joke.config)
print("Replay 结果:", replay_result)

# 4. Fork: 修改 topic 为 "databases"
print(f"\n--- Fork with modified topic ---")
fork_config = graph.update_state(
    before_joke.config,
    values={"topic": "databases"},
)
fork_result = graph.invoke(None, fork_config)
print("Fork 结果:", fork_result)
# {'topic': 'databases', 'joke': 'Why did the databases cross the road?...', 'rating': '😂😂😂'}

五、Store:跨 Thread 的长期记忆

Checkpoint 按 thread_id 隔离,不同 thread 之间无法共享数据。但很多场景需要跨对话的长期记忆——用户偏好、Agent 学习的知识、全局配置。这就是 Store 的用武之地。

5.1 Store vs Checkpointer 的区别

plaintext

┌──────────────────────────────────────────────────────────────┐
│         Checkpointer vs Store                                │
├────────────────────┬─────────────────────────────────────────┤
│    Checkpointer    │              Store                      │
├────────────────────┼─────────────────────────────────────────┤
│ 按 thread_id 隔离   │ 按 namespace 组织,跨 thread 共享       │
│ 每个 superstep 快照 │ 手动 put / get / search                │
│ 自动保存            │ 主动写入                                │
│ 短期记忆(会话级)   │ 长期记忆(跨会话级)                    │
│ 适合: 对话上下文     │ 适合: 用户偏好、Agent 学习结果          │
│ 实现: InMemorySaver │ 实现: InMemoryStore                    │
│       PostgresSaver │       PostgresStore                    │
│       RedisSaver   │       RedisStore                       │
└────────────────────┴─────────────────────────────────────────┘

plaintext

  Thread 1 (alice)           Store (跨 Thread)           Thread 2 (alice)
  ┌─────────────┐           ┌─────────────────┐         ┌─────────────┐
  │ Checkpoint  │           │ namespace:       │         │ Checkpoint  │
  │ (对话上下文) │           │  (alice, prefs)  │◀───────│ (对话上下文) │
  │             │           │  key: theme      │────────▶│             │
  │             │──写入────▶│  value: "dark"   │◀─读取──│             │
  └─────────────┘           │                  │         └─────────────┘
                             │  (alice, mems)   │
                             │  key: uuid-1     │
                             │  value: "喜欢披萨"│
                             └─────────────────┘

5.2 代码示例

python

"""
Store 跨 Thread 长期记忆示例
"""
import uuid
from langgraph.store.memory import InMemoryStore
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import StateGraph, START, END, MessagesState
from langgraph.store.base import BaseStore
from langchain_core.runnables import RunnableConfig
from langchain.chat_models import init_chat_model


# 初始化 Store + Checkpointer
store = InMemoryStore()
checkpointer = InMemorySaver()

model = init_chat_model("openai:gpt-4o-mini")


def chat(state: MessagesState, config: RunnableConfig, *, store: BaseStore):
    user_id = config["configurable"]["user_id"]
    namespace = (user_id, "memories")

    # 1. 搜索相关记忆
    memories = store.search(namespace, query=state["messages"][-1].content, limit=3)
    mem_text = "\n".join([m.value.get("text", "") for m in memories]) if memories else "暂无记忆"

    # 2. 构建系统提示(含历史记忆)
    system = f"你是一个贴心的助手。关于用户的已知信息:\n{mem_text}"

    # 3. 如果用户说"记住xxx",就存入 Store
    last_msg = state["messages"][-1].content.lower()
    if "记住" in last_msg or "remember" in last_msg:
        to_remember = last_msg.split("记住", 1)[-1].strip() if "记住" in last_msg else last_msg.split("remember", 1)[-1].strip()
        store.put(namespace, str(uuid.uuid4()), {"text": to_remember})

    # 4. 调用 LLM
    response = model.invoke([{"role": "system", "content": system}, *state["messages"]])
    return {"messages": [response]}


builder = StateGraph(MessagesState)
builder.add_node("chat", chat)
builder.add_edge(START, "chat")
builder.add_edge("chat", END)

# 编译时同时传入 store 和 checkpointer
graph = builder.compile(store=store, checkpointer=checkpointer)


# ---- 测试 ----

# Thread 1: 用户告诉偏好
config1 = {"configurable": {"thread_id": "thread-1", "user_id": "alice"}}
graph.invoke(
    {"messages": [{"role": "user", "content": "记住我喜欢深色主题和披萨"}]},
    config1,
)

# Thread 2: 新对话,但 Store 共享
config2 = {"configurable": {"thread_id": "thread-2", "user_id": "alice"}}
result = graph.invoke(
    {"messages": [{"role": "user", "content": "我喜欢什么主题?什么食物?"}]},
    config2,
)
# → "你喜欢深色主题和披萨!"(从 Store 中检索到跨 Thread 记忆)

5.3 生产环境 Store

python

from langgraph.store.postgres import PostgresStore
from langchain.embeddings import init_embeddings

# 带 embedding 的 PostgresStore(支持语义搜索)
emb = init_embeddings("openai:text-embedding-3-small")

with PostgresStore.from_conn_string(
    DB_URI,
    index={
        "dims": 1536,
        "embed": emb,
        "fields": ["text"],  # 嵌入哪些字段
    },
) as store:
    store.setup()  # 首次使用建表

    # 写入
    store.put(("users", "alice"), "prefs", {"text": "喜欢深色主题和披萨"})

    # 语义搜索
    results = store.search(("users", "alice"), query="用户喜欢什么颜色", limit=3)
    for r in results:
        print(r.value)

依赖安装:

bash

pip install -U langgraph langgraph-checkpoint-postgres

六、实战:生产级审批工作流

综合 Checkpoint + interrupt + Time Travel,构建一个可审计、可回退的完整审批流。

python

"""
生产级金融交易审批工作流
特性: Checkpoint 持久化 + 人工审批 + Time Travel 审计
"""
import uuid
from typing import Literal, TypedDict
from langgraph.types import interrupt, Command
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import InMemorySaver


class ApprovalState(TypedDict):
    transaction: dict       # 交易信息
    risk_score: float       # 风控评分
    approval_result: dict   # 审批结果
    execution_result: str   # 执行结果


def validate_transaction(state: ApprovalState) -> dict:
    """交易验证"""
    txn = state["transaction"]
    print(f"📋 验证交易: {txn}")
    # 模拟验证逻辑
    if not txn.get("amount") or txn["amount"] <= 0:
        return {"execution_result": "invalid: amount must be positive"}
    return {}


def risk_assessment(state: ApprovalState) -> dict:
    """风控评估"""
    amount = state["transaction"]["amount"]
    risk = min(amount / 1000, 100)
    print(f"🔍 风控评分: {risk:.1f}")
    return {"risk_score": risk}


def human_review(state: ApprovalState) -> Command[Literal["execute", "cancel"]]:
    """人工审批(动态中断)"""
    txn = state["transaction"]
    risk = state["risk_score"]

    # 低风险自动放行
    if risk < 20:
        print("✅ 低风险,自动放行")
        return Command(goto="execute")

    # 中高风险 → 人工审批
    approval = interrupt({
        "type": "approval_required",
        "transaction": txn,
        "risk_score": risk,
        "suggestion": "拒绝" if risk > 70 else "建议批准",
    })

    if approval.get("approved"):
        return Command(goto="execute")
    else:
        return Command(goto="cancel")


def execute_transaction(state: ApprovalState) -> dict:
    """执行交易"""
    txn = state["transaction"]
    print(f"💰 执行: {txn['type']} ¥{txn['amount']} → {txn['to']}")
    return {
        "approval_result": {"approved": True, "reviewer": "system"},
        "execution_result": "success",
    }


def cancel_transaction(state: ApprovalState) -> dict:
    """取消交易"""
    print(f"❌ 取消交易")
    return {
        "approval_result": {"approved": False, "reviewer": "human"},
        "execution_result": "cancelled",
    }


# ---- 构建图 ----
builder = StateGraph(ApprovalState)
builder.add_node("validate", validate_transaction)
builder.add_node("risk_assess", risk_assessment)
builder.add_node("review", human_review)
builder.add_node("execute", execute_transaction)
builder.add_node("cancel", cancel_transaction)

builder.add_edge(START, "validate")
builder.add_edge("validate", "risk_assess")
builder.add_edge("risk_assess", "review")
builder.add_edge("execute", END)
builder.add_edge("cancel", END)

checkpointer = InMemorySaver()
graph = builder.compile(checkpointer=checkpointer)


# ---- 运行场景 ----

def run_scenario(name: str, txn: dict, approve: bool = True):
    """运行一个完整场景"""
    print(f"\n{'='*60}")
    print(f"场景: {name}")
    print(f"{'='*60}")

    config = {"configurable": {"thread_id": str(uuid.uuid4())}}

    # 第一步:运行到 interrupt 或完成
    result = graph.invoke(
        {"transaction": txn, "risk_score": 0.0, "approval_result": {}, "execution_result": ""},
        config,
    )

    # 检查是否需要人工审批
    state = graph.get_state(config)
    if state.interrupts:
        print(f"⚠️  等待审批: {state.interrupts[0].value}")
        # 模拟审批
        decision = {"approved": approve, "reason": "审核通过" if approve else "风险过高"}
        result = graph.invoke(Command(resume=decision), config)

    print(f"最终结果: {result.get('execution_result', 'N/A')}")

    # 审计:打印完整历史
    print(f"\n📊 审计记录:")
    for snapshot in graph.get_state_history(config):
        step = snapshot.metadata.get("step", -1)
        next_nodes = snapshot.next
        print(f"  Step {step}: next={next_nodes}")

    return result


# 场景1: 小额交易(自动放行)
run_scenario("小额转账", {"type": "转账", "amount": 500, "to": "Alice"})

# 场景2: 大额交易(人工批准)
run_scenario("大额转账-批准", {"type": "转账", "amount": 50000, "to": "Bob"}, approve=True)

# 场景3: 超大额交易(人工拒绝)
run_scenario("超大额转账-拒绝", {"type": "转账", "amount": 200000, "to": "Unknown"}, approve=False)


# ---- Time Travel 审计 ----
print(f"\n{'='*60}")
print("Time Travel 审计演示")
print(f"{'='*60}")

# 重新运行一个大额交易
config_audit = {"configurable": {"thread_id": "audit-demo"}}
graph.invoke(
    {"transaction": {"type": "转账", "amount": 80000, "to": "Charlie"}, "risk_score": 0.0, "approval_result": {}, "execution_result": ""},
    config_audit,
)

# 找到风控评估后的 checkpoint,尝试修改风控分数
history = list(graph.get_state_history(config_audit))
after_risk = next((s for s in history if s.next == ("review",)), None)

if after_risk:
    print(f"\n找到风控后的 checkpoint: step={after_risk.metadata.get('step')}")
    print(f"当前 risk_score: {after_risk.values.get('risk_score')}")

    # Fork: 把风控分数改成 10(低风险),看看是否会自动放行
    print("\n--- Fork: 修改 risk_score 为 10 ---")
    fork_config = graph.update_state(
        after_risk.config,
        values={"risk_score": 10.0},
    )
    fork_result = graph.invoke(None, fork_config)
    print(f"Fork 结果: {fork_result.get('execution_result')}")
    # → 低风险自动放行,无需人工审批

七、踩坑记录

坑1: MemorySaver 在分布式部署下的坑

MemorySaver 的数据全在进程内存。如果你的服务跑在 K8s 多 Pod / 多 Worker 上,不同实例的 MemorySaver 完全隔离,A Pod 写的 checkpoint,B Pod 读不到。

python

# ❌ 分布式部署用 MemorySaver
graph = builder.compile(checkpointer=InMemorySaver())  # 多 Pod 各管各的

# ✅ 分布式部署用 PostgresSaver
with PostgresSaver.from_conn_string(DB_URI) as cp:
    graph = builder.compile(checkpointer=cp)  # 所有 Pod 共享同一个 DB

结论:开发用 InMemorySaver,上线必须换 PostgresSaver。

坑2: 状态膨胀问题

别把大对象塞进 State。每个 superstep 都存完整快照,一个 10MB 的 PDF base64 字符串在 State 里,跑 100 步就是 1GB 的 checkpoint 数据。

python

# ❌ 大对象直接存 State
class BadState(TypedDict):
    pdf_content: str  # 10MB base64 → 每步存 10MB
    messages: list

# ✅ 大对象走外部存储,State 只存引用
class GoodState(TypedDict):
    pdf_url: str       # 只存 URL 引用
    pdf_summary: str   # 存摘要
    messages: list

# ✅ 或者用 Store 存大对象
store.put(("docs",), doc_id, {"content": large_pdf_base64})

坑3: interrupt 恢复时节点重新执行的陷阱

恢复时整个节点从头重新执行,不是从 interrupt() 那一行继续。如果你的 interrupt() 前面有副作用代码(发邮件、调 API),恢复时会重跑

python

# ❌ interrupt 前有副作用
def bad_node(state):
    send_email(state)        # 恢复时重跑 → 重复发邮件!
    decision = interrupt("审批?")  # 恢复时不再暂停
    return {"decision": decision}

# ✅ 副作用放 interrupt 之后,或放单独节点
def safe_node(state):
    decision = interrupt("审批?")  # 恢复时返回 resume 值,继续往下
    if decision:
        send_email(state)    # 只在确定审批后才发
    return {"decision": decision}

黄金法则interrupt() 前的代码必须幂等(多次执行结果一致),副作用放 interrupt() 后面。

坑4: thread_id 泄漏问题

thread_id 是 Checkpoint 的索引键。如果用自增 ID 或用户 ID 直接当 thread_id,可能导致:

  • 跨业务状态混淆:不同业务用同一个 thread_id,checkpoint 链交错

  • 数据残留:用户 A 的对话残留,被用户 B 继承

python

# ❌ 简单用 user_id
config = {"configurable": {"thread_id": "user-123"}}  # 该用户所有业务混在一起

# ✅ 业务隔离的 thread_id
config = {
    "configurable": {
        "thread_id": f"trade-{trade_id}",  # 按交易 ID 隔离
    }
}

# ✅ 或者复合 ID
config = {
    "configurable": {
        "thread_id": f"{user_id}:{business_type}:{session_id}",
    }
}

总结

plaintext

┌───────────────────────────────────────────────────────────────────┐
│              LangGraph 持久化 & 人在回路 知识地图                    │
├───────────────────────────────────────────────────────────────────┤
│                                                                   │
│  Checkpoint (短期记忆)           Store (长期记忆)                   │
│  ├── InMemorySaver (开发)       ├── InMemoryStore (开发)           │
│  ├── SqliteSaver (单机)         ├── PostgresStore (生产)           │
│  ├── PostgresSaver (生产)       └── RedisStore (缓存)             │
│  └── RedisSaver (缓存)                                           │
│                                                                   │
│  Human-in-the-Loop              Time Travel                       │
│  ├── 静态中断 (调试)            ├── get_state_history (历史追溯)   │
│  │   interrupt_before/after     ├── Replay (回放)                 │
│  └── 动态中断 (生产)            └── Fork (分叉, What-if 调试)     │
│      interrupt() + Command                                         │
│                                                                   │
│  三大模式:                       踩坑:                             │
│  1. 批准/拒绝                   1. MemorySaver 分布式不共享       │
│  2. 编辑状态                    2. 大对象别塞 State               │
│  3. 获取输入                    3. interrupt 恢复 → 节点重跑      │
│                                  4. thread_id 要业务隔离          │
│                                                                   │
└───────────────────────────────────────────────────────────────────┘

选型速查

  • 开发调试 → InMemorySaver + InMemoryStore

  • 单机小项目 → SqliteSaver

  • 生产分布式 → PostgresSaver + PostgresStore

  • 高频缓存 → RedisSaver(社区维护)

  • 人工审批 → interrupt() + Command(resume=...)

  • 审计追溯 → get_state_history() + update_state() Fork

核心心法:Checkpoint 让 Agent "可恢复",interrupt 让 Agent "可控",Time Travel 让 Agent "可审计"。三者组合,才是生产级 Agent 的标配。

参考文档

0
  1. 支付宝打赏

    qrcode alipay
  2. 微信打赏

    qrcode weixin

评论区