ARTICLE DETAIL

资讯详情

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

从零搭建OpenRig:多智能体持久化协作编排系统架构与实践

从零搭建OpenRig:多智能体持久化协作编排系统架构与实践 1. 先从一个让人头疼的协作场景说起如果你和我一样手里同时维护着好几个专精的 AI Agent——一个负责 SQL 生成一个做数据可视化一个写周报——大概很快就会撞上同一个问题单打独斗的 Agent 干不了复杂的协作活而把一堆离散的 Agent 编成一套能持久化协作的系统比想象中要难得多。这也是我做 OpenRig 这个多智能体编排项目的直接原因。项目名字是 open 和 rig 的组合rig 有装配、搭建的意思我想做的正是把离散的 AI Agent 像搭钻井平台一样装配成一个持续运转的协作系统。理论讲起来云里雾里但落到代码和部署上都是实打实的工程问题状态怎么存、消息怎么传、并发怎么扛、崩了怎么恢复。这篇文章把我从零搭建 OpenRig 到生产环境稳定跑完两个月的完整过程拆开讲一遍包括架构设计、核心实现、压测数据和避坑经验适合正在做多 Agent 编排、或者准备从单体 Agent 升级到协作系统的朋友参考。先说一个背景。我最开始用 LangChain 搭单体 Agent 时觉得一切都很美好一个 Agent 能调用工具、能规划、能执行。但做着做着就发现不对劲——当任务变复杂之后单体会员会遇到三个绕不开的问题。第一是上下文膨胀Agent 的上下文窗口有限让它同时扮演需求理解 任务拆分 代码生成 结果校验 文档输出五个角色注意力会被稀释还会出现角色错乱。第二是资源没法复用我有专门做 SQL 生成的、做数据可视化的、写周报的多个 Agent如果它们是三个离散的独立进程每次新任务都要重新把它们串起来串接逻辑散落在各个脚本里时间一长就变成没人维护的屎山。第三是没有状态今天跑一个任务明天另一个任务要接着用它昨天的结果如果 Agent 之间没有共享的、持久化的协作状态一切都得从头再来。OpenRig 就是为了解决这三个问题而生的。它不是一个聊天机器人框架而是一套面向多个 Agent 协同完成复杂业务任务的编排系统核心关注点就四个注册、编排、持久化、并发控制。2. OpenRig 的核心编排模型注册中心、任务拓扑与消息总线2.1 注册中心让每个 Agent 学会自我介绍要让离散的 Agent 协作第一步是让系统知道有哪些 Agent 可用、各自能做什么。我在 OpenRig 里做了一个轻量级注册中心每个 Agent 接入系统时都要做一次自我介绍包括Agent 的 ID、能力描述、输入输出 JSON Schema、超时时间、最大并发数。这里有个容易忽略的点Agent 的能力描述不能用人话要用机器可读的 schema。我在早期版本里让开发者自由填写 capability 字段结果有人写擅长分析数据有人写帮我处理 excel 表格。到编排器要做路由决策时这些描述根本没法精确匹配。后来我强制要求能力描述必须带一个 capability_id配合输入输出的 JSON Schema路由决策才变得可靠。注册信息存储在 PostgreSQL 中同时启动时加载一份到内存缓存路由热点路径直接走缓存只有注册和注销才写数据库。注册信息是低频写、高频读的数据没必要每次路由都查库。2.2 任务拓扑用 DAG 描述协作关系而不是一个个 if-else多智能体编排最常见的错误写法是在控制代码里用 if-else 把 Agent 一个个串起来。比如如果用户要数据分析就先调用 agent_sql再调用 agent_viz最后调 agent_report。这种写法最大的问题在于协作逻辑被硬编码在代码里无法持久化无法可视化更无法动态调整。OpenRig 的做法是把每次任务的协作关系建模成一张有向无环图DAG图的节点是 Agent 调用或工具调用边是数据依赖关系。编排器只负责按照 DAG 的拓扑序调度节点而不是硬编码每一步顺序。我定义了一套简洁的编排 DSL用 YAML 描述任务拓扑。写过几轮之后我发现 DSL 的设计需要克制太复杂没人用太简单又不够表达。OpenRig 的 DSL 只保留了四个核心原语node节点、edge依赖边、branch条件分支、retry重试策略。下面是一个典型的三 Agent 协作拓扑示例——用户需求 → SQL 查询 → 数据解释 → 报告生成流程version: 1 name: sql_insight_pipeline nodes: - id: intent_parser agent: llm_router input: $user_input output: { parse_result: $parse_result } - id: sql_executor agent: sql_agent input: { query_intent: $parse_result } output: { query_result: $query_result } - id: insight_writer agent: insight_agent input: { data: $query_result, intent: $parse_result } output: { insight_doc: $insight_doc } - id: report_compiler agent: report_agent input: { insight: $insight_doc } output: { final_report: $final_report } edges: - from: intent_parser to: sql_executor condition: $parse_result.type sql_query - from: sql_executor to: insight_writer - from: insight_writer to: report_compiler这套 DSL 的关键设计是数据流而非控制流。每个 Node 只声明自己的输入来自哪个上游节点的输出编排器自动建立依赖关系并计算拓扑序不需要开发者手动安排执行顺序。条件分支也收敛为 Edge 上的 condition而不是在代码里反复判断。2.3 消息总线Agent 之间不直接握手全部走邮件在 OpenRig 架构里Agent 之间不会直接互相调用。所有消息都通过内部消息总线转发这条总线在实现上就是 Redis Stream 加 PostgreSQL 落库。为什么这么设计直接调函数的架构最简单但一旦 Agent 分布到不同机器或者某个 Agent 执行失败要重试函数调用栈根本没法处理。消息总线让 Agent 之间的通信变成了异步事件传递上游 Agent 产出结果后把消息发布到总线编排器收到消息后根据 DAG 判断哪些下游节点可以解锁。这和日常生活中邮件的工作方式很像发件人不需要等待收件人马上回复邮件到了信箱收件人自己有空再来取。异步解耦之后某个 Agent 执行 5 分钟也不影响其他 Agent 并行推进。消息体统一为 JSON必带字段包括message_id、task_id、from_node、to_node、payload、timestamp。每条消息在 Redis Stream 里流转并同步写入 PostgreSQL 的 message_log 表。写库的目的是支撑持久化协作——如果某个节点失败需要回放消息可以从 message_log 找到完整记录。3. 持久化协作的秘密状态机、任务快照与崩溃恢复3.1 每次协作都是一条状态机流水线持久化协作系统与临时脚本最大的区别在于系统必须知道现在这个任务进行到哪一步了。OpenRig 将每次协作任务建模为一条状态机流水线状态包括 PENDING、RUNNING、NODE_COMPLETED、COMPLETED以及 FAILED 和 CANCELLED 两个终止态。每个 Node 从 PENDING 变 RUNNING 时系统会写入一条事件记录到 task_event 表。这个设计在线上帮了大忙。有一次 insight_agent 因为大模型服务商限流导致超时任务卡在 NODE_COMPLETED 之后迟迟没有推进。我打开 OpenRig 的 Dashboard一眼就看到任务停留在insight_writer 已完成等待 report_compiler 启动这个阶段而不是整条任务直接显示失败。配合消息总线里的消息记录我很快定位到是下游节点没被解锁原因是事件回调偶发丢了一个消息。3.2 任务快照让协作可以断点续传持久化协作系统的另一个硬需求是任务执行到一半进程崩了重启之后要不要接着跑我的答案是必须能接着跑否则就谈不上持久化。OpenRig 每隔一段时间会为运行中的任务生成一份快照记录当前 DAG 的节点完成情况、各节点的输出摘要、消息队列中的待处理消息 ID。快照本身存 JSON放 PostgreSQL 的 task_snapshot 表。进程重启后编排器扫描所有处于 RUNNING 状态的任务加载最新快照重新计算哪些节点已完成、哪些节点需要重新执行。这里有一个重要的工程决策已完成节点的输出不存全量结果只存摘要和结果在对象存储里的路径。原因是 LLM 的输出往往很大一个 insight_doc 可能有几十 KB 甚至几百 KB如果每个快照都备份全量输出快照体积会无法控制地膨胀。只存摘要加结果指针快照体积能控制在 KB 级恢复时按需加载即可。task_snapshot 表的核心字段大致是这样设计的字段类型说明snapshot_idbigint快照 ID主键task_idvarchar关联的任务 IDversionint快照版本乐观锁用node_statesjsonb各节点状态PENDING/RUNNING/COMPLETEDoutput_refsjsonb节点结果摘要与对象存储路径pending_message_idsjsonb队列中待处理消息 IDcreated_attimestamptz快照生成时间3.3 恢复策略幂等重放与人工接管崩溃恢复最怕的是什么是同一个 Agent 在一个任务里被执行了两次。LLM 接口调用不是免费的而且有些 Agent 会触发外部副作用——比如发邮件、写数据库。所以 OpenRig 要求所有 Agent 实现幂等语义每个任务节点都有一个唯一的 execution_idAgent 在产生副作用前先检查这个 execution_id 是否已经处理过。我在实际实现中发现完全依赖 Agent 自己保证幂等是不现实的。更稳妥的做法是编排器在重放节点之前做一次结果确认向 Agent 发出一个 probe 请求询问execution_id 为 xxx 的任务你处理过吗结果是什么。如果 Agent 能返回结果就直接复用不重新执行只有确认没有结果时才重新调度。这个 probe 机制对外部系统也算友好因为它只查状态不触发副作用。4. 扛并发OpenRig 的调度器、限流与水平扩展设计4.1 为什么 Agent 服务天然容易被打爆最近社区里AI Agent 怎么扛并发的热度非常高我的观点很明确Agent 服务和普通 Web 服务的并发模型完全不一样。普通 Web 服务扛并发靠的是无状态加水平扩展每次请求独立处理完就结束。Agent 服务则不同一个任务要来回调度多个 Agent每个 Agent 都要调 LLM而 LLM 接口的响应时间是秒级甚至分钟级。这意味着一个任务占用的资源时间长、链路长简单的线程池模型根本扛不住。更麻烦的是Agent 的调用不是均匀的。用户可能在上午 10 点集中提交任务其他时间很闲。如果按峰值容量部署成本会非常难看。OpenRig 的设计目标不是让单个实例扛住所有并发而是让并发被正确排队、被正确扩散、被优雅地削峰。4.2 三级队列与背压机制OpenRig 的调度器采用了三级队列设计任务接收队列存 Redis记录所有待处理的任务 ID。就绪节点队列只有 DAG 中所有上游节点都完成的节点才会进入这个队列。Worker 执行队列每个 Worker 会从就绪队列拉取可执行的节点。这三级队列的核心价值是解耦。上游任务提交的速度和下游 Worker 执行的速度不匹配时队列就是缓冲带。当 Redis 中的待处理任务数量超过阈值时调度器会向 API 层返回 503 背压信号客户端收到后可以选择延迟重试而不是让请求继续灌进来。限流方面我对每个 Agent 注册了 max_concurrency 参数调度器会为每个 Agent 维护一个信号量。sql_agent 配置 max_concurrency3那同时最多 3 个 sql_agent 节点在跑。其他节点即使全部就绪也要等待信号量释放。这个机制有效防止了某一类 Agent 因为过度调用被外部 API 限流。4.3 水平扩展Worker 无状态化OpenRig 的 Worker 可以水平扩展因为它们不保存任何任务状态——状态全部在 PostgreSQL 和 Redis 里。要增加处理能力只需要多启动几个 Worker 进程它们去 Redis 就绪队列竞争任务即可。这里的关键在于竞争要保证不重复使用 Redis 的队列原子操作比如 LPOP 或 BRPOPLPUSH 来保证同一个节点只会被一个 Worker 取到。我在压测中遇到过一个问题多个 Worker 同时竞争任务时任务完成后要标记节点状态大家并发去写 PostgreSQL 的同一行记录会导致锁冲突。后来我引入了一套执行租约机制——Worker 取任务时写入一个租约记录过期时间 60 秒节点状态更新采用乐观锁版本号对比冲突时直接丢弃重试。这样 30 个 Worker 并发执行 1000 个节点的压测场景冲突率能控制在非常低的水平。4.4 实际压测数据并发从 5 提升到 50 的调优过程给一个真实数据参考。最初我的部署架构是单进程加内置线程池并发测到 5 个任务时 LLM 调用就开始超时因为线程都被长时间阻塞在外部 API 调用上。后来我把 Worker 拆出来独立进程数量按 CPU 核心数两倍部署我机器是 8 核所以起了 16 个 Worker并把任务接收用 FastAPI 异步接口实现并发直接提升到 30 左右。再往后加了三级队列和信号量限流调整为 50 并发稳定跑完P95 任务耗时只比单机低并发时多 18%。这个成绩对多智能体场景来说已经够用。5. 生态集成实践FastAPI 网关、LangGraph 对比与 Django 场景5.1 用 FastAPI 做统一网关OpenRig 对外暴露的是 FastAPI 服务接口包括提交任务POST /v1/tasks、查询任务状态GET /v1/tasks/{task_id}、取消任务POST /v1/tasks/{task_id}/cancel、查看 Agent 注册列表GET /v1/agents。为什么选 FastAPI一是异步支持好从请求进入到返回 task_id整个过程没有任何阻塞调用适合高并发接入二是自动生成 OpenAPI 文档前端和外部系统对接非常省事。但要注意FastAPI 只是网关不是执行引擎。任务提交后接口立刻返回 task_id所有 Agent 执行都在后台异步进行。客户端通过轮询或 Webhook 获取最终结果。这个模式叫提交即返回是多智能体服务必须养成的习惯——如果一个任务要跑 3 分钟你绝不能让 HTTP 请求一直挂着等 3 分钟。5.2 和 LangGraph 的关系编排层互补而非替代我在设计 OpenRig 时被人问得最多的问题是你不用 LangGraph 吗LangGraph 不是已经能做多 Agent 编排了吗说实话我一开始也纠结过。LangGraph 在单 Agent 内部的状态管理和工具调用链路上做得确实不错它的图模型和 OpenRig 的 DAG 有相似之处。但我最终决定自研编排层的原因是LangGraph 更多解决的是一个 Agent 内部的复杂控制流而 OpenRig 解决的是多个独立 Agent 之间的持久化协作。这两件事虽然看着像其实维度不同。LangGraph 执行流活在内存里进程一挂任务就没了OpenRig 的所有协作状态都落库可以跨进程、跨机器、跨天恢复。我实际生产方案是两者结合每个 OpenRig 的节点内部如果需要复杂的工具链调用就内嵌一个 LangGraph 工作流。这样 OpenRig 管 Agent 之间的协作LangGraph 管 Agent 内部的思考链各司其职。5.3 Django 场景与内容自动发布场景的集成注意点社区里很多人问用 AI Agent 开发 Django 项目怎么做还有人把 Agent 用在内容发布场景比如自动在社交平台发消息。这些需求的本质都是同一个把 Agent 从聊天玩具变成生产执行者。Django 场景下我建议的集成方式是这样的Django 做业务数据和 Web 管理后台OpenRig 做任务引擎。Django 的 ORM 模型里保存用户提交的业务需求转换成 OpenRig 的任务参数后调用网关接口。Agent 执行完成后的结果通过 Webhook 回调写回 Django 业务表。这样的好处是业务代码和 AI 编排隔离Django 端完全不需要知道 Agent 怎么工作只需要关心业务字段不会出现把 Celery 和 Agent 任务混在一起导致 worker 卡死的情况。至于自动发布场景核心要点就一条凡是要对真实外部系统产生副作用的 Agent都必须走人工确认或灰度执行。自动发布内容看起来省事但发错一条内容的代价比省下的那点人工成本高得多。OpenRig 在编排 DSL 里支持在某个节点前插入 approval_gate节点状态会停在 AWAIT_APPROVAL直到有用户确认才继续。这个机制在内容发布场景中是保命的。6. 一个真实任务的完整跑通过程从 YAML 编写到结果回收6.1 准备阶段定义三个 Agent 并注册我拿一个周报自动生成的任务来完整演示一遍。任务需求从项目数据库中读取本周任务记录调用大模型生成周报摘要然后发送到企业微信机器人。整个过程三个 Agent 协作。第一步写三个 Agent 的注册信息# register_agents.py from openrig import AgentRegistry, AgentSpec registry AgentRegistry() registry.register(AgentSpec( agent_idtask_reader, capability_idread_task_records, input_schema{type: object, properties: {week_start: {type: string}}}, output_schema{type: object, properties: {records: {type: array}}}, max_concurrency3, timeout_secs60, )) registry.register(AgentSpec( agent_idsummary_writer, capability_idgenerate_weekly_summary, input_schema{type: object, properties: {records: {type: array}}}, output_schema{type: object, properties: {summary: {type: string}}}, max_concurrency2, timeout_secs120, )) registry.register(AgentSpec( agent_idwecom_sender, capability_idsend_message, input_schema{type: object, properties: {content: {type: string}}}, output_schema{type: object, properties: {message_id: {type: string}}}, max_concurrency1, timeout_secs30, ))注册接口本质上就是向 PostgreSQL 的 agent_registry 表插入记录同时更新内存缓存。第二步编写编排 DSL描述三个节点之间的依赖关系name: weekly_report_pipeline nodes: - id: read_task agent: task_reader input: week_start: $task.week_start - id: gen_summary agent: summary_writer input: records: $read_task.records - id: send_wecom agent: wecom_sender input: content: $gen_summary.summary edges: - from: read_task to: gen_summary - from: gen_summary to: send_wecom6.2 执行阶段提交任务、观察状态流转用 curl 提交任务curl -X POST http://localhost:8000/v1/tasks \ -H Content-Type: application/json \ -d {pipeline: weekly_report_pipeline, inputs: {week_start: 2024-11-18}}返回内容是一串 task_id。随后我在 Dashboard 上观察状态流转read_task 从 PENDING 变 RUNNINGtask_reader 从 PostgreSQL 读取任务记录后发布消息编排器解锁 gen_summary。summary_writer 调用 LLM 生成摘要后消息继续传给 send_wecom最终整个任务 COMPLETED。整个执行过程我在 Worker 日志里能看到每个节点的执行耗时read_task 约 320ms数据库查询gen_summary 约 8.4sLLM 调用send_wecom 约 1.2sHTTP 推送。一个完整任务大约 10 秒其中绝大多数时间都花在 LLM 调用上这也是为什么调度器必须用异步和队列而不是同步等待。6.3 异常处理伪造一个 LLM 超时看看系统怎么反应为了验证持久化恢复能力我故意把 summary_writer 的超时时间改成 5 秒正常情况下需要 8 秒重新提交任务。任务在 gen_summary 节点超时后先进入重试队列重试两次仍然超时节点标记为 FAILED_RETRIABLE。根据 DSL 中的 retry 策略系统放弃该节点将任务置为 FAILED同时给管理员发送告警。关键观察点是任务失败后read_task 的结果已经持久化在消息日志里gen_summary 的输入结构还在。我修复了 summary_writer 的超时配置后通过 OpenRig 的重跑接口只重跑 gen_summary 和后续节点read_task 直接复用之前的查询结果。这个局部重跑能力在真实生产中特别有用省掉了大量 LLM 调用成本。7. 避坑实录OpenRig 落地过程中最值得说的几个坑7.1 坑一Agent 输出不规范导致下游节点频繁解析失败我在调试多智能体协作时最先遇到的问题是上游 Agent 输出的是一个 Markdown 文档里面描述了查询结果结果专门负责解析 JSON 的下游 Agent 怎么都解析失败。后来查清楚原因——不是大模型的错是我没有在 prompt 里强制约束输出格式更没有用结构化输出JSON mode / function calling。解决方式是两个一是在 Agent 的 prompt 里增加严格的格式约束二是优先使用大模型服务商提供的 JSON 模式或函数调用能力让输出直接落成 schema 定义的 JSON。如果你的 Agent 还是输出纯文本建议在编排器层面加一层输出校验器schema 校验失败就直接触发该节点的重试。我在 OpenRig 里为每个节点配置了 output_validator用 JSON Schema 校验校验不过的节点计为失败。这个机制倒逼所有 Agent 的输出慢慢收敛到规范格式。7.2 坑二Redis 消息堆积导致下游任务延迟飙升有一次线上任务 P95 耗时突然从 20 秒涨到 90 秒。排查发现是某段时间内任务提交量暴增Redis 的就绪节点队列不断堆积而 Worker 数量没有增加。每个 Worker 取到任务后要花很长时间处理队列里积压了大量消息下游节点都等着排队。这里面最值得反思的不是加机器而是为什么到 90 秒才发现。之后我加了三层监控第一层是队列长度指标超过阈值就告警第二层是每个节点在队列里的等待时长第三层是 Agent 信号量使用率。事后复盘如果只盯着任务成功率完全看不出问题因为成功率一直是 100%只是慢了。多智能体系统里延迟比失败更容易被忽视也更容易累积成事故。7.3 坑三幂等实现不彻底重试导致重复发消息幂等这个坑我踩得最狠。send_message 这个 Agent 在第一次运行时已经成功把周报发到企业微信但因为网络抖动编排器没有收到完成确认于是触发了重试。第二次运行时 Agent 又发了一条消息于是同事在群里看到两条一模一样的周报。这个问题的根子在于执行和确认没有分开。我后来给所有会产生外部副作用的 Agent 引入了两阶段模式第一阶段是 prepareAgent 生成消息内容并暂存不真正发送第二阶段是 commit编排器确认节点可以完成后Agent 才真正发送消息。如果 prepare 阶段已经完成而 commit 丢失重试时只需要再次确认 commit不会重复发送。这个模式虽然增加了一点实现复杂度但我可以负责任地说凡是涉及真实业务副作用的多 Agent 系统都值得上两阶段提交。7.4 坑四LLM 服务商限流是整个系统最大的不稳定因素跑了两个月稳定性的最大瓶颈不在代码而在大模型服务商的 API 限流。一次压测中summary_writer 的并发从 2 调高到 5结果 10 分钟内服务商开始返回 429。由于重试策略配置得比较激进429 触发重试后雪球越滚越大整个任务的失败率反而上升。最终方案是在 OpenRig 里实现一个令牌桶限流器每个 Agent 按服务商分配的 RPM 配额配置令牌桶调度器在分发节点前先申请令牌。这样做虽然会让任务排队时间变长但能保证不会触发服务商层面的限流惩罚。顺手还做了一个动态退避收到 429 后重试等待时间按指数退避并增加随机抖动避免所有重试请求在同一时刻打过去。这也是我想提醒所有做 Agent 应用的人一件事你在框架层面做再多优化最后真正卡你的大概率是外部 LLM 接口的限制一开始就要把限流、重试、熔断当成一等公民来设计而不是等爆了再修。8. 最后聊聊这套架构的边界和几个扩展方向OpenRig 不是万能的。它解决的是多个离散 Agent 如何持久化协作这个问题但有几个场景它并不擅长。首先是实时性要求极高的交互场景比如用户和多个 Agent 实时聊天每个消息都要秒回这种场景更适合流式图编排而不是任务队列表。其次是 Agent 数量特别少两三个、协作链路固定的场景用 OpenRig 反而有点杀鸡用牛刀直接写个脚本串起来更划算。OpenRig 最适合的场景是Agent 数量中等5 到 50 个协作关系复杂且会演化任务需要长时间运行、需要断点续跑同时要求系统有可观测性。后续我打算做几个方向的扩展。一个是为任务拓扑生成可交互的运行时可视化不只展示 DAG 定义还要实时展示每个节点的执行状态、消息流转延时类似分布式追踪系统。另一个是引入基于成本的调度策略比如同一个节点有多个 Agent 可以实现不同模型、不同价格调度器可以根据当前任务的预算自动选择用贵但强的模型还是便宜但稍弱的模型。还有一个是把快照机制对接对象存储的版本管理让任务可以回溯到任意历史版本的状态。从我个人的实践体会来说多智能体编排当前最大的痛点不是模型能力而是工程基础设施。模型能力可以靠外部 API 解决但状态管理、消息可靠投递、并发控制、幂等恢复这些工程问题没有任何现成框架能一步到位只能在实际项目里一点点打磨。OpenRig 是我在这条路上的一次尝试如果你也在做类似的事希望这篇文章能帮你少走几个弯路。最后分享一个建议不要一开始就追求架构宏大先用最简单的两个 Agent 串一个真实业务跑通再把状态持久化、消息总线和恢复机制逐步加进去你会发现比闭门设计半天下笔要靠谱得多。
返回列表