云计算百科
云计算领域专业知识百科平台

LangGraph 持久化详解

LangGraph 持久化详解


一、持久化概述

持久化(Persistence)是 LangGraph 的核心能力之一,它允许图状态在执行过程中被保存和恢复。LangGraph 提供了两种互补的持久化机制:

机制作用存储范围记忆类型
Checkpointer 保存图状态快照 单个线程(Thread) 短期记忆
Store 保存应用数据 跨线程 长期记忆

核心价值:

  • 人机交互:支持中断-恢复工作流
  • 对话记忆:跨多轮对话保持上下文
  • 时间旅行:回溯到历史状态重新执行
  • 容错恢复:故障后从检查点恢复执行

二、Checkpointer(检查点记录器)

2.1 核心概念

Thread(线程)

线程是检查点的组织单元,每个线程有一个唯一的 thread_id。线程包含了一系列运行的累积状态。

config = {"configurable": {"thread_id": "user-123"}}
graph.invoke({"input": "hello"}, config)

Checkpoint(检查点)

检查点是在每个超级步骤(Super-step)保存的图状态快照,由 StateSnapshot 对象表示。

超级步骤说明:

  • 对于顺序图 START → A → B → END,会创建4个检查点
  • 每个节点执行后都会保存一个检查点
  • 并行节点在同一超级步中执行

2.2 内置实现

InMemorySaver(内存检查点)

from langgraph.checkpoint.memory import InMemorySaver

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

# 使用
config = {"configurable": {"thread_id": "1"}}
result = graph.invoke({"foo": ""}, config)

注意: InMemorySaver 将检查点存储在 RAM 中,进程重启后数据丢失,仅适用于开发和测试。

SqliteSaver(SQLite 检查点)

from langgraph.checkpoint.sqlite import SqliteSaver

# 同步版本
checkpointer = SqliteSaver.from_conn_string("checkpoints.sqlite")
graph = builder.compile(checkpointer=checkpointer)

# 异步版本
from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver

async with AsyncSqliteSaver.from_conn_string("checkpoints.sqlite") as checkpointer:
graph = builder.compile(checkpointer=checkpointer)
result = await graph.ainvoke({"foo": ""}, config)

PostgresSaver(PostgreSQL 检查点)

from langgraph.checkpoint.postgres import PostgresSaver
from psycopg_pool import ConnectionPool

DB_URI = "postgresql://postgres:postgres@localhost:5432/postgres"

# 使用连接池
with ConnectionPool(conninfo=DB_URI, max_size=20) as pool:
checkpointer = PostgresSaver(pool)
checkpointer.setup() # 首次使用需要初始化数据库

graph = builder.compile(checkpointer=checkpointer)
result = graph.invoke({"foo": ""}, config)

# 使用连接字符串
with PostgresSaver.from_conn_string(DB_URI) as checkpointer:
graph = builder.compile(checkpointer=checkpointer)

2.3 检查点器接口

所有检查点器都实现 BaseCheckpointSaver 接口:

方法说明
put 存储检查点及其配置和元数据
put_writes 存储中间写入(待处理写入)
get_tuple 获取检查点元组,用于 graph.get_state()
list 列出符合配置的检查点,用于 graph.get_state_history()

2.4 StateSnapshot 字段

graph.get_state() 返回的 StateSnapshot 包含以下字段:

字段类型说明
values dict 此检查点处的通道状态值
next tuple 下一个要执行的节点名称,空元组表示图已结束
config dict 包含 thread_id、checkpoint_ns、checkpoint_id
metadata dict 执行元数据:source、writes、step
created_at str 检查点创建时间(ISO 8601)
parent_config dict | None 上一个检查点的配置
tasks tuple 此步骤要执行的任务列表

三、状态管理

3.1 获取状态

# 获取最新状态
config = {"configurable": {"thread_id": "1"}}
state = graph.get_state(config)

# 获取特定检查点
config = {
"configurable": {
"thread_id": "1",
"checkpoint_id": "1ef663ba-28fe-6528-8002-5a559208592c"
}
}
state = graph.get_state(config)

3.2 获取状态历史

config = {"configurable": {"thread_id": "1"}}
history = graph.get_state_history(config)

# 遍历历史
for checkpoint in history:
print(f"Step: {checkpoint.metadata['step']}")
print(f"Values: {checkpoint.values}")
print(f"Next: {checkpoint.next}")

3.3 更新状态

使用 update_state 编辑图状态,会创建新的检查点:

config = {"configurable": {"thread_id": "1"}}

# 更新状态
graph.update_state(
config,
values={"foo": "updated_value"},
as_node="node_a" # 可选:指定更新来源节点
)

关键点:

  • 更新不会修改原始检查点,而是创建新检查点
  • 如果字段有 reducer 函数,值会通过 reducer 处理
  • as_node 参数影响接下来执行哪个节点

四、时间旅行(Time Travel)

时间旅行允许从历史检查点重新执行图,是 LangGraph 的强大调试能力。

4.1 重放(Replay)

从特定检查点重新执行后续步骤:

# 从 checkpoint_id 重放
config = {
"configurable": {
"thread_id": "1",
"checkpoint_id": "1ef663ba-28fe-6528-8002-5a559208592c"
}
}
result = graph.invoke(None, config) # 传入 None 触发重放

重放行为:

  • 检查点之前的步骤被跳过(已保存)
  • 检查点之后的步骤重新执行
  • LLM 调用、API 请求、中断会重新触发

4.2 分叉(Fork)

通过 update_state 创建状态分叉:

# 获取历史
history = list(graph.get_state_history(config))
oldest_checkpoint = history[1]

# 从旧检查点分叉
fork_config = oldest_checkpoint.config
graph.update_state(
fork_config,
values={"foo": "forked_value"},
as_node="node_a"
)

# 继续执行(会创建新的执行路径)
result = graph.invoke(None, fork_config)


五、人机交互(Human-in-the-Loop)

持久化是实现人机交互工作流的基础。

5.1 中断点配置

在编译图时配置中断点:

graph = builder.compile(
checkpointer=checkpointer,
interrupt_before=["human_review_node"], # 进入节点前中断
interrupt_after=["ai_generate_node"] # 离开节点后中断
)

5.2 中断-恢复流程

# 1. 执行图,会在中断点暂停
config = {"configurable": {"thread_id": "1"}}
result = graph.invoke({"input": "user query"}, config)

# 2. 检查状态
state = graph.get_state(config)
print(state.next) # ('human_review_node',)

# 3. 人工审核并更新状态
graph.update_state(
config,
values={"approval": True, "feedback": "looks good"},
as_node="human_review_node"
)

# 4. 恢复执行
result = graph.invoke(None, config)

5.3 待处理写入(Pending Writes)

当超级步中的节点执行失败时,LangGraph 会保存成功节点的写入:

# 恢复时,已成功节点的写入会被保留
# 无需重新执行成功的节点
result = graph.invoke(None, config)


六、Store(存储)

Store 用于跨线程持久化应用数据,实现长期记忆。

6.1 基本使用

from langgraph.store.memory import InMemoryStore

store = InMemoryStore()

# 编译图时同时使用 checkpointer 和 store
graph = builder.compile(
checkpointer=checkpointer,
store=store
)

6.2 在节点中访问 Store

from langgraph.store.base import BaseStore

def my_node(state, config, *, store: BaseStore):
user_id = config["configurable"]["user_id"]
namespace = ("user_preferences", user_id)

# 读取存储的数据
memories = store.search(namespace, query="favorite color")

# 写入数据
store.put(
namespace,
"preference_1",
{"data": "blue is favorite color"}
)

return {"response": f"Hello {user_id}!"}

6.3 跨线程记忆示例

from langgraph.checkpoint.memory import InMemorySaver
from langgraph.store.memory import InMemoryStore

checkpointer = InMemorySaver()
store = InMemoryStore()

# 节点函数
def chatbot(state, config, *, store: BaseStore):
user_id = config["configurable"]["user_id"]
namespace = ("memories", user_id)

# 获取用户记忆
memories = store.search(namespace)
memory_text = "\\n".join([m.value["data"] for m in memories])

# 检查是否需要记住新信息
last_message = state["messages"][1].content
if "remember" in last_message.lower():
store.put(namespace, str(uuid.uuid4()), {"data": "User prefers dark mode"})

# 使用记忆生成回复
response = model.invoke([
{"role": "system", "content": f"User memories:\\n{memory_text}"},
*state["messages"]
])
return {"messages": [response]}

# 编译图
graph = builder.compile(checkpointer=checkpointer, store=store)

# 线程1:用户告诉机器人记住偏好
config1 = {"configurable": {"thread_id": "thread-1", "user_id": "user-123"}}
graph.invoke({"messages": [("user", "Remember that I like blue")]}, config1)

# 线程2:仍然可以访问用户记忆
config2 = {"configurable": {"thread_id": "thread-2", "user_id": "user-123"}}
graph.invoke({"messages": [("user", "What's my favorite color?")]}, config2)


七、序列化

检查点器需要对状态值进行序列化和反序列化。

7.1 JsonPlusSerializer

LangGraph 默认使用 JsonPlusSerializer,支持多种类型:

  • LangChain 和 LangGraph 基本类型
  • 日期时间和枚举
  • Pydantic 模型
  • 自定义类

from langgraph.checkpoint.serde.jsonplus import JsonPlusSerializer

serde = JsonPlusSerializer()
checkpointer = InMemorySaver(serde=serde)

7.2 EncryptedSerializer(加密序列化)

支持加密持久化状态:

from langgraph.checkpoint.serde.encrypted import EncryptedSerializer

# 从环境变量读取密钥
encrypted_serde = EncryptedSerializer.from_pycryptodome_aes()

# 或直接指定密钥
encrypted_serde = EncryptedSerializer.from_pycryptodome_aes(key="your-secret-key")

checkpointer = InMemorySaver(serde=encrypted_serde)


八、Checkpointer 对比

特性InMemorySaverSqliteSaverPostgresSaver
存储位置 RAM 本地文件 数据库
持久性 ❌ 进程重启丢失 ✅ 文件持久 ✅ 数据库持久
异步支持
并发性能
适用场景 开发测试 本地开发 生产环境
安装命令 内置 pip install langgraph-checkpoint-sqlite pip install langgraph-checkpoint-postgres

九、最佳实践

9.1 选择合适的检查点器

# 开发环境
checkpointer = InMemorySaver()

# 本地开发/小型项目
checkpointer = SqliteSaver.from_conn_string("checkpoints.sqlite")

# 生产环境
checkpointer = PostgresSaver.from_conn_string("postgresql://…")
checkpointer.setup()

9.2 管理检查点大小

长期对话会导致检查点无限增长,定期清理:

# 获取状态历史
history = graph.get_state_history(config)

# 删除旧检查点(保留最近N个)
for checkpoint in history[10:]: # 保留前10个
checkpointer.delete_thread(config["configurable"]["thread_id"])

9.3 使用 Reducer 处理状态更新

from typing import Annotated
from operator import add

class State(TypedDict):
messages: Annotated[list, add] # 消息列表使用 add reducer
context: str # 普通字段直接覆盖

# update_state 时,messages 会追加而不是覆盖
graph.update_state(config, values={"messages": [new_message]})

9.4 子图持久化

为子图配置独立的 checkpointer:

# 子图使用独立 checkpointer
subgraph_builder = StateGraph(SubState)
subgraph = subgraph_builder.compile(checkpointer=True)

# 或共享父图的 checkpointer
subgraph = subgraph_builder.compile(checkpointer=parent_checkpointer)


十、实战示例:多轮对话机器人

from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.store.memory import InMemoryStore
from langchain_core.messages import BaseMessage, HumanMessage, AIMessage
from langchain_openai import ChatOpenAI

class ChatState(TypedDict):
messages: Annotated[list[BaseMessage], "append"]

# 模型
model = ChatOpenAI(model="gpt-4")

# 节点
def chatbot(state: ChatState):
response = model.invoke(state["messages"])
return {"messages": [response]}

# 构建图
builder = StateGraph(ChatState)
builder.add_node("chatbot", chatbot)
builder.add_edge(START, "chatbot")
builder.add_edge("chatbot", END)

# 持久化配置
checkpointer = InMemorySaver()
store = InMemoryStore()

graph = builder.compile(
checkpointer=checkpointer,
store=store
)

# 多轮对话
config = {"configurable": {"thread_id": "user-001"}}

# 第一轮
result1 = graph.invoke(
{"messages": [HumanMessage(content="Hi, I'm Alice")]},
config
)
print(result1["messages"][1].content)

# 第二轮(自动保持上下文)
result2 = graph.invoke(
{"messages": [HumanMessage(content="What's my name?")]},
config
)
print(result2["messages"][1].content) # "Your name is Alice"

# 查看状态历史
history = list(graph.get_state_history(config))
print(f"Total checkpoints: {len(history)}")


十一、总结

LangGraph 的持久化机制为构建强大的 AI 应用提供了基础:

  • Checkpointer:保存线程级别的图状态快照

    • InMemorySaver:开发测试
    • SqliteSaver:本地开发
    • PostgresSaver:生产环境
  • Store:跨线程持久化应用数据

    • 用户偏好和记忆
    • 共享知识库
  • 状态管理:

    • get_state:获取当前状态
    • get_state_history:查看执行历史
    • update_state:编辑状态并分叉
  • 高级特性:

    • 时间旅行:从历史状态重放
    • 人机交互:中断-恢复工作流
    • 容错恢复:待处理写入保留
  • 赞(0)
    未经允许不得转载:网硕互联帮助中心 » LangGraph 持久化详解
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!