ARTICLE DETAIL

资讯详情

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

多智能体协作中的触达层设计:Agent-Reach 路由与消息派发实践

多智能体协作中的触达层设计:Agent-Reach 路由与消息派发实践 我最近在整理多智能体协作项目的时候遇到一个非常实际的问题六个 Agent 互相调用彼此之间既有 HTTP 接口又有 WebSocket 长连接还有一些走了消息队列看起来每个都通了但新加一个 Agent 之后至少有四个地方要改。后来我把所有“谁能被找到、如何被找到、请求怎么过去”的逻辑抽出来做成了一个叫 Agent-Reach 的独立服务。它是多智能体系统里的“触达层”负责把请求按照能力语义路由给正确的 Agent同时统一处理超时、重试和失联问题。这篇博文我会把这个项目的设计思路、实现细节、以及我在实测中踩过的一些坑写清楚。特别是如果你的系统里已经有很多工具型 Agent调用关系开始变成一团乱麻那这部分经验应该能帮你省不少事。1. 为什么需要一层 Agent-Reach直接互调的痛点拆解1.1 直接互调带来的四个坑先说结论Agent 之间不是不能直接调而是当数量超过三四个以后维护成本会指数上升。我最初的设计里每个 Agent 都保存了一份“通讯录”里面有其他 Agent 的地址、端口、鉴权 token、参数格式这就像是公司里每个人手里都攥着一堆其他同事的私人手机号看起来很方便但一旦有人换号通知所有人就是一场灾难。第一个坑是网状依赖。A 调用 B、B 调用 C、C 又回过来调 A这个调用图很快就变成一团乱麻。新来的人想搞清楚谁依赖谁基本只能靠猜。更麻烦的是任何一个 Agent 上线或者下线都会牵连一批上游调用方改完这个漏那个。第二个坑是接口协议五花八门。有的 Agent 暴露的是 REST 接口有的是 gRPC还有的是基于 WebSocket 的流式接口。上层如果要做一个统一的调用入口就得为每一种协议写适配层。一开始只有两个 Agent 的时候写两个 adapter 还能忍到后面每加一个 Agent 都要写一套粘合代码纯属浪费精力。第三个坑是没有全局调度策略。每个 Agent 自己决定调用谁、怎么调用结果就是系统整体行为完全不可预测。比如某个 Agent 本身只应该处理低优先级的异步任务但上游一有请求就直接打过来高峰期直接把它的资源吃满反而影响了更重要的任务。第四个坑最容易被忽略失败处理落在各个 Agent 里重试策略各写各的。有的 Agent 失败后连试三次每次等 1 秒有的重试一次就放弃还有的根本不重试直接丢数据。最后排查问题的时候你也不知道这个失败到底是网络抖动、业务异常还是对方 Agent 压根没上线只能靠日志猜。1.2 Agent-Reach 解决什么核心问题Agent-Reach 做的事情可以概括成一句把“找到一个合适的 Agent 并把任务可靠地送过去”这件事从各个 Agent 手里收回来集中到一个地方统一处理。它不是业务编排引擎不决定任务应该拆成几步也不管 Agent 内部怎么做决策。它只关心三件事可发现性你有哪些 Agent 在线、各自有什么能力、现在能不能接活。可路由性给定一个请求应该派给哪个 Agent 才最合适。可靠送达消息发出去了Agent 处理成功还是失败超时了怎么处理重试几次之后应该放弃。我把它理解成多智能体系统的高速公路调度中心。每个 Agent 只需要做一件事到 Agent-Reach 登记自己的能力和联系方式然后等在旁边接单就行。调用方不需要知道任务最后由谁执行只需要把请求丢给 Agent-Reach并声明自己要什么样的能力。下面是直接互调和引入 Agent-Reach 之后的对比这张表我后来做技术分享时也经常用。对比项直接互调Agent-Reach 触达层新增 Agent 的影响可能有 N 个调用方需要改代码只需注册一次路由规则自动生效接口协议每个 Agent 一套协议需各自适配统一为消息格式适配层下沉到触达层失败重试每个 Agent 自定义策略混乱统一重试、退避、死信策略可观测性调用链分布在所有 Agent 日志里所有请求统一留痕方便追踪全局调度无全靠 Agent 自己决定可以按优先级、负载、健康度调度这样一对比思路就很清楚了多 Agent 系统的复杂度不应该藏在各个 Agent 内部而应该收拢成一个基础设施能力。2. 整体设计和技术选型一个可以抄作业的参考架构2.1 三个核心模块注册中心、语义路由引擎、消息派发器我在设计 Agent-Reach 时没有把功能做得特别花哨只保留了三个核心模块每个模块解决一类问题。第一个模块是注册中心Registry。所有 Agent 启动之后会向 Agent-Reach 发起注册上报自己的 agent_id、能力标签、别名、调用地址、版本号等信息。注册中心会保存这些信息并且通过心跳机制维护每个 Agent 的在线状态。Agent 挂掉之后过一段时间没有被标记为“失联”路由引擎就不会再把新任务派给它。第二个模块是语义路由引擎Router。调用方提交请求时需要附带一个能力描述或者说意图声明。路由引擎会根据这个声明去注册中心里筛选候选 Agent然后按照匹配度、健康度、当前负载、优先级进行打分选出一个最合适的接收者。这个匹配不一定是死板的精确匹配可以用关键词、别名和甚至大模型辅助打分。第三个模块是消息派发器Dispatcher。路由结果确定之后派发器会把原始请求包装成一个统一的任务消息投递到消息队列里再由消费者把任务转给目标 Agent。这里我刻意做成异步模式而不是直接同步转发理由在下一小节细说。三个模块之间的关系很简单注册中心提供数据基础路由引擎产生决策派发器负责把决策落地。三者都通过同一个内部数据模型协作我用的是 JSON 消息体字段统一为request_id、intent、capability、payload、timeout、priority。2.2 同步调用改成异步消息派发我为什么这么选早期版本里调用方请求 Agent-Reach 之后Agent-Reach 去同步调用目标 Agent然后把结果原样返回。这种做法代码最简单但有两个非常致命的问题。第一个问题是调用超时不可控。LLM 场景下的 Agent 执行任务经常不是一两秒就能完成的。比如一个 Agent 要调用外部 API 拿数据再做一轮总结耗时经常在 10 秒以上。如果用同步 HTTP 长连接网关很容易打到超时上限调用方那边也可能已经断了结果任务其实在后台还在跑最后数据就丢了。第二个问题是中间节点重启会丢消息。同步转发模式下请求是带在连接上的Agent-Reach 只要一重启所有正在转发中的请求全部断掉。对内部工具型场景来说这种可用性是不可接受的。所以我把整个链路改成了异步消息模式。调用方提交任务后Agent-Reach 立刻返回一个request_id任务进入 Redis Streams 消息队列。消费者拿到消息之后再去分发到目标 Agent。这样即使 Agent-Reach 短暂重启队列里的消息也不会丢消费者恢复后可以继续消费。当然异步会带来一个代价调用方拿不到实时结果必须通过轮询或者回调接口查结果。这个代价在实际使用中可以接受因为大多数调用方本身就是一个更上层的编排 Agent它本来就要等待多个子任务完成异步提交再统一汇总反而更自然。2.3 技术栈对比与选择为什么选了 Redis Streams任务队列这块我做过几个方案的对比最终选了Redis Streams。原因很直接部署成本低、消息持久化能力够用、自带消费者组机制适合 Agent-Reach 这种几十上百个 Agent 的中小型场景。方案优势劣势适用场景Redis Streams轻量、部署简单、支持消费者组、消息可持久化不支持复杂的消息路由和分区策略中小规模任务分发、重试、延迟队列RabbitMQ支持多种路由模式、消息确认机制完善部署和运维相对重复杂业务路由、多种交换机需求Kafka吞吐量大、分区有序、生态成熟组件多、运维成本高大规模日志、数据管道、海量消息直接 HTTP 回调无中间件不可靠、重试困难内部快速调通原型我用 Redis Streams 还有一个原因它还支持给消息设置延迟投递配合一小段 Lua 脚本就能实现延迟重试队列。重试场景非常重要Agent 处理失败、目标 Agent 暂时不可达、代码发布期间消息需要延后处理都需要延迟队列。如果这个东西还要额外引一套 RabbitMQ项目复杂度立刻上去。整个 Agent-Reach 服务本身是用 FastAPI 写的Python 3.10 的类型系统够用redis-py 对 Redis Streams 的支持也比较成熟。数据库方面Agent 注册信息存在独立的 SQLite 里方便做持久化查询在线状态和路由打分的中间结果直接放 Redis 内存里性能足够。3. 核心实现注册、路由、派发三件套3.1 注册中心的数据结构与心跳机制注册中心是整个 Agent-Reach 的数据基础。每个 Agent 注册时我定义了一个比较完整的注册模型核心字段如下from datetime import datetime from typing import List, Dict, Optional class AgentRegistration: def __init__( self, agent_id: str, name: str, capabilities: List[str], aliases: Optional[List[str]] None, address: Optional[str] None, priority: int 100, max_concurrent_tasks: int 5, version: str 1.0, metadata: Optional[Dict] None, ): self.agent_id agent_id self.name name self.capabilities capabilities self.aliases aliases or [] self.address address or redis://default:agent-reach-worker self.priority priority self.max_concurrent_tasks max_concurrent_tasks self.version version self.metadata metadata or {} self.registered_at datetime.now().isoformat() self.last_heartbeat self.registered_at self.status alive # alive / lost / busy self.current_tasks 0priority字段我用来做同能力多 Agent 的场景。比如有两个 Agent 都能做天气查询一个是从免费 API 拿数据另一个是企业内部的高精度气象接口那我会给后者配更高的优先级。max_concurrent_tasks和current_tasks用来做简单的负载控制防止任务全部砸到同一个 Agent 上。心跳机制我用的很简单Agent 启动后每 15 秒上报一次心跳注册中心每次收到心跳就刷新last_heartbeat后台线程每 30 秒扫描一次如果某个 Agent 超过 45 秒没心跳就把它标记为lost。标记为lost之后路由引擎就不会再选它同时消息派发器会把已经发过去但还没确认的任务重新进入重试队列。注册接口长这样from fastapi import FastAPI, HTTPException from pydantic import BaseModel class RegisterRequest(BaseModel): agent_id: str name: str capabilities: List[str] aliases: List[str] [] address: str priority: int 100 max_concurrent_tasks: int 5 version: str 1.0 app.post(/register) def register_agent(req: RegisterRequest): if req.agent_id in registry_store: raise HTTPException(status_code409, detailagent already exists) registration AgentRegistration( agent_idreq.agent_id, namereq.name, capabilitiesreq.capabilities, aliasesreq.aliases, addressreq.address, priorityreq.priority, max_concurrent_tasksreq.max_concurrent_tasks, versionreq.version, ) registry_store[req.agent_id] registration return {status: registered, agent_id: req.agent_id}这里我把注册地址(address)直接设计成 Redis Streams 里的 worker 队列名而不是 HTTP 地址。这样做的好处是Agent 消费任务和上报结果都走同一个队列通道消息派发器不需要关心 Agent 暴露的是 HTTP 还是 gRPC。Agent 那边只需要启动一个 worker 进程订阅自己的队列拿到消息处理完再把结果写回结果通道。3.2 语义路由引擎的实现思路注册中心解决的是“有哪些 Agent 可用”路由引擎解决的是“这个请求该给谁”。我的实现分了三个层级精确匹配、别名匹配、语义匹配。每一层命中后都会加权最后取分数最高的 Agent。from typing import List, Dict, Optional import re def compute_score(candidate: AgentRegistration, capability: str, hint: str) - float: score 0.0 if capability in candidate.capabilities: score 100.0 if hint: alias_lower hint.lower() for alias in candidate.aliases: if alias.lower() in alias_lower: score 30.0 if candidate.status ! alive: score - 200.0 if candidate.priority 0: score min(candidate.priority / 10.0, 20.0) load_factor candidate.current_tasks / max(candidate.max_concurrent_tasks, 1) score - load_factor * 15.0 return score def route_request(capability: str, hint: str , top_k: int 1) - List[Dict]: candidates list(registry_store.values()) if not candidates: return [] scored [] for agent in candidates: s compute_score(agent, capability, hint) if s 0: scored.append({agent_id: agent.agent_id, score: s}) scored.sort(keylambda x: x[score], reverseTrue) return scored[:top_k]打分规则里最核心的是capability精确匹配直接给 100 分别名命中再给 30 分优先级最高加 20 分负载最高扣 15 分。这样设计是为了让“能力是否匹配”成为第一决定因素而不是让一个负载很低的低能力 Agent 抢走任务。刚开始我觉得语义匹配可以直接上大模型后来发现其实很多场景用别名和规则就够用了。只有遇到那种“用户用自然语言描述了一个模糊需求”的时候我才会让路由引擎调用一次大模型追问或者做意图改写把模糊文本转成标准能力标签再走上面的打分逻辑。这样做的好处是大部分请求走低延迟的规则路径只有少部分模糊请求走 LLM 路径成本和延迟都能控制住。3.3 消息派发与重试机制用 Redis Streams 保证不丢消息路由决策出来之后消息派发器把请求包装成标准任务消息推送到 Redis Streams。每个 Agent 对应一个 stream流名称就是 agent_id。派发器还会维护一个pending列表用来记录哪些任务还在等待 Agent 确认。任务消息格式如下{ request_id: req-001, capability: weather.query, payload: { city: 杭州, date: 2025-01-18 }, retry_count: 0, timeout: 120, created_at: 2025-01-18T10:30:00Z }派发核心代码def dispatch_task(request: dict, target_agent_id: str): request[retry_count] 0 stream_name fagent:{target_agent_id}:tasks redis_client.xadd( stream_name, request, maxlen5000, ) pending_requests[request[request_id]] { target_agent: target_agent_id, status: dispatched, updated_at: now(), }消费者 Worker 侧的逻辑def agent_worker_loop(agent_id: str): stream_name fagent:{agent_id}:tasks while True: messages redis_client.xread({stream_name: }, count10, block5000) for stream_msg_id, msg in messages: task msg_to_dict(msg) try: result execute_agent_task(task) report_success(task, result) except Exception as exc: report_failure(task, exc)这里最关键的就是失败上报之后的处理。Agent-Reach 的主服务如果收到任务失败的消息会判断retry_count是否已经达到上限。如果没到就把消息重新放回延迟队列等待退避时间后再次投递。如果到了就把任务放进死信队列同时推送一条告警到企业微信。def handle_failure(task: Dict, error_message: str): task_id task[request_id] retry_count task.get(retry_count, 0) max_retries task.get(max_retries, 3) if retry_count max_retries: task[retry_count] retry_count 1 delay_seconds min(2 ** retry_count * 10, 300) redis_client.zadd(retry_tasks, {json.dumps(task): time.time() delay_seconds}) pending_requests[task_id][status] retrying else: redis_client.xadd(dead_letter_tasks, task) notify_alert(task, error_message) pending_requests[task_id][status] dead_letter这里我用了 Redis ZSet 做延迟队列score存的是可执行时间点。主服务每隔 5 秒扫描一次 ZSet把到期任务重新投递到目标 stream。相比直接time.sleep()这种方式不会阻塞 worker 线程也让重试节奏可控。3.4 可配置参数速查表实际运行时有几个参数非常影响系统表现我列成一个表方便对照调整。参数默认值说明HEARTBEAT_INTERVAL15 秒Agent 心跳上报频率HEARTBEAT_TIMEOUT45 秒超过该时间未心跳则标记为失联TASK_TIMEOUT120 秒单次任务最长执行时间MAX_RETRIES3 次任务失败最大重试次数BACKOFF_MULTIPLIER10 秒重试退避基数默认 10 秒、20 秒、40 秒递增DELAY_QUEUE_SCAN_INTERVAL5 秒延迟队列扫描周期MAX_PENDING_TASKS1000单个 Agent 的待确认任务上限如果业务场景是内部工具型 Agent超时时间可以拉长到 5 分钟因为很多 Agent 要调外部接口如果是纯计算型 Agent超时时间可以压缩到 15 秒避免无意义的等待。4. 从零接入三个 Agent 的实测记录4.1 我实际接入的三个 Agent 与注册配置我在本地搭了一套完整环境接入了三个典型 Agentweather-agent天气查询、sql-agent数据库查询、report-agent生成 Excel 报表。这三个 Agent 正好覆盖了短任务、中长任务和重任务三种场景。注册配置如下我摘录关键部分[ { agent_id: weather-agent, capabilities: [weather.query, weather.forecast], aliases: [天气, 气温, 降雨, 气象], address: redis://localhost:6379/0, priority: 100, max_concurrent_tasks: 10 }, { agent_id: sql-agent, capabilities: [database.query, sql.execute], aliases: [数据库, 查询, SQL], address: redis://localhost:6379/0, priority: 100, max_concurrent_tasks: 3 }, { agent_id: report-agent, capabilities: [report.generate], aliases: [报表, Excel, 汇总], address: redis://localhost:6379/0, priority: 100, max_concurrent_tasks: 2 } ]实际发送请求时调用方只需要像下面这样提交任务完全不用关心最终是哪个 Agent 在执行{ capability: weather.forecast, payload: { city: 成都, days: 3 }, timeout: 60 }Agent-Reach 返回了一个request_id然后任务进入路由环节。路由引擎根据capability精确匹配到weather-agent检查它的负载之后投递到对应的 stream。整个过程从提交到消息入队本地延迟大概是 3 毫秒左右路由决策耗时几乎可以忽略。4.2 压测记录与结果分析我用一个简单的压测脚本模拟了 50 个并发请求混合三种任务类型。跑了三轮取平均值指标结果总请求数150成功数144成功率96%平均任务耗时41.2 秒P95 任务耗时94.7 秒Agent-Reach 本身平均转发耗时5.8 毫秒死信数量2成功的 96% 里大部分是异步队列处理带来的收益。失败的那 6 个请求中2 个是report-agent生成 Excel 时文件写入路径权限问题4 个是因为sql-agent在执行一条特别复杂的查询时后端连接主动断开触发了超时重试但最终错过了三次重试上限。这次压测让我意识到一个点sql-agent这种依赖外部数据库连接的 Agent很容易因为后端连接超时导致任务失败。如果注册中心能感知到它正在处理的任务已经超出了TASK_TIMEOUT就应该立刻标记为可疑状态把后续请求切到别的可用 Agent 上而不是继续往同一个 Agent 堆消息。4.3 现场事故复盘一次把消息跑丢的排查过程压测过程中我遇到过一次比较严重的消息丢失任务显示已经dispatched但 Agent 那边一直没收到。查了半天发现是注册时填写的address是 Redis Streams 的完整 key 名但 worker 订阅时用了一个带前缀的错误 key。最后结果是消息其实进对了 stream但消费者监听错了位置。这个问题暴露了一个设计缺陷地址配置应该由 Agent-Reach 统一生成而不是让 Agent 自己填完整地址。修改之后Agent 只需要在注册时传agent_idAgent-Reach 自动生成agent:{agent_id}:tasks作为它的消费流从根上杜绝了配置不统一的问题。另一个现场问题是重试风暴。某个 Agent 临时下线时注册中心要在 45 秒后才会把它标记为lost这 45 秒内消息仍然不断被路由过去。等到 Agent 恢复后积压的任务一下子同时触发重试直接把它的 Redis 连接池打满。后来我加了一个“下线快速探测”接口Agent 在正常退出时会主动调/offline接口注册中心立即把它标记为lost这样就绕开了心跳超时窗口。5. 必须记住的四个坑和后续扩展方向5.1 经验教训这几个坑我花了不少时间才填平第一个坑是幂等消费。Agent 处理完任务但还没来得及上报结果就被重启任务重新投递后Agent 会重复执行一次。如果这个任务里有写数据库或者发通知的逻辑就会产生重复数据。解决办法是在 Agent 端维护一张已处理任务表用request_id做唯一键重复消息直接透传旧结果。第二个坑是延迟队列的原子性。最开始我用 ZSet 做延迟队列时取出任务和重新投递这两个操作并不是原子的进程刚好崩在中间就会丢任务。后来我把取任务、投递、删除这几个动作封装成一段 Lua 脚本保证执行原子性才算彻底解决。第三个坑是路由打分里的“健康度”不能只看心跳。心跳正常不代表 Agent 没有在忙。我发生过一次sql-agent连接池被打满心跳照常上报但所有任务进去之后全部排队超时。后来我给注册中心增加了“当前任务数”和“最近失败率”两个指标路由打分时如果最近失败率大于 20%直接扣 50 分这类不健康 Agent 会自然被冷落。第四个坑是语义路由不能过度依赖大模型。有一次我把路由改成了“所有请求先让 LLM 判断一遍再派发”结果延迟从 5 毫秒暴涨到 1.5 秒而且 LLM 判断错一次后续任务全被带到错误的地方。现在我的策略是规则优先只有规则无法匹配时才调用 LLM而且 LLM 输出的能力标签还会被缓存避免重复请求。5.2 后面我打算继续加的功能方向Agent-Reach 第一版做下来整体框架算是稳了。我接下来想再加几个能力第一个是多租户隔离。目前注册中心只有一个全局命名空间所有 Agent 都在一起。如果把 Agent-Reach 给多个团队共用不同团队的 Agent 之间需要做权限和可见性隔离我计划引入namespace字段路由查询时强制带上租户上下文。第二个是路由规则热加载。现在打分规则硬编码在代码里改一个权重就要重新发布服务。我准备把规则抽成配置文件放到数据库里后台监听变更自动刷新让运维同学不用动代码就能调整路由策略。第三个是路由结果的可解释性。现在返回给调用方的只有一个agent_id别人根本不知道为什么选了这个 Agent。我想在响应里带上匹配分数和命中原因比如“精确能力匹配 100 分、负载扣 5 分、最终得分 95 分”。这样排查问题时就不需要靠猜了。我自己的体会是做这种基础设施型项目最忌讳一开始就堆功能。Agent-Reach 的核心价值只有一个让多 Agent 之间找到彼此这件事变得更简单。只要抓住这个点后面加再多扩展都是锦上添花。如果一开始就把路由、编排、记忆、工具调用全部揉在一起这个项目大概率会变成一个别人看不懂也维护不了的巨型怪物。
返回列表