ARTICLE DETAIL

资讯详情

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

多Agent协作编排引擎Agent-Reach架构设计与落地实践

多Agent协作编排引擎Agent-Reach架构设计与落地实践 1. 项目定位与整体架构Agent-Reach 到底解决什么问题1.1 多 Agent 协作的困境2024 年到 2025 年只要你在搞 AI 应用落地大概率会发现一个现象单个大模型 Agent 的能力边界很快就会被捅破。你可以让一个 Agent 写周报、查资料、做数据分析但一旦业务流程变长比如用户提交工单 → 自动分类 → 检索知识库 → 生成答复 → 发邮件通知用户 → 回写 CRM 系统单靠一个 Agent 就很难扛下来。不是你 Prompt 写得不好而是这条路本质上需要多个角色分工协作一个负责理解意图一个负责查数据一个负责写内容一个负责调用外部 API。把这么多职责塞进同一个 Agent会让系统变成一团乱麻你甚至分不清一次失败到底是因为模型幻觉、接口超时还是提示词冲突。我最初也是踩了这个坑。当时为了让一个 Agent 全流程处理客服工单硬是把十几个工具的调用说明写进同一个 System Prompt结果上下文动辄十几万 token响应延迟从 3 秒飙到 20 秒而且每次加一个新工具老工具的表现就会波动。后来想明白一件事Agent 不是越大越好而是应该像团队一样分工。Agent-Reach 就是在这样的背景下做的——它的定位不是让你多一个大模型应用而是解决多个 Agent 之间怎么可靠地互相触达、怎么协同工作这一层问题。1.2 Agent-Reach 的核心抽象从名字拆解Agent-Reach 的核心就两个字触达。要让 Agent A 的能力被 Agent B 使用要让一条任务链能在 N 个 Agent 之间平滑流转要让外部的工具、数据源、业务系统都能成为 Agent 可以伸手够到的资源。实现的思路是引入一个轻量级的编排层所有 Agent 不再点对点直连而是通过统一的中枢来完成注册、发现、路由和消息传递。你可以把它理解成公司里的前台你不必知道财务部的小王电话是多少你打给前台前台根据你的需求把电话转给财务部。如果小王休假前台还会帮你转给小李。这个抽象带来的最大好处是让业务逻辑和 Agent 的物理位置解耦。每个 Agent 只需要关心自己会干什么、能提供什么能力不需要知道其他 Agent 的地址、接口、调用方式。当你要新增一个 Agent只需要做两件事注册能力、订阅感兴趣的任务类型。其余的一切包括消息路由、重试、超时、降级、链路追踪都由 Agent-Reach 统一处理。整套架构分四层接入层负责暴露统一的 SDK 和 API 给各业务方编排层是大脑负责意图识别、任务规划、路由决策执行层是具体的 Worker Agent处理实际任务并返回结果基础设施层则是消息总线、状态存储和监控组件。这四层各司其职项目迭代起来非常清楚。2. 核心模块设计注册、路由与通信2.1 能力注册与发现Agent 的服务目录Agent-Reach 里很重要的一个设计是每个 Agent 在启动时必须先到协调中心注册自己的能力。注册信息不是简单写个名字就完事我给 Agent 设计了一套能力描述结构本质上是把我会什么变成一个机器可读、可匹配的元数据文件。字段包括agent_id、service_name、capabilities、input_schema、output_schema、endpoint、timeout还有priority和load_limit。其中input_schema和output_schema是用来描述参数结构的我直接用 JSON Schema 来定义这样路由层做参数校验时就不需要写一堆硬编码逻辑。举个实际例子一个做邮件触达通知的 Agent它的能力描述大概是这样的{ agent_id: agent-mailer-01, service_name: notification.mail, version: 1.2.0, capabilities: [ { name: send_email, description: 发送单封邮件通知, input_schema: { type: object, required: [to, subject, body], properties: { to: {type: string, format: email}, subject: {type: string}, body: {type: string}, cc: {type: array, items: {type: string}} } }, output_schema: { type: object, properties: { message_id: {type: string}, status: {type: string, enum: [sent, queued, failed]} } } } ], endpoint: http://agent-mailer.internal:9100, load_limit: 50, priority: 5 }为什么要用 JSON Schema因为我可以在调度时做提前校验不满足条件的请求直接在路由层拦截而不是把垃圾请求发给 Agent 让它报错。很多时候 Agent 失败的根因根本不是模型能力而是上游调用方给了不合法参数。这件事在注册阶段设计好后面能省下大量排查时间。我还给注册中心加了一个 TTL 租约机制Agent 要每隔 30 秒发一次心跳续租如果连续 3 次心跳丢失协调中心就会把这个 Agent 标记为不可用不再把新请求路由过去。这样避免了一个 Agent 挂了之后上游还在傻傻等待的窘境。2.2 路由决策把请求交给谁路由是整个系统里最有技术含量的部分。一开始我尝试过简单粗暴的规则匹配——根据请求中的intent字段查表命中哪个 Agent 就发给哪个。但很快发现真实场景没有那么听话用户表达同样一个意思可能用完全不同的措辞比如给客户发一封道歉邮件和告知用户处理结果这两个请求在规则表里如果严格匹配可能就找不到同一个 Agent。后来我把路由改成了两层第一层是语义意图识别用一个轻量级分类模型或者大模型的函数调用能力把自然语言请求映射到标准化的intent比如intent: send_email_notification第二层是能力匹配打分根据意图、参数结构、Agent 负载、历史成功率等因素给每个候选 Agent 计算一个综合得分。打分我把它定义成一个加权计算过程完全透明的规则score intent_similarity * 0.4 capability_match * 0.3 reliability_factor * 0.2 (1 - load_factor) * 0.1其中intent_similarity是意图和 Agent 能力名的语义相似度capability_match是输入 Schema 的字段匹配率reliability_factor是过去 24 小时该 Agent 的成功率load_factor是当前负载和load_limit的比值。这四个维度加权之后选得分最高的作为目标。之所以要引入负载维度是因为真实场景中如果某个 Agent 已经被打满继续把请求塞给它只会让系统更慢不如分流给速度稍慢但空闲的备用 Agent。这里的权重不一定要完全固定我给每项权重都设置成了可配置项不同业务域可以微调比如对金融场景更看重可靠性对实时交互场景更看重负载均衡。2.3 通信协议让消息有一致性骨架Agent 之间通信我采用了两条通道一条是同步通道用于简单请求-响应场景基于 HTTP/2 的 gRPC另一条是异步通道用于任务链上的消息传递基于消息队列。很多 Agent 协作场景是长任务比如先检索资料再写报告最后发邮件如果全程用同步 RPC 串起来任何一个环节卡住整个请求都会拖着一直不放来一个请求就占住一个线程服务很快就没法并发了。异步消息就能把任务中间态暂时存下来让调用方先返回等执行到后续步骤再把结果继续推下去。Agent-Reach 在异步通道上定义了一个统一消息信封关键字段有msg_id、trace_id、conversation_id、from、to、intent、payload、priority。消息在队列里流转时协调中心会持久化两个 IDtrace_id是整次业务追踪链路共用的排查问题就靠它串起所有日志conversation_id是某一次完整会话的上下文隔离边界确保不同会话之间的消息不会互相污染。给一个实际的消息体示例{ msg_id: m-8f2a9b1c, trace_id: tr-20250107-00123, conversation_id: conv-20250107-0051, from: router/agent-reach-core, to: agent-analyzer/02, intent: generate_data_insight, payload: { report_id: RPT-2025-0107, target_metrics: [revenue, active_users] }, priority: 3, created_at: 2025-01-07T10:24:00Z }这样做的收益是消息语义标准化。你不用为了对接不同 Agent 去读不同接口文档只要遵守这个信封结构往里填内容就行。任何 Agent 的输入输出在协议层都是同构的换来的是整个系统在新增角色时几乎零沟通成本。在实现时我把这条消息总线包了一层客户端 SDKAgent 接入时只需要调用send(task)之类的接口SDK 会自动填充msg_id、trace_id这些字段开发人员基本上不用关心底层通信细节。3. 编排引擎的实现要点3.1 核心数据模型编排引擎是 Agent-Reach 的心脏它负责接收上游请求、拆分任务、编排依赖关系、然后按顺序或并发地把子任务分发给不同的 Agent。我先定义清楚几类核心数据模型Task代表一次可调度的最小执行单元Plan是多个 Task 组成的有向无环图AgentNode是对线上 Worker Agent 的抽象封装RouteResult是路由决策产出。在 Python 里的骨架我用 dataclass 来实现保持代码直观from dataclasses import dataclass, field from enum import Enum from typing import Any, Dict, List, Optional class TaskStatus(str, Enum): PENDING PENDING RUNNING RUNNING SUCCEEDED SUCCEEDED FAILED FAILED SKIPPED SKIPPED TIMED_OUT TIMED_OUT dataclass class Task: task_id: str intent: str payload: Dict[str, Any] agent_id: Optional[str] None status: TaskStatus TaskStatus.PENDING retry_count: int 0 timeout_ms: int 5000 depends_on: List[str] field(default_factorylist) result: Optional[Dict[str, Any]] None error: Optional[str] None dataclass class Plan: plan_id: str conversation_id: str tasks: Dict[str, Task] root_task_ids: List[str] current_status: str IN_PROGRESS这里的depends_on是任务依赖关系的关键它让编排引擎可以把一个复杂需求拆成有向无环图按拓扑顺序执行相互依赖的任务同时把互不依赖的任务并发跑最大程度压榨系统吞吐。状态机的处理逻辑我单独写了一个模块每次任务状态变更都会触发一次Plan级别的状态评估一旦所有root_task_ids后续的叶子节点都完成整个Plan就标记为成功。这个设计让复杂流程的异步推进变得可控后续要加人工审核节点、定时任务也都方便只要往图里加节点就行。3.2 编排主流程代码骨架下面这段代码是编排引擎的核心处理逻辑我在实际项目中给它取了个名字叫RouteThenExecute。逻辑不复杂第一步路由第二步执行第三步处理异常。真正的工作量花在容错细节上比如任务失败时需要决定是重试还是跳过依赖它的下游任务要不要继续。这里我给出一个可读的骨架版本读者可以按自己的技术栈迁移import asyncio import time from typing import Optional class ReachCoordinator: def __init__(self, router, executor, state_store, bus): self.router router self.executor executor self.state_store state_store self.bus bus async def submit_plan(self, plan: Plan) - Plan: await self.state_store.save_plan(plan) ready_tasks [ t for t in plan.tasks.values() if not t.depends_on and t.status TaskStatus.PENDING ] await asyncio.gather(*[self._dispatch(task) for task in ready_tasks]) return plan async def _dispatch(self, task: Task) - None: if task.status ! TaskStatus.PENDING: return route await self.router.route(task.intent, task.payload) if route is None: task.status TaskStatus.FAILED task.error no suitable agent found await self.state_store.save_task(task) return task.agent_id route.agent_id task.status TaskStatus.RUNNING await self.state_store.save_task(task) try: result await self.executor.execute_with_timeout( agent_idroute.agent_id, intenttask.intent, payloadtask.payload, timeout_mstask.timeout_ms ) task.result result task.status TaskStatus.SUCCEEDED await self._release_downstream(task.task_id) except asyncio.TimeoutError: await self._handle_failure(task, errortimeout) except Exception as exc: await self._handle_failure(task, errorstr(exc)) finally: await self.state_store.save_task(task) async def _handle_failure(self, task: Task, error: str) - None: task.error error if task.retry_count 2: task.retry_count 1 task.status TaskStatus.PENDING await asyncio.sleep(1 * task.retry_count) await self._dispatch(task) else: task.status TaskStatus.FAILED await self._fail_downstream(task.task_id) async def _release_downstream(self, completed_task_id: str) - None: for task in self._running_plan().tasks.values(): if completed_task_id in task.depends_on: task.depends_on.remove(completed_task_id) if not task.depends_on and task.status TaskStatus.PENDING: await self._dispatch(task) async def _fail_downstream(self, failed_task_id: str) - None: for task in self._running_plan().tasks.values(): if failed_task_id in task.depends_on: task.status TaskStatus.SKIPPED await self.state_store.save_task(task) await self._fail_downstream(task.task_id)这段代码里我刻意把_fail_downstream做成递归目的是让失败传递到所有下游节点形成快速失败机制。真实项目里这一步很重要我见过很多做 Agent 编排的同学上游子任务失败后下游还在继续执行最后生成一份缺数据的报告还给用户这种故障很难排查因为出错的根本不是下游 Agent而是编排依赖没有处理干净。3.3 一个完整的业务场景演练用客服邮件触达场景完整走一遍用户提交工单我要投诉请给我回复处理进展。上游服务把这个需求转给 Agent-Reach 后编排引擎先调用意图识别模块得到拆解后的Plan如下Task IDIntent依赖目标 AgentT1classify_ticket无agent-classifierT2search_knowledgeT1agent-ragT3generate_replyT1, T2agent-writerT4send_emailT3agent-mailerT5update_crmT3agent-crmT1 和 T2 一个负责分类工单、一个可以并行准备知识库检索虽然 T2 不依赖 T1 的分类结果但这个场景里 T2 依赖的是工单文本本身不需要等分类完成。执行时 T1 和 T2 可以同时跑等两者都完成后 T3 开始生成回复之后 T4 发邮件、T5 写回 CRM这两个也是并行的。整个流程是一个有向无环图最理想情况下耗时大概是T1 T3 并发(T4,T5)而不是五个任务时间相加。Agent-Reach 执行完之后把每条任务的状态、耗时、Agent 节点、错误信息都写进状态存储整个 Plan 的执行轨迹可以完整回放这个能力在排障时价值巨大。4. 性能调优与配置实践4.1 并发模型与线程参数Agent-Reach 的编排引擎我最终选择的是异步事件循环 有界线程池的混合模型。为什么不是纯异步因为有些 Agent 的执行器底层调用的第三方 SDK 是同步阻塞的比如某些邮件服务 SDK、数据库驱动把它们直接丢进事件循环会卡住所有协程。所以我让异步层负责消息路由和状态流转真正执行 Agent 调用的部分放进一个固定大小的线程池用信号量控制最大并发。这个设计也许不极客但在真实生产环境中非常稳。线程池参数我推荐按任务类型区分短任务池的核心线程数设为 CPU 核心数乘以 2队列容量设为 512长任务池的核心线程数设为 CPU 核心数队列容量设为 128。不能把所有 Agent 调用混在一个池子里否则一个跑 30 秒的邮件发送任务占满线程之后一个只需要 200 毫秒的缓存查询也会排队到天荒地老。关于超时设置单个 Agent 调用的默认超时我不建议超过 5 秒短任务可以压到 3 秒。不要以为超时设得越大成功率越高事实恰恰相反设大超时只会让故障恢复变得更慢因为请求线程全被卡住后续请求就全部堆积了。4.2 消息队列与状态存储选型消息总线的选型早期我图省事直接上了 Redis 的 Stream后来放弃了因为 Redis Stream 的消费组机制在消费者扩容时有消息重复消费的风险而 Agent 协作场景一旦出现重复触达比如邮件发了两封就很尴尬。最终我换了 RabbitMQ开启 publisher confirm 和 consumer ack配合prefetch_count设置为 10这样既保证消息不丢也不会让消费者被积压消息打爆。这里想给一个建议如果消息量没有达到每秒数万条不要急着上 Kafka。Kafka 的优势是大吞吐、长日志保留但它的消费语义是至少一次天然会有重复需要业务层做幂等而 RabbitMQ 在中小规模下语义更清晰、运维也更简单。技术选型不是越重越好是越匹配越好。状态存储我用的是 PostgreSQL一张task_state表主键(task_id, plan_id)字段存任务状态、Agent 节点、重试次数、耗时、错误信息。为什么不直接用 Redis因为状态存储需要持久化和事务性编排引擎要在任务变更时做原子更新Redis 做这个很别扭。有人可能会觉得每次都写数据库会很慢实测下来在单机 PostgreSQL 上每秒几百个任务变更写毫无压力。我做了批量更新的优化把同一批次的状态变更合并成一条 SQL 的upsert性能直接翻倍。下面是部分配置参数贴出来供参考配置项推荐值说明SHORT_TASK_CORE_THREADSCPU核心数 * 2短任务线程池SHORT_TASK_MAX_THREADSCPU核心数 * 4短任务最大线程数SHORT_TASK_QUEUE_SIZE512短任务队列容量LONG_TASK_CORE_THREADSCPU核心数长任务如邮件/导出SHORT_TASK_TIMEOUT_MS3000短任务超时LONG_TASK_TIMEOUT_MS30000长任务超时ROUTE_MAX_CANDIDATES3路由候选Agent数AGENT_HEARTBEAT_INTERVAL30sAgent心跳间隔AGENT_HEARTBEAT_MISS_TOLERANCE3心跳丢失容忍次数5. 可靠性建设超时、熔断与可观测性5.1 重试与幂等Agent 协作系统的可靠性一大半是靠失败恢复撑起来的。但重试是有代价的如果只重试不设计幂等一个请求被重复执行就会造成重复发邮件、重复扣费、重复写库后果比重试之前更糟。Agent-Reach 里的做法是给每个 Task 强制分配task_id并且要求所有可重入的 Agent 在执行前把task_id作为幂等键写入自己的存储层。重试时协调中心把同一个task_id发过去Agent 查询到已经处理过就直接返回上一次结果不再重新执行副作用操作。这套逻辑看起来简单但它是很多 Agent 系统从 demo 走向生产的一道大坎。没有幂等前任何网络抖动引发的重试都是一次事故。重试策略不能一刀切。我用的策略是超时类错误重试 2 次第一次延迟 1 秒第二次延迟 2 秒业务逻辑错误比如参数非法、内容审核不通过不重试直接标记失败Agent 节点不可达时先把任务改路由到备用 Agent如果没有备用节点才走重试。这样分类处理失败恢复效率高也不会把资源浪费在注定失败的任务上。5.2 链路追踪与故障定位Agent-Reach 中多条 Agent 链路交错执行时没有可观测性等于闭眼开车。我把链路追踪做到完全透明SDK 会自动生成trace_id并塞进日志上下文任何 Agent 打印日志都会带上这个 ID。排查问题时只需要拿到一次用户请求的trace_id就能在所有组件日志中筛选出相关记录按时间线还原整条链路的运行过程。链路追踪的日志最好是结构化 JSON 格式不然排查会非常痛苦。每个节点上报的数据包括节点名、任务 ID、耗时、状态码、输入输出摘要。这些数据同时汇入 Prometheus 做指标监控关键指标有reach_task_success_rate、reach_task_avg_duration_ms、reach_route_miss_count、reach_agent_load_factor。我设置了一个告警规则任务成功率低于 95% 持续 5 分钟就触发告警。别小看这种基础指标它往往是系统劣化的第一个信号。有一次我排查线上问题看到reach_route_miss_count突然飙升然后定位到是一个 Agent 因为配置错误注册失败用户请求全部路由不到目标没有这个指标光靠用户反馈来发现问题至少滞后半小时。6. 落地过程中踩过的坑6.1 上下文无限膨胀第一次上线时我把整个会话的所有历史消息全部传给每个 Agent让它们自己筛选有用信息。结果跑了两周响应越来越慢token 消耗越来越高某些长会话甚至直接把模型输入上限打爆。后来改成按需裁剪每个 Agent 只接收它真正需要的字段完整历史上下文统一存在状态存储里除了少数确实需要全局记忆的 Agent默认都不带全量上下文。这里有一个经验值一个 Agent 接收的上下文不应超过它生成内容所需信息量的 1.5 倍超过就是浪费。信息压缩是 Agent 协作系统里长期要优化的命题不太可能一劳永逸但至少可以通过上下文裁剪规则把浪费控制在合理范围。6.2 Agent 之间的死循环还有一个坑是 Agent 互相触发导致的死循环。A 处理完任务后发了一个事件给 BB 处理后反过来给 A 发了一个新任务两个 Agent 就无限互相调用下去消息队列积压暴涨最后把 RabbitMQ 都拖垮了。根本原因是我在编排层没有做环节去重。解决办法是所有消息进入编排引擎前必须检查trace_id同一个trace_id下消息经过的节点集合会被记录下来如果超过设定的最大节点数比如 15 个节点新消息直接拒绝并标记整个链路为异常。这个机制能兜住大部分循环调用场景避免像无头苍蝇一样来回打转。另外编排引擎发现某个intent在同一个trace_id中出现次数超过 3 次就自动将后续同类任务转入人工审核队列。有时候循环调用不是死循环而是业务策略上产生了一个递归流程自动硬终止又太粗暴所以设置一个人工干预舱让负责人来决定要不要继续。这个设计救过我两次一次是订单状态流转写错了状态机的转移条件一次是促销活动的一个规则导致消息反复触发。6.3 路由打分不收敛在调路由算法时还遇到过一个问题两个 Agent 能力非常接近语义相似度得分也几乎一样导致请求在两个 Agent 之间来回切换一会儿走 A一会儿走 B上游用户观察到的系统行为不稳定一段时间结果风格完全不一样。为了处理这个场景我引入了亲缘性策略路由结果会缓存到 Redis如果同一个conversation_id之前已经路由到某个 Agent那么后续同类请求在同一会话内有 85% 的权重偏向该 Agent只有当它的负载超过 80% 时才强制切换。这样既保证了会话内的一致性体验又不至于把一个 Agent 打死。这是一个很典型的机器学习之外的工程调优没有太玄的技术原理但对用户体验的提升非常明显。另外路由打分要加入冷却时间的概念。某个 Agent 刚刚出现过失败不应该立刻被再次选中我给每个 Agent 维护了一个失败时间戳打分时对最近 60 秒内失败过的 Agent 施加一个 0.6 的衰减系数。这个机制让路由层天然规避了刚挂掉又被派活的尴尬场景也减少了无谓的重试。7. 一些想法跟后续扩展方向Agent-Reach 做到现在感触最深的是多 Agent 协作本质上是把分布式系统的老问题换了一套新外壳。过去我们做微服务关心服务发现、路由、超时、熔断、幂等、链路追踪现在做 Agent 系统这些问题一个不少全都要重新面对只不过多了一层语义路由多了一些模型层面的不确定性。如果你正在做一个多 Agent 的应用我建议不要一上来就堆复杂框架先用消息队列加一张任务状态表把消息能正确地从一个节点走到另一个节点这件事跑通再逐步加功能。如果你问我会不会把 Agent-Reach 继续做下去答案是肯定的。我最近在琢磨两个扩展方向一个是给路由模块引入强化学习根据历史执行结果动态调整打分权重而不是靠人工配置另一个是把任务的执行记录反过来用于 Prompt 自动优化让 Agent 能根据失败样本来调整行为。不过这两个方向都还在验证期等有稳定产出再分享细节。最后还是那句老话系统越复杂越要在基础设施上做减法把确定性留给框架把不确定性留给模型。
返回列表