ARTICLE DETAIL

资讯详情

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

多Agent协作中间层Agent-Reach:服务发现与路由调度实践

多Agent协作中间层Agent-Reach:服务发现与路由调度实践 1. 项目起底Agent-Reach 到底解决什么问题1.1 多Agent协作的混乱现状先说个最直观的场景。你团队里同时跑着六个七个AI Agent一个负责客服工单分类一个做数据分析报表一个挂在运维群里盯告警还有一个专门处理文档抽取。单独看每个都很能干一旦要让它们协同干活立刻变成噩梦。我见过太多团队倒在这一步。最开始大家觉得多Agent协作嘛让Agent之间互相调用就行了。于是A调BB调CC再回调A业务代码里写满了彼此的IP和端口。最崩溃的一次某个Agent因为流量超支被限流连带的是一整条链路上的任务全部超时。排查的时候还得逐个翻日志搞清楚到底哪个环节把任务弄丢了。那种痛苦做过一次就不想再来第二次。问题的根源在于Agent之间缺乏标准的发现和路由机制。每个Agent都像一个独立的微服务但比微服务更难搞的是Agent对外提供的是“智能能力”它可能处理的是自然语言任务、工具调用、多轮会话而不是简单的HTTP接口。如果继续用静态配置文件把Agent绑死设备一旦扩容、升级、替换整个拓扑就全乱了。Agent-Reach这个项目就是冲着解决这个问题去的。它的定位不是另一个Agent框架而是一个轻量级的Agent接入与调度中间层解决三个最基础也最头疼的问题系统的所有Agent在哪里各自具备什么能力业务方想按“能力”调用Agent而不是按某个具体实例IP去调用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_typeAgent类型比如“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。心跳的时序流程是这样Agent启动调用注册接口携带能力描述中心返回租约IDAgent在TTL/2时间内调用心跳接口续约中心每次收到心跳刷新该Agent的TTL超过TTL未续约中心把状态置为OFFLINE通知相关调用方该实例不可用。TTL我的建议默认设60秒心跳周期25到30秒。TTL设太短网络抖动会频繁造成Agent被误下线TTL设太长又不灵敏。60秒是我压过的比较舒服的值。2.3 路由策略按能力匹配按分数打分路由是Agent-Reach的核心动作。业务方发起调用时只告诉中心“我要调用能parse_pdf的Agent”剩下的选择权全部交给路由模块。路由算法并不复杂但考虑了四类因素能力匹配候选Agent的capabilities必须包含请求的能力标签这是硬门槛可用性检查Agent状态必须为ONLINE最近心跳必须在有效期内负载评分Agent上报的当前并发任务数结合它的并发上限算出一个负载率粘滞偏好如果请求带了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(keylambda 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(titleAgent-Reach Gateway) r aioredis.from_url(redis://localhost:6379/0, decode_responsesTrue) app.post(/v1/agent/register) async def register_agent(info: AgentInfo): key fagent:{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(fcap_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(fcap_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(code4001, reasonduplicate_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_idcs_agent_c011, agent_typecustomer_service, capabilities[ticket_classify, sentiment], registry_urlhttp://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_URLredis://redis:6379/0 - ROUTE_TIMEOUT15 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、上限20Agent 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_REQUEST5. 效果验证与实操心得系统稳定运行之后我做了一轮效果观察。原来团队里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自己决定调度层只负责保证“至少一次投递”不负责“精确一次”。做中间层明确边界比追求大而全重要得多。
返回列表