
最近我把手里七零八落的Agent脚本重新收拾了一遍做成了一个小框架取名OpenRig。起因很朴素项目里已经有了选题Agent、资料Agent、写作Agent、校对Agent单独跑都很乖可一旦想让它们配合完成一件完整的活就发现谁也不知道谁在干嘛。干到一半进程崩了只能从头再来烧掉的Token让人肉疼。OpenRig要解决的就是把这种离散的AI Agent编排成一张可以持久化恢复的协作网。我会把设计思路、核心代码、踩坑实录都拆开讲适合正在做多智能体应用、被状态同步和崩溃恢复折磨过的朋友参考。1. 为什么要把离散Agent编成一张持久化的网1.1 单独调Agent容易协同运作难现在大家口中的AI Agent大部分是一个“带记忆和工具的LLM封装”给定目标它能自己规划步骤、调用搜索、读写文件、最终给出结果。单个Agent的调用链路其实很简单输入Prompt → 模型推理 → 工具调用 → 输出中间再加一点重试逻辑一个人用起来已经挺顺手。可一旦你同时维护五个、十个Agent问题就来了。第一个问题是谁来调度。Agent之间不是简单的水管串联而是随时可能出现分支。比如内容生产场景资料Agent要把检索结果交给写作Agent写作Agent又要触发配图Agent这里头有先后依赖也有可并行的分支。没有编排层的话你只能在业务代码里硬编码一堆if/else去串流程。加一个新Agent就要改一遍主流程牵一发动全身后面会越来越乱。第二个问题是上下文一致性。每个Agent的输入输出都基于当前任务上下文但多个Agent并行时大家各改各的共享状态非常容易乱套。我最初的做法是把所有中间结果都塞进一个Python字典后来发现A写完的数据被B覆盖B读到的版本和C要处理的对不上光是排查“这个字段为什么变成None了”就花掉一整个晚上。所以说多Agent真正考验的不是单点模型能力而是怎么把一堆各有脾气的“员工”组织进同一条流水线。每个环节都要有明确的输入输出出错时能快速定位是哪一个环节的问题截断在哪一段。这就是编排层存在的意义。1.2 持久化不是可选项而是长任务刚需有人可能会说编排层有了状态放内存里不就行了为什么非要持久化我第一次做多Agent系统时也这么想后来连续踩翻了好几个场景才把这条铁律立起来。最直观的是成本压力。LLM按token计费token是最小文本单位一个中文词大概对应一到两个token。一个跨Agent任务动辄几十次模型调用如果运行到第35步时进程被kill掉下次从第一步重来重复调用烧掉的钱足够买好几杯咖啡。持久化能让运行中的快照保存下来恢复时从断点继续省下的不只是时间是真金白银。其次是长任务。批量生成文章、数据清洗、定时巡检这类工作经常要跑几十分钟甚至几小时。中间会有部署升级、机器重启、进程被OOM Kill如果没有持久化所有Agent进度跟着内存一起消失这种系统根本没法放到生产环境里跑。还有审计需求。项目里总要回答“这个结论是哪个Agent、哪一步、基于什么Prompt产出的”。消息队列里的历史事件就是最好的审计日志前提是它被持久化了而不是只存在于内存。所以OpenRig的第一条设计铁律是所有关键状态必须落盘运行中任何时刻杀掉进程都能从最近一个完整状态恢复继续跑。2. OpenRig的核心架构设计三个关键取舍2.1 编排层与Agent解耦写第一个版本时我把编排逻辑直接写进了Agent基类里Agent自己知道下一步该调用谁。结果Agent之间耦合越来越重业务逻辑和调度逻辑搅成一锅粥。重构OpenRig时才想明白Agent应该只关心“输入什么、输出什么”至于谁在等它、下一步去哪个节点是编排层该管的事。因此OpenRig把系统拆成两个抽象层。业务节点只实现一个execute方法在里面调用LLM、工具、搜索并把产出返回。它相当于流水线上的一名工人不做任何路由决策。路由边负责决定当前节点完成后下一步进入哪个节点。边会读取整个共享状态按规则返回下一个节点ID。这个设计的直接好处是新增一个Agent完全不用动老代码只需要在编排描述文件里加一个节点和几条边。Agent甚至不需要用同一种语言写只要通过统一协议收发状态就行。我因为项目栈是Python直接用Python实现了Node基类但协议层面已经和语言解耦。这也是OpenRig这个名字的由来它更像一个可装配的钻台把不同工具拧紧在同一个平台上。2.2 用Redis做状态底座快照与事件双写编排系统最核心的组件是状态存储。我对比过几种方案。内存方案开发最快但一崩全丢直接出局。关系型数据库很稳但要把每个Agent的产出结构化成表动态字段多了之后迁移非常痛苦。文件系统存JSON简单可并发读写和分布式场景都不友好。最终选了Redis。原因不是它快而是它的数据结构很适合表达状态和事件持久化机制也足够成熟。OpenRig采用了“快照 事件日志”双写策略。每隔N个节点执行完把整个会话状态序列化成一个JSON快照写入Redis的一个Key同时触发一次持久化。每个节点完成后的详细结果和处理事件以追加方式写入Redis Stream供审计和重放。这里要特别提一下Redis持久化里的RDB和AOF。RDB是定期全量快照恢复快但可能丢最后一次快照后的数据。AOF是命令追加日志数据丢失少但文件重放慢。我的选择是两者同时开AOF配everysec策略保证最多丢一秒数据RDB用于快速重启恢复。这样既保证顺滑恢复也不会因为AOF文件过大导致加载卡顿。这个组合在单机规模下够用如果跨地域多节点部署就要再考虑分布式存储方案。2.3 基于版本号的乐观锁解决状态冲突多Agent并行的最大麻烦是共享状态覆盖。比如资料Agent和配图Agent同时向状态里写入各自的产出如果只是读改写后写的会把先写的覆盖掉整个上下文就坏了。OpenRig没有用重量级分布式锁而是给状态挂了一个版本号。每次Node完成写入时必须带上它读取时的版本号。Redis端用Lua脚本原子地执行对比当前版本号等于传入版本号才允许更新并把版本号加一。如果版本号对不上说明有别的Agent抢先改过当前这次的写入就失败由调度层决定是重试还是丢弃。这个乐观锁的思路很像多人协同编辑不锁整个文件只检测冲突。实际测试下来在大多数Agent协作场景中比如不同Agent写不同的顶层字段冲突率很低。乐观锁的代价远小于全局锁带来的等待吞吐量要高不少。3. 核心机制落地节点、边、消息、调度3.1 节点与边的描述方式OpenRig的编排图不散落在代码里而是用一个Python字典描述这样业务同事也能看懂全局流程。一个最简单的图长这样GRAPH { nodes: { start: {agent: planner, params: {视角数量: 3}}, research: {agent: researcher, params: {max_results: 10}}, write: {agent: writer, params: {style: blog}}, review: {agent: reviewer, params: {strict: True}}, }, edges: [ {from: start, to: research, condition: always}, {from: research, to: write, condition: always}, {from: write, to: review, condition: always}, {from: review, to: write, condition: review_reject}, {from: review, to: end, condition: review_pass}, ], }这里每个节点对应一个Agent类名params是Agent的固定参数。每条边带一个condition指向调度器中的一个判断函数。condition函数接收整个状态返回True或False调度器根据结果决定走哪条边。这种描述方式的优势是图结构一目了然而且可以通过函数动态生成。我在实际项目中甚至把图定义存进了数据库运营同学直接在后台改流程配置Agent代码一行不动。这是编排层解耦带来的额外红利也让我后面维护多个不同流程时轻松了很多。3.2 会话状态与消息路由每个任务对应一个session_id它是整个生命周期里所有状态的容器。Redis里用三个关键Key组织数据session:{id}:state —— 当前最新状态JSONsession:{id}:version —— 当前版本号session:{id}:stream —— Redis Stream保存所有事件节点之间不直接通信而是把产出写入state同时把事件追加到stream。下一个节点读取最新的state作为输入。这样消息路由天然被状态驱动不需要额外维护消息总线。比如写作Agent写完初稿把draft字段写进state校对Agent启动时读state[draft]执行完把review_result写回state。这里有个很实际的细节状态字段的命名规范非常重要。如果让不同Agent自由发挥命名很快会出现同一个含义对应三种key的灾难。我引入了一个共享的字段字典文件集中定义所有跨Agent字段。Agent只能通过get_field(article.title)这类访问器读写避免字符串硬编码。后期想加字段时只需要在字段字典里补一条规则运行时会自动校验数据格式。3.3 并发控制与超时重试不少朋友问我“AI Agent怎么扛并发”OpenRig在调度器里内置了一个轻量级并发池用信号量控制同时运行的节点数import asyncio from asyncio import Semaphore class Scheduler: def __init__(self, max_concurrency4): self.semaphore Semaphore(max_concurrency) async def run_node(self, node, state): async with self.semaphore: try: return await node.execute(state) except Exception as e: # 交给上层重试策略 raise e并发数默认是4不是越大越好。多数LLM API都有每分钟调用限制并发过高只会让429错误暴涨重试等待反而拖慢整体任务。实际调优时我从1开始往上压观察错误率和吞吐量的交叉点最后在项目使用的供应商上把并发设为5。超时重试同样关键。每个Node的execute都要支持传入timeout参数超时后抛出TimeoutError调度器根据预设的max_retries决定重试还是整条边进入人工审核。重试时有个容易忽略的坑有些外部API不支持幂等重复调用会产生两条记录。我的解决办法是给每次execute生成一个request_idAPI请求带上它供应商侧就可以做去重。如果供应商不支持就在业务代码里记录已调用的参数指纹重试前先查指纹。4. 实操复盘四Agent内容生产线4.1 场景拆解与编排定义我用一个真实场景完整演示OpenRig怎么跑起来搭建一条“选题-研究-写作-校对”的内容生产线。四个Agent分别是planner分析用户给的选题方向拆出3个具体标题和关键角度researcher根据标题检索资料输出事实列表和引用链接writer基于资料写初稿reviewer对照事实列表校对初稿输出通过或不通过判断这个场景覆盖了串行依赖、条件分支和失败恢复非常适合演示。writer必须等researcher拿到资料才能启动reviewer不通过时需要带着修改意见回到writer通过后才可以结束任务。这个编排图跟上一节的GRAPH示例基本一致唯一的重点是review到write的回边条件函数。判断条件函数的实现很简单它读取reviewer写入的review_resultdef review_reject(state): result state.get(review_result, {}) return result.get(pass) is False def review_pass(state): result state.get(review_result, {}) return result.get(pass) is True这样reviewer节点只要在自己的输出里写一个“pass”标志路由完全交给边上的判断逻辑Agent自身不关心下一步去哪。4.2 核心代码实现要点Node基类的最小实现如下class BaseNode: agent_name: str async def execute(self, state: dict, request_id: str) - dict: # state是完整会话状态节点只读取自己需要的字段 # 返回dict调度器会把它合并进state并追加事件 raise NotImplementedError调度器执行一个节点的流程是锁定状态版本 → 调用node.execute → 用版本号写回state → 追加事件到stream → 判断边 → 进入下一个节点。其中状态合并策略很重要。我默认采用“深合并”如果返回的dict里有一个嵌套dict就递归合并数组则整体替换。这样不同Agent写不同的顶层key不会互相覆盖同一个key的子字段也能增量更新。举个例子researcher返回{facts: [...], sources: [...]}writer后续返回{draft: ...}两者会自然合并成包含facts、sources、draft三个顶层字段的状态。写回state并追加事件的原子操作用Redis的Lua脚本保证-- KEYS[1]: state key -- KEYS[2]: version key -- ARGV[1]: expected_version -- ARGV[2]: new_state_json -- ARGV[3]: new_version if redis.call(get, KEYS[2]) ARGV[1] then redis.call(set, KEYS[1], ARGV[2]) redis.call(set, KEYS[2], ARGV[3]) return 1 else return 0 end实际开发时我会把整个编排定义放进一个workflow.py文件里面注册节点类、图和condition函数。部署时只需要配置好Redis地址和Agent的API Key就能通过一个统一入口启动整个编排任务。4.3 崩溃恢复验证过程代码写完只算完成一半真正考验是故障恢复。我在测试时故意在writer节点执行到一半时kill掉整个进程然后重新启动OpenRig观察它能不能从最近的快照恢复。恢复的核心逻辑是启动时对每个未完成任务从Redis读取state和version检查stream里最后一条事件对应的节点id把游标定位到下一个未执行节点继续调度。这里有一个关键要求节点结果的状态写入和事件追加必须同时成功否则恢复时状态和事件对不上。我用Redis的MULTI/EXEC把这两个写操作包在一起确保要么都写要么都不写。实测下来从断点恢复平均耗时不超过两秒。但要注意如果崩溃发生在快照上传之后、节点结果写入之前恢复后会重复执行这个节点。所以每个Agent都必须设计成幂等的这个特点我下一节重点讲。5. 从坑里爬出来的排查经验5.1 常见问题速查表多Agent编排的报错往往不是单一原因我整理了一张速查表基本覆盖了前期迭代里碰到的典型问题症状可能原因排查与解法任务卡住不动并发池被耗尽某个节点长时间未返回调低单个节点超时检查外部API是否挂起状态被莫名覆盖两个节点并行写同一个顶层key拆分字段或依赖版本号检查避免共享可变字段恢复后重复调用节点不是幂等的增加参数指纹重试前查重AOF重放后状态不对RDB和AOF互相混用确保RDB快照后统一从AOF追加重放Redis内存暴涨事件日志没设长度上限用XADD MAXLEN限制比如保留最近10万条API返回429并发设置过高从并发1开始压测找到稳定上限还有一个容易踩的坑是Redis连接池。任务量上来后如果每个节点都新建连接很容易把Redis连接数打满。我在调度器里复用一个连接池并且在每次操作后立即归还连接这个问题就消失了。5.2 幂等性设计避坑幂等性是多Agent系统里最容易被忽视、又最容易致命的问题。Agent执行可能买了单、发了邮件、写了数据库你重跑一遍就会产生双倍副作用。我给Node加了一个可选接口def dedup_key(input_state) - str。调度器在执行前把这个key记到Redis去重集合里如果发现同一个key已经执行过就直接返回上次的输出不再调用Agent。对可重复查询的Agent来说这个机制能省下大量成本和等待时间。但要注意dedup_key必须包含所有影响输出的输入字段。如果只用一个标题字段做key改了个要求还是会命中旧结果。我建议用输入字段的稳定哈希至少包含topic、template、max_tokens这几个关键参数。另外去重集合要设置过期时间避免长时间运行后Redis里塞满了历史key占内存不说还容易误伤正常的重复请求。5.3 恢复与审计的边界最后说一个架构层面的体会持久化状态和审计事件要分开。RDB快照和Redis Stream承担的是恢复职责流数据会不断滚动淘汰。如果客户需要完整的审批链路审计我会把关键事件同步一份到对象存储或普通数据库而不是依赖Redis Stream。因为Redis Stream的MaxLen和淘汰策略是为性能设计的不是为长期审计设计的。另外普通日志和审计事件要明确区分。Agent的prompt、输出、异常堆栈应该进日志系统而“哪个Agent在哪个时间读写了哪个字段”才需要进事件流。如果两者混在一起事件流很快会被调试信息淹没真正要排查问题的时候反而发现关键节点被冲散了。最后再分享一个我认为很值钱的细节给每个会话加上parent_session_id。一旦任务失败需要人工介入运营可以顺着这条链查出是哪一步出的问题甚至手动调整中间状态后再恢复。这个字段只加一行但后续排查效率能翻好几倍。OpenRig在我的项目里已经跑过几十条内容生产线中间经历过几次kill和重启恢复逻辑一直很稳。做多Agent协同先把状态落盘和幂等性这两件事想透再谈并行和优化也不迟。