分布式Agent异步架构
公众号名称:柠遇AI纪元
作者名称:柠遇AI纪元
发布时间:2026-06-01 00:56
在上一篇文章中,其实也提到了,现有的开源agent很多都是针对单体架构来说的,也就是任务、调度器和结果都在同一个运行时内(或者同一个进程内),可以使用 Future、Promise、callback、async/await、线程池或协程队列来完成这些异步任务。对于分布式agent大家也在摸索中前行。

其实分布式agent异步架构无非是要满足:高并发的用户请求?如有任务执行失败该如何处理?对于很多任务的不确定性该如何保证稳定输出?长任务在分布式系统中该如何去编排或者拆分任务?其实还有很多人问题,需要去慢慢解决。总的来说,分布式agent异步系统是一个高并发、长周期、具备自愈能力且面向不确定性(LLM 推理时延、人类干预、网络波动)的分布式事件驱动系统。
一、 分布式异步 Agent 架构总体介绍
架构的核心思想是“将 Agent 与调度器 Orchestrator、状态 State 完全解耦”。

底层架构由以下五个核心层级组成:
【接入层】异步网关Gateway
-
**职责:**接收客户端请求,生成全局唯一的
Session_ID和Trace_ID。 -
**原理:**网关采用 非阻塞IO模式。收到复杂任务请求后,它一直等待 Agent 执行完毕,而是向事件总线投递一个
Task_Created事件,然后立即向客户端返回205 Accepted状态码和任务 ID。后续通过 WebSockets、SSE或长轮询(Long Polling)异步推送结果。
【通信与编排层】分布式事件总线
-
**职责:**负责 Agent 之间、网关与 Agent 之间的消息传递。
-
**原理:**采用发布-订阅模式。常见的底层组件为 Kafka 、RabbitMQ或 Redis Streams。
-
控制agent任务状态:可以使用一个中枢大脑的 Workflow Agent 订阅状态,通过异步的方式向子agent发送专门任务,最后通过sse的方式返回给中心大脑;还有一种不是由中心大脑控制,而是agent自己从订阅的当前事件的消息队列中,自动拉取对应的数据。
【执行层】无状态 Agent 算子集群
-
**职责:**真正执行 LLM 推理、Tool/MCP/Skills 调用的实体。
-
原理:Agent 被设计为无状态的。它们以微服务、容器Docker或 Serverless 节点形式存在。当有事件到来时,任何一个空闲的 Agent 节点都可以把任务抢过去执行,执行完后将状态写入持久化层,自己继续保持无状态;若没有空闲节点,则等待空闲节点来拉取当前任务。
【状态与记忆层】分布式状态存储
-
**职责:**维持长周期任务的生命。
-
原理:对于短期状态的数据,其中包括当前chat对话的历史和 Agent 执行到第几步的状态机信息,可以使用Redis中间件来存储,从而agent可以根据chat上下文,输出用户所需要的内容;而对于长期记忆的数据,使用向量数据库和关系型数据库,存储知识库、用户画像及审计日志。注意:现在很多短期状态的数据,可以把它压缩存入当前work-tree的Agent.md文件下,但是对于分布式系统来说,不能存入系统层面的work-tree,而是用户层面work-tree,方便进行上下文的压缩和总结。
二、 接入层
具体来说,接入层是整个分布式 Agent 系统的“流量闸门”与“非阻塞分流器”。
主要职责
-
**协议转换与统一接入:**统一接收来自 Web、Mobile、CLI 或第三方 Webhook 的请求(支持 HTTP/REST, WebSockets, gRPC)。
-
**请求非阻塞化:**将上游同步的“请求-响应”解耦为异步的“事件派发”。
-
**全局唯一标识注入:**为每个进入系统的请求打上分布式唯一标识(
Session_ID标识会话,Trace_ID串联生命周期)。
可能遇到的问题
【问题】**Agent长连接雪崩与 I/O 阻塞。**当数万个客户端同时在线,且每个 Agent 任务需要运行数分钟时,同步网关会因为线程耗尽而崩溃。
【解决办法】通常在分布式环境中可以基于异步 I/O 网络模型来解决,有如下方式:
-
采用 Epoll 模型构建无阻塞网关。
-
接入即断开机制:网关在收到请求的瞬间,完成参数校验与 Schema 检查,立即调用分布式 ID 生成器(如雪花算法 Snowflake)生成
Trace_ID。接着,网关直接向通信层投递一个Task_Submitted事件,并向客户端返回Accepted状态,随后网关立即释放当前连接。 -
**异步状态推流:**客户端如果需要获取实时进度,网关会通过 HTTP Long-Polling或 **WebSockets (gRPC Streaming)**建立一条轻量级的状态监听通道,通过多路复用技术,用极少的系统线程撑起百万级并发的长连接
三、通信与编排层
这一层是分布式系统的“神经系统”与“交通枢纽”,负责 Agent 间消息的无损、高吞吐传递。
主要职责
-
**事件路由:**负责将消息精准投递给目标 Agent 或订阅了特定主题的 Agent 集群
-
**拓扑控制:**定义多 Agent 协作的流向
可能遇到的问题
【问题】有大量的请求同时调用同一个模型API-Key,会造成大量的大模型 API 请求出现429限流,从而整个系统的引发的雪崩
【解决办法】可以采用分布式反压与弹性队列方法
-
基于 NATS JetStream或 Kafka实现带优先级的分布式令牌桶队列
-
可以效仿一下litellm,将每一个大模型供应商(如 OpenAI、Anthropic、claude等)在通信层都被抽象为一个“虚拟令牌池”。当下游执行层因为大模型限流(RPM/TPM 耗尽)而回传
RateLimit_Signal时,通信层启动反压机制。队列会主动放慢消息向该类 Worker 投递的速度,将消息安全地积压在持久化队列中,并采用指数退避算法异步重试,保证上游业务不中断。
【问题】在多agent和多用户请求的架构下,多agent消息乱序与重发
【解决办法】可以采用事件顺序或者确认机制来解决
-
采用支持 Raft 共识协议的消息底座。每个事件都带有严格递增的单调序列号
-
采用 确认机制。Agent 只有在完整处理完事件并将新状态写入状态层后,才会向通信层发送
ACK。若 Agent 闪退,消息会在超时后自动重新投递给其他节点,实现故障无感转移
四、 执行层
执行层是承载 Agent 运行(LLM 推理、Tool 调用、代码执行)的集群。
主要职责
-
**LLM 编排与调用:**执行 Prompt 拼接、大模型异步调用及 Token 流式解析
-
**工具沙箱执行:**异步调用外部 API 或在隔离环境中运行 Agent 生成的代码
可能遇到的问题
【问题】在多agent调用llm、tool、skills或者mcp的时候,高并发请求场景下,可能会把内存撑爆或者cpu使用率急剧飙升至100%
【解决办法】当 Agent 异步调用(时间比较长如60s) LLM 并等待其返回,或者等待用户的权限放开时,执行层会触发钝化。该 Agent 的内存上下文会被立刻 dump 到状态层,执行层释放该 Actor 所占用的所有 CPU 和内存资源。当大模型返回或人类点击审批时,通信层发送唤醒事件,任何一个空闲的 Worker 节点从状态层拉取上下文,原地激活该 Actor 继续执行。
另外需要说明下:** 对于有 Tool 调用或 Code Interpreter 需求的 Agent,执行层通过集成 **WebAssembly (WASM)**或 Docker 隔离容器作为执行宿主,通过异步进程池管理,防止恶意代码拖垮宿主机。
五、 状态层
状态层是实现 Agent“长周期运行”与“断点续传”的生命线,解决的是分布式环境下的数据一致性问题。
主要职责
-
**状态机维护:**实时记录 Agent 当前处于什么生命周期阶段(如
Init,Running,Suspended_HITL,Completed,Failed) -
**执行轨迹持久化:**记录 Agent 迈出的“每一步”
可能遇到的问题
【问题】Agent 长周期运行中途宕机,导致任务死锁或重复执行扣处积分或者token
【解决办法】通过借鉴 Temporal的工业级设计理念。状态层基于 Redis(存储活跃状态)+ PostgreSQL(存储历史轨迹)。不保存覆盖状态,只保存事件序列,
举个例子说明
❝
Step 1: Save_Status(Agent_Called_LLM)
Step 2: Save_Status(LLM_Returned_Tool_Call_XYZ)
Step 3: Save_Status(Tool_XYZ_Executed_Successfully)
当某台 Worker 服务器在 Step 2 和 Step 3 之间由于硬件故障突然宕机时,分布式调度器会启动自愈。新的 Worker 接管该 Agent,它不需要重新调用大模型(避免了二次消耗 Token 和时间),而是向状态层“回放”之前的事件日志,发现 Step 2 已经成功,它会直接拿到 Step 2 的 LLM 结果,异步发起 Tool 调用。
六、 记忆层
记忆层是 Agent 表现出“智能”的源泉,负责管理海量并发下的上下文检索与知识沉淀。
主要职责
-
**短期记忆管理:**当前 Session 的 Token 窗口内对话历史的管理与自动总结。
-
**长期记忆检索:**跨 Session 的用户画像、行业知识库、企业私有数据的向量化存储与异步检索。
可能遇到的问题
【问题】高并发下向量检索与 Embedding 计算的高延迟,拉长了 Agent 的整体响应时间。
【解决办法】记忆层采用 **Redis + 分布式向量数据库(如 Milvus)**的双层设计
-
异步检索:当接入层网关收到用户的 Prompt 时,记忆层会和通信层并发启动。在通信层还在为 Agent 调度节点、打包事件的同时,记忆层已经异步将用户的输入提交给 Embedding 模型,并在向量数据库中完成了 Top-K 的知识检索。当 Agent Worker 被激活、准备拼接 Prompt 的那一刻,相关的长期记忆和短期记忆已经作为“就绪附加件(Attachment)”推到了 Worker 的本地缓存中。
-
异步记忆向量化存储:当一次对话结束,Agent 产生的中间知识、新学到的用户偏好,不会同步写入向量库。而是由一个后台的记忆存储订阅
Session_Finished事件,在后台异步地进行文本清洗、总结、Embedding 计算并存入向量库。这让前台的 Agent 推理流程完全脱离了数据库 I/O 写入的羁绊。
七、代码设计
为了更形象地说明,可以用高并发事件驱动的伪代码(基于订阅/分发模式)来说明底层核心骨架:
// 1. 定义通用事件结构
interface AgentEvent {
eventId: string;
sessionId: string;
traceId: string;
sender: string;
eventType: 'TASK_ASSIGNED' | 'LLM_COMPLETED' | 'TOOL_REQUIRED' | 'HUMAN_APPROVAL_REQUIRED';
payload: any;
timestamp: number;
}
// 2. 抽象的分布式 Agent 基类
abstractclass DistributedAgent {
protected agentId: string;
protected eventBus: DistributedEventBus; // 封装了 Redis/Kafka 的底层连接
constructor(agentId: string, eventBus: DistributedEventBus) {
this.agentId = agentId;
this.eventBus = eventBus;
}
// 初始化时异步订阅属于自己的事件/Topic
async init() {
awaitthis.eventBus.subscribe(`topic.agent.${this.agentId}`, async (event: AgentEvent) => {
awaitthis.handleEvent(event);
});
}
// 核心的异步事件处理机
abstract handleEvent(event: AgentEvent): Promise;
}
// 3. 一个具体的 Agent 实现
class ResearchAgent extends DistributedAgent {
async handleEvent(event: AgentEvent) {
if (event.eventType === 'TASK_ASSIGNED') {
// Step 1: 异步更新状态机为 "Running"
await StateManager.updateStatus(event.sessionId, 'RUNNING');
// Step 2: 异步调用大模型推理(非阻塞)
const llmResult = await LLMService.callAsync(event.payload.prompt);
// Step 3: 推理完成后,异步投递新事件,由框架决定下一步给谁
const nextEvent: AgentEvent = {
eventId: generateId(),
sessionId: event.sessionId,
traceId: event.traceId,
sender: this.agentId,
eventType: 'LLM_COMPLETED',
payload: { content: llmResult },
timestamp: Date.now()
};
awaitthis.eventBus.publish('topic.orchestrator', nextEvent);
}
}
}
八、 整体流程梳理
通过一个用户请求的生命周期,来看这五个层级是如何完美异步协同的

行文最后说明:这周更新实在是太忙了,主要自己做了一个Vibe-coding项目,花了很长时间在调试前后端以及前端页面上,后续会开源出来,还望各位大佬帮忙提提意见,继续改进项目。最后我也想写出比较好的文章,但是需要花费很长的时间来搜索相关技术、整理技术文档、生图prompt编写以及文章整体规划,还望大家海涵。
内容效果不满意?点此反馈