最近我接手了一个内部工具平台的改造,发现手头十几个 AI Agent 都在各自为战:有做需求拆解的,有做代码生成初审的,有跑测试用例推荐的,还有做发布说明汇总的。表面上都是 Agent,实际上协作靠人肉搬运,会话上下文各自维护,服务一重启全部“失忆”。折腾了三个礼拜,我把这套东西重构成了一个叫 OpenRig 的多智能体编排系统——核心就一句话:把离散的 AI Agent 编织成持久化协作系统。这篇博文就是这次实践的完整记录,从架构设计讲到 Redis 持久化机制,从编排器代码写到排查实录,适合正在做 AI Agent 平台、需要解决多 Agent 协作和状态丢失问题的开发者参考,也可以当作入门多智能体编排的一次落地示范。
1. 先聊聊一个实际问题:为什么单 Agent 撑不住协作场景
1.1 单 Agent 的能力边界与真实业务诉求
很多团队做 Agent 应用,一开始都是“一个 Agent 包打天下”:你直接跟它对话,它调工具、搜资料、写代码、最后给结果。在小范围 demo 里这个模式很爽,但放到真实业务里,马上会遇到几个绕不开的问题。
第一是上下文窗口的物理限制。你不可能把“整个项目的代码库 + 全部需求历史 + 实时监控指标”一次性塞进一个 Agent 的上下文里,Token 成本先不谈,模型的处理质量和响应速度都会断崖式下跌。第二是职责混淆。让同一个 Agent 既做需求分析又做代码审查还做测试生成,它的 prompt 会变得越来越臃肿,输出风格越来越不稳定,你甚至不知道它在哪个环节开始出错。第三是状态维护。单 Agent 跑一次任务,状态都在内存里,偶发重启、网络中断、超时重试,都可能让整个任务从头再来。
真实业务的多智能体协作场景,其实很像一个项目组:产品经理负责拆需求,开发负责写代码,测试负责验证,发布负责人最后汇总。每个人专职做一件事,通过文档、会议、消息传递来协同,而且每个人都有自己的工作记录——这就是“持久化”的原始需求。OpenRig 的核心思路就是把这种组织方式搬进系统里,让每个 Agent 只干自己擅长的一小块,用编排器把结果串起来。
1.2 “离散 Agent”是怎么凑到一起的
所谓“离散 Agent”,指的是那些已经能独立完成特定任务的 Agent 实例。它们往往并不是为了协作而生的:可能是团队里某个人单独开发的对话机器人,可能是某个供应商提供的 API 封装,也可能是一个跑批任务脚本加了层模型调用。它们没有统一的协议,没有共同的上下文存储,各自维护自己的会话状态,甚至部署在不同的服务器上。
我这次要整合的“离散 Agent”包括:一个负责用户意图识别的轻量 Agent,一个负责方案设计的代码理解 Agent,一个负责执行安全检查的扫描 Agent,以及一个负责生成交付文档的写作 Agent。这四个 Agent 都是现成的,直接丢弃重写不现实,硬把它们塞进同一个进程更不现实。OpenRig 的定位是“编排”而非“重写”——它不关心每个 Agent 内部怎么实现,只关心怎么让它们像一个整体那样工作。
1.3 OpenRig 到底在解决什么
OpenRig 解决的,是三个层面的问题。
任务层面的“接力赛”:一个复杂任务被拆成多个步骤,每个步骤由一个 Agent 完成,上一个 Agent 的输出成为下一个 Agent 的输入,而不是每个人都要从头理解全貌。这让每个 Agent 的提示词保持精简,输出质量稳定。
状态层面的“记忆”:整个协作过程的所有消息、中间产物、任务快照、运行状态,都要沉淀到持久化存储里。任何一个 Agent 重启、网络波动、编排器升级,业务都可以从最近一次快照恢复,而不是白干一场。
管理层面的“可控性”:多个 Agent 并行执行时,谁先跑、谁后跑、失败重试几次、并发上限是多少,都需要有一个清晰的编排规则。OpenRig 里我用“状态机 + 事件总线”来实现这一层,后面会详细拆。
2. 整体设计与技术选型:不是“编排框架”,而是“系统编织器”
2.1 两种主流架构:中心化编排与去中心化消息总线
多智能体编排,行业里基本沿着两条路线在走。
一种叫“中心化编排”,类似于 LangGraph 的典型用法:一个 orchestrator 节点决定下一步调用哪个 Agent,所有 Agent 都挂在它下面,数据流和状态流都经过中心节点。这种方式的优势是逻辑清晰、容易调试,适合流程相对固定、状态流转明确的场景;劣势是中心节点容易成为瓶颈,而且每个 Agent 都要适配它的调用协议。
另一种是“去中心化消息总线”,类似 Actor 模型:每个 Agent 是独立的消息消费者,通过一个共享的消息队列发布和订阅事件,彼此之间不直接耦合。优势是扩展性好、组件独立性强,适合动态拓扑和大量并行任务;劣势是全局状态难以追踪,流程一长就容易变成“蜘蛛网”。
OpenRig 没有刻意二选一。我采用的是“中心化状态机 + 去中心化消息传递”的混合结构:业务流程和状态流转由编排器统一控制,但 Agent 之间的数据交换通过 Redis Streams 这类消息总线异步完成。这样既有中心控制的确定性,又有事件驱动的松耦合。
2.2 为什么用事件总线 + 状态机 + 持久化存储三层结构
先看一个典型的 OpenRig 执行流:
用户提交一个“分析这个仓库并生成发布说明”的任务后,编排器创建一条流水线记录,状态置为pending。编排器把任务事件发布到 Redis Stream,消息里带着 task_id 和阶段名称。负责需求解读的 Agent Worker 消费到事件,执行分析,把结果写回 Redis,同时向结果 Stream 发布一条analysis_completed事件。编排器订阅到该事件,更新状态机为analysis_done,再发布下一个阶段的事件给代码审查 Agent。如此接力,直到最终完成。
这里面三个层次各司其职:
事件总线(Redis Streams)负责“搬运消息”,它让生产者不关心消费者是谁,让消费者不关心消息从哪来,天然支持解耦、异步和并行;状态机(编排器内实现)负责“决定下一步做什么”,它从事件总线里获取事实,按照预定义的状态转移规则推进流程;持久化存储(Redis + PostgreSQL)负责“记住所有东西”,消息、状态、快照、结果全部落盘,服务重启后可以基于历史恢复上下文。
为什么用这个组合而不是直接“代码里 if-else 一步步调”?因为真实的协作流程一定会变。需求 Agent 多了一个分支判断、测试 Agent 需要等待一个异步回调、某个 Agent 想并行跑两个实例……只要用消息驱动 + 状态机,这些变化都只是增减事件和转移规则,不用改 Agent 之间的调用关系。
2.3 技术栈选型理由
我最终选定的关键组件如下:
Redis(含 Streams 与持久化配置):承担事件总线、任务队列、部分运行状态缓存。选择 Redis 不是因为它块,而是因为它同时提供了 Stream 数据结构、阻塞读取命令、TTL 过期机制和 RDB/AOF 持久化,非常适合做“轻量消息中间件 + 状态暂存区”。用 Redis Streams 而不是简单地在 K8s 里起 Kafka,是因为这个项目的消息量远没到需要 Kafka 的水平,而 Redis 部署运维成本低得多,团队已有成熟运维经验。
PostgreSQL:存放最终需要长期保留的业务数据——任务流水线、Agent 执行轨迹、交付文档元数据。Redis 里的消息和缓存可以被清理,但业务核心数据不能丢,所以一热一冷,各有分工。
FastAPI:编排器的 HTTP 入口。它天然支持异步,配合 Redis 的异步客户端可以轻松支撑大量并发请求,而且生态里自带 OpenAPI 文档,对团队协作友好。
LangGraph 借鉴但未全局依赖:我参考了 LangGraph 的图式状态机设计思想,但没有把所有 Agent 都硬塞进 LangGraph 框架。原因是项目里有不少 Agent 是老旧 Python 服务甚至外部 HTTP 接口,统一改造为 LangGraph 节点成本太高。OpenRig 只在编排器核心用了一个轻量状态机实现,外部 Agent 全部通过消息总线接入。
提示:如果你是从零开始的绿地项目,团队又愿意统一技术栈,直接全线用 LangGraph 这类框架会更省事。我这次的场景是存量系统混合改造,所以选择了自研编排器 + 消息总线的路线。技术选型没有绝对最优,只有结合上下文才知道哪个最合适。
3. 持久化的核心机制:让 Agent 协作不再“断电失忆”
3.1 持久化到底要存什么:会话、状态、消息、任务快照
说“持久化”之前,得先拆清楚到底有哪些东西需要落盘。我在 OpenRig 里把持久化对象分成了四类。
第一类是会话元数据。一次人工交互活动(比如用户反复调校一个 Agent 的结果)的所有属性:会话 ID、关联的用户、创建时间、Agent 清单、当前所在的流程阶段。没有这份数据,前端刷新一次页面,后端就不知道这个对话属于哪个流水线。
第二类是运行状态。每个任务当前处于什么状态:pending、running、awaiting_agent、completed、failed。状态机要能恢复,就必须把这些状态从内存搬到持久化存储,否则编排器重启后所有任务都会变回初始态。
第三类是消息与事件流。Agent 之间传递的每一条消息、编排器发布的每一个事件,都是可回溯的过程证据。这些数据既用于恢复执行现场,也用于事后审计和问题定位。
第四类是任务快照。某个 Agent 在某个阶段的输入输出,尤其是中间产物的引用(比如生成的文件路径、临时结果的 Key)。有了快照,失败重试时不需要整个流水线重新跑,可以从最近的快照点继续。
3.2 Redis 持久化机制详解与实践配置
Redis 本身是一个内存数据库,但 OpenRig 把它当作消息总线和运行状态存储,所以必须开启真正的数据落盘。Redis 提供两种持久化机制:RDB 快照和 AOF 追加日志。
RDB 的做法是:按照配置的频率(比如每 60 秒如果有 1000 次写操作)生成一份当前数据集的二进制快照文件。RDB 的优势是恢复速度快、文件紧凑、对 Redis 性能影响小;劣势是快照间隔期间的数据会丢失,最坏情况可能丢失上一次快照之后的所有写入。
AOF 的做法是:每次写操作都追加到一个日志文件里,重启时通过重放日志恢复数据。AOF 默认的 everysec 策略每秒落盘一次,最多丢一秒数据,但文件会持续变大,需要定期重写(rewrite)压缩。
OpenRig 在 Redis 持久化上用的是 RDB + AOF 混合模式。关键配置如下:
# redis.conf 片段 appendonly yes appendfilename "appendonly.aof" appendfsync everysec no-appendfsync-on-rewrite yes # 混合持久化:AOF 重写时生成 RDB 前缀,后续增量用 AOF 追加 aof-use-rdb-preamble yes # RDB 触发频率:至少 100 次写操作且距离上次快照 60 秒 save 60 1000 # 关联类命令直接禁掉,避免误操作清掉任务队列 rename-command FLUSHALL "" rename-command FLUSHDB ""混合模式的好处,一句话概括就是“兼顾恢复速度与数据安全性”。AOF 重写时会直接把当前数据集生成一个 RDB 格式的前缀,重启时加载 RDB 部分一步到位,之后再有新写入就增量重放 AOF 部分。
注意:
appendfsync everysec是稳妥的折中。如果业务允许丢一点数据且追求极致性能,可以改用no;如果业务绝不允许丢数据,可以改成always,但写性能会明显下降。从我的实测来看,OpenRig 这种消息编排场景用 everysec 已经完全够了,一万条消息的瞬时峰值也没造成明显瓶颈。
3.3 RDB 与 AOF 的取舍,以及“一热一冷”双存储策略
很多初学者会问:是不是把 Redis 持久化打开就万事大吉?我的答案是:不可能。Redis 的持久化只是防止 Redis 自身重启丢数据,它替代不了业务数据库,也救不了误删和逻辑错误。OpenRig 采取的是“一热一冷”策略。
Redis 作为“热存储”,保存的是短期有效、要求高吞吐的数据:事件流、任务状态、待处理消息。这些数据允许被定期清理(比如 Streams 设置MAXLEN裁剪长度),因为它们在流水线跑完后已经转化为 PostgreSQL 里更结构化的业务记录。
PostgreSQL 作为“冷存储”,保存的是长期有效、需要支持检索的业务数据:每个任务的完整生命周期记录、Agent 执行结论、交付文档内容。这里的数据是不可再生的,Redis 里丢了可以靠重放重建,PostgreSQL 里丢了就是事故。
所以我在编排器代码里会有明确的“双写”逻辑:Agent 返回结果后,内存缓存更新 + Redis 状态更新 + PostgreSQL 归档。Redis 服务于运行期,PostgreSQL 服务于回溯期。
3.4 消息总线怎么实现持久化和重放
Redis Streams 是这个项目消息总线的核心。每个 Stream 可以看作一个只追加日志,消费者用游标(ID)来标记自己消费到了哪里。不同于普通 List,Streams 天然支持故障恢复和消费组,非常适合编排场景。
整个 OpenRig 的消息总线分了三个 Stream:openrig:events存放编排事件,openrig:tasks存放发送给 Agent 的任务,openrig:results存放 Agent 返回的结果。
关键的点是消费组机制。每个 Agent Worker 都有自己独立的消费组,多个 Worker 实例可以挂在同一消费组下,Redis 会自动把消息分发给不同实例,实现横向扩展。如果某个 Worker 崩溃且没有确认消费(XACK),Redis 会把未确认消息放回 Pending Entry List,另一个实例可以读取并继续处理。这给了系统天然的“至少一次投递”语义。
为了保证任务不丢,OpenRig 还做了一个“重放”工具:扫描openrig:results里缺失的 task_id,找到原始任务消息的 ID,从openrig:tasks中按 ID 范围重新读取,重新发布给 Worker。这个工具平时不用,但在排查消息丢失或编排器 bug 导致状态不一致时,是救命稻草。
4. 实操:从零搭建 OpenRig 编排系统(可直接抄作业)
4.1 基础设施:Docker Compose 快速拉起
不想在一堆细节上浪费时间,直接上编排套件。先搭建基础设施,一条命令拉起 Redis 和 PostgreSQL,本地验证阶段也可以不开持久化,但生产环境务必配置。
# docker-compose.yml services: redis: image: redis:7.2 container_name: openrig-redis command: ["redis-server", "/usr/local/etc/redis/redis.conf"] volumes: - ./redis.conf:/usr/local/etc/redis/redis.conf - redis-data:/data ports: - "6379:6379" postgres: image: postgres:16 container_name: openrig-postgres environment: POSTGRES_USER: openrig POSTGRES_PASSWORD: openrig_secret POSTGRES_DB: openrig volumes: - pg-data:/var/lib/postgresql/data ports: - "5432:5432" volumes: redis-data: pg-data:Redis 的配置文件在上一节已经给出,放在./redis.conf即可。PostgreSQL 内部自动启用 WAL,无需额外配置就支持崩溃恢复,这也是我选择它的一个原因——冷存储的持久化基本“白嫖”。
4.2 定义 Agent Worker
接着是 Agent Worker。每一个 Agent 独立进程,通过 Redis 消费任务。下面是一个最小可用的需求分析 Agent 示例:
# agent_demand_analysis.py import asyncio import json from redis.asyncio import Redis from openai import AsyncOpenAI REDIS_DSN = "redis://localhost:6379/0" STREAM_TASKS = "openrig:tasks" STREAM_RESULTS = "openrig:results" GROUP = "demand_analysis_workers" # 同一类 Agent 的消费组 async def handle_task(redis: Redis, task: dict): prompt = f"你是需求分析师,请对以下需求进行拆解:\n{task['payload']}" # 这里可以换成任何一个模型服务,不影响编排逻辑 client = AsyncOpenAI() resp = await client.chat.completions.create( model="gpt-4o-mini", messages=[{"role": "user", "content": prompt}], temperature=0.2, ) result = {"task_id": task["task_id"], "stage": "demand_analysis", "result": resp.choices[0].message.content} await redis.xadd(STREAM_RESULTS, {"data": json.dumps(result)}) async def worker(): redis = Redis.from_url(REDIS_DSN) # 如果消费组不存在就创建,游标从 $ 开始代表只消费新消息,这里用 0 便于重放 try: await redis.xgroup_create(STREAM_TASKS, GROUP, id="0", mkstream=True) except Exception: pass while True: # 阻塞读任务流,最多等待 5 秒 entries = await redis.xreadgroup(GROUP, "worker-1", {STREAM_TASKS: ">"}, count=10, block=5000) for stream, messages in entries: for msg_id, data in messages: task = json.loads(data["data"]) try: await handle_task(redis, task) await redis.xack(STREAM_TASKS, GROUP, msg_id) except Exception as e: # 失败消息不 ack,留在 Pending 里等待重试 print(f"task {task['task_id']} failed: {e}") await asyncio.sleep(0.1)这段代码里有三个运维关键点。第一,消费组创建时id="0",意味着从头开始读,这方便调试重放;生产环境如果是新组,可以考虑id="$"只消费新消息,具体看业务需求。第二,失败的消息不 ack,Redis 会把它标记为 Pending,方便后续定位;但如果代码一直崩溃,同一个消息会反复被同一个 Worker 读到,需要结合死信处理来规避。第三,一个 Worker 实例处理完所有消息后主动sleep(0.1),避免空转。
4.3 实现编排器核心逻辑
编排器是 OpenRig 的“大脑”,决定了整场接力赛的走向。我实现了一个轻量状态机,核心配置全部用字典描述,这样做的好处是流程调整只需改配置、不用改代码。
# orchestrator.py import asyncio import json from redis.asyncio import Redis STAGES = { "demand_analysis": { "next": ["code_review"], "required_agent_group": "demand_analysis_workers", "on_complete": "demand_analysis_done", }, "code_review": { "next": ["security_scan"], "required_agent_group": "code_review_workers", "on_complete": "code_review_done", }, "security_scan": { "next": ["doc_generation"], "required_agent_group": "security_scan_workers", "on_complete": "security_scan_done", }, "doc_generation": { "next": [], "required_agent_group": "doc_generation_workers", "on_complete": "doc_generation_done", }, } class OpenRigOrchestrator: def __init__(self, redis: Redis): self.redis = redis self.state_key = "openrig:task_states" async def create_pipeline(self, task_id: str, payload: dict): state = {"task_id": task_id, "stage": "demand_analysis", "status": "pending", "payload": payload} await self.redis.hset(self.state_key, task_id, json.dumps(state)) await self.redis.xadd("openrig:tasks", {"data": json.dumps({ "task_id": task_id, "stage": "demand_analysis", "payload": payload })}) return task_id async def on_event(self, event: dict): task_id = event["task_id"] raw = await self.redis.hget(self.state_key, task_id) if not raw: return state = json.loads(raw) current_stage = stage_config = STAGES[state["stage"]] # 只处理当前阶段事件 for next_stage in stage_config["next"]: await self.redis.xadd("openrig:tasks", {"data": json.dumps({ "task_id": task_id, "stage": next_stage, "payload": state["payload"] })}) state["stage"] = next_stage if stage_config["next"] else "done" state["status"] = "done" if not stage_config["next"] else "running" await self.redis.hset(self.state_key, task_id, json.dumps(state))这里有一个值得展开的设计取舍:为什么中间产物不直接塞进 state?我一开始是把所有阶段结果都放进一个哈希字段,结果一个稍大点的分析报告就导致 Redis 哈希读取传回几百 KB JSON,性能和可读性都很差。后来改成只存任务快照的引用地址(比如结果的 Redis Key 或 PostgreSQL 主键),真正的数据放在结果表里。这样状态结构保持轻量,回溯时按引用拉取就行了。
4.4 关键的一步:接入持久化存储
编排器跑通之后,要正式接入持久化。我按四个快照点做了落盘策略。
任务创建时:写 PostgreSQL 流水线主记录,同时写 Redis 状态哈希。Agent 进入每个阶段前:发布事件到 Redis 事件流,用于回溯,同时把阶段变更写 PostgreSQL 的 stage_log 表。Agent 返回每个阶段结果后:结果写入 PostgreSQL 的 agent_results 表,Redis 只保留最近 N 条缓存,用EXPIRE设置 24 小时过期。流水线完成时:把最终结果归档到交付物表,Redis 任务状态清掉。
# persistence.py import asyncpg class PostgresStore: async def save_pipeline(self, conn, task_id: str, payload: dict): await conn.execute( "INSERT INTO pipelines (task_id, status, payload, created_at) VALUES ($1, $2, $3, NOW())", task_id, "pending", json.dumps(payload), ) async def append_stage_log(self, conn, task_id: str, stage: str, status: str): await conn.execute( "INSERT INTO stage_log (task_id, stage, status, created_at) VALUES ($1, $2, $3, NOW())", task_id, stage, status, )为什么要在 Redis 和 PostgreSQL 之间做“双写”?我当时的想法是:Redis 负责运行时快速读取和事件流转,PostgreSQL 负责审计查询和故障恢复。比如用户在前端查看一个进行中的任务,直接读 Redis 哈希,毫秒级返回;要追溯某个 Agent 几次失败重试的具体轨迹,依靠 PostgreSQL 的 stage_log 做结构化查询。两层存储职责分明,反而比一个“全能库”更清晰。
4.5 端到端验证与压测观察
把整个链路搭好之后,我写了一个简单的验证脚本,模拟 200 个任务并发提交,观察流水线是否按预期完成,同时重启编排器验证恢复能力。
关键观测结果:
- 200 个任务并发提交,50 个需求分析 Worker 实例同时消费,Redis 单个 Stream 的写入 P99 稳定在 8ms 左右,没有触发明显的热点问题。
- 重启编排器后,正在进行的任务状态从 Redis 哈希恢复,未完成阶段的任务事件从 Streams 的 Pending 队列里重新读取并重放,实现了“从上次断点继续跑”。
- AOF everysec 策略下,模拟 Redis 进程
kill -9,重启后丢失的数据只有进程被强杀前最后一秒内的一小段事件流,其余状态完整恢复。
验证脚本的核心逻辑很简单:提交任务后轮询 Redis 哈希状态,直到全部显示 done,同时人工每隔几秒观察 PostgreSQL 的 stage_log 表,确认每个阶段都有记录。
5. 常见问题与排查技巧实录
5.1 状态不一致:编排器重启后 Agent “失忆”
我在联调阶段遇到的最典型问题:编排器重启后,通过hgetall openrig:task_states能查到任务状态,但 Agent 端却像没接到任务一样长时间不响应。后来定位发现,原因是编排器重启前发布了阶段任务到 Redis Stream,但 Agent Worker 在消费后还没来得及 ack 时就断了连接。消息回到 Pending 状态,Worker 恢复后确实能继续读,但我的 Worker 代码里xreadgroup的游标用了>,只读新消息,Pending 消息根本不会被读到。
排查思路:先从 Redis 里检查消息消费组的情况,用XPENDING看 Pending 消息列表和消费者归属,确认消息是“没人处理”还是“处理中”。修正方法有两种:一种是把 Worker 改成先读 Pending,再读新消息,但要注意避免同一个消息被同一个 Worker 重复处理;另一种是加一个定时巡检任务,把超时未 ack 的消息重新派发。最终我用的是巡检派发方式,代码逻辑更简单,也不会在读取路径上引入双语义。
提示:在写事件驱动的编排系统时,一定要把“至少一次投递”作为默认前提,不要预设“每条消息只会消费一次”。消息重复、消息乱序不是 bug,是分布式系统的基本属性,所有业务逻辑都要按幂等来设计。
5.2 消息重复消费导致的重复执行
有一次我发现测试环境里的某个 Agent 执行了两次,交付文档莫名其妙追加了两遍。排查过程是这样的:任务消息在 Redis Stream 中被 Worker A 消费,但 Worker A 在执行外部 API 调用时超时,消息被重新派发给 Worker B,Worker B 成功执行并 ack。此时 Worker A 其实也完成了调用,只是返回结果时网络中断,于是它把结果再次发布到结果流,导致同一阶段出现两份结果。
解决思路是“幂等字段 + 去重表”。我在每个阶段结果里带上task_id + stage_name作为唯一键,在 PostgreSQL 里加了一个agent_results_unique约束,重复插入直接报错,再由编排器忽略失败。更彻底一点的做法是在业务逻辑里支持幂等重放,Agent 执行前检查该阶段是否已经有成功结果,有就直接返回,无才真正执行。
再造一个系统时,我建议在设计 Agent 接口的第一天就把request_id参数设计进去,让 Agent 支持“同一个请求 ID 重复提交返回同一结果”。这是加一个参数的成本,却能在后面省下无数排查重复执行的时间。
5.3 并发场景下的幂等设计与死锁规避
多 Agent 并行是编排器的常态能力,但它引出了一个很接地气的坑:多个 Worker 同时操作同一个任务状态时的竞态问题。我的编排器用 Redis 哈希保存任务状态,如果两个 Worker 同时读到同一个任务的旧状态、各自推进到不同阶段再写回,后写的人会直接覆盖前写的人,导致任务状态直接错乱。
规避手段是“乐观锁 + 版本号”。在任务状态的哈希里加一个version字段,每次读取时拿到版本号,写回时使用 Lua 脚本原子比较版本号再更新。如果版本号不匹配,说明状态已经被别的线程改过,当前线程需要重新拉取最新状态再执行后续逻辑。
-- update_state.lua local state = redis.call('HGET', KEYS[1], ARGV[1]) local parsed = cjson.decode(state) if parsed.version ~= tonumber(ARGV[2]) then return 0 end parsed.stage = ARGV[3] parsed.status = ARGV[4] parsed.version = parsed.version + 1 redis.call('HSET', KEYS[1], ARGV[1], cjson.encode(parsed)) return 1很多人在本地单线程测试时永远发现不了这类问题,一上并发就炸。提前用 Lua 脚本保证原子性,是在多智能体编排这个场景里必须遵守的纪律。
5.4 排查工具与方法清单:Redis 命令行三板斧
遇到问题不要盲猜,先把现场取证做扎实。下面是我高频使用的 Redis 排查命令:
XINFO STREAM openrig:tasks:查看 Stream 的整体长度、消费组数量、每个组的 Pending 数量、消费者数量,一秒钟判断消息是堆积还是枯竭。XPENDING openrig:tasks demand_analysis_workers:查看待确认消息,能看到哪个消费者拿走了消息、已经停留了多久,快速判断哪个 Worker 卡住或崩溃。XACK openrig:tasks demand_analysis_workers <msg-id>:确认某条消息处理完成,如果测试时手动改坏了状态,可以谨慎使用它来清理 Pending。HGETALL openrig:task_states:一次性查看所有任务当前状态,配合 grep 定位异常状态的任务。SLOWLOG GET 10:检查有没有慢命令,比如大键读取、KEYS命令误用等,提前发现性能隐患。
有了这几条命令,大部分编排问题都能在三步内定位:先确认消息是否发出,再确认消息是否被消费,最后确认状态是否更新。排查思路本身比单条命令更重要,但工具是思路的落点。
在整个 OpenRig 的落地过程中,我最大的感受是:多智能体编排的核心难点从来不在单个 Agent 的智能程度,而在系统层面的协作纪律与故障恢复能力。每次踩坑复盘后我都发现,问题几乎都出在“状态没存好”或“消息没管好”这两个基础环节,而不是模型选型不够聪明。如果你也正在尝试类似架构,建议先把持久化和消息确认机制做扎实,再考虑更花哨的编排策略。最后再分享一个小技巧:上线前一定要做一次“编排器进程被强杀 + Redis 容器重启”的演练,很多你以为“应该没问题”的环节,都会在这个演练里现出原形。