☰
多Agent协作编排引擎Agent-Reach架构设计与落地实践
2026/10/6 5:11:51 网站建设 项目流程

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。消息在队列里流转时,协调中心会持久化两个 ID:trace_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 的输入输出在协议层都是同构的,换来的是整个系统在新增角色时几乎零沟通成本。在实现时我把这条消息总线包了一层客户端 SDK,Agent 接入时只需要调用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_factory=list) 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_id=route.agent_id, intent=task.intent, payload=task.payload, timeout_ms=task.timeout_ms ) task.result = result task.status = TaskStatus.SUCCEEDED await self._release_downstream(task.task_id) except asyncio.TimeoutError: await self._handle_failure(task, error="timeout") except Exception as exc: await self._handle_failure(task, error=str(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依赖目标 Agent
T1classify_ticket无agent-classifier
T2search_knowledgeT1agent-rag
T3generate_replyT1, T2agent-writer
T4send_emailT3agent-mailer
T5update_crmT3agent-crm

T1 和 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 处理完任务后发了一个事件给 B,B 处理后反过来给 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 能根据失败样本来调整行为。不过这两个方向都还在验证期,等有稳定产出再分享细节。最后还是那句老话:系统越复杂,越要在基础设施上做减法,把确定性留给框架,把不确定性留给模型。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询