ARTICLE DETAIL

资讯详情

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

多智能体编排实战:用持久化状态管理突破Agent协作上限

多智能体编排实战:用持久化状态管理突破Agent协作上限 1. 先把结论放在前面单 Agent 的能力上限不在模型而在状态管理写这篇文章的时候我刚刚把一个跑了将近一整天的多智能体编排任务恢复到断点继续往下执行。系统没有异常也没有丢失任何中间结论。这个名叫 OpenRig 的编排层是我过去几个月反复推倒重写之后确定下来的方案核心目标就一句话把若干个互不认识的 AI Agent组织成一个有记忆、有分工、能容错的持久化协作系统。先交代清楚我为什么要碰这件事。最近大半年我一直在做 AI Agent 相关的业务落地从扣子平台上的智能体应用到基于 FastAPI LangChain LangGraph 的自研链路都试过。一个非常明显的感受是单 Agent 在真实业务里的天花板往往不是模型能力而是上下文和状态管理。任务一长对话历史一多模型就顾此失彼中间一旦出现网络抖动、服务重启、某个工具调用超时整个流程就得从头再来。这种体验放在“让 AI 真的下地干活”的场景里基本是不可接受的。OpenRig 不是从零训练或者微调模型的框架它做的是编排。所谓编排就是给一堆功能单一的 Agent 定义明确的角色、通信方式和任务交接规则让它们像一个小团队那样协作。比起一个什么都干的“全栈 Agent”一组各司其职、通过持久化状态协作的小 Agent在可控性、可维护性、可观测性上都要好得多这也是我在这个项目里最核心的取舍。这篇文章会把 OpenRig 的设计思路、核心抽象、关键实现、完整实操过程以及踩坑记录都写出来给正在做 Agent 搭建、或者被多 Agent 并发和状态丢失问题困扰的同行一个参考。无论你用的是 LangGraph、Spring AI 还是自研链路第 2 章到第 5 章里的大部分经验和结论应该都能直接用上。2. 整体设计与核心抽象把“协作”变成可持久化的数据结构2.1 四个核心抽象Agent、Task、Message、Session我在设计 OpenRig 的时候给自己定了一条硬规矩系统里所有概念必须能落成数据库里的表。任何只存在于内存、进程内部、或者某个 Agent 的上下文窗口里的东西都不允许成为协作的一部分。这条规矩听起来简单但执行起来非常难因为它直接决定了后面所有的技术选型。围绕这条规矩OpenRig 定义了四个核心抽象。Agent是能力的边界一个 Agent 只做一件具体的事。比如“搜索情报的 Agent”“写初稿的 Agent”“做格式校验的 Agent”。Agent 本身不关心业务流程它只接收输入任务执行然后返回结果。Agent 的注册信息存在表里包括名称、能力描述、支持的输入输出 schema、超时时间、重试策略。Task是工作单元表示一个需要被某个 Agent 执行的具体工作。Task 有状态、有依赖、有负责人。一个业务流程会被拆成一个有向无环图DAG每个节点是一个 Task每个 Task 的数据进出都会被记录。为什么强调 DAG 而不是允许任意循环因为只有无环的流程才方便持久化、恢复和追踪循环逻辑可以靠编排层的“重新触发”来实现不需要在依赖关系上引入环。Message是 Agent 之间的通信载体。与直接函数调用不同Agent 之间不直接调对方的方法而是通过 Message 交换数据。每条 Message 都带 sender、receiver、message_type、payload、context_id 这些字段。这个设计模仿了消息总线目的是让通信本身可记录、可回放、可审计。多 Agent 系统最容易出问题的就是“谁在什么时候给谁传了什么”Message 全量落库之后这个问题就变成了一个简单的 SQL 查询。Session是贯穿始终的协作上下文。一个 Session 对应一次完整的业务目标从触发开始到最终结果交付结束。Session 这条记录是系统的主线它的状态决定了现在有多少个 Task 处于可执行状态、哪些 Agent 已经被召回、哪些上下文已经被写入。Session 的设计直接决定了系统能不能“断点续跑”所以在 2.3 节我会单独展开。这四个抽象之间的关系可以理解成Session 是一个项目的工单Task 是工单拆出来的工序Agent 是完成工序的工人Message 是工序之间传递的半成品和交接单。每个环节都有记录每个环节都能回溯这样协作才不会是黑盒。2.2 编排模型星型、链式还是图式多 Agent 编排不是只有一种模型。我在早期原型里试过三种拓扑最后明确了一个结论没有最好的模型只有按场景选择后封装成统一 DSL 的模型。链式编排是最直观的。Agent A 的输出直接作为 Agent B 的输入一个接一个执行。适合流程固定的场景比如“翻译 - 校对 - 排版”。缺点是中间只要有一个节点失败后面就断掉而且一旦某个 Agent 需要并行处理多个子任务链式模型就表达不了。星型编排是典型的主管-下属模式supervisor pattern。一个主管 Agent 负责任务理解、拆解、分派和结果整合下面的执行 Agent 各自独立。这种模型适合意图开放、子任务不固定的场景比如“用户说了一个模糊需求系统先让规划 Agent 拆解再分给多个执行 Agent”。OpenRig 对星型模型支持得最好因为它和 Task 依赖表天然兼容主管产生 Task执行 Agent 消费 Task结果回收后由主管继续决策。图式编排把整个流程建模成 DAG节点之间可以并行、汇聚、条件分支。这是表达能力最强的方案也是 LangGraph 主推的方案但同时也是对持久化要求最苛刻的方案。图式模型里任何一个节点的执行状态、上下文版本、分支选择都必须持久化否则重启恢复就是一句空话。OpenRig 的最终做法是底层统一用 DAG 表达上层提供链式和星型的 DSL 封装。用户如果只想写一个顺序流程直接写链式声明如果需求复杂再下沉到 DAG 层手动定义节点和依赖。这个分层让新手写得快也让复杂场景有表达空间。需要说明的是这是我在项目实践中为了兼顾易用性和通用性做的取舍并不是说 DAG 是最优的只是它作为底层模型最不容易被业务场景卡住。2.3 持久化优先为什么 Session 必须成为系统的脊柱早期版本里我把 Session 当做一个普通的运行时对象用内存字典保存每隔一段时间手动快照一次。结果在一次进程重启后四个 Agent 协作到一半的场景全部丢失人工排查花了两个多小时。那次教训之后我把持久化从“附加功能”改成了“第一优先级”。具体来说Session 必须解决三个问题。第一上下文不断增长怎么办。Agent 协作过程中会产生大量中间文本不可能每次全量塞进模型上下文。OpenRig 的做法是分层记忆原始 Message 全量保存在存储里但每个 Session 维护一个“压缩后的上下文摘要”只有摘要和最近 N 条关键 Message 会被送进后续 Agent 的提示词。这个思路和我用 LangChain 做长对话记忆时的做法一脉相承只是从单 Agent 的 memory 扩展成了多 Agent 的共享记忆。第二恢复点怎么定义。不是每条 Message 都值得作为恢复点OpenRig 只在两类时机写检查点一个 Task 状态变为 success/failed 的时刻以及一个 Agent 完成全部工作并返回结果给编排器的时刻。这两个时机都是业务上的稳定节点记录的上下文也是完整且自洽的恢复起来不会出现“执行到一半引用了尚未存在的数据”这种问题。第三失败之后怎么续跑。恢复逻辑很简单加载 Session 最近的检查点找到所有处于 pending 和 in_progress 状态的 Task把 in_progress 的 Task 重置为 pending因为它的执行结果未必可靠然后重新调度。这里有一个隐含的重要约定Agent 必须被设计成幂等的也就是说同一个 Task 被重复执行两次结果不会产生副作用叠加。这一点我会在第 3 章展开讲。持久化优先的设计让 OpenRig 在真实环境里表现得非常稳。数据库里有全量消息、有任务状态、有检查点重启、扩容、甚至换一台机器跑调度器都不会影响协作进度。可以说Session 就是整个系统的脊柱其他所有模块都是围绕它长的肉。3. 关键模块实现调度、通信、状态恢复一个都不能少3.1 消息协议与路由Agent 之间到底怎么说话Agent 之间的通信协议我吃了不少苦头才定下来。最开始的版本里Agent 之间直接传递 Python 对象爽是爽但一旦 Agent 跑在独立进程或者另一台服务器上这套设计就完全崩了。后来我把协议收敛成 JSON 格式的消息字段固定payload 部分兼容任意合法 JSON。一个典型的 Message 长这样{ message_id: msg_01HZX..., session_id: ses_01HZW..., sender: search_agent, receiver: writer_agent, message_type: task_result, correlation_id: task_01HZY..., timestamp: 2025-01-18T10:23:11Z, payload: { status: success, data: {topic: OpenRig, materials: [...]} } }几个字段值得单独解释。message_type用来区分是任务结果、请求补充信息、还是错误回报correlation_id把 Message 和触发它的 Task 关联起来这一步非常关键因为没有它你没法在消息日志里还原“这条结果到底是响应哪个任务的”。payload只装业务数据不装系统控制字段避免业务和框架耦在一起。路由逻辑在实现上其实很轻。每个 Session 内部维护一张“接收者表”记录当前处于等待状态的 Task 对应的 Agent。当一条 Message 到达时调度器根据correlation_id找到对应的 Task再把 payload 写入 Task 的结果字段最后推进依赖状态。这里完全不需要一个全局的消息中心因为所有的路由信息都能从持久化的 Task 表里查出来。用 SQL 做路由是我这个项目里最满意的一个决定它让消息系统变得极其简单可靠。3.2 调度引擎与并发控制多 Agent 怎么扛得住并发“AI Agent 怎么扛并发”是很多人问我的问题。OpenRig 的做法并不神秘核心就是三点一是 Session 级隔离二是 Task 级并行三是全局限流。Session 级隔离的意思是不同 Session 之间完全不共享可变状态每个 Session 有自己的上下文、自己的任务列表、自己的执行线程。隔离保证了一个业务方的高负载不会污染另一个业务方的上下文。Task 级并行是整个系统吞吐量的来源一个 Session 内部如果有多个互不依赖的 Task调度器会同时执行它们。比如“情报收集”阶段要查三个来源那就起三个并行 Task分别交给三个 Agent 实例执行。实现上我用了一组工作线程池线程数可以通过配置调整。调度器每次从所有未完成 Session 里挑出当前可执行的 Task放进一个优先级队列然后由固定数量的 worker 消费。这里最重要的细节是同一个 Session 内可并行的 Task 可以同时执行但同一个 Session 的状态提交必须串行。我用了每个 Session 一把独立锁Task 完成后worker 需要先抢到 Session 锁再更新状态和写检查点。这个设计避免了“两个并行任务同时修改上下文导致覆盖”的经典问题。全局限流则是为了照顾模型 API 的速率限制。OpenRig 里维护了一个简单的令牌桶每个 Agent 调模型之前都要申请令牌。令牌桶的速率按模型供应商和模型名分开配置比如 GPT-4 类慢速模型限流严一点本地小模型可以放宽。没有这个限流高并发下最先挂掉的一定不是你的服务而是上游的模型 API。3.3 状态快照、检查点与恢复重启不再是灾难先看快照的数据结构。每次检查点保存三样东西Session 的元信息、所有 Task 的状态快照、以及从上一个检查点以来新增的 Message。这三样合起来足够在任意时刻重建整个协作现场。恢复的代码逻辑很直接核心流程如下def resume_session(session_id: str) - Session: session load_session(session_id) checkpoint load_latest_checkpoint(session_id) tasks load_tasks(session_id) # 把执行中的任务重置为待执行等重新调度 for task in tasks: if task.status TaskStatus.IN_PROGRESS: task.status TaskStatus.PENDING task.assignee_id None save_task(task) # 重建消息索引和路由表 rebuild_message_index(session_id) scheduler.wake_up(session_id) return session恢复完成后调度器会像正常流程一样把所有 PENDING 且依赖已满足的 Task 继续派发给对应 Agent。这里有个容易踩的坑如果 Agent 本身不是幂等的重置 IN_PROGRESS 任务会导致同一步被执行两次。比如一个“扣款”Agent 被重复调用用户就会被扣两遍钱。所以 OpenRig 官方推荐所有 Agent 的副作用操作都走后置补偿或者幂等键。对多数内容生成、检索类 Agent 来说天然就是幂等的这一点不太用担心。另一个值得注意的小细节是检查点的存储频率。写检查点本身有 I/O 成本频率太高会影响吞吐。我的经验值是按 Task 粒度写入不要按 Message 粒度写入。一个 Task 可能会产生 10 条 Message但只有它完成时才有必要落一次检查点。这样既保证了恢复粒度足够细又不至于把数据库写入变成瓶颈。3.4 可观测性与问题定位多 Agent 系统的救星多 Agent 系统最大的噩梦不是写不出来而是出了问题不知道去哪里查。两个 Agent 来回传了几轮信息最后结果不对你连是哪一步开始错的都找不到。OpenRig 的可观测性设计是围绕“链路追踪”展开的。每条 Message、每个 Task、每次 Agent 调用都会带上同一个trace_id。这个 trace_id 会在 Session 创建时生成贯穿全程。当业务方来投诉时我只需要拿到 trace_id就能把完整链路从数据库里拉出来先查 Session再查这个 Session 下所有 Task 的执行顺序再查每个 Task 消耗的 Message最后定位到具体是哪一步产生了错误输出。在实现上我写了一个简单的查询接口按时间线输出每个阶段的输入摘要、输出摘要、耗时和 Token 消耗。摘要不会存全文只存前几百字符和关键统计字段避免日志表膨胀得太快。早期我遇到过“排查问题要翻几万条原始记录”的情况加了摘要层之后90% 的问题一眼就能定位。还有一个小工具值得做Agent 输出对比。当同一个 Task 被重试多次时系统会把每次重试的输入参数和输出结果并排展示。这个功能在排查“为什么 Agent 这次回答和上次不一样”时特别好用因为多数时候不是随机性问题而是重试时上下文里混进了新的消息导致 Agent 决策变了。4. 实操用 OpenRig 搭一个“情报收集 内容生成”协作流程4.1 初始化工程与环境理论讲再多不如动手跑一遍。下面用一个完整的例子演示 OpenRig 的使用流程任务是把一个技术主题做成一篇结构化的调研简报。整个过程分三步先让“情报收集 Agent”搜集素材然后让“内容规划 Agent”整理大纲最后让“写作 Agent”生成全文。三个 Agent 之间通过持久化 Session 协作任何一步失败都可以从断点恢复。先初始化一个 Python 工程安装依赖mkdir openrig_demo cd openrig_demo python -m venv venv source venv/bin/activate pip install openrig fastapi uvicorn pydantic2 sqlalchemyOpenRig 默认使用 SQLite 起步零配置跑通之后再换 PostgreSQL 也不费劲。配置好环境变量后启动一个最小的 FastAPI 服务作为宿主进程OpenRig 会把编排调度和 API 服务都挂在这个进程里。4.2 定义两个角色 AgentAgent 的定义非常简单核心就是一个函数加上注册信息。以情报收集 Agent 为例from openrig import Agent, register_agent def search_materials(topic: str, max_results: int 5) - dict: # 这里是实际检索逻辑可以是搜索 API、内部知识库或本地文档 results fake_search(topic, max_resultsmax_results) return { topic: topic, materials: [{title: r.title, snippet: r.snippet} for r in results], count: len(results) } search_agent Agent( namesearch_agent, description负责检索主题相关的资料、文章和案例返回结构化素材列表, executesearch_materials, input_schema{topic: string, max_results: integer}, output_schema{topic: string, materials: array, count: integer}, timeout_seconds60, retry_policy{max_retries: 3, backoff: 2.0} ) register_agent(search_agent)这里面最容易被忽略的字段是description。OpenRig 里的规划 Agent 靠这个描述来决定任务分给谁描述写得越具体分派准确率越高。我见过很多人写“负责搜索”就完事了结果规划 Agent 把写作任务也分给它。描述至少要包含职责边界、擅长和不擅长的内容。写作 Agent 的写法类似输入是主题和素材列表输出是 Markdown 文本。为了让 Agent 幂等我要求它在开头带上session_id作为标识同一个 Session 重复生成时直接复用已有草稿不会产生重复内容。4.3 创建 Session 并编排任务依赖接下来是核心的编排逻辑。OpenRig 提供一套声明式的编排接口用户可以定义一个工作流指定 Agent 之间的依赖关系from openrig import Workflow, SessionManager, task workflow Workflow(material_report) task(assigneesearch_agent, depends_on[]) def search_step(ctx): # ctx 里自动带上了 Session 上下文包括用户输入 topic return search_agent.execute(ctx.inputs[topic]) task(assigneewriter_agent, depends_on[search_step]) def write_step(ctx): # 这里能直接拿到 search_step 的结果因为持久化上下文自动注入了 materials ctx.result(search_step) return writer_agent.execute(ctx.inputs[topic], materials) workflow.add_task(search_step) workflow.add_task(write_step) manager SessionManager() session manager.create_session(workflow_idmaterial_report, inputs{topic: OpenRig 多智能体编排}) manager.start_session(session.id)这段代码最值得解释的是ctx.result(search_step)。OpenRig 允许后续 Task 按前置 Task 名称拉取结果实现上是直接从持久化存储读的而不是从内存变量读。这意味着即使写作 Agent 被调度到另一台机器上它依然能拿到前置结果。跨进程、跨机器的数据依赖就靠这一步打通了。4.4 运行验证与效果检查跑完上面的流程可以检查一下 Session 的最终状态openrig-cli session show ses_01HZW...输出会列出 Session 下每个 Task 的执行状态、耗时、Token 消耗、输入输出摘要。正常情况应该看到两个 Task 都是success写作结果已经写入了 Session 的最终产物字段。为了确认协作过程确实落库我习惯直接查一下消息表看看 search_agent 发送给 writer_agent 的那条 Message 是否包含完整的素材列表。这一步能验证消息路由是否真的把数据传到了正确的地方。在本地跑通后这个流程可以用容器镜像部署成常驻服务FastAPI 会暴露一个简单的 HTTP 接口外部业务系统通过 POST 请求创建 Session、查询进度、获取结果。OpenRig 本身不绑定特定前端扣子、微信机器人后端、企业内部工作台都可以通过这个接口接进来。4.5 断点恢复与并发压测演示完正常流程我再手动模拟一次故障。跑完搜索步骤之后直接把服务进程 kill 掉然后重新启动调用恢复接口manager.resume_session(session_idses_01HZW...)观察日志会发现search_step 的状态从 success 被保留而尚未执行的 write_step 被重新调度并运行。整个恢复过程在秒级完成没有丢失搜索结果也没有重复执行搜索 Agent。这就是“持久化协作系统”和普通脚本最本质的区别脚本挂了就是挂了协作系统挂了还能接着干。压测方面我模拟了 50 个并发 Session每个 Session 内部又拆出 3 个并行子任务。在 SQLite 环境下吞吐量大概到每分钟 300 个 Task 就开始有写入竞争切到 PostgreSQL 之后同样配置下每分钟能跑到 1800 个 Task 以上。这个数据说明瓶颈主要在后端存储而不是编排层生产环境建议直接上 PostgreSQL连接池调到 20 起步。这部分压测结果是我在自己项目里的实测数据不同机器上会有差异但结论方向是一致的存储先行别先优化 Agent 调用那多半不是瓶颈。5. 常见问题与排查速查5.1 高频问题实录与解决方案多 Agent 编排系统在真实环境里的问题往往集中在这几类Agent 输出格式失控、任务卡死、上下文串台、恢复后行为不一致。下面是我实战里遇到的高频问题整理成速查表。问题现象根因分析解决方案Agent 返回的 JSON 解析失败模型输出掺杂了 Markdown 或额外文字输出层加结构化解码器先提取代码块再 parse同时提示词里强调“只输出 JSON”多个 Task 并行执行后上下文互相覆盖Session 状态提交没有串行引入 Session 级锁状态更新走统一提交入口禁止各 worker 直接写上下文日志里出现大量超时重试系统吞吐骤降上游模型 API 被限流按模型供应商和模型名配置令牌桶超时重试用指数退避避免雪崩断点恢复后Agent 重复执行了副作用操作Agent 不满足幂等性所有副作用操作加幂等键执行前查重无法幂等的操作挪到编排层做后置补偿两个会话之间的上下文互相污染全局变量或缓存 key 没有带 session_id所有缓存 key 强制拼接 session_id禁止使用裸的业务 ID 作为缓存键Agent 明明收到了正确输入却答非所问提示词里保留了过期上下文升级上下文管理只送摘要 最近 N 条关键消息过期消息归档后不参与生成其中 Agent 输出格式失控是最常见也最隐蔽的。模型不是不会输出 JSON而是它会很自然地“顺带”输出解释文字。我的经验是第一层用正则抽取 JSON 片段抽不到就调一次“修复 Agent”专门整理格式交给“修复 Agent”前把原始输出原样保存下来方便定位是哪里开始坏的。卡死问题也值得多说一句。OpenRig 默认给每个 Task 设置了超时时间但单纯超时往往不够很多任务卡住是因为依赖关系里有隐式的环。比如 A 等 B 的结果B 又在等 A 的确认两个 Task 都处于 waiting 状态谁也不会先动。排查这种问题我的工具是把所有 waiting 任务的依赖表导出来做一次拓扑排序有环就是编排定义错了没有环那就是超时配置太小。5.2 性能与并发调优从能跑到跑得快如果你的多 Agent 系统已经从“能跑”进入“跑得快”阶段下面的调优经验可以直接抄。第一把模型调用和业务逻辑分离。Agent 的函数体内不要一边调模型一边写数据库。模型调用放在独立的执行器里方便统计耗时和 Token业务逻辑放在编排层方便做事务控制。这个分离对性能排查有决定性的帮助出问题时你能立刻分清是模型慢还是代码慢。第二上下文压缩要跑在后台。长会话的摘要生成本身也要调模型如果放在关键路径上会拖慢整个链路的响应。我的做法是Task 完成后先落库、先标记状态、先返回给上层摘要压缩放到异步队列里慢慢算。上层用户感知到的延迟就只包括真正必要的 Agent 调用耗时。第三高频读取走缓存状态写入走数据库。Session 的当前状态读取频率很高每次调度决策都要读一次但写入频率低。所以我把 Session 的实时状态放在 Redis 里做缓存数据库只保存可恢复的检查点。注意这里有个一致性代价缓存和数据库之间可能有秒级延迟但多 Agent 协作系统天然容忍这种延迟因为最终一致性足够满足业务需求。实测下来这个改动把调度器的决策开销从 30ms 降到了 3ms 以内收益非常明显。第四预留一个“降级通道”。业务高峰期模型 API 不稳定我给 OpenRig 加了一个全局开关可以把所有 Agent 的执行切换到“模拟模式”直接返回缓存的历史结果。这样在联调环境和上游故障时系统依然能完整跑完整个编排流程。这个模式在生产排障时救过我很多次强烈建议保留。6. 一些个人体会最后分享一点个人感受不算总结就是些实在话。多智能体编排这件事最难的地方其实不是技术而是克制。我见过很多团队一上来就搭了十几个 Agent看起来热闹实际上职责重叠、通信混乱最后连自己都说不清哪个 Agent 干了什么。OpenRig 的设计让我学会了一件事Agent 的粒度宁愿偏大一点也要职责清晰。一个 Agent 从搜索到筛选到总结全都包了虽然在概念上“不够优雅”但在排障和调优上的收益远大于概念上的损失。另一个体会是持久化不是“为了不丢数据”这么简单它实际上是整个系统的调试工具。有了全量消息、任务状态和检查点你就拥有了多 Agent 系统的黑匣子。我最近处理的一个线上问题就是靠查 Message 历史定位到是某个 Agent 在第三轮协作时把价格信息写进了错误字段进而导致下游决策偏差。没有持久化日志这种问题基本只能靠猜。如果你正在做自己的 Agent 项目我的建议是先别急着上编排框架先把“状态打出来、日志留下来、断点续跑跑通”这三件事做好。这三件事做好了哪怕你用的是最简单的顺序调用系统的健壮性也会超过大多数花哨的编排方案。等确实需要并行、需要分工、需要跨进程协作的时候再引入 OpenRig 这样的编排层你会发现自己对“系统需要什么”有了远比之前清晰的认识。最后再分享一个小技巧给每个 Agent 取名的时候把它的 Router 路径和职责缩写加进去比如search_agent_v2。这个小习惯在日志检索的时候能省下大量时间排查问题靠grep的时候你才会感谢当初命名的用心。
返回列表