☰
Agent-Reach:多Agent协作的消息路由与注册中心设计实战
2026/10/7 11:54:53 网站建设 项目流程

最近手头跑着调度、搜索、摘要、代码生成好几个Agent,想让它们互相配合干点正事,结果发现最难的居然不是模型本身,而是让这些Agent能互相找到对方、把话递过去。刚开始我直接在各Agent之间互相写调用,两个月后代码乱成一锅粥——A要调B,B要调C,C又要找A,每次加个新Agent都得改一圈配置。Agent-Reach就是在这个背景下折腾出来的:给这堆Agent加了一层“触达层”,统一管注册、发现、消息路由,让每个Agent只关心自己擅长的事,别人找它就来活,它需要别人帮忙就发一条消息出去。

这个方案解决的核心问题,一句话概括就是:把多Agent协作从“人肉维护的点对点调用网”改造成“自动注册与路由的消息网”。如果你也在本地跑了好几个Agent、想让它们协作完成复杂任务、又不想每次加节点都重写一遍调用逻辑,那这篇文章值得你花十分钟看完。我会把Agent-Reach的设计思路、消息协议、最小实现对一遍,再把我联调过程中踩过的坑原原本本列出来。

1. 先说清楚Agent-Reach到底解决什么问题

1.1 多Agent协作的真实痛点

咱们先还原一个典型场景。你有一个代码生成Agent,一个文档搜索Agent,一个需求分析Agent。想让它们仨协作完成“根据需求文档生成项目骨架并输出说明”这个任务,最直接的做法是什么?代码生成Agent自己调用搜索Agent拿到文档,再自己调用需求分析Agent拿到结论。

这套方案看起来没毛病,跑几次就发现问题了:需求分析Agent如果换了接口,代码生成Agent得跟着改;搜索Agent挂了,代码生成Agent要处理异常;再塞进来一个测试Agent,代码生成Agent又要多写一个调用。这就是典型的点对点调用网——每个Agent都要知道所有协作方在哪、怎么调、参数长什么样。

更麻烦的是,消息从A到B到C再到D要走好几跳,每一跳都可能超时、失败、丢消息,你不做全链路追踪,出了问题根本不知道卡在哪。我那时候排查一个问题,愣是在日志里翻了半天才找到是哪两个Agent之间断了。这种复杂度随着Agent数量指数上涨,到七八个Agent的时候基本就管不动了。

1.2 Agent-Reach给出的解法:加一个触达层

Agent-Reach的思路很朴素,套用生活里的例子就是:一个项目群里几十个人,没人会去记每个人的手机号,需要谁帮忙就在群里喊一嗓子,或者通过群主分派。Agent-Reach就是那个“群主”——它维护了一张“谁能干什么”的表,谁想发消息就往它这里投递,它负责找到正确的接收方、把消息送到、再把回执带回来。

这样每个Agent只需要做两件事:启动时向Agent-Reach注册自己的能力和地址,运行时通过Agent-Reach收消息、发消息。Agent之间不需要知道对方任何细节,不知道对方的IP,不知道对方用什么语言,不关心对方是不是同机部署,甚至不需要双方同时在线——消息先交给Agent-Reach,接收方上线了再取走。

这种模式在分布式系统里叫消息中间件或者服务网格,但Agent-Reach针对AI Agent的场景做了几个专门的优化:能力匹配不靠固定接口签名,而是靠语义描述;消息格式里带了跳数限制防止Agent之间互相调用来回调去形成环;内置了重试和超时机制,不用每个Agent自己写一遍容错逻辑。

2. 核心设计:注册、发现与消息路由

2.1 消息协议怎么定最顺手

Agent-Reach里跑的消息,我定义为一段JSON,字段不多但每个都有讲究:

{ "msg_id": "8f2a9c1e-4b7d-4b2a-9e31-5d8f2a6c1b3d", "sender": "agent-scheduler", "target_skill": "web_search", "target_agent": "", "payload": {"query": "AI Agent框架最新动态", "limit": 5}, "deadline_ms": 15000, "max_hops": 3, "reply_to": "agent-scheduler", "timestamp": 1734567890123 }

逐字段说下我的设计逻辑:

  • msg_id:全链路唯一ID,排查问题全靠它串联各Agent的日志。这个务必用UUID,别用自增数字,多Agent并发时自增ID绝对会重复。
  • sender:来源Agent名,Agent-Reach拿它做回执路由,也用来做最简单的鉴权(判断这个Agent是否注册过)。
  • target_skill:目标能力标签,告诉Agent-Reach这条消息要发给“具备这个能力”的Agent。这是Agent-Reach区别于传统消息队列的核心点——MQ只认队列名,Agent-Reach认能力标签,粒度更粗但更贴合Agent协作的语境。
  • target_agent:选填。如果你明确知道要找某个具体的Agent,填上它就能精确路由,跳过能力匹配那一步。
  • payload:业务内容,结构随意,每个Agent对payload有自己的约定。
  • deadline_ms:绝对超时时间(毫秒)。我特意用“绝对时间戳”而不是“相对超时”,防止消息在中转过程中等待太久,导致超时判断失真。
  • max_hops:允许经过多少跳。Agent A发消息给B,B处理完可能又要找C,这个字段就是给这种链式调用准备的。等于3意味着A是第1跳,B是第2跳,C是第3跳,到了C就不能再往下传了——这是防环的关键设计。
  • reply_to:回执地址。B处理完结果要发回给谁,就靠这个字段,实际上就是sender的副本,单独拎出来是因为可能有人希望回执发给另一个Agent而不仅是原始发送者。

注意:deadline_ms和max_hops是我踩了坑之后才加上的。第一版没有这两个字段,结果某次两个Agent逻辑写岔了,互相发消息形成死循环,把所有Agent的CPU跑满了。这两个字段就是保险丝,宁可设大一点,但不能没有。

2.2 注册中心与Agent能力声明

Agent-Reach的核心数据是一张“Agent能力表”。每个Agent启动后,往注册中心上报四样东西:Agent唯一ID、能力标签列表、连接地址、当前状态。

我用的注册协议长这样:

{ "action": "register", "agent_id": "agent-search", "skills": ["web_search", "news_fetch", "url_check"], "endpoint": "ws://127.0.0.1:9002", "meta": { "description": "负责搜索互联网信息,返回URL列表和网页摘要", "model": "qwen2.5-14b", "version": "0.3.1" } }

meta.description是这个Agent的自然语言能力描述,这一段很有用。Agent-Reach内部有一个语义匹配模块,当某个target_skill和所有Agent的能力标签都匹配不上时,就会拿这个描述去做语义相似度计算,看看是不是能力标签取名的差异导致没匹配上。

注册之后Agent还要维持心跳,默认每5秒发一次,如果连续3个心跳周期没收到,Agent-Reach就把这个Agent标记为offline,不再往它那里路由消息。心跳消息很简单:

{"action": "heartbeat", "agent_id": "agent-search", "status": "busy"}

有个细节值得提:心跳里可以带上当前状态idle或busy。我做过一个优化,Agent-Reach优先把消息路由给空闲的Agent——同一种能力挂了多个实例时,这个字段就能起到简单的负载均衡作用。

2.3 路由策略:精确匹配 + 语义兜底

消息来了之后,Agent-Reach按三步走:

第一步,精确匹配。遍历所有在线Agent,谁的skills列表里包含target_skill,谁的地址就被列为候选。如果有多个候选,优先选status=idle的,再随机选。

第二步,语义兜底。精确匹配一个都没命中,就把target_skill和所有在线Agent的meta.description做embedding相似度计算,取相似度最高的那个,阈值设0.75,高于阈值才路由过去。这个阈值我调了好久,设太低容易路由到明显不对的Agent,设太高又失去兜底意义。实测0.75在大多数场景下比较平衡。

第三步,找替补。候选Agent如果发送失败,Agent-Reach会尝试选下一个候选,而不是直接把错误返回给发送方。这个逻辑模拟了人的行为——你找一个同事帮忙他没空,你自然会去找另一个能干的同事,而不是直接跟领导说“没人干活”。

我对比过纯精确匹配和带语义兜底的效果。在一个六个Agent的测试环境里,精确匹配的成功率大概是82%,因为总有那么几次我记错了能力标签的拼写;加上语义兜底之后成功率到了96%,剩下4%主要是描述本身太模糊,比如某个Agent写“处理各种乱七八糟的文本”,这种描述谁来了也匹配不准。

3. 实操:搭一套Agent-Reach最小可用版本

3.1 技术选型与理由

选型之前我纠结过一阵:用现成的消息队列(RabbitMQ、Redis Streams)、还是用gRPC搞服务发现?最后都没选,理由如下:

不用RabbitMQ/Kafka:这类系统是为海量消息设计的,功能强,但部署和运维成本对一个个人Agent项目来说太重了。消息是发给具体Agent的,不是广播给订阅者的,MQ的发布订阅模型和Agent通信模型有个根本错位。

不用gRPC服务发现:gRPC那套健康检查、负载均衡、拦截器,对Agent之间这种轻量协作来说显得臃肿。Agent间消息频率低(一秒几条到几十条),但每条消息的结构都不同,JSON更灵活。

最终选型:Python 3.11 + asyncio + aiohttp WebSocket。模块通讯用WebSocket而不是HTTP,因为Agent之间是双向通信——A发消息给B,B可能过几秒才回,WebSocket天然支持这种模式,HTTP就得靠轮询了。aiohttp比websockets库多带了HTTP服务能力,注册中心本身也要跑一个简单的HTTP接口,用aiohttp一个库全搞定。

另一个机型点是存储用SQLite而非Redis。Agent能力表的数据量极小(几十个Agent算很多了),SQLite持久化文件既不用额外起服务,又能保证重启不丢注册信息。Redis还得考虑内存淘汰、持久化配置,对这个小项目来说是杀鸡用牛刀。

3.2 核心模块代码实现

Agent-Reach拆成三个文件:registry.py——注册中心,router.py——路由引擎,agent_sdk.py——给Agent用的客户端库。

先看registry.py的核心逻辑:

# registry.py import asyncio import json import sqlite3 import time class AgentRegistry: def __init__(self, db_path="agent_reach.db"): self.conn = sqlite3.connect(db_path, check_same_thread=False) self.conn.execute(""" CREATE TABLE IF NOT EXISTS agents ( agent_id TEXT PRIMARY KEY, skills TEXT, endpoint TEXT, description TEXT, status TEXT DEFAULT 'idle', last_seen REAL, registered_at REAL ) """) self._lock = asyncio.Lock() async def register(self, agent_data): async with self._lock: self.conn.execute( "INSERT OR REPLACE INTO agents VALUES (?, ?, ?, ?, ?, ?, ?)", ( agent_data["agent_id"], json.dumps(agent_data.get("skills", [])), agent_data["endpoint"], agent_data.get("meta", {}).get("description", ""), "idle", time.time(), time.time() ) ) self.conn.commit() async def heartbeat(self, agent_id, status="idle"): async with self._lock: self.conn.execute( "UPDATE agents SET status=?, last_seen=? WHERE agent_id=?", (status, time.time(), agent_id) ) self.conn.commit() async def get_online_agents(self): return self.conn.execute( "SELECT * FROM agents WHERE last_seen > ?", (time.time() - 20,) ).fetchall() async def mark_offline(self, agent_id): async with self._lock: self.conn.execute( "UPDATE agents SET status='offline' WHERE agent_id=?", (agent_id,) ) self.conn.commit()

这个类做的事情很简单,但有几个点我特别说明:INSERT OR REPLACE专门处理Agent重启后重新注册的情况——同一ID重复注册,直接用新数据覆盖旧数据,不用先删再插;last_seen和当前时间差20秒作为在线判定,对应前面说的5秒心跳×3次容忍,超时了自动视为离线,不用等Agent主动发offline消息;asyncio.Lock()是防并发写库,多个Agent同时注册时SQLite会报database is locked,这个锁能直接避免。

然后是路由器的核心:

# router.py import asyncio import json import time import uuid import random class AgentRouter: def __init__(self, registry, websocket_connections): self.registry = registry self.connections = websocket_connections # {agent_id: ws} self.semaphores = {} # 每个Agent限流用 async def route(self, message): # 跳数检查 if message.get("max_hops", 3) <= 0: return {"error": "max_hops exceeded", "msg_id": message.get("msg_id")} message["max_hops"] -= 1 # 能力精确匹配 target_agent = message.get("target_agent") if not target_agent: target_agent = await self._match_skill( message["target_skill"], message["sender"] ) if not target_agent: return {"error": "no agent found for skill", "skill": message["target_skill"]} # 检查目标Agent是否在线、是否连接着 ws = self.connections.get(target_agent) if not ws: return {"error": "target agent offline", "agent": target_agent} # 限流:每个Agent同时处理的消息不超过10条 sem = self.semaphores.setdefault( target_agent, asyncio.Semaphore(10) ) if sem.locked(): return {"error": "target agent busy", "agent": target_agent} async with sem: await ws.send_json(message) return {"ok": True, "msg_id": message["msg_id"], "target": target_agent} async def _match_skill(self, skill, sender): agents = await self.registry.get_online_agents() candidates = [ a for a in agents if skill in json.loads(a[1]) and a[0] != sender ] if candidates: return random.choice(candidates)[0] # 这里省略了embedding语义匹配的具体实现 return None

这段代码包含了我反复调过的几个细节。message["max_hops"] -= 1是借消息的跳数检查——我把减一跳的操作放在路由入口做,这样Agent收到消息之后如果再转给下一个Agent,用的就已经是减过一跳的计数,每个环节都不用重复减。semaphore.locked()检查用的不是获取锁而是直接看状态,目的是一旦Agent忙就立刻返回错误,不要让发送方干等着排队——对我们这个体量,告诉发送方“对方忙”比让它排队体验更好。

最后是Agent侧的SDK:

# agent_sdk.py import asyncio import json import uuid import aiohttp class AgentClient: def __init__(self, agent_id, skills, endpoint, registry_url): self.agent_id = agent_id self.skills = skills self.endpoint = endpoint self.registry_url = registry_url self.ws = None self.pending_replies = {} async def start(self): async with aiohttp.ClientSession() as session: # 连接Agent-Reach的WebSocket网关 self.ws = await session.ws_connect(f"{self.registry_url}/ws/{self.agent_id}") # 注册自己 await self.register() # 开启心跳任务 asyncio.create_task(self._heartbeat_loop()) # 处理消息循环 async for msg in self.ws: if msg.type == aiohttp.WSMsgType.TEXT: data = json.loads(msg.data) if "reply_to" in data: # 回复消息,唤醒等待的调用方 fut = self.pending_replies.pop(data["reply_to"], None) if fut: fut.set_result(data) else: # 普通业务消息,交给用户处理函数 asyncio.create_task(self.on_message(data)) async def send_request(self, target_skill, payload, timeout=15): msg = { "msg_id": str(uuid.uuid4()), "sender": self.agent_id, "target_skill": target_skill, "payload": payload, "deadline_ms": timeout * 1000, "max_hops": 3, "reply_to": self.agent_id, } loop = asyncio.get_event_loop() fut = loop.create_future() self.pending_replies[msg["msg_id"]] = fut await self.ws.send_json(msg) try: return await asyncio.wait_for(fut, timeout=timeout) except asyncio.TimeoutError: self.pending_replies.pop(msg["msg_id"], None) return {"error": "timeout", "msg_id": msg["msg_id"]}

AgentClient里最关键的是pending_replies这个字典加上asyncio.Future的组合——发送请求时创建一个Future丢进字典,等待期间如果Agent-Reach转来了这个请求的回复,就立刻把结果set给Future,唤醒等待中的调用方。这个模式比用队列简单得多,而且超时器能直接取消等待,不用额外清理。

3.3 三个Agent联调实录

为了验证Agent-Reach真的能用,我搭了一个最小验证环境:一个调度Agent(负责拆解任务并分发),一个搜索Agent(负责查资料),一个文本处理Agent(负责做摘要)。跑通后再把代码生成Agent加进来做第二期验证。

流程是这样的:用户在终端对调度Agent说“查一下MCP协议最近几个版本的更新内容”。调度Agent收到后,生成一条target_skill=web_search的消息通过Agent-Reach发给搜索Agent;搜索Agent执行搜索返回URL列表;调度Agent再把URL列表打包成target_skill=content_summarize的消息发给文本处理Agent;处理Agent返回摘要;调度Agent把最终结果打印到终端。

跑起来后我在Agent-Reach日志里看到的完整链路是这样的:

[10:21:33.221] route msg=8f2a9c1e target=web_search from=agent-scheduler hops=2/3 [10:21:33.652] route result ok target=agent-search latency=431ms [10:21:35.104] route msg=9b4c7d2f target=content_summarize from=agent-scheduler hops=2/3 [10:21:35.587] route result ok target=agent-nlp latency=483ms [10:21:37.022] task complete total_time=3.80s

这个结果看起来顺畅,但中间其实翻过几次车。第一次跑的时候,搜索Agent收到消息后格式错了,直接报JSON解析异常回了段乱码,Agent-Reach把乱码当正常消息转发给了调度Agent,调度Agent一脸懵。后来我在Agent-SDK里加了一层消息校验,格式不对的回复直接在Agent-Reach端拦截掉,不往业务层传——这个坑我放在后面“常见问题”里细说。

4. 踩坑实录:我在联调中遇到的5个典型问题

4.1 问题速查表

现象根因解决方法
消息路由过去,对方没反应,等满15秒超时Agent绑定了127.0.0.1,Agent-Reach用局域网IP连不上endpoint统一填0.0.0.0或实际可路由的IP
一个Agent处理慢,其他Agent全被堵住没做消息并发控制,所有消息挤在一起处理为每个Agent加Semaphore并发上限,超限直接返回busy
两个Agent逻辑写岔,互相发消息死循环没有跳数限制,消息永远传不完加max_hops字段,每过一跳减1,减到0就丢弃
重试风暴:Agent超时后重发,又把其他消息挤掉重试没有退避策略,超时立刻重发指数退避+随机抖动,初始重试间隔2秒,逐次翻倍
消息发错Agent,且发送方还没发现语义匹配阈值设太低,相似度0.6就路由了阈值提到0.75,低于阈值返回not_found让发送方自己判断

这张表按我踩坑的频率排序,前三个都是上线第一天就遇到的,后两个是运行一周后逐渐暴露的问题。

4.2 每个问题的详细排查过程

第一个问题特别隐蔽:新加的摘要Agent注册成功,心跳正常,Agent-Reach显示在线,但路由过去就是没响应。我查了半天,最后发现Agent进程里aiohttp服务绑定的是127.0.0.1:9002,而Agent-Reach在同一台机器上,本应能连通——问题出在Agent-Reach作为服务端保存的连接地址来自Agent注册时上报的endpoint字段,我测试时手滑把endpoint写成了另一台测试机的内网IP。检查注册信息才发现上报地址和实际监听地址不一致。Agent-Reach对endpoint做了约束:注册时必须从本机检测实际的监听地址,未匹配直接拒绝注册。

第二个问题暴露得比较晚:并发任务一多,某个处理能力较弱的Agent(本地小模型推理,单条要好几秒)接收新消息时还在处理旧消息,积压越来越多,最终把整个链路拖到全线超时。Agent-Reach为每个Agent加了并发上限,超过上限的消息立即返回busy,由发送方决定是稍后重试还是换个Agent。

第三个问题的排查过程是我这个月最难受的一个凌晨。现象就是CPU飙满、整个系统卡死,起初以为是模型推理负载太高,一查进程发现是两个Agent在疯狂互发消息——A要求B验证数据,B发现数据不符要求A重新生成,两边都没设置最大处理次数,消息就无限循环起来了。幸好数据量不大,几轮循环就触发了日志刷屏,顺着日志才定位到这两个Agent。

第四个重试风暴问题紧跟在第三个后面暴露:我加完max_hops限制后,循环停了,但A发给B的消息一旦超时,发送方立刻重发,B刚好在处理大批量任务,每条消息都超时,每条都立刻重发——这等于把本来只有几Mbps的消息量放大到几十Mbps。给SDK内置重试机制后更糟,每层App都重试,最终消息量翻了几十倍。后来在Agent-Reach中央节点做统一重试控制,只允许Agent发消息时标记retry_count,最多3次,且必须退了再试。

第五个问题是有一次搜索结果莫名其妙指向了另一个Agent——细查发现语义匹配用了相似度0.6的阈值,把“搜索”和“分析数据”看成了一回事。提高阈值并加了“匹配结果必须附带相似度分数”的机制,发送方能看到路由的依据,误路由就好查了。

5. 压测数据、扩展方向与常用工具链

5.1 简单的性能摸底

写完之后我好奇这套轻量方案到底能承受多大压力,就在本地跑了轮简单的压测:三个Agent实例,消息大小约1KB,连续发了1000条路由消息,记录耗时。

结果如下:

指标数值说明
总耗时约14.2秒1000条消息全部路由完毕
平均处理耗时约14.2ms/条包含消息校验、能力匹配、WebSocket发送
P99延迟42ms最慢的那1%,能感受到延迟
最大并发数10条/Agent受Semaphore限制,再高就返回busy

坦白讲,这个数据在“性能”上并不亮眼——没有做批处理优化,没有消息压缩,也没有多进程并行路由。但Agent协作场景本身不需要高性能消息通道,更需要的是灵活的消息格式和语义匹配能力。吞吐几十条每秒已经远超“几个Agent互相说会话”的需求了。

真遇到高吞吐场景,我建议这么改:消息批量拉取而不是单条发送;payload用MessagePack压缩而不是JSON;WebSocket的收发循环拆成独立进程跑。这三板斧做完,吞吐量提升一个数量级没问题。

5.2 还能怎么扩展:从单机到多机的Agent触达网络

Agent-Reach这个名字里带了个Reach,本意是“触达”——现在这套方案只做到了单机触达,几台机器之间的Agent协作还没覆盖。我的扩展方向有三个:

方向一:多机部署。现在Agent-Reach的WebSocket网关绑定的是单机端口,换成多机之后需要把注册中心单独抽出来,用gRPC或HTTP暴露给各台机器。每台机器跑一个本地Agent-Reach的副本,本地路由走后端注册中心。这个改动预计要动registry.py的存储层,把SQLite换成PostgreSQL或者接口化的存储后端。

方向二:向MCP(Model Context Protocol)方向靠拢。越来越多Agent开始支持MCP标准,Agent-Reach的语义路由能在这上面发挥价值——每个Agent暴露成MCP server,Agent-Reach做MCP host,统一处理工具调用、上下文传递、跨Agent请求。这个方向是和行业标准对齐的最好路径,不用自己造协议生态。

方向三:给Agent-Reach加一个“Agent市场”的概念。注册时不仅能声明能力,还能声明计费方式、质量等级。路由的时候,用户指定目标能力和预算上限,Agent-Reach根据历史成功率和延迟做路由选择——这实际上就是把触达网络升级成交易市场,Agent之间从协作关系变成服务关系。

我现在已经实现了本地多Agent触达的核心逻辑,下一步想先把多机联通做了。按这几轮踩坑的经验,多机改造的重点不是网络打通,而是注册信息的同步时效和心跳检测的时限设计——单机心跳快慢无所谓,多机环境里网络抖动是常态,心跳超时阈值得重新调。

关于Agent-Reach的一点个人体会

项目做到这里,我对“Agent协作需要什么”这件事的看法改变了不少。最初以为难点在模型能力、提示词工程这些“上层建筑”上,做完Agent-Reach才意识到,底层的消息触达、发现、路由才是地基——地基不牢,上面堆再多Agent也是各行其是,碰到复杂一点的协作任务就乱成一团。

如果让我给同样在折腾多Agent项目的人一句建议,我会说:别急着上K8s、别引入复杂的服务网格,先用几百行代码把消息层控制住,等你真的遇到瓶颈了再升级。Agent协作的核心痛点是“谁有本事干这活、怎么把话递过去”,这是逻辑层面的事情,跟基础设施的规模没多大关系。

个人体会最深的一点:Agent之间的通信协议一定得允许“不知道对方是谁”。点对点调用是具体思维,路由网络是系统思维——前者适合两个Agent的临时任务,后者才是多Agent协作的常态。Agent-Reach这套东西骨架只花了一个下午就搭出来了,但前前后后填坑的时间估摸着得翻十倍。但这个过程值,至少现在我再加一个新Agent,只需要写一个配置文件然后启动,其余的事Agent-Reach全接管了。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询