作者: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:"什么上次?" │
│ │
│ 审计追溯 ──▶ 无历史记录,出问题只能猜 │
│ 运维:"第三步到底发生了什么?" │
│ 系统:"不知道,反正现在坏了" │
│ │
└──────────────────────────────────────────────────────┘
持久化的核心价值就三个词:
表格
二、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 的标配。
参考文档:
评论区