Agent开发第9步:持久化与 Checkpoint
第1步:抽象模型接口
第2步:定义工具抽象
第3步:实现 Agent Loop
第4步:定义 RunEvent
第5步:实现 Run Store
第6步:实现 SSE 编解码
第7步:创建订阅与打断接口
第8步:开发简易客户端
第9步:持久化与 Checkpoint
第10步:异常处理与测试
进入到“生产化(Production)”阶段,我们在第5步中编写的 InMemoryRunStore 显然是不够用的:只要服务器一重启,所有的对话历史和执行状态就会灰飞烟灭。
为了让 Agent 能够稳定运行并支持长线任务,我们需要接入真正的数据库,并引入**Checkpoint(检查点)**机制。
1. 为什么需要 Checkpoint?
在长程推理任务(如写代码、做深度研究)中,Agent Loop 可能会循环几十次。如果在这期间,服务器由于 OOM(内存溢出)被系统杀掉,或者发生了停机更新,会导致整个任务失败。
Checkpoint 机制是指:在 Agent Loop 的每一个关键步骤(例如思考结束、工具调用完成),将当前的状态快照序列化并保存到持久化存储(数据库)中。
如果进程崩溃,重启后调度器可以读取最新的 Checkpoint,从断点处继续执行。
2. 数据库选型与架构
持久化存储通常需要保存两类数据:
- 结构化数据:如 Run 的状态(pending/running/completed)、时间戳等。适合关系型数据库(如 PostgreSQL)。
- 半结构化/文档数据:如对话历史 (messages)、事件流序列 (events)、状态快照。适合文档型或 NoSQL 数据库(如 MongoDB, Redis, 或者是 PostgreSQL 的 JSONB 字段)。
在这里,我们以一个抽象的 DatabaseStore 和 PostgreSQL/JSONB 的思路来进行演示。
3. Python 代码实现:持久化 Store
import json
from typing import Optional
# 假设我们使用了 SQLAlchemy 等 ORM
# from sqlalchemy.orm import Session
# from models import RunRecord
class PostgresRunStore:
def __init__(self, db_session):
self.db = db_session
def create_run(self, run_id: str, initial_message: str):
"""在数据库中插入一条新的运行记录"""
# record = RunRecord(
# id=run_id,
# status="pending",
# messages=json.dumps([{"role": "user", "content": initial_message}]),
# events="[]"
# )
# self.db.add(record)
# self.db.commit()
pass
def save_checkpoint(self, run_id: str, messages: list, current_state: dict):
"""
保存检查点
更新数据库中的 messages 列表和当前的内部状态
"""
# record = self.db.query(RunRecord).filter_by(id=run_id).first()
# record.messages = json.dumps(messages)
# record.state_snapshot = json.dumps(current_state)
# self.db.commit()
print(f"💾 [DB] 已保存 Run {run_id} 的 Checkpoint。")
def load_checkpoint(self, run_id: str) -> Optional[dict]:
"""加载最新的检查点用于恢复"""
# record = self.db.query(RunRecord).filter_by(id=run_id).first()
# if record and record.status != "completed":
# return {
# "messages": json.loads(record.messages),
# "state_snapshot": json.loads(record.state_snapshot)
# }
return None
4. 将 Checkpoint 接入 Agent Loop
在 Agent 的循环中,我们在每一次 Think 和 Observe 之后调用 save_checkpoint。
class ResilientAgent:
def __init__(self, model, tools, store: PostgresRunStore):
self.model = model
self.tools = tools
self.store = store
def resume_or_run(self, run_id: str):
# 1. 尝试加载 Checkpoint
checkpoint = self.store.load_checkpoint(run_id)
if checkpoint:
print(f"🔄 正在从 Checkpoint 恢复任务: {run_id}")
messages = checkpoint["messages"]
# 恢复其他内部状态...
else:
print(f"🚀 开始新任务: {run_id}")
messages = [{"role": "system", "content": "..."}] # 初始化
# 2. Agent Loop
for step in range(5):
# ... 思考 ...
response = self.model.chat_with_tools(messages, self.tools.declarations)
# 思考结束,保存快照
messages.append(response.message_dict)
self.store.save_checkpoint(run_id, messages, {"step": step, "phase": "thought"})
if response.tool_calls:
for tool_call in response.tool_calls:
# ... 执行工具 ...
result = self.tools.execute(tool_call)
messages.append({"role": "tool", "content": result})
# 工具执行完毕,保存快照
self.store.save_checkpoint(run_id, messages, {"step": step, "phase": "tool_executed"})
continue
# 完成任务
self.store.update_status(run_id, "completed")
break
5. 分布式调度(拓展)
结合 Checkpoint,如果你有多个后端节点集群,还可以引入类似于 Celery 或 RabbitMQ 的任务队列。
如果 Node A 在处理任务时宕机,任务队列在超时后会将任务重新分配给 Node B。Node B 接手后,直接从数据库拉取最新的 Checkpoint 即可继续执行,实现真正的高可用 (High Availability)。
总结
持久化与 Checkpoint 是区分“玩具项目”与“工业级框架”的分水岭。通过将状态落盘,Agent 具备了穿越重启与宕机的生命力。
但这还不够,在复杂的现实世界里,大模型会胡言乱语,第三方 API 会超时报错。在最后一篇中,我们将探讨异常处理与测试,为 Agent 加上坚固的护甲。
更多推荐


所有评论(0)