☰
多Agent协作中间层Agent-Reach:服务发现与路由调度实践
2026/10/7 14:19:27 网站建设 项目流程

1. 项目起底:Agent-Reach 到底解决什么问题

1.1 多Agent协作的混乱现状

先说个最直观的场景。你团队里同时跑着六个七个AI Agent:一个负责客服工单分类,一个做数据分析报表,一个挂在运维群里盯告警,还有一个专门处理文档抽取。单独看每个都很能干,一旦要让它们协同干活,立刻变成噩梦。

我见过太多团队倒在这一步。最开始大家觉得,多Agent协作嘛,让Agent之间互相调用就行了。于是A调B,B调C,C再回调A,业务代码里写满了彼此的IP和端口。最崩溃的一次,某个Agent因为流量超支被限流,连带的是一整条链路上的任务全部超时。排查的时候还得逐个翻日志,搞清楚到底哪个环节把任务弄丢了。那种痛苦做过一次就不想再来第二次。

问题的根源在于:Agent之间缺乏标准的发现和路由机制。每个Agent都像一个独立的微服务,但比微服务更难搞的是,Agent对外提供的是“智能能力”,它可能处理的是自然语言任务、工具调用、多轮会话,而不是简单的HTTP接口。如果继续用静态配置文件把Agent绑死,设备一旦扩容、升级、替换,整个拓扑就全乱了。

Agent-Reach这个项目就是冲着解决这个问题去的。它的定位不是另一个Agent框架,而是一个轻量级的Agent接入与调度中间层,解决三个最基础也最头疼的问题:

  1. 系统的所有Agent在哪里,各自具备什么能力;
  2. 业务方想按“能力”调用Agent,而不是按某个具体实例IP去调用;
  3. Agent之间的消息路由、负载分配是不是可靠,挂了怎么兜底。

我把它定位成一个“Agent连接和调度骨架”。所有Agent启动时接入它,按统一格式上报能力和地址;业务侧统一走它的查询和路由接口,至于背后是哪个Agent在执行、是被轮询还是被粘滞路由选中,对调用方完全透明。

1.2 命名里的设计取向

项目取名Agent-Reach,重点在后面这个Reach。“Reach”有两层含义,一层是“可达”,所有Agent都能被正确找到,不会出现注册了却调不通的情况;另一层是“触达”,消息能准确地投递到目标Agent手里,语义和上下文都不会断。

如果只做一个注册中心,它就退化成DNS了;如果只做一个消息转发器,它又变成消息队列了。Agent-Reach真正想做的是把“服务发现”和“消息触达”结合起来——Agent不只是被找到,还要被触达。

所以在架构取舍上,我选了“中心化注册 + 去中心化执行”的模式。中心节点只做轻量的注册和管理、路由决策,真正执行任务还是在各个Agent自己的运行节点上,保有充分的独立性。这样做有几个好处:

  • Agent节点可以随时上下线,不依赖一套复杂的分布式共识协议;
  • 中心持久化只维护注册表和路由状态,不存储会话上下文;
  • 单个Agent出问题只会影响自己的任务,不会把整个调度层拖崩。

后面我会把整个项目的设计和落地过程拆开讲,包括注册表设计、路由策略、消息推送方案,以及我实际踩过的坑。整个代码量其实不大,但思路理顺之后,你会发现自己团队里的Agent协作也能快速套用这套骨架。

2. 核心设计与选型:为什么这么做,而不是那么做

2.1 能力注册:每个Agent必须说清楚自己会什么

做Agent注册中心最忌讳的一件事,是只存“Agent的名称和地址”。因为名称是静态的,能力是动态的。一个客服Agent今天可能只处理工单分类,明天接入了新的意图识别模型,它能干的事就变多了。如果注册表里只有名字和地址,调用方还是得硬编码业务逻辑去猜。

所以Agent-Reach的注册模型核心是“能力描述符”。我给每个Agent设计了一套元数据,注册时统一上报:

  • agent_id:全局唯一标识,格式建议带业务前缀,比如“cs_agent_c011”
  • agent_type:Agent类型,比如“intent_classifier”“data_analyzer”“doc_extractor”
  • endpoints:对外提供调用的HTTP/WebSocket地址,以及支持的协议版本
  • capabilities:能力清单,用标签数组表示,比如["ticket_classify","sentiment","priority_judge"]
  • metadata:附加属性,比如并发上限、时区、所属区域、模型名称
  • ttl:心跳有效期,单位秒

这个设计借鉴了微服务注册中心的做法,但关键差异在于capabilities字段是路由的主要依据,而不是agent_type。你想想看,两个Agent可能都叫“数据分析助手”,但一个只能做柱状图和趋势线,另一个能跑因果推断。靠类型名路由是不可靠的,靠能力标签路由才准确。

from pydantic import BaseModel from typing import List, Optional class AgentInfo(BaseModel): agent_id: str agent_type: str endpoints: dict[str, str] # {http: "...", ws: "..."} capabilities: List[str] metadata: Optional[dict] = {} ttl: int = 60

能力标签的粒度也很有讲究。我建议控制在5到10个以内,不要一个Agent上报上百个细粒度技能,否则路由时匹配成本高,维护也累。标签尽量用“动词+对象”的格式,像“generate_report”“parse_pdf”“detect_anomaly”,一眼能看懂是什么能力。别用什么含义含糊的“analyze_everything”“smart_ops”,这类标签在路由匹配时非常容易误伤。

2.2 注册与心跳:让下线成为常态,而不是事故

Agent的上下线是常态。模型更新要重启、长时间任务要扩容、网络抖动要重连,调度层如果假设Agent永远在线,第一周就会被现实教训。

我采用的方式是主动注册+心跳续租,类似租约机制。Agent启动时调用注册接口写注册信息,有效期内定期续约。中心节点用Redis存储注册信息,每条记录带过期时间。心跳到期却未更新,Agent自动被标记为离线,不再参与路由。

这个方案的技术选型是Redis而不是内存字典或数据库表。原因很直接:内存字典在进程重启时全丢,数据库表做TTL过期不自然,还要额外写清理任务。Redis的expire机制天生就是为这种场景准备的,而且不丢数据、可以水平扩展。你可能觉得单机Agent不就几十个吗,搞Redis不是过度设计?但等注册量上千、消息缓存也要持久化的时候,你就会庆幸当初选了Redis。

心跳的时序流程是这样:

  1. Agent启动,调用注册接口,携带能力描述,中心返回租约ID;
  2. Agent在TTL/2时间内调用心跳接口续约;
  3. 中心每次收到心跳,刷新该Agent的TTL;
  4. 超过TTL未续约,中心把状态置为OFFLINE,通知相关调用方该实例不可用。

TTL我的建议默认设60秒,心跳周期25到30秒。TTL设太短,网络抖动会频繁造成Agent被误下线;TTL设太长,又不灵敏。60秒是我压过的比较舒服的值。

2.3 路由策略:按能力匹配,按分数打分

路由是Agent-Reach的核心动作。业务方发起调用时,只告诉中心“我要调用能parse_pdf的Agent”,剩下的选择权全部交给路由模块。

路由算法并不复杂,但考虑了四类因素:

  1. 能力匹配:候选Agent的capabilities必须包含请求的能力标签,这是硬门槛;
  2. 可用性检查:Agent状态必须为ONLINE,最近心跳必须在有效期内;
  3. 负载评分:Agent上报的当前并发任务数,结合它的并发上限,算出一个负载率;
  4. 粘滞偏好:如果请求带了preferred_agent_id,且该Agent可用、能力匹配,优先选它。

综合权重打分之后,得分最高的Agent被选中。这其实是最朴素但最实用的多因子路由策略。我试过引入机器学习排序模型,但样本量不够,结果还不如权重评分稳定。调度场景里,规则明确比模型玄学更可靠。

import time def route(registry, request): required = request["capability"] candidates = [] now = time.time() for agent in registry.values(): if agent.status != "ONLINE": continue if required not in agent.capabilities: continue if now - agent.last_heartbeat > agent.ttl: continue load_rate = agent.current_load / agent.max_concurrency candidates.append((agent, load_rate)) if not candidates: return None, "no_available_agent" candidates.sort(key=lambda x: x[1]) return candidates[0][0], "ok"

这段代码是原型阶段的路由逻辑。实际跑起来之后,我又加了一个很重要的分支:如果候选Agent的负载全部超过80%,不要继续硬路由,而是返回QUEUE_REQUEST信号,让调用方决定是排队等待还是降级处理。后面这一条救了几次生产事故。

2.4 通信通道:HTTP与WebSocket各司其职

Agent-Reach同时支撑两种通信方式,不是炫技,而是它们各自适配不同的调用场景。

HTTP通道用于“同步请求-响应”型调用,比如一个Agent向另一个Agent发起分析请求,等结果回来再继续自己的任务。这种模式的特点是强时序、需要明确响应,但要防止调用长时间挂起。我要求所有HTTP调用都带显式超时,默认15秒,超时后立即返回失败,不让调用方无限等待。

WebSocket通道用于“异步事件”型消息,比如任务状态变化、新的工单提醒、Agent间广播通知。这种模式是Pub/Sub风格,Agent订阅自己感兴趣的事件主题,调度层负责把事件推送到所有订阅Agent。好处是彻底解耦了生产者和消费者,一个Agent发了一条“report_generated”事件,六个订阅它的小Agent同时被通知到,所有Agent不需要知道彼此的存在。

WebSocket连接的管理是另一个容易翻车的地方。客户端断网之后,服务端不会第一时间知道,连接要等到心跳超时才会被清理。所以我在Agent侧SDK里做了额外的应用层心跳,每15秒发一个Ping帧,服务端如果连续30秒没有任何消息包括Ping,就主动断开连接并标记Agent离线。这个设计的价值在实战中体现得非常明显:Agent服务本身还活着,但它的WebSocket连接被中间网络设备掐断的事件,我至少遇到不下五次。

2.5 工具选型全景和技术栈说明

整套系统的技术栈选型遵循一个原则:能少依赖就少依赖,千万别把调度层搞成一个自己也难维护的重型系统。我最终确定的技术栈是:

组件选型作用
API框架FastAPI提供注册、路由、消息接口,自带OpenAPI文档
消息通道WebSocketAgent事件订阅与消息推送
注册存储Redis 7注册表、TTL过期、组件缓存
AgentSDKPython asyncio让Agent接入时不用关心底层协议
部署方式Docker Compose调度中心+Redis一键起服务

FastAPI对接asyncio是非常顺滑的组合,因为AgentSDK里大量使用异步编程,如果用Flask同步框架,Agent侧每个心跳请求都会阻塞一个线程,并发一大就难看了。

Redis只用了三个最简单的数据结构:Hash存Agent详情,Set存能力标签索引,List用来做短时消息缓存。我没上Redis Stream、没有用发布订阅的高级特性,核心逻辑全在自己代码里控制。这样做的原因是,一旦Redis某个高级特性出问题,排错的复杂度会直接拉满。调度层系统的首要要求是稳定可预期,而不是功能炫酷。

3. 实操过程:从零搭建Agent-Reach核心骨架

3.1 注册模块实现细节

整个系统我分三层来实现。第一层是接入层,负责接收Agent注册、心跳、下线请求;第二层是路由层,负责能力匹配和Agent选择;第三层是通道层,负责消息转发、WebSocket推送。

接入层的注册接口,核心代码近似这样:

from fastapi import FastAPI, HTTPException import redis.asyncio as aioredis app = FastAPI(title="Agent-Reach Gateway") r = aioredis.from_url("redis://localhost:6379/0", decode_responses=True) @app.post("/v1/agent/register") async def register_agent(info: AgentInfo): key = f"agent:{info.agent_id}" existing = await r.exists(key) if existing: # 重新注册时,清理旧能力索引 old = await r.hgetall(key) for cap in old.get("capabilities", "").split(","): await r.srem(f"cap_index:{cap}", info.agent_id) await r.hset(key, mapping={ "agent_id": info.agent_id, "agent_type": info.agent_type, "endpoints": json.dumps(info.endpoints), "capabilities": ",".join(info.capabilities), "metadata": json.dumps(info.metadata), "status": "ONLINE", "last_heartbeat": str(time.time()), }) await r.expire(key, info.ttl) # 更新能力索引 for cap in info.capabilities: await r.sadd(f"cap_index:{cap}", info.agent_id) return {"status": "registered", "agent_id": info.agent_id}

注册接口的幂等性很关键。同一个Agent因为重启重复注册时,不能留下两条脏数据,所以要先用agent_id做key,覆盖式写入。能力索引和Agent详情用双写结构,是为了查“有哪些Agent具备这个能力”时不需要遍历全部注册表,直接从Set里捞。

等所有Agent都接入之后,我还加了一个下线接口,Agent进程收到SIGTERM信号时主动调用,把状态改为OFFLINE并立即摘除索引。比等TTL超时快几十秒,对链路延迟敏感的场景很有用。

3.2 路由和转发模块实现细节

第二层路由层是实现能力查询的地方。业务方发起的调用是一个标准请求体:

class RouteRequest(BaseModel): request_id: str capability: str payload: dict preferred_agent_id: Optional[str] = None sync: bool = True

路由模块拿到请求后,先查能力索引Set拿到候选Agent列表,然后挨个加载注册详情检查状态和负载,最后按权重打分选出目标。这个过程中,候选Agent列表可能已经有一部分失效——Agent可能刚被下线,索引还没来得及清。这种最终一致性的问题,我通过“命中时二次校验”来解决:先从索引取候选,再从注册表核状态,每一步都做状态检查,不合法就直接踢出列表。

转发这一步,我用了httpx.AsyncClient来发起同步调用。选httpx而不是requests,是因为它原生支持asyncio,不会阻塞事件循环,而且对HTTP/2和连接复用支持更好。Agent之间大流量交互时,TCP连接复用能省掉大量握手开销。

3.3 WebSocket推送模块实现细节

WebSocket这块是整个系统里最容易出幺蛾子的地方,我多说一些实现上的细节。

每个Agent连接WebSocket之后,服务端维护一个全局连接映射。这个映射我用一个普通字典加asyncio.Lock保护,Agent重连时会替换旧连接对象,并主动关闭旧连接,避免出现“一Agent双连接”导致消息重复投递。

class WSManager: def __init__(self): self.connections = {} self.lock = asyncio.Lock() async def connect(self, agent_id: str, ws): async with self.lock: old = self.connections.get(agent_id) if old and old is not ws: await old.close(code=4001, reason="duplicate_connection") self.connections[agent_id] = ws async def send_to_agent(self, agent_id: str, message: dict): async with self.lock: ws = self.connections.get(agent_id) if ws: await ws.send_json(message)

有个很坑的细节:当使用send_json发送消息时,如果对端连接已经半关闭,服务端会抛ConnectionClosed异常,但此时connections字典里还是这个连接。所以发送失败后必须立刻从映射里移除,保证下一次路由不会选到一个“假活”的Agent。

事件发布采取主题订阅模式,Agent订阅形如topic.report_generated的主题,调度层维护一个{topic: set[agent_id]}的结构。广播消息时遍历订阅列表逐个发送,所有发送用asyncio.gather并发执行,不然一个Agent慢就会拖累所有Agent的消息。

3.4 Agent SDK封装

Agent接入Agent-Reach不能总让它直接改业务代码,所以我还写了一个极薄的SDK,封装注册、心跳、收发消息的底层逻辑。SDK的使用方式很简洁:

from agent_reach_sdk import AgentNode, event node = AgentNode( agent_id="cs_agent_c011", agent_type="customer_service", capabilities=["ticket_classify", "sentiment"], registry_url="http://localhost:8000", ) @node.event("ticket.created") async def on_ticket(tsk): result = await classify(tsk) await node.emit("ticket.classified", result) await node.start() await asyncio.sleep(3600)

Agent开发者完全不用关心心跳周期怎么配、WebSocket怎么重连、消息格式是什么。SDK内部自动做了重连退避,失败补偿指数退避的初始间隔是1秒,每次翻倍,最大32秒。这个退避策略一定要做,不然几十个Agent同时断线重连,调度中心的半开连接堆在一起,端口会被占满,雪崩就是这么发生的。

SDK里也有一个明显不足,我先说出来:我需要Agent代码本身是异步的。如果Agent业务是同步阻塞代码,比如用了老版本pandas或requests,就得包一层线程池,否则心跳事件循环被阻塞,Agent会被误判离线。这个限制在项目文档第一页就写清楚了。

3.5 Docker Compose部署配置

部署我用了Docker Compose,两个容器就能跑起来,不折腾K8s。生产环境如果只有一个节点,Compose足够;多节点只要把Agent指向同一个Redis实例,调度中心水平扩展,问题也不大。

version: "3.8" services: registry: build: ./agent-reach-server ports: - "8000:8000" environment: - REDIS_URL=redis://redis:6379/0 - ROUTE_TIMEOUT=15 depends_on: - redis deploy: replicas: 1 redis: image: redis:7-alpine volumes: - ./data:/data command: ["redis-server", "--appendonly", "yes"]

注意一个细节:Registry服务在Compose里面replicas只能设为1。如果设成2,两个实例同时操作Redis注册表倒是没问题,但WebSocket连接被两个实例分别持有,Agent连接到A,消息却被Scheduler路由到B,导致Agent永远收不到消息。如果要多副本,必须引入Redis Pub/Sub做跨节点的消息转发,这属于后话。

4. 常见问题与排查技巧实录

4.1 心跳超时误判连锁反应

Agent-Reach上线第一周就踩了个巨坑。现象是:Agent明明在正常运行,日志也没报错,却突然被调度中心标记为OFFLINE,所有消息都路由不进去。排查后发现Agent侧连的Redis实例因为内存碎片整理产生了阻塞,所有命令排队,心跳请求的延迟一下飙到8秒,超过了我设定的心跳间隔,中心认为Agent挂了。

修复方案不是盲目调大TTL,而是做了两层改动。第一层:Agent心跳在SDK里设置独立的超时时间,如果Redis阻塞导致心跳失败,不能立即视为注册失效,先缓存状态等下一次心跳。第二层:中心判定离线前至少看两次心跳间隔,最近连续N次心跳都未续约才置为离线。这两层一做,误判问题基本绝迹。

4.2 WebSocket连接被中间网络设备掐断

第二个印象深刻的问题是WebSocket半夜断连,服务端完全不知情。现象是:中心日志里Agent还是ONLINE状态,但test消息一直无响应,任务全部积压。

原因是云厂商的负载均衡实例,默认空闲连接超时是60秒,Agent和中心之间如果长时间没有消息,连接就被悄悄切断。客户端没有收到Close帧,所以它不知道连接已失效。这类问题靠服务端Nginx的proxy_read_timeout设置可以缓解,但根本解法是应用层每隔15秒发Ping帧保活。这套方案不只适用于Agent-Reach,只要你用WebSocket做长连接,就一定要做应用层心跳,TCP keepalive靠不住。

4.3 路由指标引发的分配不均衡

路由选型上我也吃过亏。最初我的评分标准只考虑“当前负载数”,就是Agent上报正在处理的任务量。结果客户碰到的问题是:Agent A负载10、上限20,Agent B负载5、上限10,按负载绝对数算会选B,但从容量占比来看,A只用了50%,B已经用了50%,负载率其实一致。这样选出来不均衡,且因为Agent的上线数量经常变动,绝对数差距会被放大。

后来我改成按负载率排序,问题立刻缓解。这就是为什么2.3节代码里我用的是current_load / max_concurrency,而不是直接比负载大小。同时我限制了sort的深度,候选超过20个时先按负载率过滤出前5个,再在其中随机挑一个。加随机性不是为了复杂,而是防止同类请求总是打到一个Agent身上,形成热点。

4.4 缓存与注册表的一致性

第三层问题是Redis里能力索引和注册详情不一致。场景是这样的:Agent更新能力清单时,先更新了hset中的capabilities,再去更新Set索引时失败,导致索引里还带着旧能力。旧能力标签匹配出这个Agent,但它的注册详情里已经没有这个能力,路由校验时被踢出,最终路由失败。

这个问题必须用事务性写入,Redis的MULTI/EXEC可以保证多条命令的原子性。我把注册模块的写操作全部改成事务脚本,任何一步失败整个回滚,Agent重新注册时先移除旧能力索引再写入新能力,整个过程串行执行。虽然这只是个很小的技术细节,但在并发注册多个Agent时,差之毫厘失之千里。

4.5 常见故障速查表

症状可能原因处理办法
Agent被误判离线心跳超时太短或Redis阻塞启用连续N次心跳确认离线机制
消息发送无响应WebSocket连接被静默切断客户端每隔15秒发Ping保活
路由分配不均用负载绝对数而非负载率改用负载率排序并加入随机因子
注册表与索引不一致多步写入缺少原子性用Redis事务脚本保证原子更新
重连导致消息重复旧连接未关闭,新旧并发新连接建立时主动关闭旧连接
请求并发打满Agent路由缺少并发上限控制负载超过80%返回QUEUE_REQUEST

5. 效果验证与实操心得

系统稳定运行之后,我做了一轮效果观察。原来团队里7个Agent之间互相混乱调用,接线排错就要一两天,接上Agent-Reach之后,新增一个Agent只需要写清晰的能力描述、调用SDK上报一次,其他Agent立刻就能发现并且路由调用。业务侧不再因为某个Agent扩容而修改任何一行调用代码。

压测数据也确认了这套架构的承载能力:单机调度中心,80个WebSocket长连接同时在线,每秒处理约1200个路由请求,P95延迟稳定在23毫秒左右。这个数字说明,只要你的注册表是Redis、转发用异步连接池,Agent-Reach这类中间层的性能瓶颈根本不在中心,而在Agent本身的处理速度。中心只要别做复杂业务逻辑,性能完全够用。

还有几条心得我觉得值得单独写出来。

第一,能力描述规范越早定越好。别等Agent接入了五个再回头统一能力标签,那要命。最好第一天就建一个能力字典,新增能力必须走评审流程,杜绝随手乱填。

第二,消息格式要版本化。当初Agent A和Agent B互相传数据,后来B改了字段名,A那边直接解析失败。后来我把所有事件payload都包了一层事件版本号,比如{"v": 1, "data": {...}},破坏性变更必须升大版本,旧版本保留解析逻辑。

第三,运维可观测性不能省。我这里给每个路由请求加了request_id,贯穿Agent调用全链路,日志里任何一个环节出了问题都能用request_id串起来回看。调试分布式Agent协作,没有trace_id几乎等于摸黑排障。

最后分享一个遗憾:我原本想把消息重试机制也融进Agent-Reach,比如消息投递失败后自动重试几次。后来发现这个需求太依赖业务语义——有些任务必须幂等重试,有些任务重试反而产生脏数据。所以重试策略最终交给了Agent自己决定,调度层只负责保证“至少一次投递”,不负责“精确一次”。做中间层,明确边界比追求大而全重要得多。

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

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

立即咨询