从智能体孤岛到触达网络是我这段时间做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 字符串。之所以不用专门的数据库是因为注册表的读操作极其频繁路由引擎每次请求都要扫一轮在线 AgentRedis 可以把该操作压到毫秒级。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(), keylambda x: x[1], reverseTrue)实际运行中我发现一个有趣的规律语义相似度这个指标不能只看意图和能力名的相关性还要结合调用方历史上成功调用过哪些 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 fagent:{agent_id} data json.dumps(req.dict(), ensure_asciiFalse) 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, reasontimeout) raise这里我刻意把投递成功和处理成功区分开。投递成功只代表请求已经送达到 Agent 网关不代表业务逻辑完成。真正的完成信号由 Agent 在处理完毕后通过回调接口上报或者由调用方在约定时间内主动查询。这个区分非常重要避免了很多误判。3.4 调用方 SDK 的一次请求生命周期最后给调用方看的是一个小 SDK。我希望业务方调用 Agent-Reach 就像调用一个本地函数from agent_reach import Client client Client(endpointhttps://reach.internal.example.com) async def get_news_summary(news_text: str): task_id await client.submit( intentsummarize_news, params{text: news_text}, timeoutnormal, ) result await client.wait_for_result(task_id, timeout15) return result[summary]submit时只需要提供意图和参数不需要指定任何 Agent。SDK 内部会自动完成注册发现、路由匹配、协议协商、投递、回执等待整个链路。如果第一个选中的 Agent 超时SDK 会自动在候选列表里选下一个重试最多重试两次。调用方拿到的只有 task_id 和最终结果这个抽象让业务代码非常干净。4. 灰度上线后踩到的坑超时雪崩、重复执行和僵尸 Agent4.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 fidem:{task.idempotency_key} acquired await self.redis.set(key, task.task_id, nxTrue, ex300) 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 发到目标 APIProcessAdapter 通过子进程方式运行本地脚本型 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 逻辑隔离互不感知另一个是沉淀路由日志让每个任务的完整链路都可以回放方便做协作质量分析和成本核算。如果你们也在搭类似的智能体协作层我建议把这两个需求在一开始就纳入设计后期返工的成本非常高。
阅读完成 · 觉得有帮助?