LangGraph Checkpointing:持久化与Human-in-the-Loop模式

引言

想象这样一个场景:你正在使用一个复杂的AI代理系统处理一笔银行转账。系统已经完成了用户身份验证、风险评估、合规检查,正准备执行转账操作——突然,进程崩溃了。用户重新发起请求时,整个流程从头开始,用户又被问了一遍“请问您的身份证号是?”。

这不仅仅是体验问题,更是架构问题。在传统的单体应用中,一次请求的生命周期通常在毫秒级,进程崩溃的概率低到可以忽略。但在LangGraph这类有状态、多步骤、长耗时的Agent架构中,一个工作流可能运行数分钟甚至数小时,期间可能涉及多次LLM调用、外部API交互、人工审批环节。任何一次崩溃都意味着上下文丢失,而LLM对话的上下文重建成本极高——不仅是金钱成本(Token费用),更是时间成本和用户信任成本。

这正是LangGraph Checkpointing要解决的核心问题:让Agent工作流在任意节点暂停、恢复、甚至分支重放,而不会丢失任何状态。

更关键的是,Checkpointing是Human-in-the-Loop(HITL)模式的基石。没有持久化,就无法实现“机器暂停,等人确认,再继续执行”的交互范式。今天,我们从源码级别拆解LangGraph的Checkpoint机制,并给出可落地的实战代码。

核心概念:从“书签”到“时间旅行”

生活类比:电子书阅读器的“书签”与“历史记录”

想象你在读一本电子书。Checkpointing就像阅读器自动保存的书签——你随时合上iPad,下次打开时,它精准地把你带回上次离开的段落,连高亮标记、翻页位置都分毫不差。

但LangGraph的Checkpoint更强。它不只是“书签”,而是完整的“历史记录”:你可以回退到第10页重新阅读(分支重放),也可以跳转到第20页看看另一条故事线(时间旅行),甚至可以对比两次阅读的差异(状态Diff)。

技术定义

LangGraph的State是线程级(Thread-level)的。一个thread_id对应一条独立的Agent执行轨迹。Checkpoint则是State在某个超级步骤(Super-step)结束时的不可变快照

# 伪代码:Checkpoint的数据结构
class Checkpoint:
    thread_id: str          # 线程标识
    step: int               # 步数
    state: dict             # 完整状态快照
    parent: Optional[CheckpointID]  # 父检查点,构成链式结构
    metadata: dict          # 时间戳、LLM调用次数等元信息

每次图执行经过一个节点后,LangGraph会调用BaseCheckpointSaver将当前状态序列化并存储。这不仅仅是简单的KV存储,而是一个支持并发写入、条件读取、分支检索的版本化状态系统。

源码深度分析:LangGraph Checkpointing如何工作

核心抽象:BaseCheckpointSaver

LangGraph的持久化接口定义在langgraph.checkpoint.base中,核心抽象如下:

class BaseCheckpointSaver(ABC):
    @abstractmethod
    def put(self, config: dict, checkpoint: Checkpoint, metadata: dict) -> None:
        """写入一个检查点"""

    @abstractmethod
    def get_tuple(self, config: dict) -> Optional[CheckpointTuple]:
        """获取指定配置的检查点(含父指针)"""

    @abstractmethod
    def list(self, config: dict, *, limit: int = 10) -> Iterator[CheckpointTuple]:
        """列出历史检查点(用于时间旅行)"""

注意get_tuple返回的是CheckpointTuple,它包含checkpointparent_checkpoint_id。这个设计是为了构建链式结构——每次状态更新都指向其父检查点,形成一个不可变的追加日志。

Pregel引擎与Checkpoint的交互

LangGraph的执行引擎是Pregel(受Google Pregel图计算模型启发)。在每次节点执行前,引擎会调用_get_checkpoint从Saver中恢复状态;节点执行完毕后,调用_put_checkpoint保存新状态。

# Pregel._run_step 简化代码
def _run_step(self, input: dict, config: dict) -> dict:
    # 1. 从持久化存储中恢复状态
    checkpoint = self._get_checkpoint(config)
    
    # 2. 将checkpoint中的状态作为起始输入
    state = checkpoint.state if checkpoint else self.graph.input_schema(**input)
    
    # 3. 执行节点(可能包含多次LLM调用)
    for task in self._tasks_for_step(state, config):
        result = task.execute()
        state = self._update_state(state, result)
    
    # 4. **原子性地**保存新检查点
    self._put_checkpoint(config, state, metadata)
    return state

关键点在于第4步的原子性。LangGraph使用InMemorySaver时,通过Python的threading.Lock保证同一线程内的原子操作;使用SqliteSaver时,通过数据库事务保证崩溃一致性。

分支与时间旅行:CheckpointList的威力

HITL模式的核心能力来自get_next_versionlist的组合。当人类审批者拒绝某个操作时,我们可以:

  1. 通过list获取所有历史检查点
  2. 选择某个历史检查点作为新起点
  3. 用新的用户输入构造config,传入invoke,实现从指定步骤重放
# 时间旅行实现
def replay_from_step(thread_id: str, step_number: int, new_input: dict):
    config = {"configurable": {"thread_id": thread_id}}
    checkpoints = saver.list(config, limit=100)
    target_checkpoint = next(cp for cp in checkpoints if cp.metadata["step"] == step_number)
    
    # 构造带有父指针的配置
    replay_config = {
        "configurable": {
            "thread_id": thread_id,
            "checkpoint_id": target_checkpoint.checkpoint["id"],
        }
    }
    return graph.invoke(new_input, config=replay_config)

这里有一个隐藏的坑:checkpoint_id必须与thread_id配合使用。LangGraph的get_tuple会优先查找checkpoint_id,如果找不到再回退到该线程的最新检查点。

实战代码:从入门到HITL

示例1:基础持久化——SqliteSaver实现跨会话恢复

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

# 定义状态结构
class AgentState(TypedDict):
    messages: Annotated[list, operator.add]
    current_step: int
    user_approved: bool

# 定义图节点
def step_a(state: AgentState):
    print(f"执行Step A,当前步数:{state['current_step']}")
    return {"messages": ["A执行完毕"], "current_step": state["current_step"] + 1}

def step_b(state: AgentState):
    print(f"执行Step B,当前步数:{state['current_step']}")
    return {"messages": ["B执行完毕"], "current_step": state["current_step"] + 1}

def step_c(state: AgentState):
    print(f"执行Step C,当前步数:{state['current_step']}")
    return {"messages": ["C执行完毕"], "current_step": state["current_step"] + 1}

# 构建图
builder = StateGraph(AgentState)
builder.add_node("a", step_a)
builder.add_node("b", step_b)
builder.add_node("c", step_c)
builder.set_entry_point("a")
builder.add_edge("a", "b")
builder.add_edge("b", "c")
builder.add_edge("c", END)

# 关键:使用SqliteSaver(内存中创建临时文件)
conn = sqlite3.connect(":memory:", check_same_thread=False)
saver = SqliteSaver(conn)
graph = builder.compile(checkpointer=saver)

# 第一次执行——完整跑完
config = {"configurable": {"thread_id": "demo-1"}}
result = graph.invoke({"messages": [], "current_step": 0, "user_approved": True}, config)
print(f"第一次执行结果:{result['messages']}")

# 模拟进程重启——重新连接同一个数据库
conn2 = sqlite3.connect(":memory:", check_same_thread=False)
saver2 = SqliteSaver(conn2)
graph2 = builder.compile(checkpointer=saver2)

# 第二次执行——不提供初始输入,从检查点恢复
config2 = {"configurable": {"thread_id": "demo-1"}}
result2 = graph2.invoke(None, config2)
print(f"恢复执行结果:{result2['messages']}")

注意:SqliteSaver的check_same_thread=False是必须的,因为LangGraph可能在不同线程调用Saver。另外,invoke(None)表示从已有检查点恢复,而不是新开流程。

示例2:Human-in-the-Loop——人工审批暂停与恢复

from langgraph.checkpoint.memory import MemorySaver
from typing import TypedDict, Literal
import time

class ApprovalState(TypedDict):
    transaction_id: str
    amount: float
    approval_status: Literal["pending", "approved", "rejected"]
    risk_score: float

def assess_risk(state: ApprovalState):
    """模拟风控评估"""
    print("🏦 风控引擎评估中...")
    time.sleep(1)
    # 模拟规则引擎打分
    risk_score = min(100, state["amount"] / 1000 * 5)
    return {"risk_score": risk_score}

def request_human_approval(state: ApprovalState):
    """暂停并等待人工审批"""
    print(f"⚠️  交易 {state['transaction_id']} 金额 ${state['amount']} 需要人工审批")
    # 注意:这里不返回任何状态更新,图会在此节点暂停
    return {}

def execute_transaction(state: ApprovalState):
    print(f"💸 交易执行成功!交易ID:{state['transaction_id']}")
    return {"approval_status": "approved"}

def reject_transaction(state: ApprovalState):
    print(f"❌ 交易被拒绝!交易ID:{state['transaction_id']}")
    return {"approval_status": "rejected"}

# 构建带条件分支的图
builder = StateGraph(ApprovalState)
builder.add_node("risk", assess_risk)
builder.add_node("human", request_human_approval)
builder.add_node("execute", execute_transaction)
builder.add_node("reject", reject_transaction)

builder.set_entry_point("risk")
builder.add_edge("risk", "human")

# 条件边:根据人工审批结果决定走向
def route_after_approval(state: ApprovalState):
    if state["approval_status"] == "approved":
        return "execute"
    elif state["approval_status"] == "rejected":
        return "reject"
    return "human"  # 保持等待

builder.add_conditional_edges("human", route_after_approval, 
                                {"execute": "execute", "reject": "reject", "human": "human"})
builder.add_edge("execute", END)
builder.add_edge("reject", END)

# 使用内存持久化(生产环境请替换为Sqlite/Postgres)
graph = builder.compile(checkpointer=MemorySaver())

# 模拟用户发起交易
initial_state = {
    "transaction_id": "TXN-2024-001",
    "amount": 15000,
    "approval_status": "pending",
    "risk_score": 0
}

config = {"configurable": {"thread_id": "approval-flow-1"}}

# 第一次调用——执行到human节点后自动暂停
result = graph.invoke(initial_state, config)
print(f"当前状态:{result['approval_status']},风险分:{result['risk_score']}")

# =========== 模拟管理员审批 ===========
# 关键:再次invoke,但只更新approval_status字段
admin_config = {"configurable": {"thread_id": "approval-flow-1"}}
new_state = {"approval_status": "approved"}  # 增量更新

# 继续执行图
final_result = graph.invoke(new_state, admin_config)
print(f"最终状态:{final_result['approval_status']}")

# 测试拒绝场景
config2 = {"configurable": {"thread_id": "approval-flow-2"}}
graph.invoke({"transaction_id": "TXN-2024-002", "amount": 5000, 
              "approval_status": "pending", "risk_score": 0}, config2)
reject_result = graph.invoke({"approval_status": "rejected"}, config2)
print(f"拒绝场景最终状态:{reject_result['approval_status']}")

核心机制:当节点返回的状态与当前状态无差异(或返回空dict)时,LangGraph会认为该节点是“终态”,暂停执行。后续invoke会从暂停点继续。这是interrupt_before的隐式用法——实际上LangGraph还提供了interrupt_before参数显式控制暂停节点。

示例3:时间旅行——回溯到任意历史状态

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

# 定义状态
class DebugState(TypedDict):
    counter: int
    log: Annotated[list, operator.add]

def increment(state: DebugState):
    new_counter = state["counter"] + 1
    print(f"计数器增加:{state['counter']} -> {new_counter}")
    return {"counter": new_counter, "log": [f"increment to {new_counter}"]}

# 构建线性图
builder = StateGraph(DebugState)
builder.add_node("inc1", increment)
builder.add_node("inc2", increment)
builder.add_node("inc3", increment)
builder.set_entry_point("inc1")
builder.add_edge("inc1", "inc2")
builder.add_edge("inc2", "inc3")
builder.add_edge("inc3", END)

conn = sqlite3.connect(":memory:", check_same_thread=False)
saver = SqliteSaver(conn)
graph = builder.compile(checkpointer=saver)

config = {"configurable": {"thread_id": "time-travel-demo"}}
result = graph.invoke({"counter": 0, "log": []}, config)
print(f"最终状态:{result}")

# =========== 时间旅行 ===========
# 获取所有历史检查点
checkpoints = list(saver.list(config, limit=10))
print(f"\n历史检查点数量:{len(checkpoints)}")

# 查看每个检查点的状态
for i, cp_tuple in enumerate(checkpoints):
    cp = cp_tuple.checkpoint
    print(f"检查点 {i}: 步数={cp['step']}, 状态={cp['state']}")

# 回退到第二步执行前(检查点索引1,即inc2执行前)
target_cp = checkpoints[1].checkpoint
replay_config = {
    "configurable": {
        "thread_id": "time-travel-demo",
        "checkpoint_id": target_cp["id"]
    }
}

# 重新执行(从inc2开始)
print("\n🔄 开始时间旅行重放...")
replay_result = graph.invoke(None, replay_config)
print(f"重放后的状态:{replay_result}")

关键洞察:时间旅行不是修改历史,而是基于历史状态创建新的分支。原有线程的记录还在,但新的执行路径会以目标检查点为父节点继续延伸。

方案对比:Checkpointing的三大实现路线

| 方案 | 存储介质 | 并发性能 | 事务性 | 适用场景 |

|------|---------|---------|--------|---------|

| MemorySaver | 内存 | 极高 | 无 | 单进程、测试、无状态原型 |

| SqliteSaver | SQLite文件 | 中(写锁) | 支持 | 单机生产、中小规模、轻量部署 |

| PostgresSaver | PostgreSQL | 高(MVCC) | 强 | 多实例、生产级、需要HA |

深度对比

MemorySaver:本质是dict + Lock。看似简单,但它的put方法使用了读写锁,读操作可并发,写操作互斥。适合快速原型验证,但进程崩溃即丢失。

SqliteSaver:利用SQLite的BEGIN IMMEDIATE事务。其put方法会执行INSERT OR REPLACE。注意:SQLite的写锁是数据库级的,高并发写入会串行化。对于Agent场景(写入频率不高,但写入体积大),SQLite反而够用。

PostgresSaver:使用SELECT ... FOR UPDATESERIALIZABLE事务隔离级别。支持多进程同时读写。其序列化格式是JSONB,配合pgvector甚至可以直接存储向量状态。

选型建议

  • 单机部署、状态<1GB:用SQLite
  • 多实例部署、需要故障恢复:用Postgres
  • 极致性能、可接受恢复重建:用Memory + 定期快照

最佳实践与避坑指南

四个必须避开的坑

坑1:在节点内部修改状态对象

# ❌ 错误:直接修改传入的state
def bad_node(state):
    state["counter"] += 1  # 这会污染检查点!
    return state

# ✅ 正确:返回新的状态字典
def good_node(state):
    return {"counter": state["counter"] + 1}

LangGraph的_update_state会深度比较新旧状态差异,直接修改可变对象可能导致检查点无法正确捕获变更。

坑2:线程ID使用不当

# ❌ 每次请求都生成新thread_id,导致无法恢复
config = {"configurable": {"thread_id": uuid4().hex}}

# ✅ 从请求上下文(如用户ID)派生稳定thread_id
config = {"configurable": {"thread_id": f"user-{user_id}-task-{task_id}"}}

坑3:忽略checkpoint_id的优先级

config中同时存在thread_idcheckpoint_id时,LangGraph优先使用checkpoint_id。这会导致“我明明更新了状态,但执行还是旧的”的错觉。

坑4:在interrupt节点返回非空状态

# ❌ 错误:interrupt节点返回了状态
def interrupt_node(state):
    state["approved"] = True  # 这会导致图不暂停!
    return {"approved": True}

# ✅ 正确:interrupt节点返回空状态
def interrupt_node(state):
    print("等待人工审批...")
    return {}  # 空状态触发暂停

三条最佳实践

  1. 定期压缩检查点:长时间运行的Agent会产生大量检查点。建议每小时合并旧检查点,只保留关键节点(如每个HITL决策点的状态)。
  1. 使用get_state而非直接访问Savergraph.get_state(config)会返回当前状态,包括图内部的next节点信息,比直接读Saver更安全。
  1. 为检查点添加业务元数据
saver.put(config, checkpoint, metadata={
    "user_id": user_id,
    "business_step": "approval",
    "llm_cost": total_tokens * price
})

这为后续的审计、成本分析、异常回溯提供了数据基础。

总结

LangGraph的Checkpointing远远超出了“保存进度”的范畴。它是Agent工作流可观测性可恢复性人机协作的基石。通过源码分析我们看到,它的核心设计是不可变检查点链 + 基于父指针的分支重放,这赋予了Agent系统类似Git的版本控制能力。

延伸思考:当你的Agent工作流需要处理多用户并发审批(如多人会签)时,单一的thread_id会失效。LangGraph 0.2+的checkpoint_id支持多级嵌套,可以解决这个问题——这将是下一篇博客的主题。在那之前,建议你动手将现有的Agent流程加上Checkpointing,体验一次“状态不丢失”的安全感。

代码仓库:github.com/yourname/langgraph-checkpointing-demo (示例代码可直接运行)