ARTICLE DETAIL

资讯详情

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

多智能体协作的轻量触达层:Agent-Reach设计与实践

多智能体协作的轻量触达层:Agent-Reach设计与实践 前段时间我们团队在搭一套多智能体协作系统最让人头疼的既不是模型本身的推理效果也不是Prompt写得多好而是那些各自独立运行的Agent之间根本没法顺畅地互相找到对方、互相调用能力。传统微服务那套注册中心和网关虽然能用但直接搬过来会发现两个问题一是太笨重二是不贴合Agent的工作方式——Agent不只是被动接收请求它得主动声明自己会什么、能碰什么数据、状态是忙还是闲。基于这些痛点我在公司内部落地了一个叫Agent-Reach的轻量触达调度层核心就三件事让Agent能被发现、让请求能路由到对的Agent、让调用过程能扛住超时和故障。这篇文章把从设计到写代码再到压测排查的完整过程都拆开讲适合正在做多Agent系统、被Agent互相调用问题折腾过的开发者参考。1. 项目背景为什么我需要一个“触达层”1.1 多Agent协作的三个典型痛点先说清楚我遇到的具体问题这样后面讲方案的时候大家才知道我在解决什么。第一个痛点是“找人难”。我们团队同时维护着七八个业务Agent有做客服工单分类的有做订单数据分析的有做自动化报表的还有两个是搞知识库问答的。每个Agent都跑在自己的独立服务里IP地址和端口都是写死在配置文件中的。后来业务一扩展新增一个Agent要通知所有人改配置下线一个Agent要紧急排查谁还在往那个死地址发请求。整个调用关系就是一张蜘蛛网没人说得清全局状态。第二个痛点是“说话难”。每个Agent都是不同同事用不同技术栈写的有的暴露REST接口有的走消息队列有的干脆只能通过命令行脚本触发。A Agent想调用B Agent的能力往往得先看B的接口文档手动构造参数处理返回值格式。一旦B升级了接口数据结构A这边就要跟着改代码。这种强耦合的调用方式在Agent数量少的时候还能忍一旦超过五六个维护成本就是指数级上升。第三个痛点是“管不住”。Agent之间互相调用的时候权限怎么卡A Agent到底有没有资格调用B Agent里涉及用户手机号的接口如果调用失败是重试还是放弃重试多少次如果不做统一管控最后就是各个Agent自己写一套重试逻辑自己维护一套调用白名单出了故障各查各的日志排障靠吼。这三个痛点凑到一起我意识到我们需要一个统一的“触达层”。它不是简单的API网关也不是单纯的消息队列而是介于Agent与Agent之间的一层调度基础设施负责注册、发现、路由、鉴权和可靠性保障。这个层我命名为Agent-Reach。1.2 选型权衡为什么不直接套用现成框架当时我们认真评估过几类方案。第一类是直接用Spring Cloud或者Dubbo那套微服务治理体系功能确实齐全但引入的组件太重——要配注册中心集群、配置中心、链路追踪运维负担不小。更重要的是微服务的服务发现模型是“服务名 IP 端口”三层而Agent的粒度是“能力 状态 上下文”比如同样是那个客服Agent它也许同时具备“工单分类”和“情绪识别”两种能力传统注册中心只能注册一个服务名然后在调用方写死具体路径没法按能力动态匹配。第二类方案是上消息队列Kafka或者RabbitMQ。消息队列擅长异步解耦但Agent之间的调用很多是需要同步返回结果的比如一个决策Agent要实时问询数据Agent“这个月的退货率是多少”这时候异步消息再加回调链路就复杂了排障也不方便。第三类方案是把网络层做成Service Mesh用Sidecar统一代理所有通讯。想法很好但对现有系统改造太大每个Agent都要额外部署代理容器而且我们的Agent大多是非容器化部署在虚拟机上SetSidecar的成本很高。综合权衡后我决定自研一个轻量触达层。核心思路是提供一个注册中心存Agent元数据提供一个路由引擎做能力匹配再提供一个执行网关统一处理超时、重试和鉴权。三块组件全部围绕Agent的实际工作模式设计而不是把传统微服务的概念硬套过来。2. 核心机制拆解Agent-Reach是如何工作的2.1 总体架构注册中心 路由引擎 执行网关Agent-Reach的整体架构分三层每层只干一件很纯粹的事。第一层是注册中心。每个Agent启动时会向注册中心上报自己的基础信息包括Agent名称、实例ID、主机地址、端口、当前状态以及最重要的一份“能力声明表”。能力声明表用JSON描述这个Agent能做什么包含动作列表、参数格式、返回格式和权限级别。注册中心把数据落库并通过TTL机制判断每个Agent是否存活。第二层是路由引擎。当调用方发起一个意图请求比如“帮我查一下上季度华东区的销售额”路由引擎先解析意图识别出这是“销售数据查询”类动作然后在注册中心里检索所有具备这类能力的Agent再结合调用方的权限范围和目标Agent的负载状态最终选出最合适的那个。第三层是执行网关。网关负责真正把请求发给选中的Agent同时统一处理三件事超时控制、重试策略、鉴权校验。调用方只跟网关交互不用关心目标Agent到底在哪个IP也不用自己写重试逻辑。这三层组合起来调用方视角的体验就是我发一个请求给Agent-Reach传一个意图和参数它给我返回结果或者一个明确的失败原因。至于请求到底被转给谁、目标Agent是死是活、中间重试了几次调用方统统不用管。2.2 服务注册与心跳机制让Agent“可被发现”注册和心跳是整个系统的地基。每个Agent在启动的时候调用一次注册接口传入自身信息随后每10秒上报一次心跳。注册中心维护一个LastHeartbeat时间戳每30秒做一次扫描把超过30秒没心跳的Agent标记为离线状态。这里有个关键细节不能一超过30秒就立刻把Agent摘掉因为网络抖动导致偶尔丢一两个心跳包是正常的。所以我把状态分为online、suspected、offline三档。连续两次心跳丢失进入suspected再丢一次才正式转offline。这个策略后来在压测中被验证非常有效下面的问题排查章节我会细讲。心跳包的格式也很简单只包含agent_id和时间戳不包含任何业务数据。这样注册中心的写入压力很小SQLite就能扛住小规模集群如果Agent数量超过几百个可以平滑替换成MySQL或者Redis。Agent下线或重启时应该主动调用注销接口但为了避免“进程被杀来不及注销”的异常场景TTL机制兜底仍然保留。如果Agent换了IP重启注册中心通过agent_name 业务标识字段判断是否同一个逻辑Agent如果是就更新地址信息而不是新建一条记录。2.3 能力声明与路由规则让请求“找到对的Agent”Agent只是报了一个“名字”是不够的路由引擎必须有办法知道“你会干什么”。所以每个Agent在注册时都附带一份能力声明表这是Agent-Reach路由的核心输入。能力声明表的结构大致如下{ agent: order-analysis, capabilities: [ { action: query-sales, params: [region, quarter, dimension], return_schema: sales_summary, scope: read, latency_ms: 200 }, { action: generate-daily-report, params: [date, format], return_schema: file_ref, scope: write, latency_ms: 15000 } ] }路由引擎收到调用意图后先做意图到动作的映射这一步我在调用方Agent的SDK里内置了一个简单的意图分类函数支持基于关键词的规则映射和基于调用的模糊匹配。比如“查一下销售额”会被映射到query-sales动作“做日报”会被映射到generate-daily-report动作。映射完成之后路由引擎在在线Agent列表中检索所有声明了对应action的Agent。如果命中了多个候选再按负载策略排序——默认是随机调用但我也实现了加权策略权重值由Agent心跳时携带的实时CPU使用率和队列长度计算。真实场景下同样的“客服意图识别”动作可能部署了两个实例一个空闲一个繁忙这时候权重策略就很有用了。2.4 执行网关超时、重试与鉴权路由选定了目标Agent接下来的活儿全交给执行网关。网关是同步HTTP调用因为大部分Agent之间的交互期望实时拿到结果异步场景暂不纳入这个版本。我希望每个动作有明确的时间预期所以给每个动作类型配置了超时时间。超时配置放在注册中心的能力声明表里由Agent自己声明一个latency_ms网关取这个值乘以一个系数作为请求超时时间。经验值读类动作超时5秒写类动作10秒文件生成类动作30秒各自乘以1.2的冗余系数。这样设置的好处是不同Agent可以自报“我需要多久”网关不用全局写死。重试逻辑我遵循两个原则只重试幂等动作只重试连接异常。调用方明确标记为幂等的动作比如query类查询可以重试最多两次非幂等动作比如订单创建只会在连接失败时重试一次避免重复创建。重试退避用指数退避基础间隔0.5秒倍增系数1.5。鉴权这块做得相对简单但必要。每个Agent在注册时签发一个token调用方发起请求时必须携带自己的token网关校验通过后再校验调用方是否有权调用目标Agent的某个动作权限来自一套简单的RBAC规则表。这样至少能把“谁都能调谁”的问题约束住。3. 实操过程从零搭一版最简Agent-Reach3.1 环境准备与依赖我选择用Python实现Agent-Reach的原型因为团队里Agent有一半是Python写的SDK集成方便。依赖很克制就三样FastAPI暴露注册、路由、网关的HTTP接口SQLite存储Agent元数据和权限规则Python内置httpx网关发起异步HTTP调用支持重试直接用pip安装pip install fastapi httpx uvicornSQLite不需要额外安装。生产环境如果Agent数量多了可以把存储层换成MySQL或Redis接口不用大改。整个项目分为registry.py、router.py、gateway.py、main.py四个模块。下面逐个写出核心代码并解释关键设计。3.2 注册中心的落地实现注册中心最重要的两个功能接收Agent的注册信息并落表定时扫描清理失联Agent。# agent_reach/registry.py import time import threading import sqlite3 import uuid class AgentRegistry: def __init__(self, db_pathreach.db): self.conn sqlite3.connect(db_path, check_same_threadFalse) self.conn.execute( CREATE TABLE IF NOT EXISTS agents ( agent_id TEXT PRIMARY KEY, name TEXT NOT NULL, host TEXT NOT NULL, port INTEGER NOT NULL, capabilities TEXT NOT NULL, status TEXT NOT NULL DEFAULT online, last_heartbeat REAL NOT NULL, created_at REAL NOT NULL ) ) self._lock threading.Lock() self._stop_scan False self._scan_thread threading.Thread(targetself._sweep_loop, daemonTrue) self._scan_thread.start() def register(self, name, host, port, capabilities: dict, ttl30): Agent启动时调用。同一个name的新实例会复用agent_id更新地址信息。 agent_id str(uuid.uuid4()) now time.time() with self._lock: cursor self.conn.execute( SELECT agent_id FROM agents WHERE name? ORDER BY created_at DESC LIMIT 1, (name,) ) row cursor.fetchone() if row: agent_id row[0] self.conn.execute( INSERT OR REPLACE INTO agents VALUES (?,?,?,?,?, online, ?, ?), (agent_id, name, host, port, json.dumps(capabilities, ensure_asciiFalse), now, now) ) self.conn.commit() return agent_id def heartbeat(self, agent_id: str) - bool: with self._lock: cursor self.conn.execute( UPDATE agents SET last_heartbeat? WHERE agent_id?, (time.time(), agent_id) ) self.conn.commit() return cursor.rowcount 1 def mark_offline(self, agent_id: str): with self._lock: self.conn.execute( UPDATE agents SET statusoffline WHERE agent_id?, (agent_id,) ) self.conn.commit() def _sweep(self): 这里用三态判断suspected - offline而不是直接离线。 cutoff time.time() - 30 with self._lock: self.conn.execute( UPDATE agents SET statussuspected WHERE last_heartbeat ? AND statusonline, (cutoff,) ) self.conn.execute( UPDATE agents SET statusoffline WHERE last_heartbeat ? AND statussuspected, (cutoff - 15,) ) self.conn.commit() def _sweep_loop(self): while not self._stop_scan: self._sweep() time.sleep(10)这段代码里有几个细节值得展开。第一同一个逻辑Agent第二次注册时我不会生成新的agent_id。这个处理很重要因为Agent可能因为发布新版本而重启如果每次都生成新ID那么注册中心里会堆满同一个Agent的历史记录路由查询时还得额外做去重。用name分组取最新一条记录保持agent_id不变等于自然实现了“实例重启但身份不变”的效果。第二扫描逻辑用了三态判定。第一次扫描发现心跳超时先置为suspected再过15秒如果还没有心跳才转offline。真实环境里偶尔的GC停顿、网络抖动导致10秒内丢一个心跳很常见直接标记离线会造成大量请求误路由到别的Agent触发不必要的重试。3.3 路由引擎的实现路由引擎的输入是一个“意图请求”输出是候选Agent列表。# agent_reach/router.py import json class Router: def __init__(self, registry: AgentRegistry): self.registry registry self.intent_map {} self.strategy weighted def register_intent(self, intent: str, action: str): self.intent_map[intent] action def resolve_action(self, intent: str) - str: return self.intent_map.get(intent, intent) def match(self, intent: str): action self.resolve_action(intent) rows self.registry.conn.execute( SELECT * FROM agents WHERE statusonline ).fetchall() candidates [] for row in rows: caps json.loads(row[4]) for cap in caps: if cap[action] action: candidates.append({ agent_id: row[0], name: row[1], host: row[2], port: row[3], capability: cap, latency: cap.get(latency_ms, 500) }) break if self.strategy weighted and candidates: total_latency sum(c[latency] for c in candidates) candidates.sort(keylambda c: c[latency]) return candidates匹配的核心逻辑很简单遍历在线Agent解析能力声明表找出声明了对应action的那些。这在小规模Agent集群中完全够用而且每轮查询不会超过几十毫秒。如果未来Agent数量增长这里可以加一层缓存把能力表加载到内存中而不是每次查SQLite。我当时的精力主要放在同步逻辑上所以先用这种最直白的查询保证正确性优先。路由策略我实现了两种随机和按延迟加权。按延迟加权就是我上面写的遍历之后按latency_ms排序把响应速度快的Agent排前面。实际跑下来发现这策略有个好处慢的Agent会自动被边缘化流量都打到快的那几台集群吞吐量反而更稳定。3.4 执行网关的代码细节网关负责最后的请求发送和失败兜底。# agent_reach/gateway.py import asyncio import json import time import httpx class Gateway: def __init__(self, retry_mapNone): self.retry_map retry_map or { query: {max_retries: 2, base_delay: 0.5}, command: {max_retries: 1, base_delay: 0.5}, } async def dispatch(self, target: dict, intent: str, payload: dict, identity: str): url fhttp://{target[host]}:{target[port]}/invoke timeout_ms target[capability].get(latency_ms, 500) timeout timeout_ms / 1000 * 1.2 meta self.retry_map.get(target[capability].get(scope, command)) max_retries meta[max_retries] base_delay meta[base_delay] headers {X-Agent-Token: identity} body {intent: intent, payload: payload} for attempt in range(max_retries 1): try: async with httpx.AsyncClient(timeouttimeout) as client: resp await client.post(url, jsonbody, headersheaders) if resp.status_code 400: raise RuntimeError(ftarget returned {resp.status_code}) return resp.json() except (httpx.TimeoutException, httpx.ConnectError): if attempt max_retries: return {error: timeout_or_unreachable, intent: intent} await asyncio.sleep(base_delay * (1.5 ** attempt)) return {error: unknown}网关的核心设计是超时时间由Capability自己上报的latency_ms乘以1.2得出重试次数由动作的scope决定——read类允许重试两次write类只允许重试一次。接收端Agent需要实现一个/invoke接口这是Agent-Reach对外的唯一协议约定。有一个小坑大家可能容易踩httpx.AsyncClient如果每次dispatch都重新创建在高并发下会有TCP连接重复建连的开销。我的做法是在Gateway初始化时创建一个全局的AsyncClient实例并复用但要注意把timeout参数移到client级别而且每次请求如果传了不同timeout复用同一个client会导致后设置的timeout全局生效。这个问题我查了好半天才定位解决方案是改用httpx.AsyncClient(timeouthttpx.Timeout(None))然后在单次请求时用timeoutxxx覆盖。代码里“简洁起见”是每次new生产版本建议做一个连接池。3.5 部署联通性测试我在三台虚拟机上各部署一个模拟Agent分别声明了query-sales、generate-daily-report、query-stock三个能力再加上一台机器跑Agent-Reach核心服务。启动Agent-Reach后依次执行以下测试注册三个Agent查询注册中心日志确认状态均为online正常心跳运行5分钟确认无suspected误报调用一个不在任何Agent能力声明中的意图验证返回明确的“无匹配Agent”错误手动kill掉一个Agent观察10秒后其状态变为suspected再过15秒变为offline在offline期间发起对应意图的调用确认请求不进入网关直接返回失败重启kill掉的Agent确认注册中心自动复用agent_id并恢复online状态压测用50并发持续发query-sales请求30秒观测丢包率和平均响应时长这轮测试做下来我发现了一个重要问题超时重试在某些场景下依然会带来重复请求比如目标Agent其实已经处理完请求只是响应慢了导致客户端超时才重试。这个问题我在第四部分结合幂等设计细讲。4. 常见问题与排查技巧实录4.1 心跳抖动导致的误杀问题上线第一周就遇到一次“Agent集体离线”事故。现象是某个Agent发出大量请求成功、但部分请求返回超时排查时发现注册中心里那个Agent的状态已经被置为offline了。分析日志发现那段时间Agent所在机器有一次短暂网络抖动持续了约20秒心跳全部没送过来。按照我最初“30秒没心跳就离线”的设计它果然被误杀了。后来我把扫描逻辑改成了两步客户端连续丢两个心跳进入suspected再丢一个才转offline相当于给“误杀”加了一道缓冲。同时把心跳上报接口的耗时压到毫秒级确保Agent侧没有因为接口慢而漏发。改造之后同样场景复测Agent只被标记suspected没有离线也没有引发大规模请求重试。这个案例给我的经验是心跳阈值不要拍脑袋定最好结合Agent部署环境的网络状况来调。像我们虚拟机环境偶尔会有几秒的延迟阈值就适当放宽如果你跑在专有网络里延迟普遍低于10毫秒那阈值可以收紧到15秒。4.2 超时重试引发的“重复执行”事故这是压测阶段暴露出来的另一个坑。有个Agent实现了“工单自动分类”能力属于写操作。我一开始把它的重试配置设成了2次结果压测时出现大量重复工单。原因是目标Agent处理分类请求需要3秒而我根据latency_ms配置的请求超时是2.4秒导致请求已经发出、目标Agent正在处理客户端就判定超时并重试于是同一个工单被分类了两次。排查方案分两步走。第一步把所有非幂等动作的重试次数降到1次或者0次。工单分类这个动作虽然没有直接写数据库但会触发下游通知重试的成本很高。第二步在Agent侧增加请求幂等校验。Agent的/invoke接口接收到请求时先看请求头里的x-request-id是否已经处理过。这里我在网关里生成了x-request-id并随请求透传Agent侧用一个Redis或者SQLite表存最近1000个已处理请求ID重复请求直接返回上一次的结果。加了这两道保险之后重复执行问题彻底消失。从结果看超时重试这种机制必须和幂等设计配套使用光靠客户端配置重试次数是堵不住的。4.3 调用链日志联调的心得Agent-Reach上线之后不同部门开始接入各自的Agent。很快遇到一个问题消息在A、B两个Agent之间转手出错一边说是对方没返回另一边说是根本没收到互相扯皮。原因是每次转发的调用链路上没有统一的标识。我在网关和注册中心之间加了trace_id的生成和透传机制调用方发起请求时如果没有附带trace_id网关就生成一个并透传给目标Agent目标Agent收到后在其所有内部日志里都带上这个trace_id。这样一条请求从入口到出口所有环节的日志可以通过trace_id一次性捞出来。日志格式统一为JSON包含trace_id、action、caller、target、status、cost_ms、timestamp几个字段。用grep或者日志平台搜索trace_id一秒就能定位是哪一环出了问题。4.4 注册中心的并发写锁问题SQLite本身是单写者的虽然我用threading.Lock保证了单进程内的互斥但Agent数量增多、心跳频率提高之后还是出现了注册接口变慢的现象。压测到200个Agent实例时注册接口的平均响应从0.5毫秒涨到15毫秒。分析发现瓶颈在INSERT OR REPLACE操作每次都要重写整行数据能力声明表如果字段特别多写放大严重。后来我把心beat和regiter拆成了两张表agent_meta存静态信息agent_heartbeat只存agent_id和last_heartbeat。心跳只更新第二张表写入量缩小了一个数量级。这个优化后压测500个Agent实例心跳写入稳定在2毫秒以内。如果你预计Agent规模过万那还是建议直接换Redis或者MySQLSQLite适合小规模场景。5. 后续扩展思路与个人经验Agent-Reach这个项目目前还在迭代中有几块内容已经在计划表里。第一块是Agent的优雅上下线。现在Agent要下线时是直接kill进程完全依赖TTL被动发现。更好的做法是支持一个健康检查接口让路由引擎在摘流之前主动探活一次做优雅摘流避免正在处理的请求被中断。第二块是联邦调度。我们有两个团队各自维护了一套Agent集群将来要做跨集群调度。目前想到的方案是Agent-Reach之间通过gossip协议相互同步“能力摘要”即只传“某个集群里有谁会做什么”不传具体的节点心跳实际请求到达后再做二级路由。这样既保留了集群自治又能实现跨集群能力触达。第三块是更细粒度的可观测性。想加入一个请求样本采样功能对每个进入网关的请求按1%的比例记录完整的协议栈耗时包括注册中心查询时间、路由匹配时间、网关网络IO时间、目标Agent处理时间。有了这些数据就能精确定位“慢”到底是哪个环节慢的。最后说几个我这几个月实操下来最深刻的体会。第一Agent调度这类基础设施核心价值是“把不确定性收拢到一个地方”。不管是目标Agent宕机、响应超时还是重试风暴只要有统一的中转层故障的范围就是可控的。分散在各自Agent里的重试和超时逻辑越多整个系统越难收敛。第二所有超时和重试参数一定要经过压测验证之后再固化。我刚开始上线时用的参数都是拍脑袋估的结果压测一上就出了重复执行问题。后来每个动作类型的latency_ms都是让Agent负责人自己上报网关再按倍率分配超时这样就实现了“各人给自己的服务设预期”。第三如果重新来一遍我会更早设计“可观测性”。战斗力和调试效率的差距往往就差一个贯穿全链路的trace_id和结构化日志。这东西越早做后面排查问题的时候越省心。Agent-Reach现在还在我们内部跑着时不时还会有新坑冒出来。但到目前为止它已经把“Agent之间怎么互相找到、怎么对话、怎么保证可靠”这个问题从“靠人肉沟通”变成了“靠统一平台解决”。如果你也在做类似的多Agent系统希望这篇拆解能给你一个足够清晰的起点。
返回列表