2026/10/6 4:11:45

多智能体编排实战:基于Redis Streams的代理通信与任务调度

多智能体编排实战:基于Redis Streams的代理通信与任务调度 刚刚做完Agent-Reach这个小项目的时候我本来没打算写一篇复盘但后来陆续有几个朋友问我多智能体编排到底该怎么落地、消息队列该选什么、为什么代理总是出现各种诡异失联问题……这些问题其实我在做这个项目的时候基本都踩了一遍。借这个机会把Agent-Reach从设计思路到具体实现、再到坑位排查完整地梳理一遍。这篇文章里我说的每一句话都来源于实际跑通的经验不是概念堆砌读完了你可以直接照着自己的需求去改造。Agent-Reach这个名字拆分来看就很有意思。Agent指智能体、执行单元Reach指覆盖范围、触达能力。合在一起它想解决的问题非常直观让一个个相互独立的AI代理能够被统一调度、按需触达把原本单个代理的能力边界向外扩展。简单说它是一个面向多AI代理的分布式编排框架你可以把它理解成一个“代理通信枢纽”加“任务路由大脑”。我构建它的核心目标不是做一个软件框架而是解决三个极其现实的痛点代理各自为战、任务难以跨单元协作、运行过程黑盒不可观测。这三件事不解决再强大的模型能力也只是停留在单次对话里根本无法形成系统性生产力。这个项目适用的场景主要是那些已经开始搭建AI应用、准备从“单Agent演示”走向“多Agent协作”的团队也包括想在公司内部落地AI自动化流水线的个人开发者和技术负责人。如果现在的你正在被以下问题困扰任务拆分不干净、多个代理抢着干活、上下文在不同代理之间丢失、故障发生之后不知道找谁复盘那么Agent-Reach的这套设计就非常适合你参考哪怕你不直接把代码拿过去用架构思路和避坑方法也有很高的复用价值。搞清楚了要解决什么问题我心里就踏实了一大半。接下来的事情就是一个个模块地拆解、选型、写代码、试错、再修直到整条链路能稳定连续跑几个小时不崩。后面我把这个过程原原本本讲给你听。1. 整体设计与思路拆解动手写第一行代码之前我先把设计目标定得尽量清晰。一个多代理编排框架最容易掉进去的坑就是用“中心化总线”把一切都串在一起表面上逻辑简单了实际上耦合度极高代理一多中心节点就变成瓶颈而且任何一个代理的失败都可能拖垮整条链路。所以Agent-Reach一开始就确立了三个设计原则代理自治、路由去中心化、状态可追踪。代理不需要知道其他代理在哪里、什么时候被调用它们只需要跟接入层通信剩下的事情全部交给编排层处理。在“路由去中心化”这件事上我做了个折中方案。框架里有一个逻辑中心叫Coordinator协调器看名字好像是中心化的但它实际上只负责注册信息和下发路由规则并不负责搬运具体的任务数据。任务数据流走的是Redis StreamsCoordinator只做轻量级的决策和转发。这样设计的好处很直接Coordinator本身不处理业务逻辑它永远不会成为瓶颈挂了也不会导致所有代理瘫痪——代理们依然可以通过Streams里已经存在的任务继续工作只是新任务的注册和调度会受影响。这套设计里还有一个很重要的概念我管它叫“触达范围”。如果咱们把每个代理理解成一个驻扎在独立工作区的小团队它们的感知半径是有限制的只能听到自己工作区里的声音。Agent-Reach要做的就是给每个代理装上一套“扩音器和天线”让它们能听到远处的需求、把成果发回总部。这套“扩音器和天线”具体实现就是一套基于主从模式的轻量级SDK代理端只需要调用几个API就能完成注册、监听、回报状态完全不需要理解底层的Streams和路由逻辑。选型上我在几个方案里徘徊过一阵子下面把这个权衡过程也写出来我觉得这部分对大家是最有价值的。2. 核心模块与通信链路的选型逻辑2.1 为什么最终选择了Redis Streams而不是RabbitMQ或Kafka做多代理通信大家第一时间想到的肯定是消息队列。我早期试过RabbitMQ确实很成熟但有一个问题让我很别扭它天然以“队列”为单位组织消息一个队列一般只对应一种消息类型代理多了之后队列数量会爆炸式增长管理成本很高。而且每个代理需要维护自己的交换机、路由键SDK封装会变得很厚重跟“轻量接入”的初衷矛盾。Kafka呢吞吐量确实大但它的设计目标是大数据流处理引入它会给一个小型多代理编排项目带来不必要的运维复杂度光是一个Topic的分区设置和消费者组管理就够写一本手册了。Redis Streams则处在一个非常合适的生态位。它的写入延迟极低单机性能完全撑得住中小规模的代理集群消费者组Consumer Group天然支持一条Stream被多个消费者竞争消费正好对应“多个代理抢一个任务”的场景再加上Redis本身是被绝大多数团队已经使用的组件不需要额外引入运维依赖。所以我的结论是在一个百级代理规模以下、以任务调度和状态同步为核心需求的系统里Redis Streams的性价比是最高的。如果你有上万级并发、海量消息堆积的需求那直接上KafkaAgent-Reach的路子就不适用了。2.2 任务数据的格式设计与上下文隔离代理之间协作最容易翻车的地方不是把任务发出去而是在任务接力的时候把前文的上下文给丢了。我在设计任务数据格式的时候坚持了一个原则每个任务都自带完整的上下文包不允许代理在执行过程中去别的Stream里“翻旧账”。一个标准任务数据长这样{ task_id: TASK-20240521-001, source_agent: planner, target_agent: coder, payload: { type: code_generation, body: { ... } }, context: { session_id: SESSION-X, summary: 前序对话摘要, artifacts: [artifact://file/src/main.py] }, timeout: 120, priority: 5, trace: [planner] }这个格式看起来简单但每一个字段都是有讲究的。source_agent和target_agent用于路由决策context字段是上下文隔离的关键它的值是从对象存储里取出来的会话快照而不是引用ID。这样即使某个代理在处理过程中崩溃重启换个实例继续跑它依然能拿到完整的上下文不会因为进程消失就导致任务中途夭折。关于timeout和priority我用实际经验告诉你这两个字段能省掉你一半的Debug时间。没有超时控制的任务会在异常场景下无限期挂起尤其当某个代理背后的模型服务抽风时一个卡死的任务会拖住整个消费线程。优先级则用于应对突发任务插队平时可能用不上但一旦线上出问题需要立即执行回滚或诊断任务优先级的价值就瞬间体现了。2.3 代理注册与心跳机制看似简单实则暗坑很多Agent-Reach里每个代理启动后的第一件事就是向Coordinator发起注册请求把自己的能力描述、订阅事件类型、健康状态报上去。这个注册信息会被写入Redis的Hash结构中key值形如agent:{name}里面存储的是代理的元数据和最近一次心跳时间。心跳机制我一开始想得比较简单代理每隔15秒上报一次状态就算完事。后来发现这样根本不够用。因为代理和Redis之间的连接是长连接如果代理所在进程GC停顿或者网络抖动心跳上报就会延迟一旦超过阈值Coordinator就会误判它失联把发给它的任务转到备份代理上去导致重复执行。我后来做了一版优化心跳间隔设为10秒失联判定阈值设为30秒也就是允许2次心跳丢失才判定失联。这个数字不是拍脑门定的而是压测后得出的经验值。压测数据表明大多数GC停顿和次要网络抖动都会在3秒内恢复30秒的阈值能过滤掉90%以上的误报。另外一个坑是代理重启后注册信息里的IP和进程ID变了但Redis里的旧信息可能还在。如果不做版本号比对Coordinator会把任务发给一个已经死亡的进程然后因为心跳超时才慢慢感知到再重新调度浪费一整轮时间。解决办法是在注册信息里加入instance_id字段每次启动生成一个新的UUIDCoordinator看到instance_id不匹配就直接把旧记录覆盖同时把发给旧实例的任务重新入队。这一个小改动让我在生产环境中的任务失败率下降了差不多一半。3. 核心机制拆解与实操要点3.1 靠“广播定向”结合的任务路由策略路由策略是最能体现一个编排框架智力水平的地方。Agent-Reach支持两种路由模式定向路由和广播路由。定向路由就是上游任务明确指定目标代理比如规划代理分析完需求后给编码代理发送一个代码生成任务。这种模式简单直接适合链路清晰的场景也是我最早实现的版本它用到的是Redis Streams里的消息ID关联能力任务的target_agent字段会被消费端用来过滤不放给其他代理。广播路由要复杂一些。它适合那些“不知道谁最适合干这个活但肯定有人能处理”的场景。我在设计时用一个独立的Stream作为广播频道所有代理都从这个频道消费消息但每个代理在拿到消息后会先做一个“能力匹配检查”匹配算法基于代理注册时上报的能力标签。比如一个需要“Python后端开发”能力的任务只会被带有python和backend标签的代理接走其他代理直接丢弃消息不处理。这个策略在单一代理拥有多个标签的情况下会同时被多个代理消费导致同一个任务被执行两遍。解决方式是在广播消息里加入lock_id代理在正式处理前会尝试用SET NX EX命令抢占一把锁只有抢占成功的代理才继续执行失败的自动放弃。我会在后面的代码实现里展示这个细节。广播路由虽然好玩但我必须提醒一句能不用就不用。它本质上是在用消息重复换取灵活性而重复必然带来幂等性的压力。实际项目里我尽量把大部分路由设计成定向的只有在任务类型确实无法预测目标时才启用广播比如每个代理在处理完用户元数据变更后需要通知“所有可能感兴趣的模块”。3.2 上下文对象存储的必要性前面我提到每个任务都携带完整上下文这个上下文到底是放在哪里的Agent-Reach里有一个对象存储服务用Redis Hash加JSON序列化实现每个会话的数据都会以ctx:{session_id}的key存起来value里包含对话历史摘要、中间产物引用、用户偏好设置等等。刚开始我尝试把所有上下文都塞进任务的context字段里一个任务消息动辄几百KBRedis Streams处理起来倒是能撑住但网络开销非常大尤其是任务频繁在中转过程中复制时整个系统的响应时间被拖慢了将近一倍。后来我改成只把对象的引用ID放进任务里代理如果需要完整上下文再按需去对象存储里拉取。这个改动让消息体积缩小到原来的十分之一系统吞吐量上了一个台阶。这个策略也带来一个代价代理必须能访问对象存储服务。所以在部署的时候我会保证所有代理和Redis实例在同一个网络平面内避免跨区域频繁拉取上下文。如果你是多区域部署那建议在每一区域内放一个Redis副本用主从复制把上下文数据同步过去网络延迟会好很多。3.3 人工审核节点的设计机器干活人还是要盯做AI流水线最怕的不是不准确而是“单项选择题变成发散填空题”。我在Agent-Reach里专门设计了一个人工审核节点Human Review Gate。在一些关键任务类型上比如代码合并、支付流程、对外发布的营销文案任务路由不会直接把结果发给下一个代理而是先进入一个专门的Review Stream由人工审核面板消费审核通过后任务才继续往下走。这个环节听着很简单实际上很考验架构设计的细致程度。审核结果有三态通过、打回重做、直接终止。打回重做时任务会附加上审核意见然后重新投递到原代理的输入Stream里原代理看到“打回”标记就知道要拿着反馈重新生成。为了防止打回和重新生成之间形成死循环我设置了最大重做次数为3超过这个次数任务自动进入终止状态并标记为“需要人工介入”。这块经验我可以拍着胸脯说没有人工审核节点的多代理系统上线后光擦屁股就足以让你崩溃。特别是代理之间层层传递、一环错环环错的时候如果中间没有人工关卡兜底整个链路的错误会被自动放大到不可思议的规模。在关键路径上保留人工审核是一种非常值得秉持的工程态度。4. 实操过程与核心环节实现有了前面思路和选型的铺垫接下来就是动手写码的过程。这部分的代码我是跑通验证过的你可以直接参考甚至可以用最小配置把它跑起来再慢慢改造。4.1 搭建环境与初始化接入层我的开发环境是Python 3.11 FastAPI Redis 7.x代理SDK的核心逻辑用Python写其他语言可以通过HTTP接口调用Coordinator来实现同样的接入效果。环境准备好后第一个任务就是写Coordinator的启动入口它负责初始化Redis连接池、创建必要的Stream、启动心跳检查的后台任务。import redis.asyncio as aioredis from fastapi import FastAPI REDIS_URL redis://localhost:6379/0 STREAMS { task.router: task.router, event.broadcast: event.broadcast, human.review: human.review, } app FastAPI(titleAgent-Reach Coordinator) app.on_event(startup) async def startup(): app.state.redis aioredis.from_url(REDIS_URL, decode_responsesTrue) for stream in STREAMS.values(): try: await app.state.redis.xgroup_create( stream, coordinator, id0, mkstreamTrue ) except Exception: # 已存在的group会报错忽略即可 pass这短短的代码里蕴含了一个容易踩坑的细节xgroup_create的mkstreamTrue参数必须在流不存在时也能自动创建流。否则你第一次启动系统时因为消息还没产生流也不存在后续所有消费端的xreadgroup会直接报错。这一点很多教程不会专门提但实际部署时非常关键。4.2 代理SDK的核心实现代理SDK要做的事情就是让代理开发者把精力集中在业务逻辑上。每个代理启动时先注册信息然后启动一个后台循环从订阅的Stream里拿消息、跑业务、回传结果。class AgentClient: def __init__(self, agent_name: str, tags: list[str]): self.name agent_name self.tags tags self.instance_id str(uuid.uuid4()) self.redis None self.running False async def connect(self): self.redis await aioredis.from_url( REDIS_URL, decode_responsesTrue ) await self.register() async def register(self): key fagent:{self.name} await self.redis.hset( key, mapping{ tags: json.dumps(self.tags), instance_id: self.instance_id, status: idle, last_heartbeat: time.time(), }, ) # 设置TTL失联后自动清理 await self.redis.expire(key, 60) async def heartbeat_loop(self, interval10): while self.running: key fagent:{self.name} await self.redis.hset(key, last_heartbeat, time.time()) await asyncio.sleep(interval)心跳循环里有一个细节注册时的TTL我设置的是60秒而心跳上报间隔是10秒。这两个数值组合起来意味着如果在60秒内代理都没有任何心跳Coordinator就会认为它彻底失联了从而触发任务重新调度。这个60秒的TTL不是随手设置的它至少需要覆盖两个完整的心跳间隔否则会因为网络稍有波动就误删注册信息反而增加不必要的调度开销。主循环里消费任务的逻辑我建议用XREADGROUP阻塞模式async def consumer_loop(self, stream: str): group_name fgroup-{self.name} consumer_name fconsumer-{self.instance_id} while self.running: # 只取当前消费者还没处理完的消息 resp await self.redis.xreadgroup( group_name, consumer_name, {stream: }, count10, block3000, ) if not resp: continue for _, entries in resp: for msg_id, data in entries: try: ok await self.process_message(data) if ok: await self.redis.xack(stream, group_name, msg_id) except Exception as exc: logger.error(f处理消息失败 {msg_id}: {exc}) # 这里先不ack让消息留在pending里等待人工处理容器组这个设计中consumer_name要用self.instance_id拼上是为了防止代理重启后之前未确认的消息被误认成新消息。如果不做这一点代理重启再启动消费者时它会把隶属于旧实例的pending消息继续拉出来二次处理可能导致重复执行。用instance_id隔离后每个实例只消费它自己的pending消息旧实例的pending消息则会由异常检测机制单独处理整个系统的执行语义就会清晰很多。4.3 任务分发与路由器的实现Coordinator的核心路由逻辑是整个系统的心脏。它会监听task.router流读取每条消息的target_agent字段然后按策略把任务转发到对应的代理接入流中。这里有个路由的细节每个代理其实有一个专用的输入流agent:{name}.inputCoordinator只需要把任务XADD到这个流里对应的代理消费端就会收到。async def route_task(task: dict): target task[target_agent] stream fagent:{target}.input await redis.xadd(stream, {payload: json.dumps(task)})定向路由看上去非常简单但真正有价值的反而是异常回退逻辑。如果target_agent指向的代理已经失联注册表里没有它、心跳超时这个任务直接扔进流里就会被无限期搁置。我的处理方式是在路由前检查代理状态如果目标失联就退回到广播路由寻找备用能力标签匹配的代理再不行就进入人工审核流。广播分发的代码里就用到前面提到的分布式锁async def broadcast_to_matching(task: dict): target_tags task.get(required_tags, []) stream event.broadcast await redis.xadd( stream, { payload: json.dumps(task), lock_id: flock:{task[task_id]}, }, ) async def try_acquire_lock(lock_id: str, ttl: int 30) - bool: # 利用SET NX EX实现分布式锁 ok await redis.set(lock_id, 1, nxTrue, exttl) return ok is not None这个锁的TTL设置为30秒对应的是AI代理处理一个任务的平均耗时上限。如果你代理跑的模型特别慢可能需要把这个值调大但同时也得有配套的超时机制来兜底不能无限期占锁。我实际跑下来发现TTL 30秒配合5秒单次推理耗时的模型非常够用几乎不会有任务因为锁过期被另一个代理二次执行。4.4 人工审核面板和状态回写人工审核是Agent-Reach里必不可少的一环。我把它实现为一个FastAPI路由前端通过HTTP拉取需要审核的任务列表审核人点击通过或驳回后状态通过Redis的ctx:{session_id}里的状态位更新然后决定把任务继续投放到下游还是打回上游。审核界面的代码不复杂核心逻辑就是根据审核动作修改状态app.post(/review/{task_id}) async def review_task(task_id: str, action: str, comment: str ): task await redis.hget(ftask:{task_id}, data) if not task: raise HTTPException(status_code404, detail任务不存在) if action approve: # 继续投递给下游代理 await redis.xadd(fagent:{task[next]}.input, {payload: task}) await redis.hset(ftask:{task_id}, status, approved) elif action reject: retry_count int(await redis.hget(ftask:{task_id}, retry_count) or 0) if retry_count 3: await redis.hset(ftask:{task_id}, status, escalated) else: await redis.hset(ftask:{task_id}, status, rejected) await redis.xadd( fagent:{task[source]}.input, {payload: json.dumps({**task, review_comment: comment})}, ) await redis.hincrby(ftask:{task_id}, retry_count, 1) return {status: ok}这段代码里有一个“升级”的动作也就是当审核打回超过3次时任务状态变成escalated。这些任务通常意味着代理在处理该需求时存在系统性不足纯靠自动反馈已经解决不了需要人类去看流程本身而不是单独看这一次任务。我把所有escalated任务单独做成了一个仪表盘列表方便集中排查这个习惯帮我发现过好几个代理因为标签匹配规则错误导致的“怪癖行为”。5. 常见问题与排查技巧实录代码能跑通只是第一关真正磨人的是各种运行时诡异现象。这里我把Agent-Reach跑起来后遇到的高频问题按发生频次排个序每一个都附上排查思路和我的解决办法这些绝对是压箱底的经验。5.1 代理彻底失联但任务还在排队第一次遇到这个问题时后台日志没有任何报错Coordinator的注册表里代理状态还是idle但任务发过去之后无人消费。排查后发现原来代理进程因为内存溢出被系统杀掉了但Redis里的注册信息和心跳时间还停留在一个有效的状态中。因为心跳是Agent主动上报进程死亡后心跳自然停止但Coordinator需要等待TTL过期才能感知。这个问题的核心教训是注册信息的TTL不能太长。我一开始为了减少注册抖动设置的TTL太长导致失联发现被严重后置。调整策略后我把心跳间隔控制在10秒、TTL控制在60秒同时把Coordinator的心跳巡检后台任务改成每15秒扫描一次所有注册代理发现心跳时间戳跟当前时间差超过45秒就直接标记失联不再等TTL自然过期。5.2 任务被多个代理重复消费前文提到过广播路由配锁的解决方案但在实际测试中我发现锁TTL设置不合理时依然会出现两个代理同时处理同一个场景。后来我从日志中发现问题出在代理的业务处理时间大于锁TTL锁过期后第二个代理抢到了锁就开始了重复执行。解决方式有两条路。一是为不同类型的任务动态设置锁TTL二是要求任务处理器支持幂等。两者在我眼里必须同时做动态TTL解决并发竞争幂等设计解决系统在异常状态下的最终一致性。幂等实现最简单的方式是在每个代理的输入段处理一个processed_tasks集合用Redis的SISMEMBER判断任务ID是否已经处理过。虽然理论上会有额外开销但对于中小型系统这点开销远小于处理重复任务带来的混乱。5.3 上下文丢失导致任务错乱Agent-Reach有一次在生产环境测试中A代理的产出作为B代理的输入进行处理结果B的回答明显缺上下文好像A刚刚生成的内容根本没传过来。排查链路时发现任务消息确实携带了A产物的引用IDB代理也调用了对象存储去拉取但拉取的时候对象存储已经因为过期把数据清理了。问题根源是我给上下文对象设置的TTL太短只给了15分钟而A到B之间的任务在人工审核节点那里堵了20分钟。等B去取数据时TTL刚好到期已经被自动删除了。修复措施是对象存储的TTL不能是一个拍脑袋的固定值必须动态计算至少是“这个会话所有任务的最大允许处理时间 额外安全余量”。我后来统一改成所有上下文对象TTL默认24小时并且任务完成时主动删除比单纯依赖过期清理更可靠。5.4 代理之间的死循环调用这是多代理系统最让人头大的问题。A任务调用BB处理完回传CC又触发需要A来处理的新任务然后A再调用B……如果链路没有设置终止条件整个系统就会一直空转白白消耗模型API的调用费用。我在测试中还真跑出过一次这样失控的循环后来灵光一闪查了一下API账单肉疼了好久。死循环的解法主要靠两点任务最大跳数Max Hops和会话级终止开关。任务数据里有一个trace字段每经过一个代理就追加一行在路由前检查如果len(trace)超过某个阈值我设置为12任务自动丢弃并写入异常日志。会话级终止开关则是在ctx里加一个status字段一旦人工审核或监控系统确认这个会话出了问题可以直接把status设置为aborted所有代理每次拿到任务先检查会话状态遇到aborted直接拒单。5.5 排查工具和调试技巧Agent-Reach最外围的诊断手段就是盯日志。但等到代理多了、日志打满整个屏幕后再去人肉翻效率太低。我后来为SDK在输出日志时增加了一些约定格式的元数据比如task_id、agent_name、session_id、trace_depth然后利用日志聚合工具单独按task_id抽取整条链路的执行轨迹。这不仅让我发现过几次具体误调度原因也给优化链路性能提供了数据支撑比如可以直接看出每个代理处理节点花了多少时间哪一环是瓶颈。如果你也想在自己的项目上排查类似问题我有几条很实用的通用建议每条关键路径上的日志必须包含同一个关联ID没有关联ID的日志再详细也无法串联分析。每个代理启动时打印自己的注册信息和订阅的流排查“消息进不来”的问题时先确认注册信息还在不在。用Redis自带的XINFO STREAM和XINFO GROUPS看消费者组的pending消息数量一旦pending数持续增长不同消费者之间必然有处理不平衡别急着改代码先看数据。尽量让代理的输入输出都是可序列化的JSON哪怕NLP模型返回的是文本外壳上也包一层JSON后续做审计和分析会方便太多。6. 实际操作中我最后想多提一嘴的体会Agent-Reach从构思到跑通花了不少没必要的弯路。最深的体会是多代理编排系统里面“编排”这个词的重点不在AI模型的思路而在工程化的节点关系、编排结构和异常处理。模型只管生成内容生成之后内容往哪走、什么情况下该重新生成、什么情况下要人工介入这些才是系统工程要解决的问题也是系统最终能不能稳定运行的核心。在做这个项目之前我一度很迷信“每个代理都有独立思考能力”这句话真正落地后发现代理的独立性最珍贵的用途不是让它们自由发挥而是让它们能被独立扩展、独立替换、独立降级。这就像你招了一群能力很强的人如果大家各自为政项目一定乱成一锅粥但如果有一个清晰的组织架构、任务分配和反馈机制每个人才能高效地协作。另外分布式锁、幂等设计、超时控制、人工审核富余量这些东西看起来都是老生常谈的工程标配但放在AI代理场景里确实有一些新的坑位需要重新思考。AI代理不仅会有常规的系统错误还会出现“看似合理却完全错误”的模型输出这是传统并行系统里极少遇到的因此验证和兜底机制要更细致。最后再分享一个小技巧。如果你要在自己的项目里推广多代理架构一开始别追求全面铺开把业务链条中一个固定、高频、结果相对容易验证的场景单独抽出来先把一条链路完整跑通。我当初就是从“产品需求摘要自动生成”这个单链路开始做起的稳了之后才逐步扩展。这样你调试难度低获得的反馈也明确整个系统的架构也能在迭代中逐渐成熟。希望Agent-Reach的这些经验能帮你在搭AI代理协同系统的时候少踩几个坑多跑通几条真正的业务链路。