Clipping 微信公众号

分布式Agent异步架构

by 柠遇AI纪元 原文 ↗
Created: 2026-06-17

公众号名称:柠遇AI纪元

作者名称:柠遇AI纪元

发布时间:2026-06-01 00:56

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

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

一、 分布式异步 Agent 架构总体介绍

架构的核心思想是“将 Agent 与调度器 Orchestrator、状态 State 完全解耦”。

底层架构由以下五个核心层级组成:

【接入层】异步网关Gateway

  • **职责:**接收客户端请求,生成全局唯一的 Session_IDTrace_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 JetStreamKafka实现带优先级的分布式令牌桶队列

  • 可以效仿一下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(存储历史轨迹)不保存覆盖状态,只保存事件序列,

举个例子说明

  1. Step 1: Save_Status(Agent_Called_LLM)

  2. Step 2: Save_Status(LLM_Returned_Tool_Call_XYZ)

  3. 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编写以及文章整体规划,还望大家海涵。


内容效果不满意?点此反馈

输入关键词开始搜索