生产环境LangGraph Agent因checkpointing策略不当导致多步任务每次从零开始;SqliteSaver在多实例场景存在缺陷,需切换至更适合生产规模的持久化方案。
最初发表于 AIdeazz——经授权在此交叉发布,保留原始链接。
我的第一个 LangGraph 生产环境智能体(Agent)在数周内静默丢弃了每一个任务。这个智能体原本应该从 Telegram 接收用户请求,将其分解为子任务,然后在多个步骤中执行,并存储中间结果。然而,它处理完第一步后,下一次调用时就会从头开始,失去所有之前的上下文。问题不在于智能体逻辑本身,而在于我尝试持久化其状态的方式。
我经历了三种截然不同的 LangGraph 检查点(checkpointing)策略,才找到一种能够可靠地为生产环境中的有状态智能体工作的方案。每一次失败都教会了我关于 LangGraph 所做的假设以及部署多步骤 AI 智能体的现实情况的一个关键教训。
我最初的做法很简单:使用 LangGraph 内置的 SqliteSaver。它对于 Oracle Cloud VM 上的单实例部署来说看起来足够稳健。智能体的图(graph)定义了一个状态,姑且称之为 AgentState,其中包含 user_id: str、request_id: str、task_list: list[str] 和 current_step: int 等字段。
class AgentState(TypedDict):
user_id: str
request_id: str
task_list: list[str]
current_step: int
# ... 后续会添加更多字段
智能体接收消息后,会初始化 AgentState 并启动图。同一个 user_id 的后续消息应该恢复现有状态。
问题始于我需要向 AgentState 添加一个新字段,比如 tool_output: dict。我更新了 TypedDict,重新部署了智能体,期望它从中断处继续运行。但它没有。现有的对话会重启,新的对话则正常工作。
我调试了好几天,追踪 SqliteSaver 的调用,直接检查数据库。SQLite 中的检查点表存储的是状态的 JSON blob。我发现的问题很隐蔽:LangGraph 的 SqliteSaver 不会执行 schema 迁移,甚至不会对 schema 不匹配发出警告。当旧的检查点(没有 tool_output)被加载到期望新 AgentState schema 的智能体中时,TypedDict 实例化会静默丢弃加载的 JSON 中不存在的任何字段,而且如果新字段没有显式处理,也不会用默认值初始化。
智能体会加载一个不完整的状态,然后假装完整地继续运行,之后因为缺少 tool_output 而在下游失败。更糟糕的是,它会因为部分加载导致 current_step 这样的关键标志被重置,从而直接重启整个流程。这不是错误,是静默的数据丢失。我的解决方案是一个手动且痛苦的过程:导出 SQLite 检查点,手动迁移 JSON,然后重新插入。这不可扩展。
在 SqliteSaver 灾难之后,我转向了自定义的 OracleCloudObjectStorageSaver。我的智能体运行在 Oracle Cloud Infrastructure(OCI)上,对象存储既便宜又高可用。我实现了一个 BaseCheckpointSaver 子类,将状态序列化为 JSON 并上传到 OCI 对象存储桶,使用 thread_id 作为对象名称。
class OracleCloudObjectStorageSaver(BaseCheckpointSaver):
def __init__(self, bucket_name: str, namespace: str, object_storage_client):
self.bucket_name = bucket_name
self.namespace = namespace
self.client = object_storage_client
def get(self, thread_id: str) -> Optional[Checkpoint]:
try:
response = self.client.get_object(self.namespace, self.bucket_name, thread_id)
data = json.loads(response.data.content.decode('utf-8'))
return Checkpoint(**data)
except Exception as e:
# 日志记录,如果对象不存在或损坏则返回 None
print(f"Error loading checkpoint {thread_id}: {e}")
return None
def put(self, thread_id: str, checkpoint: Checkpoint) -> None:
data = json.dumps(checkpoint, default=str) # default=str 处理 datetime 对象
self.client.put_object(
self.namespace, self.bucket_name, thread_id, data.encode('utf-8')
)
这看起来更稳健。我完全控制序列化和反序列化,可以为状态对象添加版本控制并显式处理迁移。
然后检查点损坏了。偶尔,智能体会无法加载其状态,报告 JSON 解码错误。get_object 调用返回有效数据,但 json.loads 会抛出异常。经检查,OCI 对象存储中的 JSON 文件被截断或格式损坏。
根本原因是并发。我的智能体本身是无状态的,作为无服务器函数或可以扩展的 VM 运行。同一个 thread_id 的多个调用可能几乎同时发生,特别是当用户发送连续快速消息时。如果两个 put 操作并发执行,一个可能会部分覆盖另一个,或者读取可能发生在写入进行中,导致 JSON 损坏。OCI 对象存储提供最终一致性,但不提供覆盖的强一致性。
我最初的修复是为 put 操作添加带指数退避的重试机制。这降低了损坏频率,但没有消除它。问题本质上是:在高度并发环境中,简单覆盖状态模型是等待发生的竞态条件。
最终稳定生产环境中 LangGraph 有状态智能体的解决方案涉及两个关键组件:显式状态版本控制和条件写入。
我不再只存储 Checkpoint 对象,而是将其包装在一个包含版本号的自定义信封中。
class VersionedCheckpoint(TypedDict):
version: int
checkpoint: Checkpoint
加载检查点时,我会读取版本字段。写入时,我会递增版本号。关键部分是条件写入。OCI 对象存储与 S3 类似,支持使用 If-Match 或 If-None-Match 头基于对象的 ETag 进行条件请求。这允许乐观锁。
我修改了 put 操作:
class OracleCloudObjectStorageAtomicSaver(BaseCheckpointSaver):
def __init__(self, bucket_name: str, namespace: str, object_storage_client):
self.bucket_name = bucket_name
self.namespace = namespace
self.client = object_storage_client
self.max_retries = 5
def get(self, thread_id: str) -> Optional[Checkpoint]:
try:
response = self.client.get_object(self.namespace, self.bucket_name, thread_id)
data = json.loads(response.data.content.decode('utf-8'))
# 期望 VersionedCheckpoint 结构
versioned_data = VersionedCheckpoint(**data)
return Checkpoint(**versioned_data['checkpoint'])
except Exception as e:
print(f"Error loading checkpoint {thread_id}: {e}")
return None
def put(self, thread_id: str, checkpoint: Checkpoint) -> None:
for attempt in range(self.max_retries):
current_etag = None
current_version = 0
# 尝试获取当前对象和 ETag 用于条件写入
try:
response = self.client.get_object(self.namespace, self.bucket_name, thread_id)
current_etag = response.headers.get('etag')
existing_data = json.loads(response.data.content.decode('utf-8'))
current_version = existing_data.get('version', 0)
except Exception as e:
# 对象可能不存在,或其他临时错误。继续而不使用 ETag。
print(f"No existing object or error getting ETag for {thread_id}: {e}")
new_version = current_version + 1
versioned_checkpoint = VersionedCheckpoint(version=new_version, checkpoint=checkpoint)
data_to_write = json.dumps(versioned_checkpoint, default=str)
try:
headers = {'If-Match': current_etag} if current_etag else {}
self.client.put_object(
self.namespace, self.bucket_name, thread_id, data_to_write.encode('utf-8'),
opc_meta={'version': str(new_version)}, # 也将版本存储在元数据中
**headers
)
return # 成功
except Exception as e:
if "412 Precondition Failed" in str(e) and attempt < self.max_retries - 1:
print(f"Precondition failed for {thread_id}, retrying (attempt {attempt+1})...")
time.sleep(2 ** attempt) # 指数退避
else:
raise # 如果达到最大重试次数或其他错误则重新抛出
raise Exception(f"Failed to save checkpoint for {thread_id} after {self.max_retries} attempts.")
这个模式确保对于给定的 thread_id,一次只有一个写入操作成功。如果多个智能体尝试并发更新同一状态,只有其 If-Match 头正确标识当前状态 ETag 的那个会成功。其他会失败并重试,最终获取新写入的状态并在其上应用更新。这有效地将并发更新串行化到同一检查点。
这种原子更新模式与显式状态版本控制和 schema 管理相结合,终于为生产环境 LangGraph 智能体提供了所需的稳定性。无论智能体是将请求路由到 Groq 获取快速初始响应还是路由到 Claude 进行复杂推理,它们现在都能在多次步骤中可靠地维护状态。OCI 对象存储的成本可以忽略不计,对于数十万个检查点通常每月不到 5 美元。
LangGraph 的 SqliteSaver 适用于单进程、不演化的状态。对于演示来说没问题,但不适用于生产环境中状态 schema 会变化或涉及并发的场景。
显式管理你的状态 schema。TypedDict 是编译时提示,不是运行时验证器或迁移器。为你的状态对象实现你自己的版本控制和迁移逻辑。
并发是有状态的杀手。任何共享状态在分布式或并发系统中都需要原子更新机制。简单覆盖会导致静默数据损坏。
利用云原生原语。具有条件写入(If-Match / ETag)的对象存储是用于乐观锁的强大且成本效益高的原语。除非绝对必要,否则不要用复杂的分布式锁重新发明轮子。
监控和记录一切。静默失败是最难调试的。对检查点加载、保存、版本和重试的广泛日志记录至关重要。
我目前的智能体,从客户支持路由到内部数据分析,都使用这个模式运行。三次重写的最初痛苦是值得的,因为它带来了稳定性和信心。
Q: 为什么不使用带行级锁的 PostgreSQL 等 proper 数据库进行检查点?
A: 完整的关系统数据库会增加运营开销(管理、备份、扩展)和成本。对于简单的键值状态,带条件写入的对象存储提供了足够的原子性,而且在规模上操作起来比这种特定用例的数据库便宜和简单几个数量级。我目前的设置每月成本不到 5 美元。
Q: 用这种方法如何处理 AgentState 的 schema 迁移?
A: 加载 VersionedCheckpoint 时,我会检查版本字段。如果加载的版本比智能体当前期望的 schema 版本旧,我会应用显式迁移函数(例如,为新字段添加默认值、转换旧字段名),然后再实例化 AgentState TypedDict。
Q: 如果智能体在 put 操作期间崩溃,留下损坏的检查点怎么办?
A: put 操作设计为幂等且有弹性的。如果在写入中途发生崩溃,下次 get 操作将要么获取上次成功写入的检查点(如果部分写入没有覆盖 ETag),要么无法解析 JSON,触发重试或全新开始。条件写入有助于防止部分写入损坏有效的先前状态。
Q: 这种方法会为每个智能体步骤增加显著延迟吗?
A: 每次 get 和 put 操作都涉及对 OCI 对象存储的网络调用。对于典型的智能体步骤,每次状态访问增加 50-200 毫秒的延迟,这对于大多数对话式 AI 应用来说是可以接受的,因为 LLM 调用主导延迟(例如 Groq 100ms,Claude 1-5s)。对于极高吞吐量、低延迟场景,可以在顶部分层带有最终一致性的内存缓存。
— Elena Revicheva · AIdeazz · Portfolio