ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

多智能体编排实战:用OpenRig构建持久化Agent协作系统

多智能体编排实战:用OpenRig构建持久化Agent协作系统 从单个 Agent 到一群 Agent最难的不是把它们接在一起而是让它们像一支队伍那样稳定地协作。做了一段时间 AI Agent 应用后你会发现单 Agent 再聪明也只是个超级个体上下文一断就失忆任务一长就漂移中间某个环节崩了整条链路就从头上演。而当我们把多个 Agent 放到同一个系统里让它们各自负责一块、彼此交接任务、共享状态、协作产出情况马上变得不一样——但与此同时一个更棘手的问题浮出水面这些离散的 Agent怎么才能被编织成一个持久化协作系统这就是我这篇文章想聊透的事。我围绕 OpenRig 这套多智能体编排实践讲清楚为什么单 Agent 扛不住真实业务、多 Agent 编排到底在编排什么、持久化为什么是协作系统的底座以及当有人问AI Agent 怎么扛并发时我们到底在解决什么问题。内容会包含完整的设计思路、可落地的架构拆解、实操代码和踩坑记录适合正在做 Agent 应用落地的开发者、架构师以及想从 demo 走向生产的团队参考。1. 先从为什么说起离散 Agent 的真实困境1.1 单 Agent 的三种崩溃方式我见过太多把 Agent 当万能接口用的项目用户提问Agent 回答任务结束。这种模式在 demo 里很完美但在真实业务里它基本会以三种方式崩溃。第一种是上下文崩溃。对话一长Agent 就开始丢三落四。LLM 的上下文窗口是有限的一旦超出早期的关键信息就被截断Agent 会一本正经地基于残缺信息做决策。你让它在第 50 轮对话时还记得第 3 轮提到的某个约束它大概率做不到。第二种是任务崩溃。真实任务往往不是一轮问答而是多步骤的流程分析数据、制定方案、写代码、测试、修订。单 Agent 做这些事时每一步之间的状态散落在对话历史里没有结构化的承接。一旦某一步出错它可能从头开始或者更糟——带着错误状态继续往下走。第三种是职责崩溃。一个 Agent 又想当产品经理又想当程序员又想当测试还想当运维。它会频繁地在不同角色之间切换导致行为不一致。今天你让它写代码它写得挺好明天同样的需求它可能先给你来一段需求分析因为它的 system prompt 太宽泛角色没有边界。所以多 Agent不是炫技而是单 Agent 模式碰到真实业务复杂度之后的必然演化。但你如果只是把多个 Agent 平铺在一起、用 if-else 调来调去那比单 Agent 更糟糕——你会发现协调成本比干活成本还高。1.2 编织意味着什么OpenRig 这个名字里最核心的隐喻是编织rig 本身有索具、装配的意思。离散的 Agent 就像散落的线每条线都有自己的能力但只有通过某种结构把它们编织起来才能成为能承重的绳索。这就需要三样东西拓扑结构谁在哪个环节、负责什么、产出给谁。协作协议Agent 之间用什么语言交流怎么交接任务怎么反馈结果。持久化底座所有协作状态、中间产物、决策记录不会因为某个 Agent 崩溃或重启而丢失。这三者缺一不可。没有拓扑Agent 就是一盘散沙没有协议消息就是鸡同鸭讲没有持久化系统就是空中楼阁——跑得再漂亮宕机一次全部归零。2. 多智能体编排的核心设计2.1 编排器和总线控制平面与数据平面分离多智能体编排最容易踩的坑是把所有消息都经过一个中央 Agent 转发。这种中心化方案看着直观但一旦 Agent 数量上来中心节点就成了瓶颈而且单点故障会拖垮整个系统。OpenRig 的设计思路是控制平面和数据平面分离。控制平面由一个编排器Orchestrator负责它只做三件事路由决策、状态管理、任务调度。数据平面则是任务总线Task BusAgent 之间通过总线交换任务和结果不直接点对点通信。这样做的好处是解耦。每个 Agent 不需要知道消息是谁发的只需要知道自己该处理什么、处理完往哪发布。新增一个 Agent 时只需要在注册表里登记它的能力和订阅规则不用改其他 Agent 的代码。举个具体的例子。一个内容生成系统里有三个 Agent选题 Agent、初稿 Agent、审校 Agent。选题 Agent 产出一批选题后把结果发布到选题就绪主题初稿 Agent 订阅这个主题拿到选题后生成初稿发布到初稿就绪审校 Agent 订阅初稿就绪审完发布到终稿就绪。整个流程是异步的、解耦的任何一个 Agent 挂掉或替换其他 Agent 不受影响。2.2 拓扑注册表Agent 的身份与能力声明多 Agent 系统要能编织首先得有个东西记录谁是谁、能干什么。OpenRig 里这块叫拓扑注册表Topology Registry本质是一张结构化元数据表。每个 Agent 启动时向注册表上报自己的信息Agent 名称、能力声明能处理什么类型的任务、订阅主题、产出主题、依赖关系启动前需要哪些就绪信号、资源需求比如需要 GPU 还是纯 CPU 推理。注册表存在的意义不只是让编排器知道有哪些 Agent更重要的是让系统有能力做动态调度和故障转移。当一个 Agent 实例崩溃时编排器可以通过注册表找到同类能力的其他实例把未完成的任务重新路由过去。文档型 Agent 如果 A 实例挂了B 实例只要声明了相同能力就能接管任务。持久化在这里也发挥作用——注册表本身要持久化。不能每次重启都让 Agent 重新注册这样编排器在启动初期会有一段盲区不知道到底有哪些能力可用。2.3 任务状态机协作的通用语言多 Agent 协作中任务状态必须统一建模。OpenRig 用一套任务状态机来定义任务的流转PENDING已创建、READY已就绪可被消费、IN_PROGRESS处理中、SUCCEEDED成功、FAILED失败、RETRYING重试中、SKIPPED跳过。这看起来很简单但它是整个协作系统的语法规则。因为 Agent 之间彼此不直接调用它们通过任务状态的变更来感知系统进度。A 产出一个任务状态是 READYB 接手后置为 IN_PROGRESSB 完成结果后置为 SUCCEEDED并产出新的 PENDING 任务给下游。状态机必须有明确的转换规则和幂等语义。编排器在更新状态时要保证幂等——同一个任务被多次标记为 SUCCEEDED不能导致重复执行下游任务。我在实现时给每个任务分配了全局唯一的 task_id并在状态流转时用 CASCompare-And-Swap保证并发安全。代码示意如下dataclass class TaskState: task_id: str manifest: dict # 任务元数据 status: str # PENDING / READY / ... created_by: str # 产生该任务的 agent assigned_to: str | None # 当前处理 agent attempt_count: int 0 payload_ref: str | None # 指向持久化产物的引用每次状态更新都是一个独立事务带版本号更新前先比对版本防止并发覆盖。这个设计在后面讲怎么扛并发时还会再提。3. 持久化协作系统的构建3.1 状态存储与事件溯源持久化的第一个层次是状态存储。任务状态、Agent 注册信息、编排器的路由表这些必须落到稳定的存储里。我用 Postgres 作为主存储同时也把事件日志写到独立表里形成事件溯源Event Sourcing的底子。这里要说明一下为什么用事件溯源多 Agent 协作系统的调试难度远比单 Agent 要高。当链路出错时你不仅要知道最终结果不对还要知道在哪个环节、由哪个 Agent、基于什么上下文产出了什么中间结果。事件溯源把每一次状态变更、每一条 Agent 间消息、每一次重试决定都追加到日志里这样你可以完整重放一次协作过程。表结构大致是这样的CREATE TABLE task_events ( id BIGSERIAL PRIMARY KEY, task_id VARCHAR(64) NOT NULL, event_type VARCHAR(32) NOT NULL, agent_id VARCHAR(64), payload JSONB, created_at TIMESTAMPTZ DEFAULT now() ); CREATE INDEX idx_task_events_task_id ON task_events(task_id);每个事件都是只追加的不改写历史。任务当前状态是从事件流聚合出来的投影。这个设计让回溯变得极其自然——出问题了把任务事件按时间排开一眼就能看到当时发生了什么。很多人问持久化和缓存的区别。这里必须讲清楚如果你只是把状态存在 Redis 里那叫缓存重启就没了不是持久化。OpenRig 做持久化有明确的层次——核心状态进 Postgres热数据和高频读写走 RedisRedis 只是缓存层允许丢丢了可以从 Postgres 重建。3.2 消息队列作为任务交接通道Agent 之间不直接通信而是通过任务总线。我在实现时选了 Redis Streams 作为总线因为它在持久化方面比普通 Pub/Sub 更可靠——消息会留在 Stream 里消费者组可以记录消费位点同一个消息不会被重复消费配合 ACK 机制。任务总线的结构是主题 消费者组。每个主题对应一类任务Agent 实例以消费者组的方式订阅主题。编排器或上游 Agent 往主题里写任务消息消费组里的某个实例取走并处理。用 Redis Streams 而不是普通消息队列的原因很实际一是部署简单和状态缓存共用 Redis 实例就行二是 Stream 天然支持读取未 ACK 消息的机制Agent 崩溃后重新上线可以拿到没有 ACK 的任务继续处理实现最基本的断点续传。XADD task:article:research * agent_id researcher-A task_id task_20250101_001 payload {...} XREADGROUP GROUP article_agents consumer_A COUNT 1 STREAMS task:article:research 这套机制解决了任务在 Agent 之间传递时的持久化问题。消息不会因为消费者临时掉线而丢失Redis 会把它保留在 Stream 里等新的消费者实例上线后重新投递。3.3 快照与会话恢复机制事件溯源给了完整性但完整重放的成本很高。如果某个任务的协作链有 500 个事件每次恢复都从头重放一遍性能上吃不消。所以 OpenRig 引入了快照机制每处理 N 个事件或每隔一定时间就把当前任务状态做一次快照。恢复时先加载最近的快照再重放之后的事件。这其实就是数据库中全量备份 增量日志的思路。快照本身也存 Postgres事件表保留足够长的窗口用于追溯。会话恢复的流程是编排器检测到某个 Agent 实例失联心跳超时。编排器从注册表找到替代实例。替代实例加载该任务最近的快照。重放快照之后的事件重建上下文。从最后一个成功的状态继续执行而不是从头开始。这里有个实操细节快照不能只存任务状态还要存 Agent 的工作记忆——比如分析 Agent 已经读过的资料摘要、初稿 Agent 已经写出的前几版内容。因为 Agent 不是无状态函数它的决策依赖之前积累的信息。如果恢复后工作记忆是空的它会失忆基于不完整的信息继续干活。这意味着快照的数据结构要包含任务元数据、中间产物引用、对话记忆摘要、上下文窗口截断策略。我通常会为每个快照预留一个 payload 字段塞进去一个结构化文档而不是把全部原始历史都存下来——选择性保留关键信息恢复效率和上下文完整性之间做个平衡。4. 扛并发从快到稳的关键4.1 Agent 并发模型不是线程是消费者组AI Agent 怎么扛并发是我在搜索热词里看到最多的问题之一。很多人的直觉是给 Agent 加线程池、上异步框架但真正的瓶颈往往不在执行层而在协作层的竞争条件。OpenRig 的并发模型是每个 Agent 类型一个消费者组组内多个工作实例并行消费。比如初稿 Agent类型注册了三个实例它们是一个消费组的三个消费者。任务从队列里被投递到某个空闲实例处理完后发布产出。这个模型天然支持水平扩展——发现任务积压给这个消费组加实例就行不用改代码。但消费者组模型有一个隐藏问题任务顺序。如果同一个任务的多个子任务被不同实例并行处理它们之间可能有依赖关系。比如写初稿和做校对虽然有先后但做校对依赖写初稿完成如果两个子任务被不同实例同时接手就可能出现校对比写稿先完成。解决办法是给任务分片——按任务来源、业务线或依赖链哈希分区同一个依赖链的任务总是进同一个消费组里的同一个消费者。我在实践里用的是任务 ID 前缀哈希同一条业务链的任务 ID 以同一前缀开头哈希后落在固定消费者上从根源上避免乱序。4.2 状态读写竞争与隔离多实例并行最大的痛点是状态竞争。两个实例同时读到待处理的任务同时尝试更新状态为处理中就冲突了。我在前面提过任务状态更新必须走 CAS。实现上有两种做法一种是数据库层乐观锁更新时带上原状态条件UPDATE tasks SET status IN_PROGRESS, assigned_to :agent, version version 1 WHERE task_id :task_id AND status READY AND version :old_version;影响行数为 0 说明别人已经抢占了当前实例应该跳过这个任务。这保证了同一个任务只会被一个实例真正处理。另一种是在 Redis 里用 setnx 做分布式锁抢到锁的实例才有权处理任务。锁要有过期时间防止死锁还要有续期机制防止长时间任务中途锁失效。我在实践中两种都用了——Redis 锁做第一道拦截快速过滤数据库乐观锁做最终裁决保证准确。锁不是越多越好但任务接手的这一下必须有锁因为它是所有并发窗口的入口。4.3 背压、超时与降级并发扛得住的标志不是无限处理得快而是在流量上来时系统依然稳定。这就要提背压和降级。每个 Agent 消费者都有队列长度上限。当队列塞满时编排器不再往里投递新任务而是返回忙信号让上游慢下来。这比让任务无限堆积在内存里靠谱得多——堆积到一定程度系统会雪崩。超时控制同样关键。每个任务从 READY 到完成都要设一个 SLA 超时。比如资料分析任务最多 2 分钟超时后编排器判定该任务失败或重新投递。这里最容易踩的坑是 LLM 调用迟迟不返回——服务端偶尔会卡在流式输出上。我的做法是给每次 LLM 调用加双层超时总超时 90 秒流式空闲超时 30 秒。任何一层触发都主动断开让任务进入 RETRYING。降级策略在另一头兜底。当某个 Agent 依赖的第三方服务不可用时它不是简单报错而是尝试走备选模型或降级方案。比如初稿 Agent 的主模型超时可以切换到配置更低的备用模型跑一版质量可能差一点但流程不会断。系统设计里流程不断和结果最好之间应该有一套明确权衡。5. 实操搭建一个持久的 Agent 协作流水线5.1 场景定义一条内容生产流水线理论讲再多不如直接搭一个。我用一个内容自动化生产流水线做示例它包含五个 Agent选题 Agenttopic_agent基于行业资讯和用户画像产出选题清单。研究 Agentresearch_agent针对选题做资料搜集和要点提炼。初稿 Agentdraft_agent基于研究结果生成文章初稿。审校 Agentreview_agent检查初稿质量、事实偏差和风格问题输出修订意见。发布 Agentpublish_agent把终稿格式化后写入内容库。这些 Agent 之间的依赖是接力式的选题 → 研究 → 初稿 → 审校 → 发布。但实际运行时选题 Agent 可以同时产出多个选题研究 Agent 可以并行处理多个选题的研究这就是天然的并发场景。5.2 注册表配置与任务流转每个 Agent 启动时向注册表登记{ agent_name: draft_agent, capabilities: [article_drafting], subscribe: [research_completed], publish_to: [draft_ready], requires: [research_result], resources: {model: gpt-4o, max_concurrency: 3} }这样编排器就知道当研究完成的事件到达时应该把任务投递给 draft_agent 的消费组。draft_agent 完成任务后把成果发布到draft_ready主题触发 review_agent。任务流转的完整代码非常直白# 编排器核心事件路由 def handle_event(event): topic event[topic] # 查找订阅此 topic 的 agent subscribers registry.query(topictopic) for agent in subscribers: # 创建任务并发布到 agent 的队列 task TaskState( task_idgenerate_id(), manifestevent[payload], statusREADY, created_byevent[source], assigned_toagent[agent_name] ) store.save(task) bus.publish(agent[subscribe_topic], task.to_dict())5.3 状态持久化与拓扑结构的核心代码任务执行过程中每一阶段做一次状态落地。以研究 Agent 为例def run_research(task): topic_id task.manifest[topic_id] # 1. 标记进行中 state load_state(task.task_id) state.status IN_PROGRESS persist_event(task.task_id, research_started, payload{topic_id: topic_id}) # 2. 调用工具收集资料 sources search_web(topic_id, limit10) summary llm_summarize(sources) # 3. 持久化中间产物 artifact_ref save_artifact(research_summary, { topic_id: topic_id, summary: summary, sources: sources[:5] }) # 4. 发布完成事件触发下游 persist_event(task.task_id, research_completed, payload{artifact: artifact_ref}) bus.publish(research_completed, { task_id: task.task_id, artifact: artifact_ref, agent_id: AGENT_ID }) # 5. 更新状态为 SUCCEEDED state.status SUCCEEDED persist_snapshot(task.task_id, state)这段代码的用意是让读者看到Agent 的核心业务逻辑只占一半代码另一半全是状态同步和事件发布。这就是多智能体编排和普通函数调用的本质区别——每次干完活都要留痕。5.4 如何设计与维护好拓扑结构我在实操中积累的一个重要经验是拓扑结构不要一开始就设计得很复杂。第一次跑通时只保留最关键的三个 Agent、一条直链就够了。等直链稳定了再往里面加并行分支、加条件路由、加重试链路。有个具体建议把拓扑配置和业务代码分离。拓扑注册表里的 JSON 是配置Agent 代码是执行单元。这样调整流程时不用改代码改配置就能重排协作链路。发布审批流想从审校后直接发布改成审校后人工确认再发布只需要在注册表里把 publish_agent 的订阅条件加一个需人工确认标签代码一行不动。6. 常见问题与排查技巧实录6.1 状态漂移Agent 以为做完了编排器以为没做完这是我遇到最频繁的问题。Agent 完成工作了也把自己的状态更新了但编排器没有收到事件或者收到了事件却没有正确落地。于是整个系统的状态视图出现分歧。排查路径通常是看 Redis Stream 里的消费位点消息是否被 ACK还是卡在 pending 列表里。看 Postgres 的 task_events 表Agent 声称的状态更新事件是否真实落库。看 Agent 的日志它在发布事件之前是否发生异常退出导致活干了一半事件没发出去。解决方案是在 Agent 端引入状态变更即事务的思想业务处理和事件发布放到同一个事务里或相近的保证不能先改状态再发事件因为两者之间如果崩溃系统就裂了。可靠做法是业务处理完成后先写事件表和产物表再发布 MQ 消息——事件表作为最终事实MQ 只是加速通知的通道。6.2 重放导致重复副作用事件溯源能恢复任务但也带来一个问题重放可能重复执行有副作用的操作。比如发布 Agent 已经把文章发布到了内容库事件记录在案恢复时如果从头重放它可能又发布了一遍。解决办法是所有 Agent 的副作用操作都要做幂等化。发布操作以 task_id 作为幂等键内容库中已经存在同 task_id 的发布记录就直接跳过。外部 API 调用比如发邮件、调第三方平台更要在调用前先做一次是否已执行的检查。我在设计 Agent 接口时加了一个约定每个 Agent 的任务处理器要支持传入 checkpoint 信息它从 checkpoint 开始恢复而不是从头跑。这样重放时跳过了已完成的部分副作用自然就不会重复。6.3 编排器本身会不会成为单点一个现实的问题是编排器挂了怎么办只把 Agent 做得高可用编排器是单点整个系统还是脆弱的。我用的方案是编排器多副本 选主机制。同一时刻只有一个活跃编排器在做路由决策其他副本处于热备。活跃节点挂了备节点通过数据库里的租约lease机制接管。租约的核心很简单活跃编排器每隔一段时间在存储里续租备节点发现租约过期就尝试抢占。抢占成功者成为新主重新加载拓扑注册表和未完成任务列表。-- 租约表 CREATE TABLE orchestrator_lease ( lease_key VARCHAR(16) PRIMARY KEY, holder_id VARCHAR(64) NOT NULL, expires_at TIMESTAMPTZ NOT NULL );选主的关键是时钟和超时的设置要合理。租约过期时间设为 10 秒心跳续约 3 秒一次。这样故障在最多 10 秒内就能被感知系统完成接管。6.4 Agent 死循环与任务风暴多 Agent 系统里会出现一种诡异的故障Agent A 产出任务给 BB 的结果又触发 A 继续产出形成一个闭环。如果双方都在生产新任务而不是最终消费任务任务数量会指数级增长直接把存储和队列打爆。我在系统里加了任务风暴防护每个任务记录来源链origin_chain追溯它由哪个根任务衍生。单根任务的最大衍生数量上限默认 50超出后拒绝创建新任务。协作链最大深度限制默认 10 层防止无限套娃。另外要在监控里盯着任务的产生/消费速率。一旦发现某主题的生产速率远大于消费速率就自动报警并暂停该主题的投递。这不是多高端的技术但很多系统是因为没人盯监控才被拖垮的。7. 工具选型与后续扩展方向7.1 存储、队列与框架的选型对照在 OpenRig 的搭建过程中我对各组件的选型做了横向对比。这部分基于我的实际踩坑经验不是标准答案仅供参考。组件选项我的选择原因主存储Postgres / MySQL / MongoDBPostgresJSONB 对事件日志友好事务能力强恢复工具生态成熟缓存/队列Redis / RabbitMQ / KafkaRedis Streams部署简单和缓存复用同一实例Stream 的消费组机制刚好匹配消费者组并发模型Agent 运行时Python / Node / RustPythonAI 生态最丰富langchain 等工具链成熟性能不满足时可用 Rust 重构热路径编排框架LangGraph / Temporal / 自研自研轻量编排层业务编排逻辑简单直接框架太重反而受限选型逻辑就是一句话能用简单方案解决的不引入重型组件。Redis Streams 在高吞吐、强一致的场景不如 Kafka但在这个系统的规模下它已经够用且省心太多。7.2 未来扩展接入其他业务域的编排OpenRig 的这套模式不只能做内容流水线。电商场景里用户下单后的订单履约流程同样非常适合多 Agent 协作库存校验 Agent、支付回调 Agent、物流调度 Agent、售后跟踪 Agent它们之间的交接逻辑和上面的内容流水线一模一样——基于任务状态机和持久化事件总线。我现在在做的一个扩展方向是把 OpenRig 的编排层独立成一个通用服务业务方只需要注册自己的 Agent 类型和拓扑配置就能复用持久化、恢复、并发的整套能力。这相当于把多 Agent 如何协作这件事从业务代码里剥离出来变成一个平台能力。另外一个值得探索的方向是人机协同接口。目前 Agent 协作是全部自动化的但真实业务里人需要参与审批、确认、纠偏。我计划在任务状态机里增加一个 AWAITING_HUMAN 状态——当任务流转到需要人工介入的节点时状态挂起等待人在 Web 界面确认后再放行到下游。这一步做完整个系统才算是真正进入生产环境可用的状态。从实践经验来看做多智能体编排最后拼的不是模型多聪明而是工程上多稳。模型能力是上限工程底座是下限——OpenRig 解决的是下限问题是让几十个 Agent 能持续、稳定、可追溯地一起干活的问题。把这些细节铺完之后我最大的感受是Agent 单个的聪明程度是天花板而编排系统的下限决定了你实际能拿到多少结果。这个道理踩过坑的团队应该都有共鸣。
返回列表