第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,从断点处继续执行。

进程突然崩溃 ❌

从断点恢复

Start

Save_Checkpoint_1

Think

Call_Tool

Save_Checkpoint_2

运维重启服务

Restore_Checkpoint_2

Return_Tool_Result

Think_Again

Final_Answer

2. 数据库选型与架构

持久化存储通常需要保存两类数据:

  1. 结构化数据:如 Run 的状态(pending/running/completed)、时间戳等。适合关系型数据库(如 PostgreSQL)。
  2. 半结构化/文档数据:如对话历史 (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 加上坚固的护甲。

Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区

更多推荐