ARTICLE DETAIL

资讯详情

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

Agent-Reach实战:多Agent注册发现与消息路由的轻量级方案

Agent-Reach实战:多Agent注册发现与消息路由的轻量级方案 最近在不少技术社群里频繁看到Agent-Reach这个词被反复提到尤其是做 AI Agent 落地的人见面几乎都会聊一句你手里这批 agent到底能不能稳定地互相“喊得通”我一开始以为这是个网络层面的连通性问题真正动手做了一轮之后才发现它远比想象中复杂但也没有复杂到需要自研一套消息中间件的地步。这篇文章把我从设计到落地、再到填坑的全过程整理出来尽量说人话给同样在搞多 Agent 协作的人一份可以直接抄作业的参考。1. Agent-Reach 到底在解决什么问题1.1 Agent 之间为什么需要“可达性”现在每个团队手里都有不少 Agent有基于 LLM 的对话代理有做自动化流程的 RPA 式 Agent有负责定时任务的调度 Agent还有嵌入到业务系统里的决策 Agent。这些 Agent 往往由不同小组、不同时间、甚至不同技术栈开发出来孤岛效应非常明显。想让 A Agent 把任务转给 B Agent第一件事不是“协议怎么定”而是“B Agent 现在活着吗它到底能干什么我怎么找到它”这就是 Agent-Reach 的核心场景为异构 Agent 提供一套注册、发现、健康检查和消息投递的轻量基础设施。它解决的核心问题有三个管理 Agent 的注册信息你是谁、在哪、能做什么判断 Agent 是否存活心跳与租约机制把消息准确投递给目标 Agent按能力或实例路由1.2 它和微服务注册中心的区别如果你熟悉 Spring Cloud Eureka 或者 Consul会觉得 Agent-Reach 有点眼熟。但严格来说它和传统注册中心有几个关键差异这也是我一开始没想明白的地方维度传统微服务注册中心Agent-Reach注册对象无状态的 HTTP 服务实例有状态、带意图和能力的 Agent发现粒度按服务名找实例按能力语义找 Agent或按 ID 精确找健康判定端口探活或心跳心跳 会话级 ACK 任务完成状态消息模型通常是请求/响应支持请求/响应也支持单向任务派发路由依据负载均衡能力匹配 实例存活 上下文亲和性简单类比传统注册中心相当于公司总机按分机号转接Agent-Reach 更像是一个带“技能标签”的协作通讯录你不仅能找到人还能知道这个人擅长什么、当前有没有在忙然后再决定要不要把活儿派给他。2. 整体架构设计注册 心跳 路由三件套2.1 注册中心的数据模型我在设计 Agent-Reach 时第一个决策是注册表里到底存什么。字段太少后面做能力匹配会很痛苦字段太多Agent 接入成本就会变高没人愿意接。最终我定下来的注册信息包含以下几部分agent_id全局唯一标识格式统一为namespace/name/uuid例如ops/email_agent/a1b2c3避免不同业务线之间冲突。endpointAgent 暴露的回调地址可以是 HTTP URL、gRPC 地址也可以是内部消息队列的 topic。capabilities能力描述列表用domain:action:version三段式表达比自由文本更利于做匹配。metadata附加信息比如当前负载、可用并发数、所属区域、优先级等。lease租约过期时间由服务端计算客户端不传这个字段防止时钟偏移。注册接口的请求体长这样{ agent_id: ops/email_agent/a1b2c3, endpoint: http://10.20.30.40:8080/callback, capabilities: [email:send:v1, email:batch_send:v1, template:render:v2], metadata: { region: cn-east-1, max_concurrency: 8, preferred_queue: high_priority } }这里有一个我踩过坑的点capabilities 的命名规范一定在第一天就定死。我们早期没有强制校验结果email/send、email.send、email_send三种写法都出现在注册表里路由模块为了兼容这三个版本写了整整一屏的匹配逻辑。后来我直接在注册入口做了正则校验只允许[a-z][a-z0-9_]*:[a-z][a-z0-9_]*:v\d这种格式不合法直接 400问题瞬间消失。2.2 心跳与租约如何判定 Agent 还活着Agent 是典型的“可能随时死掉”的进程比如 LLM 服务超时导致进程卡死、容器被 OOM Kill、网络分区导致消息送不达。如果注册表里充满了死 Agent路由层会把大量任务投进黑洞。所以我用了租约机制而非永久注册Agent 启动后调用register注册自己之后每隔固定时间发送心跳heartbeat(agent_id)服务端记录last_heartbeat_at并计算lease_expire_at last_heartbeat_at lease_timeout路由模块在选择目标时只选lease_expire_at now的 Agent一旦租约过期注册表异步清理该 Agent并触发下线通知。心跳间隔和租约超时怎么定我参考了常见的分布式系统经验给了一组适合大多数业务场景的参数参数建议初值依据心跳间隔5 秒大多数 Agent 任务耗时在秒级5 秒内能感知死亡兼顾网络开销租约超时15 秒容忍 3 次心跳丢失避免网络抖动导致误摘除清理周期1 秒定时扫描过期租约及时排除死节点重试注册间隔10 秒注册失败后不至于频繁打爆注册中心这个参数组合在 200 个 Agent 的规模下实测很稳注册中心的压力完全可以忽略。如果你的 Agent 数量到了几千这个量级心跳间隔可以根据业务容忍度适当拉长到 10 秒削减无谓的心跳 QPS。2.3 消息路由与重试策略路由是 Agent-Reach 最核心的执行路径。调用方有两种玩法按 ID 直投明确知道要找哪个 Agent只做存活检查 投递适合任务交接场景。按能力路由不关心具体是哪个 Agent 干活只要有人具备email:send:v1能力就行适合任务分发场景。我实现的是两阶段路由1. 根据消息头里的 target_modeby_id / by_capability找到候选 Agent 列表 2. 过滤掉租约过期、状态为 busy 且无队列容量的 Agent 3. 按策略排序同区域优先、负载最低优先 4. 选中最优 Agent调用其 endpoint 投递消息投递过程中的重试逻辑我做得比较保守retry_delays [0.1, 0.2, 0.4, 0.8] # 单位秒 max_attempts 5 for attempt in range(max_attempts): try: resp requests.post(endpoint, jsonmessage, timeout3) if resp.status_code 202: return accepted if resp.status_code in (404, 410): mark_agent_stale(agent_id) return agent_unreachable except requests.Timeout: pass time.sleep(retry_delays[attempt]) return failed_after_retries这里隐藏了一个容易被忽略的细节收到 202 只代表 Agent 接了消息不代表它成功处理完了。所以 Agent-Reach 只承诺“送达”任务是否执行成功需要第二层确认——由 Agent 在任务完成后回调一个task_complete事件。这个设计把消息投递和业务执行解耦避免投递系统变成一个分布式事务管理器复杂度骤降。3. 核心实现与关键代码3.1 注册中心的服务端实现我用 Python FastAPI 实现了注册中心的 MVP核心就是三张内存表加一个后台清理任务。这里为了快速验证设计是否正确没有一上来就引入 Redis 和数据库内存表配合锁在单机几千个 Agent 的场景下完全够用。from fastapi import FastAPI, HTTPException from pydantic import BaseModel, Field, validator import time, threading, re app FastAPI() lock threading.Lock() agents {} # agent_id - agent record capability_index {} # capability - set of agent_id class AgentRegistration(BaseModel): agent_id: str Field(..., min_length3, max_length128) endpoint: str capabilities: list[str] metadata: dict Field(default_factorydict) validator(capabilities, each_itemTrue) def check_capability_format(cls, v): if not re.match(r^[a-z][a-z0-9_]*:[a-z][a-z0-9_]*:v\d$, v): raise ValueError(finvalid capability format: {v}) return v app.post(/register) def register(reg: AgentRegistration): with lock: now time.time() agents[reg.agent_id] { endpoint: reg.endpoint, capabilities: reg.capabilities, metadata: reg.metadata, last_heartbeat_at: now, lease_expire_at: now 15, } for cap in reg.capabilities: capability_index.setdefault(cap, set()).add(reg.agent_id) return {status: registered, lease_expire_at: agents[reg.agent_id][lease_expire_at]} app.post(/heartbeat) def heartbeat(agent_id: str): with lock: rec agents.get(agent_id) if not rec: raise HTTPException(status_code401, detailagent not registered) now time.time() rec[last_heartbeat_at] now rec[lease_expire_at] now 15 return {status: ok} app.post(/route) def route(message: dict): target_mode message.get(target_mode) if target_mode by_id: candidates [message[agent_id]] elif target_mode by_capability: candidates list(capability_index.get(message[capability], set())) else: raise HTTPException(status_code400, detailinvalid target_mode) # 过滤过期租约 now time.time() live [aid for aid in candidates if aid in agents and agents[aid][lease_expire_at] now] if not live: return {status: no_reachable_agent} # 简单策略取第一个生产环境可改为负载优先 target live[0] return {status: routed, target: target, endpoint: agents[target][endpoint]}这个版本刻意保持简单核心逻辑全都在锁内部单线程安全。但我必须提醒一个性能隐患锁粒度太大。当心跳频率升高、路由请求变多时这把全局锁会成为瓶颈。实测中 2000 个 Agent、5 秒心跳间隔的场景下还没有出现问题但如果你准备扩容到万级 Agent建议把注册表改成分段锁或者直接换用 Redis hash sorted set 方案。3.2 租约清理线程的正确写法后台清理任务有一个常见的坑把“扫描间隔”和“租约超时”混为一谈。扫描线程每 1 秒跑一次但并不代表 15 秒一到节点立刻被清理真正的清理时机是lease_expire_at now扫描只是碰巧发现并回收。def cleanup_loop(): while True: time.sleep(1) now time.time() with lock: expired [aid for aid, rec in agents.items() if rec[lease_expire_at] now] for aid in expired: rec agents.pop(aid, None) if rec: for cap in rec[capabilities]: capability_index.get(cap, set()).discard(aid)这里要注意清理动作必须同时更新capability_index否则会留下幽灵索引路由时明明列表里有这个 Agent取到之后却查不到记录返回的还是一个错误的 endpoint。我刚开始忘了做这一步调试了半天才定位到是索引残留的问题。3.3 客户端 SDK 的心跳与优雅退出服务端写好了客户端也得配套一个 SDK否则每个 Agent 都要自己实现心跳逻辑接入成本太高。我在 SDK 里做了三件事启动时注册注册失败就按指数退避重试后台线程按固定间隔发心跳捕获 SIGTERM 信号时主动调用unregister把下线动作从被动等待租约过期变成主动通知。主动注销看起来是一件小事但作用很大。被动清理意味着 15 秒内路由层可能还会把消息发给一个已经死掉的 Agent而主动注销可以把这个时间窗口压缩到几百毫秒对任务延迟敏感的场景非常关键。import signal, threading, time, requests class AgentReachClient: def __init__(self, registry_url, agent_id, endpoint, capabilities): self.registry_url registry_url self.agent_id agent_id self.endpoint endpoint self.capabilities capabilities self._stop threading.Event() def start(self): self._register_with_retry() threading.Thread(targetself._heartbeat_loop, daemonTrue).start() signal.signal(signal.SIGTERM, self._shutdown) def _register_with_retry(self): payload {agent_id: self.agent_id, endpoint: self.endpoint, capabilities: self.capabilities} for delay in [1, 2, 4, 8, 16]: try: resp requests.post(f{self.registry_url}/register, jsonpayload, timeout3) if resp.status_code 200: return except requests.RequestException: pass time.sleep(delay) raise RuntimeError(register failed after retries) def _heartbeat_loop(self): while not self._stop.is_set(): try: requests.post(f{self.registry_url}/heartbeat, json{agent_id: self.agent_id}, timeout3) except requests.RequestException: pass self._stop.wait(5) def _shutdown(self, *args): try: requests.post(f{self.registry_url}/unregister, json{agent_id: self.agent_id}, timeout2) finally: self._stop.set()4. 实测中高频问题与排查思路4.1 Agent 注册成功但始终路由失败这个现象最迷惑人注册接口返回 200心跳也在发但路由结果一直是no_reachable_agent。我排查下来的原因有几种按照出现频率从高到低排列现象根因解法路由按能力找不到目标capability 大小写/分隔符不一致注册入口强制校验格式路由按 ID 找不到目标调用方传的 agent_id 带多余前后缀注册和调用共用同一个 ID 生成函数路由返回 no_reachable_agent租约已过期但客户端以为还在线查看服务端日志中 last_heartbeat_at 是否在更新消息发到旧 endpointAgent 重启后 IP 端口变化但未重新注册SDK 强制要求在启动流程中先注册再对外服务第四个原因特别隐蔽。容器化部署时 Agent 重启会拿到新的 IP如果注册表里还是旧地址路由照常返回目标但消息永远投不出去。我在客户端 SDK 里加了一个保护机制注册成功后返回的 lease_expire_at 如果和本地预期差太多说明注册表之前有残留记录直接先调 unregister 再重新注册。4.2 心跳丢失导致 Agent 被误摘除网络抖动是心跳机制的天敌。我当时在本地实验室跑得好好的一上生产就出现 Agent 频繁被摘除、又频繁注册回来的现象像抽风一样。排查后发现是两条链路叠加的问题生产环境某些节点的网络延迟偶发超过 5 秒Agent 的心跳发送本身有 3 秒超时超时后线程直接跳过这一轮两条叠加导致两次心跳间隔实际上达到了 8~10 秒超过 15 秒租约的容错空间吗并没有但如果是连续两次抖动叠加就会出现误判。解决办法是双管齐下心跳超时从 3 秒提到 2 秒给下一轮心跳留出更多时间租约超时从 15 秒提到 20 秒。不要小看这 5 秒的余量它让系统在不牺牲太多感知速度的前提下扛住了偶发的网络抖动。4.3 消息重复投递消费端该如何自保Agent-Reach 的投递语义是 at-least-once这就意味着极少数情况下消息会被重复发送。最典型的场景是Agent 成功处理完消息但返回的 202 响应在网络传输中丢失客户端 SDK 判断超时触发重试同一消息就被投了两次。重复消费对很多 Agent 是致命的比如发邮件、转账、调外部 API。我的方案是在消息结构里增加message_id消费端用一个去重表来判断是否见过这条消息import redis r redis.Redis(hostlocalhost, port6379, db0) def is_duplicate(message_id: str, ttl_seconds: int 3600) - bool: # 消息处理时长不会超过 ttl超过则视为新消息 result r.set(fdedup:{message_id}, 1, nxTrue, exttl_seconds) return result is False这套去重逻辑虽然简单但非常管用。唯一要注意的是 Redis 的SET NX EX在极端情况下的语义如果两个请求同时进入只有一个能成功写入另一个直接返回 False正好满足去重需求。4.4 时钟偏移与租约计算分布式系统里最经典、也最容易被新手忽略的问题就是时钟偏移。我一开始图省事让客户端把本地时间戳放进心跳报文服务端用这个时间戳计算租约。结果在生产环境出现了两个诡异现象部分 Agent 注册后立刻过期部分 Agent 租约永远不过期即使进程已经死了。根因是服务器的系统时间没有同步有的快了 30 秒有的慢了 2 分钟。快速 Agent 的时间比服务端快 30 秒它发的心跳时间戳看起来像是未来时间服务端算出的租约虽然还在但清理线程用的是服务端本地时间一对比就认为过期了。慢的 Agent 则反过来永远在“续命”。解决办法很简单客户端在心跳报文里只传 agent_id不传任何时间戳租约计算完全基于服务端本地时间。这应该是分布式系统的设计铁律时间以服务端为准客户端永远不要参与时间计算。5. 规模扩展的瓶颈与应对方案5.1 单机注册中心的极限在哪里我简单压测过单机 FastAPI 注册中心在 2000 个 Agent、5 秒心跳间隔的情况下心跳 QPS 大约是 400加上路由请求压力非常小。真正吃资源的是消息投递本身但那个环节可以水平扩展。注册中心要扩容的话第一个瓶颈是内存第二个是锁竞争。内存方面每个 Agent 的注册记录大约占用 500 字节1 万个 Agent 也就是 5MB 级别完全不是问题。锁竞争才是真正的瓶颈全局锁在路由高峰期会拖慢心跳处理。如果一定要支撑万台规模我有两个方向对agent_id做一致性哈希分片每片一个独立的注册表实例注册表数据放进 Redis Hash利用 Redis 单线程原子性替代应用层锁索引关系用SADD维护。5.2 消息可靠性与任务追踪Agent-Reach 做得再多本质上也只是个“投递层”。任务发出去了结果怎么样Agent 处理失败怎么办这些问题不能全部让业务方自己解决否则 Agent-Reach 的价值就少了一半。我在路由模块里顺手做了一个trace_id贯穿机制每条消息生成一个trace_id路由日志、Agent 请求头、任务完成回调都带上这个 ID。排查问题的时候只要拿trace_id一查整个链路的每一跳都清清楚楚。这是我在 Agent-Reach 里最庆幸做对的一个设计建议每个做 Agent 协作系统的人都把 trace 提前加上不要等线上出了问题再后悔。5.3 动态扩缩容与注册风暴最后提一个容器化时代的特殊问题Agent 批量重启时会同时向注册中心发送注册请求形成注册风暴。几十个 Agent 还好如果是几千个一起重启注册中心瞬间 QPS 可能冲到上千容易把服务打挂。我的解法是在注册入口加了一个简单的令牌桶限流import time from collections import deque class TokenBucket: def __init__(self, rate, capacity): self.rate rate self.capacity capacity self.tokens capacity self.updated_at time.time() def acquire(self): now time.time() self.tokens min(self.capacity, self.tokens (now - self.updated_at) * self.rate) self.updated_at now if self.tokens 1: self.tokens - 1 return True return False实测效果很好。注册请求被限流后返回 429客户端 SDK 会按退避策略重试不会丢数据。配合上文的客户端指数退避机制注册风暴基本能被消化掉。6. 经验总结与个人体会Agent-Reach 这个项目做下来我的最大感受是做 Agent 协作基础设施难的不是技术而是对“不确定性的容忍设计”。Agent 不是普通服务它会死、会卡、会超时、会返回非结构化结果你必须在每一层都把“它可能失败”当作默认前提来设计。注册表要有租约而不是永久记录投递要有重试而不是一次了事消费要有去重而不是指望发送端绝对可靠。另一个体会是协议约定一定要前置。capabilities 命名规范、agent_id 格式、trace_id 传递链路这些看起来都是小事但等到几十个 Agent 接入之后再想统一成本会翻好几倍。我见过太多项目栽在“先跑通再说”这四个字上最后都在协议兼容上耗费了大量精力。最后分享一个实用的小技巧在开发阶段把注册中心的管理接口开放成一个简单的 HTML 页面——能看到所有在线 Agent、手动模拟心跳丢失、强制摘除某个节点。这个页面我一开始只想用来调试结果后来变成团队排查问题最常用的工具几乎人手一个书签。如果你也在做类似 Agent-Reach 的系统建议把这步加上调试体验会好很多。
返回列表