尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

自我改进智能体的事件溯源架构:可追溯、可重放、可回滚

自我改进智能体的事件溯源架构:可追溯、可重放、可回滚 大家好我最近在一套自研 Agent 平台里做“自我改进”功能遇到了一个所有团队最终都会撞上的问题智能体确实能通过反思改自己的策略但代码库和数据库里只留下了一个“最新策略”完全看不到它为什么改、怎么改、改错了能不能还原。后来我试着用事件溯源Event Sourcing的思路重新组织整个系统发现“自我改进智能体”这个场景几乎是事件溯源的最佳拍档。这篇文章会把概念讲清楚再用一个完整可运行的 Python 项目演示如何把 Agent 的运行、反思、策略发布和回滚统一建模成一份不可变的追加事件流。如果你正在做 Agent 的 Prompt 自动优化、技能库更新、策略自学习或者只是想让智能体的行为可回滚、可审计这篇文章应该能给你一套可以落地的架构思路。1. 背景与核心概念1.1 自我改进智能体到底在改进什么在深入事件溯源之前我们需要先明确“自我改进智能体”是什么。传统 Agent 的流程通常是固定 Prompt 工具调用执行完就结束。它的策略在开发阶段被写好运行时不会主动改变自己。而自我改进智能体会在运行过程中根据执行结果、失败样例、用户反馈生成新的反思或策略再把这些经验沉淀下来影响后续的每一轮执行。这套流程近年来已经被多个知名项目验证过Reflexion 会把执行失败的反思以文本形式保存下来在下一轮尝试时重新读完再决策。Self-Refine 会让模型对自己的输出生成反馈然后迭代改进。Voyager 会把探索到的经验封装成代码技能存入技能库持续扩展能力。MemGPT 通过分层记忆让 Agent 获得更大的上下文和长期记忆。这些项目的共同点是它们都有一个“执行 - 评价 - 提炼经验 - 更新策略 - 再执行”的过程。也就是说改进不是一个瞬间赋值而是一条时间线。真正的问题在于绝大多数工程落地时把这条时间线压缩成了一个字段比如agent.policy 新版 Prompt。一旦这样设计你就永远失去了“为什么”。所以我们需要换一种持久化视角那就是事件溯源。1.2 事件溯源用流水账而不是余额来建模事件溯源是领域驱动设计里一个很经典的架构模式在金融、审计、追溯类系统里被大量使用。你应该见过银行账户余额。余额是当前状态但它本质上是“存入、取出”这一系列事件的投影。如果我们只存余额就丢失了每一笔钱的来源如果我们存下一整条存取流水随时可以重放流水算出任意时刻的余额。事件溯源正是采用后一种模型把所有改变系统状态的动作记成不可变的“事件”。事件只能追加不能修改不能删除。当前状态只是从事件流里逐步重放出来的“投影”。投影可以丢弃事件不能丢因为事件是唯一的事实来源。举个例子传统状态存储会执行UPDATE agent SET current_policy 新的策略 WHERE id agent_001;事件溯源则会记录事件1AgentRunStarted任务是计算总价 事件2ToolExecuted结果是 30缺少单位 事件3ReflectionGenerated输出需要带单位 事件4PolicyApplied新策略为“输出计算结果时带上单位”这两者最大的区别不是表结构而是你选择相信“最终值”还是相信“发生的事情”。对普通业务系统最终值往往够用。但对一个会修改自己策略的智能体过程事实才是最有价值的部分。1.3 自我改进智能体与事件溯源为什么天然契合如果我们把“智能体的一次自我改进”拆开看它本身就是一个完整的事件链。一次改进至少包含一个任务被交给 Agent。Agent 使用了某个策略去执行。工具返回结果。反思引擎根据结果生成了经验。经验被加工成策略候选。策略经过审批、生效。后续任务开始使用新策略。这条事件链就是我们需要的“流水账”。如果让它和模型中的Agent对象解耦不要再试图去维护一个“最新的策略字符串”而是每次需要策略时从事件流投影出来整个系统就会变得非常干净。接下来我们先把“为什么要用事件溯源”这个话题聊透然后直接进入实战。2. 为什么自我改进智能体需要事件溯源2.1 可追溯知道每个策略从哪里来对自我改进智能体来说最危险的操作是“自动改我们自己的大脑”。这里的“大脑”可能是 Prompt、工具描述、奖励函数、技能代码。如果只保存最终策略那么当线上效果突然变差时你看到的是一个已经变成“新策略”的文本完全不知道它是在哪个失败任务后被谁生成出来的。这就是为什么我们需要把策略候选产生时对应的任务、反思、审批过程都记录下来。事件溯源带来的是完整的因果链路哪次失败触发了哪段反思哪段反思被批准成了策略。没有这种可追溯性任何看似智能的自我改进都只是黑盒长期维护会非常痛苦。2.2 可重放拿历史事件做离线评估事件溯源给了我们一个被很多人忽略的能力重放历史。假设你现在有两个候选策略一个是当前的一个是反思引擎新生成的。你拿不准新策略是否更好最稳妥的办法不是直接上生产而是把过去七天的事件流取出来用两个策略各自重跑一遍。只要事件流里保存了当时的任务输入、执行工具和结果判定你就能离线对比新旧策略的正确率。这种“基于真实历史重放的评估”比拍脑袋相信“反思文本越长越好”要可靠得多。这在工程上非常有用。事件溯源天然支持这种“时间旅行式实验”因为你记录的是历史事实而不是当时某个时刻的临时状态。2.3 可回滚策略坏了能一秒还原自我改进智能体上线之后一定会遇到“新策略反而把所有任务搞坏”的情况。在传统状态存储方式下回滚需要你手动备份旧策略、再覆盖新策略很容易出现遗漏。在事件溯源模型里回滚只需要追加一条新事件比如PolicyRolledBack投影逻辑会让已经回滚的策略不再生效。历史事件不需要改动策略的“当前状态”自然回到上一个有效版本。这就像 Git你不需要把代码仓库“改回去”只需要 reset 到某个 commit所有历史提交还在。事件溯源给了智能体一个同样可靠的版本控制底座。2.4 可审计给智能体的自我修改装上记录仪在真实业务环境里Agent 自动修改自己行为边界会带来合规和安全问题。谁允许策略变更策略候选经过了什么检查是否有人类审批这些都是需要回答的问题。事件溯源提供的不可变事件流就是一套天然审计日志。每一次策略提案、批准、驳回、上线、回滚全部在事件表里记录得清清楚楚。配合权限系统可以确保只有经过授权的人或自动化流程才能追加策略相关事件。2.5 状态存储与事件溯源对比对比维度状态存储事件溯源存储内容当前状态例如current_policy发生过的全部事实修改方式覆盖追加能否回到过去不能能重放到任意事件位置审计能力弱强排查问题只能看最终值能看到完整链路存储成本低高需要快照控制实现复杂度低偏高但可控3. 环境准备与系统设计3.1 环境准备本文的示例全部基于 Python 标准库完成没有第三方依赖方便你直接复制运行。操作系统Windows / macOS / Linux 均可。Python 版本Python 3.9 即可示例在 Python 3.10 下验证。数据库SQLitePython 内置支持。构建工具不需要 Maven/Gradle直接运行 Python 脚本。IDE任意支持 Python 的编辑器比如 VS Code、PyCharm。你只需要新建一个目录把下面几个文件创建出来然后执行python run.py。示例完整项目结构如下event_sourced_agent/ ├── event_store.py # 事件存储只追加日志 ├── projection.py # 状态投影从事件流重建 Agent 状态 ├── agent.py # Agent 运行时与反思引擎 ├── policy_service.py # 策略审批 / 发布 / 回滚 └── run.py # 演示主流程3.2 总体架构我们把系统拆成四个核心部分Agent 运行时、事件存储、投影模块、策略服务。┌─────────────────────┐ │ Agent 运行时 │ └──────────┬──────────┘ │ 产生事件 ▼ ┌─────────────────────┐ │ Event Store │ │ (append-only) │ └──────────┬──────────┘ │ 投影 ▼ ┌─────────────────────┐ │ 状态投影/读写模型 │ └─────────────────────┘Agent 运行时负责执行任务并在执行过程中不断追加事件。Event Store 是事件日志的存储层只允许追加。投影模块从事件流重建当前 Agent 的策略状态。策略服务负责策略候选的提出、审批、应用和回滚。3.3 事件类型设计在进入代码前先定义好事件类型。事件命名使用过去式表示“已经发生的事实”。事件类型含义关键属性AgentRunStartedAgent 开始处理任务taskToolExecuted工具执行完成task, result, ok, errorReflectionGenerated反思引擎生成经验reflectionPolicyCandidateProduced生成策略候选policy_id, contentPolicyApproved审批通过policy_idPolicyRejected审批驳回policy_idPolicyApplied策略正式生效policy_idPolicyRolledBack回滚策略policy_id它们合起来能描述一次完整的“自我改进闭环”AgentRunStarted - ToolExecuted - ReflectionGenerated - PolicyCandidateProduced - PolicyApproved - PolicyApplied - AgentRunStarted - ToolExecuted - PolicyRolledBack3.4 数据库模型这里用 SQLite 的一张事件表来存储所有事件。CREATE TABLE IF NOT EXISTS events ( seq INTEGER PRIMARY KEY AUTOINCREMENT, event_id TEXT NOT NULL UNIQUE, aggregate_id TEXT NOT NULL, event_type TEXT NOT NULL, payload TEXT NOT NULL, occurred_at TEXT NOT NULL, schema_version INTEGER NOT NULL DEFAULT 1 );字段含义seq事件序列号也是事件写入顺序。event_id事件唯一 ID用于消费端幂等。aggregate_id聚合 ID在本文里表示某个 Agent 实例。event_type事件类型。payload事件的业务数据JSON 序列化。occurred_at事件发生时间。schema_version事件结构版本方便未来演进。3.5 投影规则投影就是从一个空状态开始按顺序把每个事件“作用”到状态上最终得到当前状态。我们只关心三组事件AgentRunStarted运行次数加一保存当前任务。ToolExecuted保存最近一次执行结果和成败。PolicyCandidateProduced / PolicyApproved / PolicyRejected / PolicyApplied / PolicyRolledBack维护策略候选、审批、生效和回滚关系。最终生效策略取“最后一个被应用、且没有被回滚、没有被驳回、且已审批通过”的候选。这个逻辑会在projection.py中体现。4. 完整代码实现4.1 事件存储模块 event_store.py事件存储是全部功能的地基。先定义一个Event数据类再实现EventStore。# event_store.py import json import sqlite3 import uuid from dataclasses import dataclass from datetime import datetime, timezone dataclass class Event: event_id: str aggregate_id: str event_type: str payload: dict occurred_at: str schema_version: int 1 class EventStore: def __init__(self, db_path: str agent_events.db): self.db_path db_path self._init_schema() def _connect(self): conn sqlite3.connect(self.db_path) conn.row_factory sqlite3.Row return conn def _init_schema(self): with self._connect() as conn: conn.execute( CREATE TABLE IF NOT EXISTS events ( seq INTEGER PRIMARY KEY AUTOINCREMENT, event_id TEXT NOT NULL UNIQUE, aggregate_id TEXT NOT NULL, event_type TEXT NOT NULL, payload TEXT NOT NULL, occurred_at TEXT NOT NULL, schema_version INTEGER NOT NULL DEFAULT 1 ) ) def append(self, aggregate_id: str, event_type: str, payload: dict) - Event: event Event( event_idstr(uuid.uuid4()), aggregate_idaggregate_id, event_typeevent_type, payloadpayload, occurred_atdatetime.now(timezone.utc).isoformat(), ) with self._connect() as conn: conn.execute( INSERT INTO events(event_id, aggregate_id, event_type, payload, occurred_at, schema_version) VALUES(?, ?, ?, ?, ?, ?) , ( event.event_id, event.aggregate_id, event.event_type, json.dumps(event.payload, ensure_asciiFalse), event.occurred_at, event.schema_version, ), ) return event def read_all(self, aggregate_id: str): with self._connect() as conn: rows conn.execute( SELECT * FROM events
返回列表