ARTICLE DETAIL

资讯详情

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

多Agent统一触达层设计:从接口混乱到任务标准分发

多Agent统一触达层设计:从接口混乱到任务标准分发 我手头的Agent实例已经快两位数了一个跑在Docker里的本地问答Agent一台物理机上负责凌晨生成报表的调度Agent还有一个时不时被拉起来做数据抓取的临时容器Agent。最让我崩溃的不是它们各自的功能而是每接入一个新Agent我都要重新翻接口文档、换不同的鉴权头、写一套新的调用脚本。等到第七个Agent进来的时候我终于没法忍了于是花了两周做了一个轻量级工具取名Agent-Reach。它把“找到Agent、叫它干活、收回结果”这套高频动作变成标准流程这篇文章就把它的设计思路、实现细节和我在实测中踩过的坑完整记录下来给同样被多Agent环境折磨的人一个参考。1. 越想越不对Agent多了以后维护变成一场寻宝游戏1.1 我手上的三个Agent就是三个完全不同的“方言体系”先说具体场景。我参与维护的实验环境里有三个用途完全不同的Agent常驻Agent运行环境职责暴露方式鉴权model-qaDocker容器基于开源模型跑问答HTTP接口OpenAI兼容格式API Key放在Header里report-agent内网物理机定时生成报表支持手动触发自定义JSON接口请求体里带MD5签名scraper-agent临时容器按规则抓取公开数据并结构化启动一个命令行进程通过STDIN/STDOUT通信根本没有鉴权三个Agent三种接口风格三种调用方式。想给model-qa发个问题我得写OpenAI格式的请求体想触发report-agent我得按它的字段拼JSON再算签名想用scraper-agent还得先进容器里起进程。每次切换脑子里都要重新加载一套“方言”。这种心智负担在Agent数量少的时候还能忍一旦超过五个就会变成常态性的烦躁。1.2 “再写一个调度脚本”根本不是好答案一开始我的想法很简单既然每个Agent的接口不一样那就写一个Python脚本把三个调用封装成函数统一返回。脚本确实跑通了但问题很快出现。首先是改代码成本高。每加一个Agent就要在主脚本里新写一个函数、加一个配置项、改一遍帮助文档脚本从200行膨胀到800行时我已经不敢随便动了。其次是故障感知靠缘分。report-agent曾经因为物理机磁盘满了卡住我两天后才发现中间有三次定时任务其实都失败了。调度脚本本身没有健康检查Agent挂了它照样傻乎乎地发请求。第三是结果格式各说各话。有的Agent返回JSON有的返回纯文本有的返回一个文件路径下游想统一处理得写一长串类型判断。所以我意识到我需要的不是“再包一层”而是一个专门解决触达问题的中间层它要知道每个Agent在哪里、是否活着、能干什么、怎么把任务发过去、怎么把结果收回来。这些东西如果继续散落在调度脚本里那脚本迟早会变成一团乱麻。1.3 Agent-Reach的定位一个不抢活的“总机接线员”Agent-Reach给自己划的边界很清晰不接管任何Agent的业务逻辑也不做大模型编排只做注册中心、健康探针、任务通道、结果收敛。它的工作方式可以类比成酒店总机我不关心电话那头是座机还是手机也不关心通话内容是什么我负责把你拨的号码转到正确的人那里如果对方占线就记录一条未接来电。在这个定位下Agent-Reach需要解决四个核心问题寻址我怎样才能用一个统一入口找到任意一个Agent健康怎么知道一个Agent现在到底能不能干活分发怎么把任务安全地送到目标Agent手里收敛怎么把不同Agent五花八门的返回变成统一结构把这四个问题想清楚之后后面写代码就只是体力活了。2. Agent-Reach的核心设计把“找到Agent、叫它干活、收回结果”变成标准动作2.1 统一身份与注册表每个Agent都得先“报到”在Agent-Reach里任何Agent想被管理第一步必须是注册。注册时Agent要上报一份元数据我最终定的字段是这样agent_id全局唯一ID比如agent-model-qa是整个系统里寻址的关键name给人看的名字方便控制台展示control_addrAgent-Reach访问该Agent控制接口的地址必须从Agent-Reach的视角可达这个非常重要后面踩坑章节会细说auth_token_hashAgent的访问令牌哈希Agent-Reach后续给Agent发任务时会带上它做鉴权capabilities能力标签比如[chat]、[report]、[fetch]这是做能力发现的基础versionAgent代码版本号排查问题时很有用。为什么保留capabilities而不是直接按Agent ID硬编码路由因为有了能力标签Agent-Reach可以做“按能力发现”用户不关心具体是哪个Agent在做抓取只需要给一个fetch标签系统自动找到当前健康的抓取Agent下发任务。这个设计在Agent数量上去之后价值会越来越明显。2.2 触达协议越简单越好但要有三个固定端点我一度纠结要不要上gRPC毕竟流式传输和强类型接口都很诱人。但考虑到Agent可能跑在资源受限的容器里也可能跑在完全没有安装gRPC库的旧机器上最后还是选了HTTP JSON。原因很朴素HTTP是默认搭载能力最强的协议任何语言都能处理调试也最方便。Agent侧需要实现三个固定端点端点方向请求体响应体说明/agent/reach/healthAgent-Reach → Agent空{status:ok,now:1690000000}健康检查也用于Agent-Reach启动时的回连探测/agent/reach/taskAgent-Reach → Agent{task_id:...,payload:{...}}{task_id:...,accepted:true}任务下发Agent收到后立刻返回接受表示会异步处理/agent/reach/resultAgent → Agent-Reach{task_id:...,status:succeeded,output:{...}}{ok:true}任务执行完毕后的回传接口注意第三个端点的方向是反的结果不是Agent-Reach去拉的而是Agent主动post回来的。这个设计是为了适配Agent可能长时间运行任务、中途可能需要分片回传的场景。Agent-Reach只负责登记任务状态后续调用方轮询任务ID即可。2.3 心跳机制连续三次失联才判死不误杀健康检查是Agent-Reach最重要的基础能力。如果不做心跳那所谓的“触达”就退化成一个到处撞运气的HTTP请求。我采用的方案是每个Agent默认每30秒向Agent-Reach的/heartbeat接口发一次心跳Agent-Reach记录收到心跳的时间戳。当距离当前时间超过90秒也就是连续3个心跳周期没有任何消息时才把Agent状态从online改为offline。这里用“连续3次”而不是“1次”是有讲究的网络偶尔抖动一次丢包不代表Agent挂了但是90秒还没有任何消息基本可以确定Agent进程或网络出了严重问题用滑动窗口而不是简单计数可以避免Agent在临界状态频繁上下线。另外光靠Agent主动心跳只能说明“Agent能发出请求”不能说明“Agent-Reach能访问到Agent”。所以在Agent注册成功时Agent-Reach会主动向control_addr发一次/agent/reach/health请求双方都能通才标记online。这个“双向验证”在我后面踩到网络坑时救了我很多次。2.4 为什么不做全双工长连接轮询就够了别过度设计很多朋友看这个架构的第一反应是为什么不直接让Agent和Agent-Reach之间维持一条WebSocket长连接实时性不是更好吗我评估过这条路线最终放弃了原因有三个HTTP轮询已经能满足我的场景。报表任务、问答任务、抓取任务的实时性要求都不高秒级延迟完全可接受。长连接会显著增加Agent端的心智负担。Agent可能跑在休眠节能设备上也可能被NAT挡住维护长连接的保活、重连、心跳协调本身就是一套复杂的逻辑和“让Agent轻量接入”的目标冲突。任务不是连续流而是离散批次。Agent-Reach绝大多数场景是一问一答或一任务一结果WebSocket的推送优势根本发挥不出来。所以Agent-Reach最终采用“短HTTP 任务ID轮询”的方式任务提交后立即返回task_id下游每秒轮询一次任务状态。这套组合简单、可靠、好调试等以后真的出现需要实时推送的场景再在局部引入长连接也不迟。3. 落地0.1版这套系统我是怎么一步步搭起来的3.1 技术选型Python FastAPI SQLite 就够技术栈我选得很克制Python 3.11 FastAPI Uvicorn SQLiteAgent端SDK用aiohttp和标准库asyncio。选型理由如下组件理由FastAPI自带OpenAPI文档写接口几乎不用额外模板异步支持好Uvicorn轻量、性能足够单进程就能扛住十几个Agent的心跳和任务分发SQLite单文件部署不用引入数据库服务测试环境迁移也方便aiohttpAgent端发心跳和回传结果用异步非阻塞不拖累Agent主流程有人会说SQLite是不是太寒酸了我的判断是Agent-Reach 0.1版最重要的目标是快速验证架构每天的读写量也就几千次SQLite完全撑得住。等真正需要多节点部署时再迁移到PostgreSQL也不迟数据表结构是可以平滑迁移的。3.2 注册表数据结构两张表搞定数据库里我建了两张表一张存Agent元数据一张存任务流转状态。CREATE TABLE agents ( id TEXT PRIMARY KEY, name TEXT NOT NULL, control_addr TEXT NOT NULL, auth_token_hash TEXT, capabilities TEXT, status TEXT DEFAULT online, last_heartbeat INTEGER NOT NULL, created_at INTEGER NOT NULL ); CREATE TABLE tasks ( id TEXT PRIMARY KEY, agent_id TEXT NOT NULL, payload TEXT, status TEXT DEFAULT pending, result TEXT, created_at INTEGER NOT NULL, started_at INTEGER, finished_at INTEGER, timeout_at INTEGER, retry_count INTEGER DEFAULT 0 );agents表里的last_heartbeat是判断Agent是否存活的核心字段capabilities存的是JSON字符串方便后续做能力检索。tasks表里的status字段我维护了一个状态机pending→running→succeeded/failed/timeout。其中timeout_at字段很关键它记录了任务的最晚完成时限扫描线程会定期清理超时任务。3.3 注册与心跳接口实现这是Agent-Reach最核心的代码片段注册接口负责接收Agent上报并做回连探测from fastapi import FastAPI, HTTPException from pydantic import BaseModel import time, sqlite3, json, hashlib app FastAPI() DB_PATH reach.db class AgentRegister(BaseModel): agent_id: str name: str control_addr: str auth_token: str capabilities: list[str] [] version: str app.post(/register) def register(agent: AgentRegister): token_hash hashlib.sha256(agent.auth_token.encode()).hexdigest() now int(time.time()) conn sqlite3.connect(DB_PATH) cur conn.cursor() cur.execute( INSERT OR REPLACE INTO agents (id, name, control_addr, auth_token_hash, capabilities, status, last_heartbeat, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?), (agent.agent_id, agent.name, agent.control_addr, token_hash, json.dumps(agent.capabilities), online, now, now) ) conn.commit() conn.close() return {ok: True, id: agent.agent_id, status: registered}心跳接口更简单只需要更新时间戳并把状态拉回onlineapp.post(/heartbeat) def heartbeat(agent_id: str, auth_token: str ): now int(time.time()) conn sqlite3.connect(DB_PATH) cur conn.cursor() cur.execute( UPDATE agents SET last_heartbeat?, statusonline WHERE id?, (now, agent_id) ) conn.commit() conn.close() return {ok: True, now: now}3.4 任务分发与状态回收状态机是核心任务分发是整个Agent-Reach里逻辑最重的部分。调用方通过POST /tasks提交任务Agent-Reach先查注册表确认Agent存在且状态是online再往Agent的/agent/reach/task发任务最后把task状态置为running。app.post(/tasks) def create_task(agent_id: str, payload: dict): conn sqlite3.connect(DB_PATH) cur conn.cursor() cur.execute(SELECT control_addr, status FROM agents WHERE id?, (agent_id,)) row cur.fetchone() if not row: raise HTTPException(status_code404, detailagent not found) if row[1] ! online: raise HTTPException(status_code403, detailagent is offline) task_id ftask-{int(time.time()*1000)} cur.execute( INSERT INTO tasks (id, agent_id, payload, status, created_at) VALUES (?, ?, ?, pending, ?), (task_id, agent_id, json.dumps(payload), int(time.time())) ) conn.commit() conn.close() return {task_id: task_id, status: pending}实际代码里任务下发之后会立刻尝试调用Agent的/agent/reach/task如果调用成功状态从pending改为running如果调用失败直接标记failed并记录失败原因。Agent执行完之后会通过POST /result回传结果我把这个接口做成了幂等更新同一个任务ID可以重复回传以后做任务重放也方便。3.5 Agent端三分钟接入一个装饰器就够为了让已有Agent快速接入我写了一个特别简单的Python SDK核心是一个装饰器def reach_agent(agent_id: str, auth_token: str, controller_url: str): def decorator(handler): from fastapi import FastAPI, Request app FastAPI() app.post(/agent/reach/health) async def health(): return {status: ok, now: int(time.time())} app.post(/agent/reach/task) async def run_task(req: Request): body await req.json() task_id body[task_id] result await handler(body[payload]) # 异步回传结果避免阻塞Agent主流程 import aiohttp async with aiohttp.ClientSession() as session: await session.post( f{controller_url}/result, json{task_id: task_id, status: succeeded, output: result} ) return {task_id: task_id, accepted: True} return app return decorator这样接入的成本就是Agent原有业务逻辑写成一个异步函数再用reach_agent(...)包一层即可获得健康检查、任务接收、结果回传三件套。理论上三分钟就能把老Agent接进来。4. 三台真实Agent实测跑通之后才发现这三个坑4.1 测试环境一台控制器三台异构Agent我搭了一个最贴近真实环境的测试网络一台宿主机上跑Agent-Reach控制节点另外三个Agent分别以三种不同角色接入Agent ID运行位置control_addrcapabilities备注agent-model-qaDocker容器172.17.0.2:9000[chat]容器网络NAT隔离agent-report内网物理机192.168.1.10:9100[report]物理机网络agent-scraper临时K8s Pod10.0.0.8:9200[fetch]动态端口整个测试流程分三步启动控制节点、注册三个Agent、逐个下发任务并观察状态流转。前两步都很顺利三个Agent都成功注册并显示online第4.2节我放一张观测表。4.2 从注册到第一次任务分发看起来一切顺利我依次给三个Agent各发了三个不同任务Agent-Reach的控制台输出很清晰[10:01:22] task-a-001 - agent-report: dispatched [10:01:26] task-a-001 - agent-report: running [10:01:29] task-a-001 - agent-report: succeeded (took 7s) [10:01:35] task-b-001 - agent-model-qa: dispatched [10:01:38] task-b-001 - agent-model-qa: running [10:01:42] task-b-001 - agent-model-qa: succeeded (took 4s)观测项agent-reportagent-model-qaagent-scraper注册耗时0.8s1.2s0.6s心跳到首次任务下发延迟0.2s0.3s0.1s任务平均耗时7s4s11s任务成功率9/99/99/9一切顺利的感觉最容易让人放松警惕后面的坑都是真实跑了一周才暴露出来的。4.3 踩坑一心跳线程把异步事件循环拖垮现象Agent-Reach跑了两天后注册接口和心跳接口的响应时间从1毫秒劣化到5秒以上但CPU占用并不高看起来像死锁又不像死锁。排查过程我先看数据库锁排除SQLite写锁问题接着看Uvicorn访问日志发现大量请求堆积在/heartbeat上最后翻Agent端SDK的心跳实现发现问题出在心跳线程上。我给Agent端的心跳是用threading.Timer实现的每30秒在线程里调用asyncio.get_event_loop().create_task()往主事件循环塞心跳任务和Agent主流程的异步任务形成了竞争。根因asyncio的事件循环不是线程安全的跨线程向同一个事件循环提交协程轻则调度延迟重则引起未知的竞态。我当时文档没读透管它叫“线程和协程一起用”其实是不规范的混用。修复Agent端SDK的心跳逻辑彻底改成同步实现用requests.post直接发心跳请求不碰asyncio事件循环Agent-Reach服务端则用心跳接口的幂等更新。另外我养成了一个习惯所有接口都加上响应时间日志这样下次再出现类似问题一眼就能看出是哪个环节慢。4.4 踩坑二endpoint填的地址控制端根本触达不到现象agent-model-qa注册时control_addr我填了容器内的172.17.0.2:9000。从Agent-Reach控制节点所在的宿主机上访问死活不通但Agent-Reach日志里显示的却是注册成功、心跳正常。排查过程最先怀疑是防火墙但关了还是不通然后用curl从不同网络位置分别测试这个地址发现只有容器所在宿主机能通局域网内都不通。原因很清楚容器网络是NAT隔离的172.17.0.2这个地址只在Docker内部有效换一个网络视角就变成了不可达地址。根因Agent注册时上报control_addr的角度错了。Agent自己觉得“我用这个地址访问自己是通的”但Agent-Reach需要的是“从Agent-Reach的视角访问Agent是通的”。这两个地址在NAT环境下完全不同。修复我在注册流程里增加了一个“回连探测”机制Agent-Reach收到注册请求后主动向control_addr发起GET /agent/reach/health只有探测成功才标记online否则拒绝注册并返回错误。同时在文档里用粗体写明control_addr必须是Agent-Reach网络视角下可路由的地址。这个机制后来帮我在每次新Agent接入时都能第一时间发现网络配置错误。4.5 踩坑三一个超时任务差点卡死整条分发队列现象某天下午我给agent-report下发了一个报表任务Agent端因为物理机负载过高迟迟没有响应。结果接下来所有新提交的任务全部停在pending状态控制台看起来像整个Agent-Reach都挂了。排查过程先重启Agent-Reach重启后前几个任务能跑但过几分钟又全部堵死。看日志发现分发队列里堆了几十个任务它们都在等同一个Agent的回复。问题的根源是我在任务下发代码里没有设置HTTP调用的超时上限Agent不响应await就一直挂着消息队列前面排队的任务全部被堵住。根因任务级超时和请求级超时都没配置。一个慢Agent拖垮了所有Agent的任务分发典型的“单点拖累全局”。修复我做了三个改动。第一所有HTTP调用统一增加超时参数服务端默认120秒超时第二capabilities里带slow标签的Agent走独立的慢任务队列不占用普通任务队列的并发槽第三增加一个后台扫描线程定期清理timeout_at已过期的任务并标记为timeout释放队列占用。这三板斧下去之后再没出现过因为单个Agent卡死导致全局瘫痪的情况。5. 从“能跑”到“能用”安全、观测与扩展的实战补课5.1 安全加固令牌、签名与最小权限0.1版的安全模型很简陋Agent上报的token直接存哈希但任务下发时的身份验证做得不够严格。要往生产环境走至少需要补齐四件事令牌最小化每个Agent单独分配一个token只允许它注册自己对应的agent_id不允许跨ID操作消息签名Agent-Reach向Agent下发任务时请求体用HMAC-SHA256签名Agent端验签后才执行防止接口被局域网内其他进程恶意调用传输加密HTTP换成HTTPS内网环境可以用自签CA但必须做证书校验审计日志所有注册、注册注销、任务下发、结果回传操作都记录审计日志方便事后追踪。安全这块我不建议一步到位但要尽早补。Agent本身可能就是执行敏感任务的入口Agent-Reach作为统一触达层如果被攻破等于一次性拿到了所有Agent的控制权。5.2 观测性没有指标你根本不知道系统在变慢很多自建系统死就死在“看起来正常其实已经在崩溃边缘”。Agent-Reach现在强制要求每个核心接口都暴露三类指标请求延迟P50/P95哪个接口慢了立刻能看出来任务状态分布pending、running、failed、timeout各自有多少动态看出队列健康度心跳间隔统计Agent心跳到达时间的抖动情况提前预判网络质量。我直接用一个/metrics端点输出文本格式的指标配合现有的监控系统做抓取。代码实现不复杂就是在一个全局字典里累计计数器METRICS {} def inc(metric_name, value1): METRICS[metric_name] METRICS.get(metric_name, 0) value def observe(metric_name, value): METRICS.setdefault(metric_name, []).append(value)这比事后翻日志高效得多。第4.3节那个心跳线程问题如果第一天就上了延迟指标根本不用靠猜。5.3 从单机到多节点演进路径要想清楚但不要过早做0.1版是单节点SQLite架构上有一个明确的演进路径当Agent数量超过几百个、心跳频率超过每秒几十次时SQLite会先成为瓶颈。届时可以分两步走存储层替换把SQLite换成PostgreSQL或MySQLagents表增加索引tasks表按时间做分区任务总线替换把asyncio.Queue换成Redis队列或轻量消息中间件Agent-Reach支持水平扩展出多个实例共享同一个任务队列。但现阶段我不建议做这些。过早引入分布式组件只会让系统排错变难单节点跑得好好的就先把单节点榨干再说。5.4 我接下来计划做的四件事Agent-Reach 0.1版已经跑了一个多月稳定支撑了我手头所有Agent的日常调度。接下来我给自己列了一个优先级明确的TODO插件化协议适配器目前SDK只支持Python以后可能让Agent接入时声明“我支持HTTP/JSON”让Agent-Reach自动选择协议适配器扩大接入范围任务重放与幂等消费Agent执行任务超时后Agent-Reach应该能重新投递同一个task_id但必须保证Agent侧幂等执行防止重复报账或重复扣款之类的问题人工审批门禁对于高风险任务Agent-Reach支持先挂起交给人工确认再下发这个能力在敏感操作场景非常关键结果Schema校验每个Agent在注册时声明response_schemaAgent-Reach在回收结果时强校验不符合立即标记失败而不是让脏数据流到下游。整个Agent-Reach做下来我最深的体会是很多看似复杂的问题不是需要更复杂的系统而是需要把“发现、触达、回收”这三个高频动作从“每次手动做”变成“一次定义、长期复用”。现在我再接一个新Agent最花时间的反而是准备好Agent自己的业务逻辑接入部分十分钟就能搞定。如果你也被同类问题逼疯过建议先别想着上什么重型平台照着这个思路做一个几十KB的轻量触达层可能就够用了。
返回列表