从"智能体孤岛"到"触达网络",是我这段时间做Agent-Reach最核心的感受。团队里同时跑着七八个智能体,有做舆情摘要的、有写周报草稿的、有盯代码评审提醒的、还有处理客服工单分类的,单个拿出来都能干活,但让它们协作起来就非常痛苦:任务不知道怎么分给谁、某个 Agent 挂了没人知道、调用方为了等一个响应愣是把超时设成了 60 秒。最后我干脆做了一层统一调度与触达基础设施,把"发现智能体、把请求送过去、确认对方有没有收到、拿到反馈"这四件事一次性解决。这东西适合任何手头有多个 Agent、又不想靠人肉编排的团队参考,接下来的内容就是完整的设计思路和踩坑记录。
1. 从"智能体孤岛"到"触达网络":Agent-Reach 要解决的真实问题
1.1 团队场景里最常见的分工困境
先描述一个大概率你也会遇到的画面。公司内部陆续上线了几个 AI Agent,每个都是独立项目、独立部署、独立 API Key,甚至超时和重试策略都不一样。运维那边有个"故障诊断 Agent",客服那边有个"话术推荐 Agent",数据分析组有个"报表解读 Agent"。这些 Agent 各自服务好各自的业务方,看似没什么问题,但如果出现跨部门任务,比如客服人员希望直接调用运维知识库来答复用户,就麻烦了。
你作为调用方,需要自己查文档找到运维 Agent 的地址,自己处理对方 API 的鉴权方式,自己设置超时时间,还得自己处理对方频繁升级导致的接口变化。更麻烦的是,你根本不知道对方 Agent 当前是繁忙还是空闲,是活着还是已经悄悄崩溃。所谓"触达",在大多数团队里就是靠企业内部聊天工具把一堆接口文档传来传去,这就是典型的"智能体孤岛"。
我做Agent-Reach的出发点非常简单:让调用方不再关心"目标 Agent 是谁、地址是什么、协议是什么",只需要描述"我要什么",剩下的发现、连接、确认、反馈全部由这一层完成。说得直白一点,它是一张把各种智能体连接起来的触达网络,而不是一个个单独运维的孤岛。
1.2 为什么消息队列解决不了路由问题
有人会问,这用消息队列不就行了吗?Kafka、RabbitMQ 都能做异步解耦,Agent 把消息往队列里一扔,另一端订阅消费,不也算触达吗?说实话,最初我也想直接套消息队列,但很快发现两者思考问题的层次完全不同。
消息队列解决的是"消息怎么从 A 可靠地到达 B",它不关心 B 是不是有能力处理这条消息,也不关心现在是否有另一个更合适的 C 在线。你仍然要在生产端硬编码队列名,本质上是人肉指定路由。一旦 Agent 上下线、能力升级、负载变化,队列名不会跟着变,调用方代码得跟着改。
Agent-Reach把"路由语义"提到了核心位置。它不看队列名,而是看意图、能力和状态。比如调用方说"我要做新闻摘要",Reach 层去看注册表里有哪些 Agent 声明了summarize_news能力,再结合当前心跳、负载、历史成功率选择一个最合适的。队列只负责传输,不负责决策;而我要做的这一层,恰恰把决策放在传输之前。
1.3 触达不是通信,是四件事的闭环
很多设计失败的系统,问题出在只做了"通信"这个动作,没有形成闭环。我在最开始就把触达拆成四个环节,每个环节单独设计:
| 环节 | 要回答的问题 | 失败的表现 |
|---|---|---|
| 发现 | 谁能做这件事? | 不知道调谁 |
| 连接 | 怎么联系上它?协议对吗? | 地址错误、鉴权失败 |
| 确认 | 它收到并开始处理了吗? | 消息丢了没人管 |
| 反馈 | 处理结果是什么?成功还是失败? | 调用方干等 |
任何一个环节断了,对整个业务来说就是一次触达失败。举个例子,假设你成功把请求发给了某个 Agent,对方也确实收到了,但处理到一半崩溃了,没有返回任何结果。这种事情用轮询接口根本察觉不到,需要有一层机制去发现"发出去的任务没有按期完成",然后触发重试或降级。所以Agent-Reach表面看是个调度器,实际上它是个"带状态感知和闭环确认的任务分发系统"。
2. Agent-Reach 的三个关键抽象:注册表、意图路由、连接保鲜
2.1 注册表:每个智能体的"身份证"与能力描述
注册表是整个系统最基础也最重要的抽象,它决定了一个 Agent 如何被"知道"。我设计的注册信息不是简单存一个地址,而是一份结构化描述,每个接入的 Agent 需要上报以下字段:
{ "agent_id": "agent_ops_diagnosis", "name": "运维故障诊断助手", "version": "2.1.0", "capabilities": [ { "name": "diagnose_fault", "params_schema": "https://schema.example.com/diagnose_fault.json", "routing_weight": 80 }, { "name": "analyze_log", "params_schema": "https://schema.example.com/analyze_log.json", "routing_weight": 60 } ], "protocols": ["http", "mcp"], "endpoint": "http://10.21.33.45:8080", "heartbeat_interval": 15, "max_concurrency": 20, "tags": ["运维", "故障", "日志"] }这里有几个设计细节值得展开。第一,capabilities数组是路由的核心依据,它不是给人看的描述,而是机器可读的能力声明,每个能力都带一个params_schema指向参数校验的 JSON Schema。第二,protocols字段非常关键,它告诉 Reach 层这个 Agent 能接受哪几种协议,后面路由时会根据调用方的能力做协议协商。第三,routing_weight是给路由引擎做同能力多 Agent 排序用的,但注意它只是初始权重,实际运行时分数还会被心跳状态、失败率动态调整,后面我会专门讲。
注册表底层我用 Redis Hash 存储,key 是agent:{agent_id},字段就是上面对应的 JSON 字符串。之所以不用专门的数据库,是因为注册表的读操作极其频繁,路由引擎每次请求都要扫一轮在线 Agent,Redis 可以把该操作压到毫秒级。
2.2 意图路由:从"调接口"到"说目的"
传统接口调用是"我知道你是谁,我要调你某个方法",但 Agent 协作不是这样。业务方通常只知道目的,比如"把这段日志分析一下",他并不知道该找诊断 Agent 还是日志分析 Agent。意图路由要解决的就是这个转换。
我的做法是两层匹配。第一层是标签与关键词匹配,属于粗筛。系统会把调用方的意图文本做分词,和注册表的name、tags、capabilities.name做向量相似度匹配,返回一个候选列表。第二层是结构化参数匹配,属于精筛。调用方提交任务时如果附带结构化参数,Reach 会拿这些参数逐一校验候选 Agent 的params_schema,校验不过的直接淘汰。
匹配逻辑用 Python 写核心部分大约长这样:
def match_candidates(intent: str, params: dict, registry: list) -> list: score_by_agent = {} for agent in registry: if not is_agent_healthy(agent["agent_id"]): continue capability = get_best_capability(agent, intent, params) if capability is None: continue score = capability.routing_weight * 0.6 + \ semantic_similarity(intent, capability.name) * 100 * 0.4 score_by_agent[agent["agent_id"]] = score return sorted(score_by_agent.items(), key=lambda x: x[1], reverse=True)实际运行中我发现一个有趣的规律:语义相似度这个指标不能只看意图和能力名的相关性,还要结合调用方历史上成功调用过哪些 Agent。所以我加了一个"最近 7 天成功率"作为额外因子,权重不高,但可以避免每次都因为初始权重高而把请求打给一个总是超时的 Agent。
2.3 连接保鲜:心跳、探活与状态感知
注册信息只是静态快照,运行中 Agent 的状态是不断变化的。连接保鲜这一层的目标就是保证路由引擎不会把请求发给一个已经失联的节点。我同时启用了主动心跳和被动探活两条链路。
主动心跳很容易理解,每个 Agent 每 15 秒调用一次 Reach 的心跳接口,上报自己的存活状态和当前 backlog。心跳延迟超过 45 秒,也就是连续 3 个周期没收到,注册表就标记该 Agent 为SUSPECT,此时路由权重直接降到 0,不会分配新任务,但已分配的任务仍允许完成。再过一个周期仍未恢复,就标记DEAD,从活动列表移除。
被动探活是我后来加上的补丁。因为有些 Agent 进程存活、心跳正常,但内部线程池已经阻塞,任何任务进来都是超时。这种"僵尸 Agent"靠心跳发现不了。我的办法是绕开心跳,由 Reach 侧定期发送一个轻量级探测任务,比如让 Agent 花 100 毫秒算一下ping并返回当前队列深度,如果连续两次探测响应时间超过 5 秒,就强制摘除。
这里我踩过一个大坑:探活频率太高会把 Agent 压垮。最初我设成每 5 秒探测一次,结果某个 Java 写的 Agent GC 压力暴涨。后来改成每 60 秒探测一次,只探出了 30% 的僵尸,但整体稳定性好很多。这个频率最终调成 30 秒,并且不同 Agent 可以配置不同的探测间隔。
3. 我把第一版跑通用的最小实现:代码骨架与实测表现
3.1 技术选型:为什么是 Python + FastAPI + Redis Streams
先说结论,这套选型适合中小规模,团队在 10 人以下、日请求量在百万以内的场景。如果你们的 Agent 数量在 50 个以上、单日请求千万级,我建议直接上 Kubernetes 事件驱动架构,但这篇文章先聊最小可用版本。
选 Python 是因为团队里多数 Agent 就是 Python 写的,统一语言能减少接入成本。FastAPI 的好处是自带 OpenAPI 文档,每个 Agent 接入时可以很直观地看到一个 POST 接口长什么样,团队协作效率会高很多。Redis Streams 则是我对比了 Kafka 之后的选择,原因有三个:
第一,运维成本几乎为零,我们已经有 Redis 了,不需要再额外维护 Kafka 集群。第二,Redis Streams 自带消费者组和 pending entries 机制,天然适合"任务分发 + ACK 确认 + 重新投递"这个场景。第三,千万级以下的消息量 Redis Streams 完全扛得住,瓶颈根本不在传输层,而在每个 Agent 的实际处理能力。
3.2 注册接口与能力描述 Schema
接入 Agent-Reach 第一步就是注册。我给每个 Agent 暴露一个注册接口,Agent 启动时调用,之后每 15 秒通过心跳续约。核心代码如下:
@app.post("/agent/register") async def register_agent(req: AgentRegisterRequest): agent_id = req.agent_id # 用 Redis Hash 存储注册信息,TTL 设为 60 秒 reg_key = f"agent:{agent_id}" data = json.dumps(req.dict(), ensure_ascii=False) await self.redis.hset(reg_key, "info", data) await self.redis.expire(reg_key, 60) # 设置心跳续约的回调地址 await self.redis.hset(reg_key, "last_heartbeat", time.time()) resource = { "agent_id": agent_id, "capabilities": req.capabilities, "legacy_route": f"/agent/{agent_id}/invoke", } return {"status": "ok", "resource": resource}能力描述 Schema 我单独抽了一个 Pydantic 模型:
class Capability(BaseModel): name: str description: str = "" params_schema: dict routing_weight: int = 50注意我加了一个legacy_route字段,这是为了兼容那些不想改造自身协议、只提供 REST 接口的老系统。注册表允许一个 Agent 同时声明多种协议,实际调用时由适配器去翻译。
3.3 路由引擎的匹配逻辑
路由引擎是独立于 API 服务的一个消费者进程,从 Redis Streams 里读任务,经过匹配、选路、投递三步完成一次触达。选路的核心我上面已经给了match_candidates函数,但投递环节还需要一个包装层。
async def deliver(agent_info: dict, task: AgentTask): protocol = agent_info["protocols"][0] adapter = get_adapter(protocol) try: result = await adapter.invoke(agent_info["endpoint"], task) # 投递成功后,写回执 await task_stream_ack(task.task_id, result) except AgentTimeoutError: await task_stream.nack(task.task_id, reason="timeout") raise这里我刻意把"投递成功"和"处理成功"区分开。投递成功只代表请求已经送达到 Agent 网关,不代表业务逻辑完成。真正的完成信号由 Agent 在处理完毕后通过回调接口上报,或者由调用方在约定时间内主动查询。这个区分非常重要,避免了很多误判。
3.4 调用方 SDK 的一次请求生命周期
最后给调用方看的是一个小 SDK。我希望业务方调用 Agent-Reach 就像调用一个本地函数:
from agent_reach import Client client = Client(endpoint="https://reach.internal.example.com") async def get_news_summary(news_text: str): task_id = await client.submit( intent="summarize_news", params={"text": news_text}, timeout="normal", ) result = await client.wait_for_result(task_id, timeout=15) return result["summary"]submit时只需要提供意图和参数,不需要指定任何 Agent。SDK 内部会自动完成注册发现、路由匹配、协议协商、投递、回执等待整个链路。如果第一个选中的 Agent 超时,SDK 会自动在候选列表里选下一个重试,最多重试两次。调用方拿到的只有 task_id 和最终结果,这个抽象让业务代码非常干净。
4. 灰度上线后踩到的坑:超时雪崩、重复执行和僵尸 Agent
4.1 坑一:慢响应 Agent 把调用方线程池占满
第一版里所有请求默认超时 30 秒,看起来合理,但上线第三天就出了事故。有个数据分析 Agent 在最坏情况下要跑 40 秒,于是调用方侧大量线程阻塞在等待响应上,线程池被占满,连那些本来可以 200 毫秒就返回的请求也全被拖死。
这个问题的根子在于我用了"一刀切"的超时策略。后来我把任务按预期耗时分成三类,分别配置线程池和超时,这个我下一节详细讲。触发这次事故之后我才意识到,超时策略不是顺手填一下的参数,而是触达系统最核心的可靠性设计之一。
4.2 坑二:重试导致业务重复执行
重试本身不是坏事,坏的是重试没有幂等意识。有一回客服工单分类 Agent 在处理一条工单时,因为网络抖动导致某个响应包在最后丢了,SDK 自动重试了一次,结果同样的工单被分类了两次,导致客服系统里出现了两条重复工单记录。
我们的修复方案很标准,但值得分享出来:调用方 SDK 在创建任务时自动生成一个idempotency_key,作为请求头随任务一起发送。Reach 收到后用 Redis 的SETNX记住这个 key,如果相同 key 3 秒内重复到达,直接返回第一次的结果,不再投递到下游 Agent。
async def submit_with_idempotency(task: AgentTask): key = f"idem:{task.idempotency_key}" acquired = await self.redis.set(key, task.task_id, nx=True, ex=300) if not acquired: # 重复请求,直接返回已保存的 task_id existed = await self.redis.get(key) return {"task_id": existed, "replayed": True} return await self._normal_submit(task)这个设计不仅防了重试,还防了调用方自己在代码里的循环误调,因为同样的请求 5 分钟内提交多次都会直接被去重。
4.3 坑三:僵尸 Agent 占着路由表
有个内部工具 Agent,用 Go 写的,进程稳定运行,心跳一直正常。但它内部调用了一个第三方封闭 SDK,有天这个 SDK 的某个连接池泄漏,导致所有请求进入后都排队阻塞,平均响应时间从 100 毫秒一路涨到 20 秒。注册表认为它活着,路由还继续给它派活,结果就是大量任务卡住。
前面我说过用被动探测治僵尸,这里补充一个细节:探测任务不能是单纯的ping/pong,我得让它去实际执行一个极轻量但带业务逻辑的任务,比如"返回最近一次成功处理的时间戳和当前请求队列长度"。只有这样,才能暴露 Agent 内部的真实处理能力,而不是进程层面的存活。
4.4 三个坑的解决思路汇总
| 问题现象 | 根本原因 | 解决方案 |
|---|---|---|
| 调用方线程池被占满 | 所有任务共用 30s 超时,慢任务阻塞线程 | 超时分级、线程池隔离 |
| 重复工单记录 | 重试未做幂等校验 | idempotency_key + Redis SETNX |
| 路由派给半死 Agent | 心跳无法反映处理能力 | 业务探活 + 连续失败摘除 |
5. 让触达真正可靠:超时分级、降级策略与熔断摘除
5.1 超时分级:Fast / Normal / Long 三类任务策略
我最终把任务按耗时预期分成三档,每一档都有独立的线程池和超时策略,互不干扰。
| 任务类型 | 预期耗时 | 客户端超时 | 线程池 | 典型场景 |
|---|---|---|---|---|
| fast | ≤3秒 | 5秒 | 独立大线程池 | 实时查询、风险判断 |
| normal | ≤15秒 | 20秒 | 独立线程池 | 内容生成、工单处理 |
| long | 分钟级 | 不等待,异步回调 | 少量长连接 | 批处理、深度分析 |
这个分组不仅改变了超时数值,还改变了调用方式。long 类任务我建议调用方不要傻等,而是提交后立刻拿到 task_id,后续通过 Webhook 或主动查询获取结果。这样的好处是,即便某个 Agent 需要跑 5 分钟,也不会占用调用方的请求线程。
5.2 降级:没有可用 Agent 时的兜底方案
路由匹配不到 Agent,或者匹配到但连续调用失败时,系统必须有降级路径,不能让调用方直接面对异常。我的降级分三个层次:
第一层是无匹配降级。当没有 Agent 声明对应能力时,Reach 返回一个 503 响应,但响应体里带上一个suggestion字段,告诉调用方"可以考虑哪些 Agent",比如"当前没有新闻摘要 Agent,但有文档摘要 Agent 可以用于 PDF 摘要"。这个对业务方非常重要,因为通常只是能力命名不一致,换个意图描述就能匹配上。
第二层是单 Agent 失败降级。第一候选 Agent 调用失败后,自动路由到候选列表里分数第二的 Agent。注意这里不是简单的重试,而是带着相同的idempotency_key去重试,防止多个 Agent 同时处理同一任务。
第三层是全部失败降级。如果所有候选 Agent 都失败,任务会进入死信队列,同时给注册表里的on_failure_callbacks发告警。这个回调通常是内部运营群的 webhook,让值班人员看到消息能手动介入。
5.3 熔断与淘汰:连续失败自动摘除
路由时要把连续失败的 Agent 排除出候选列表。我给每个 Agent 维护一个滑动窗口,记录最近 10 次调用的成功/失败状态。如果失败率超过 60%,这个 Agent 的 routing_weight 动态调整为 0,不再被选中。
但全部靠失败率有个问题:一个服务的流量是波动的,不足 10 次请求的情况下统计意义不大。所以我加了一个最小样本数,至少要有 5 次调用记录,失败率统计才生效。不足 5 次的按初始权重处理。
淘汰不是永久性的。摘除一个 Agent 后,我会每 30 秒给它发送一个单飞探测请求,如果探测成功且当前失败率回落到 20% 以下,就自动恢复路由资格。这个机制我称之为"半开状态",可以避免一个好 Agent 因为一次瞬时故障被永久打入冷宫。
6. Agent-Reach 与协议生态:统一接入层设计
6.1 为什么不能只支持 HTTP
最初我以为所有 Agent 都可以提供 HTTP 接口,后来发现现实很骨感。有团队内部的旧脚本 Agent,是几十个 Shell/Python 命令的组合,根本没有常驻服务,更别说对外开放 HTTP。还有一些 Agent 跑在边缘设备上,网络环境不允许直接开放端口,只能通过消息网关单向通信。如果接入层只认 HTTP,这些 Agent 就等于被排除在协作网络之外。
所以 Agent-Reach 的接入层必须支持多种协议。我设计了一个适配器机制,每种协议一个适配器,统一把请求转换成内部标准的InternalEnvelope,再交给路由引擎处理。
6.2 适配器模式:把每种 Agent 翻译成统一内部格式
适配器接口定义非常简单,核心只有一个invoke方法,它负责把统一格式的请求翻译成目标协议能识别的形式:
class ProtocolAdapter: async def health_check(self, endpoint: str) -> AgentHealth: pass async def invoke(self, endpoint: str, envelope: InternalEnvelope) -> InvokeResult: pass目前我实现了三个适配器:HTTPAdapter 是最常见的,把 envelope 序列化成 JSON 发到目标 API;ProcessAdapter 通过子进程方式运行本地脚本型 Agent,适合那种"传参数、跑命令、取结果"的场景;MQTTAdapter 适合在受限网络的边缘 Agent,通过发布订阅模式传任务和收结果。
每个 Agent 在注册时声明的protocols数组,就是告诉 Reach"我支持哪几种适配器"。路由时,Reach 会看调用方声明的可用协议和 Agent 声明的协议,取交集,按优先顺序选。比如调用方只走 HTTP,而某个 Agent 只支持 MQTT,那么直接跳过该 Agent,而不是尝试一条走不通的路。
6.3 协议协商与版本兼容
协议协商还有一个容易被忽略的细节:同一协议也有版本差异。Agent 的 MCP 接口可能是 v1,也可能是 v2,语义甚至完全不同。为了让升级不破坏已有集成,我要求每个 Agent 在注册信息里声明协议版本,并且路由时做兼容性判断。如果内部已经有agent:{agent_id}:protocols的映射表,Reach 就会优先选择双方都支持的协议版本。
协议升级时,不要直接在原 Agent 上改,而是新版本注册成一个新的 agent_id,比如agent_ops_diagnosis_v3,两个版本并行运行一段时间。等旧版本流量全部切换完毕且观察稳定,再下线旧版本。这样协议升级基本不影响线上稳定性。
这套统一接入层的好处是,业务方永远不感知协议的差异,对他们来说就是一次client.submit。协议适配和协商的复杂度全部被吞在了接入层内部,这也是我认为 Agent-Reach 最值得复用的一部分设计。
7. 最后:实测中的几点体会与后续扩展
跑了大半年,Agent-Reach从最初的分发工具逐步变成团队里 Agent 协作的事实标准。最想分享的个人体会是,注册表要尽量薄,不要一开始就把它做成配置中心或数据仓库,否则每个 Agent 接入成本都会急剧上升,大家就不愿意接了。薄的注册表只要做到"能发现、能连接、能感知状态"就够了,其余的业务配置应该留在各自 Agent 内部。
还有一个实用小技巧:给每个接入 Agent 增加一个带业务语义的/healthz拦截请求,返回体里除了{"status": "ok"}之外,至少带上当前队列深度、最近 10 分钟失败率、最后一个任务耗时。Reach 侧的被动探活直接消费这个接口。这样你排查一个 Agent 是否需要摘除时,不必临时翻日志,直接看探活数据就能定位。
后续我计划给Agent-Reach加两个能力:一个是对多租户的支持,让不同部门的 Agent 逻辑隔离,互不感知;另一个是沉淀路由日志,让每个任务的完整链路都可以回放,方便做协作质量分析和成本核算。如果你们也在搭类似的智能体协作层,我建议把这两个需求在一开始就纳入设计,后期返工的成本非常高。