JOTO
Contact us
← AI 智库
大语言模型

Apache RocketMQ 面向 AI 演进:LiteTopic 支撑百万级多 Agent 会话协作

2026 年 8 月 27 日

Apache RocketMQ 推出 LiteTopic 新模型,专为 AI 场景设计,解决 AI 应用在通道拓扑动态性、任务时长不可预测、GPU 成本敏感等约束下对消息系统的全新要求。其通过 Parent Topic/LiteTopic 两层结构、会话级隔离、RocksDB 索引、Ready Set 投递等机制,在 Qoder Cloud Agents 和百炼网关两个真实业务中实现了百万级会话的有序协作、故障恢复与精准限速。

AI 应用对异步协作的新约束

过去一年,阿里云内部有两类 AI 系统在 RocketMQ 上完成了生产落地。一类是编程智能体 Qoder Cloud Agents,单个任务的执行周期可达数小时甚至数天;另一类是百炼平台的模型推理网关,需要在百万级租户规模下完成请求调度。两者的业务形态差异很大,但对消息系统提出的诉求是一致的:会话状态、任务交接和故障恢复需要由基础设施可靠承载,而不是由应用层反复实现。

本文说明 AI 应用给异步协作带来的新约束,RocketMQ 为此进行的一次面向 AI 场景的演进(LiteTopic),以及它在上述两个业务中解决的具体问题。

AI 应用与传统应用在通道拓扑、任务时长、成本结构三个维度的对比示意图
AI 应用与传统应用在通道拓扑、任务时长、成本结构三个维度的对比示意图

相比传统应用,AI 应用在通道拓扑、任务时长和成本结构三个维度都发生了根本变化。

通道拓扑的变化。传统应用的业务流在设计期即已确定,链路上的节点、消息的流向在系统上线后基本不变,因此消息通道的数量和拓扑可以预先规划。AI Agent 的执行路径由运行期决策产生,任务被自主拆解、动态派生,通道的数量与拓扑无法在设计期枚举。

任务时长的变化。传统异步消费的处理耗时在毫秒到百毫秒量级,AI 任务包含多轮推理和工具调用,单次执行常达分钟级,长会话可跨越数天,且耗时高度不可预测。同一任务内的交互轮次显著增加,中间状态需要被频繁读写,状态管理与会话管理由此成为必须在基础设施层解决的问题。

成本结构的变化。GPU 的单位算力成本远高于 CPU,因此在 AI 场景中,任何形式的算力空转都会直接反映在成本上。这一约束贯穿后文的多处设计选择。

在这些约束下,传统的事件驱动模型出现了三处不适配。

第一是积压斜率的变化。入口流量不变,但单条消息的处理耗时从毫秒级上升到分钟级,队列积压的增长速度完全不同,原先通过增加消费者即可吸收的流量抖动,现在需要成倍的算力才能消化。

第二是队头阻塞。传统模型中一个队列由一组消费者共享,其隐含前提是队列内的消息同质且可互换。在 AI 场景中,一条消息对应一个具体用户的一次具体会话,其耗时长是任务自身特征,并非系统异常。在共享队列模型下,该消息会阻塞其后所有任务,导致无关会话被动等待。缓解方式是细化队列粒度,而真正有效的粒度是会话级;传统 Topic 的创建与元数据成本决定了"一个会话一个 Topic"无法成立。

第三是重复推理的成本。单次推理的算力开销已经发生,若因进程重启、实例调度或网络异常导致中间结果丢失,则需重新计算。反之,若历史执行结果能够可靠持久化并在恢复后被读回,即可避免重复调用模型。消息与事件的持久化能力在此不仅关系可靠性,也直接关系成本。

由此可以归纳出消息系统需要补齐的三项能力:支撑海量、动态生成且彼此独立的高速通道将隔离粒度从 Topic 级下沉到会话级,使每个会话具备独立的顺序和故障边界;识别消息与会话的归属关系,从而提供会话级的有序性与亲和性。这三项能力共同指向一个新的 Topic 模型。

RocketMQ LiteTopic 的设计与实现

RocketMQ LiteTopic 的核心思路是将每个 AI 会话映射为一个独立的轻量级 Topic,而非让所有会话共享同一个业务 Topic。

将粒度切至会话级会引入两个必须先行解决的工程问题:百万级通道同时存在时,消费队列索引如何组织;百万级订阅关系下,如何确定当前哪些通道有消息可投。LiteTopic 对应做了两处改造,存储层以 RocksDB KV 引擎替代基于文件的消费队列索引,投递层以事件驱动替代长轮询。

需要说明的是,这不是对 RocketMQ 存储内核的重写。顺序写、零拷贝、主从复制、刷盘策略等经过大规模生产验证的能力保持不变,LiteTopic 只在其上增加一层索引以及一套新的订阅与投递语义,这决定了它的成熟度和接入成本。

▍Topic 模型:Parent Topic 与 LiteTopic 两层结构

LiteTopic 挂载在 Parent Topic 之下,构成两层结构。

Parent Topic 承担命名空间与控制边界的职责,需要预先创建,权限、路由、配额等管控能力挂载在这一层。其数量很少,通常一个业务域或一个 Worker 池对应一个即可。

LiteTopic 是运行期的动态单位,具备有序、隔离、可重放三个属性。它无需预创建,第一条消息到达时自动创建,空闲后按 TTL 自动回收。业务侧需要做的仅是在发送消息时标明该消息所属的会话。

在这一模型下,逻辑通道的创建与销毁成本降至可以按会话使用的程度,通道生命周期与会话生命周期对齐,运维侧无需为通道的增减进行任何操作。

▍消费模型:订阅粒度从 Group 下沉到 Client

RocketMQ 原有集群消费模型与 LiteTopic 消费模型对比图
RocketMQ 原有集群消费模型与 LiteTopic 消费模型对比图

Topic 模型的变化要求消费模型同步调整。RocketMQ 原有的集群消费以 Group 为单位组织,同组内所有实例订阅完全相同的一组 Topic,并在组内做负载均衡。该模型适合同质任务,但无法表达实例之间的差异,也就无法指定某条消息由某个实例处理。

LiteTopic 将订阅粒度细化为同时支持 Group 级与单个 Client 级,消息可以被路由到特定实例,并保证一个 LiteTopic 在任一时刻仅被一个 Client 独占消费。由此,同一会话的消息不会被拆分到多个实例并发处理,避免了顺序错乱与状态竞争。同时,同一消费组内的不同实例可以订阅不同的 LiteTopic 子集,使负载能够按会话维度摊开。

▍存储结构:单份写入,多路索引

RocketMQ 存储分层结构及 LiteTopic 索引构建示意图
RocketMQ 存储分层结构及 LiteTopic 索引构建示意图

RocketMQ 的存储是分层的,消息体统一顺序追加至 CommitLog,ConsumeQueue 仅是指向物理位置的轻量索引。LiteTopic 沿用该结构,改动集中在索引层:后台 dispatch 线程并行构建两套索引,一套是按 Parent Topic 组织的标准 ConsumeQueue,一套是按 LiteTopic 组织的 LiteCQ。消息本身仍只写入一份,写放大受控,消费侧可按 LiteTopic 精确拉取所属通道的消息。

▍索引引擎:以 RocksDB 承载消费队列索引

原有的 ConsumeQueue 采用文件实现,一个队列对应一个独立文件。在数十至数百个 Topic 的规模下不存在问题,但百万级 LiteTopic 意味着百万级小文件,文件句柄、内存映射与元数据管理的开销将非线性增长,写入也会退化为大量随机 I/O。

改造后,所有 LiteCQ 索引统一收敛到一个 KV 引擎中,Key 由 Topic、QueueID、Offset 组成,Value 为消息在 CommitLog 中的物理偏移量。百万级 LiteTopic 共享同一套 KV 存储,单条索引查询为 O(1),元数据规模不再随通道数量膨胀。RocksDB 的 LSM 结构适合高频写入与海量小记录,与该场景的访问特征一致。从“每队列一个文件”转为“共享 KV、O(1) 元数据”,是百万通道共存的前提条件

▍投递机制:以 Ready Set 替代全量轮询

LiteTopic Ready Set 投递机制工作流程图
LiteTopic Ready Set 投递机制工作流程图

存储层解决了通道的容纳问题,投递层需要解决活跃通道的定位问题。长轮询的工作方式是消费者周期性发起请求,Broker 扫描相关队列。当队列数量达到百万级,周期性全量扫描的 CPU 开销显著上升,而其中大部分扫描无效,因为同一时刻真正有消息的通道只占很小比例。

LiteTopic 引入了两个结构:Lite Subscription Set 保存订阅关系,Ready Set 聚合当前确实有消息可投的通道。消息写入时触发事件匹配,命中订阅关系的消息被聚合进 Ready Set,消费者仅从 Ready Set 获取任务。这一改造与 epoll 相对 select 的改进思路一致,调度成本由 O(订阅总数) 降为 O(活跃数):存量会话可以是百万级,而单位时间内活跃的通道仅为千级,系统只需为活跃部分付出调度开销。

▍业务侧能力:按需创建、自动过期、差异化订阅、消费挂起

上述机制在业务侧表现为四项能力,分别对应第一节提出的问题。

  • 按需创建:第一条消息到达时通道自动创建,无需预先申请或在配置中登记,解决海量通道的资源供给问题。

  • 自动过期:空闲通道按 TTL 自动回收,无需业务显式销毁。按需创建必须配套自动回收,否则历史会话的残留元数据会持续累积。

  • 差异化订阅:消息可精确投递至单个 Client,同一消费组内的实例可订阅不同的 LiteTopic 子集,实现会话级隔离。

  • 消费挂起:消费者返回 Suspend(N) 时,仅当前 LiteTopic 暂停 N 毫秒,其余通道不受影响。这一能力用于替代消费线程内的 sleep 限速:sleep 期间消费线程被占用,其本可处理的其他会话只能等待;Suspend 由 Broker 记录该通道下次可投递的时间,消费线程立即返回线程池继续服务其他会话。在 GPU 成本占主导的场景中,这一差异直接影响资源利用率。

▍接入方式:改造集中在会话标识和限速方式两处

生产端在构造消息时增加一个会话维度的标识:

// 生产:按会话维度写入,通道由 Broker 自动创建
Message msg = provider.newMessageBuilder()
       .setTopic(PARENT_TOPIC)
       .setLiteTopic("session-" + sessionId)
       .setBody(payload)
       .build();
producer.send(msg);

消费端订阅 Parent Topic 下的全部通道,需要限速时返回 Suspend:

// 消费:订阅父 Topic 下所有通道,按通道级别调速
consumer.subscribe(PARENT_TOPIC, "*");
// 仅挂起当前通道,不阻塞消费线程
return ConsumeResult.suspend(Duration.ofMillis(50));

以上为示意代码,实际 API 以所用版本为准。原有的编程模型、客户端依赖与运维体系均无需更换,改造集中在消息的会话维度标识和限速方式两处。

生产落地

本节以 Qoder Cloud Agents 为主,它较为完整地使用了上述能力;百炼网关代表另一类问题,放在后半部分。

▍案例一:Qoder Cloud Agents

核心难点:脑手分离架构下的有序交接与故障恢复

Qoder Cloud Agents 是运行在云上的编程智能体服务,接收需求后自主完成代码读取、修改、测试与迭代修正,单个任务的执行周期较长。其架构为脑手分离,Agent Runtime 负责决策,Sandbox Worker 池负责执行,整体事件驱动。会话状态不绑定在具体实例上,因此可以做到无损伸缩,并支撑万级规模的并发推理实例。

需要澄清的是,分布式架构本身不会缩短单次推理的耗时,其收益体现在资源利用率、并发能力与成本效率上。该系统的主要难点在于执行状态的可恢复性、任务交接的有序性、Agent 循环中复杂语义的表达,以及故障恢复与版本演进,这些均属于消息与事件层面的问题。

基于 LiteTopic 的实现:请求/响应双通道,每个会话一条 lane

Qoder Cloud Agents 架构中 Agent Runtime 与 Sandbox Worker 通过 LiteTopic 进行双向通信示意图
Qoder Cloud Agents 架构中 Agent Runtime 与 Sandbox Worker 通过 LiteTopic 进行双向通信示意图

Agent Runtime 与 Sandbox Worker 之间使用两个 Parent Topic 完成双向异步通信,一个承载 request,一个承载 response。Worker 仅订阅 request 方向,Runtime 仅订阅 response 方向,两端均不阻塞在对方的处理过程上。每个会话在两个 Parent Topic 下各拥有一条 lane,即以会话标识命名的 LiteTopic,随首条消息创建,会话结束后按 TTL 回收。

流经消息系统的是任务指令、执行结果与状态变迁等事实与事件,不包含庞大的中间状态。执行进度由 Checkpoint 承载并持久化至存储,大对象与产物同样存放在存储侧,消息中仅保留引用。这里的边界需要明确:消息系统负责有序、可靠地传递事件,而任务当前进展的权威来源是 Checkpoint,不通过消息的 ACK 时机推断。因此实例重启、会话迁移与事件重放都不会导致进度失真。

在这一结构下,会话内消息有序而会话间并行执行,顺序性与并发度可以同时满足。流量突增时由 Broker 承接积压,Worker 按自身处理能力拉取,形成的是通道级的局部背压,单个会话变慢不会传导为全局过载,故障影响也被约束在单条通道内。隔离层次上,Parent Topic 隔离不同的 Worker 池,LiteTopic 隔离业务单元,脑与手可以各自独立伸缩。

Qoder Cloud Agents 中四类复杂场景的统一 LiteTopic 表达示意图
Qoder Cloud Agents 中四类复杂场景的统一 LiteTopic 表达示意图

四类复杂场景:并行、打断、Webhook、挂起恢复

1. 多 Agent 并行

一个会话内会派生多条执行线,任务以 LiteTopic = session-id 有序下发,Worker 抢占后并行执行,结果由调度侧通过 Mailbox 与 Barrier 聚合,从而同时获得分发的确定性、执行的并发度与汇合语义。

2. 打断处理

Agent 执行期间收到用户的追加指令时,新消息先进入该 Agent 的 Mailbox,在执行完成前的下一个安全边界注入;需要中止正在执行的任务时,通过独立的控制通道下发停止指令,随后经 LiteTopic 恢复。数据面与控制面在此必须分开。

3. Webhook 回调

以 LiteTopic = endpoint-id 组织,保证单个 endpoint 的回调有序、不同 endpoint 之间并行,某个 endpoint 变慢或故障时仅阻塞其自身通道。

4. 挂起与恢复

Agent 在等待用户输入、工具返回或子 Agent 完成时,将待处理状态持久化并释放 Worker,等待期间不占用算力;事件返回后,后续任务段经 LiteTopic 下发,由任意空闲 Worker 接手继续执行。传统架构中会话与 Worker 固定绑定,会带来连接中断导致上下文丢失、等待期间实例空转、扩容时会话迁移困难三个瓶颈,改为按需接入后这三点均得到解决。

上述四类场景在同一套 LiteTopic 架构下统一表达,无需为每类场景引入独立机制。

▍案例二:百炼网关

核心难点:限流问题的性质变化

传统网关限流与百炼网关基于 LiteTopic 的漏桶阵列限流对比图
传统网关限流与百炼网关基于 LiteTopic 的漏桶阵列限流对比图

传统网关限流判断调用方的请求速率是否超出配额,超出即拒绝,其前提是后端容量可弹性扩张,限流仅用于防止极端流量打穿系统。模型推理服务不满足该前提,GPU 是当前最稀缺的生产要素,后端处理能力存在硬上限。限流因此从放行与拒绝的二元判断,变为还需确定以何种节奏将请求送往后端。

在百万级租户规模下,单纯的阈值限流会同时放大两类误差。一类是误拒,请求实际可被后端接纳,但瞬时速率超过阈值而被直接拒绝,调用方收到 503,此时后端并未真正饱和;另一类是过载放行,多个维度的流量在同一时刻叠加,各维度均未超过阈值,合并后压垮后端。其原因在于传统限流只有放行与拒绝两种状态,缺少先接纳请求、再按合适节奏送往后端的第三种状态。

基于 LiteTopic 的实现:漏桶阵列按 user + model 精准调速

重构后,网关自身仅保留一层固定窗口的硬上限用于兜底极端情况,请求按 user + model 维度调用 setLiteTopic 写入消息系统,实际调速由 LiteTopic 构成的漏桶阵列承担。漏桶按需创建、按 TTL 回收,无需预先为租户与模型的组合规划资源。消费侧使用统一消费组以 subscribe(parent, "*") 匀速消费,需要调速时返回 Suspend(N ms),精确控制单个 user + model 组合的输出速率而不影响其他组合。

效果上,限流维度从单一维度扩展到数千个维度,限流比降低至原先的十分之一左右,即在相同后端算力下可接纳的有效请求明显增加,被误拒的请求明显减少;调速精度达到毫秒级,可跟随后端实时水位动态调整输出速率。该方案目前运行于公有云百炼网关,服务百万级租户规模。

两个案例的问题形态不同,一个是有序交接,一个是算力分配,但依赖的是同一组能力:会话粒度的通道、按需创建与自动回收、以及不阻塞消费线程的调速。

社区与学术进展:合入 Apache 主干,论文入选 FSE 2026

▍社区进展:合入 master,斩获 OSCAR 与 InfoQ 奖项

社区方面,LiteTopic 已合并进 Apache RocketMQ 的 master 分支,成为主干能力,RocketMQ-A2A 的实现也已开放给社区,目前正在探索 MCP 传输层的实现,将该传输能力延伸至 Agent 与工具之间的通信场景。相关工作获得中国信通院 2025 年度 OSCAR 开源 + AI 尖峰案例奖与 InfoQ AI 开源明星项目奖。

Apache RocketMQ 面向 AI 演进:LiteTopic 支撑百万级多 Agent 会话协作 配图 8

JOTO 企业落地观察

  • 企业部署 AI 智能体系统时,若沿用传统消息队列的 Topic 模型,将面临会话级隔离缺失、通道爆炸性增长、GPU 算力空转等刚性瓶颈;LiteTopic 提供的会话级通道与自动生命周期管理,是支撑百万级并发会话的基础设施前提。
  • 这类系统的取舍在于:是否将状态管理与会话生命周期的复杂性下沉至消息中间件。LiteTopic 的价值不在于替代应用层状态,而在于将“有序性”与“故障边界”的保障责任从应用层移交至基础设施,从而降低智能体工程的整体复杂度。
  • 对于 RAG 知识工程中的多轮问答会话、Agent 工具调用链路追踪等场景,LiteTopic 提供的会话级有序性与可重放性,可直接作为轻量级知识上下文流转载体,避免在应用层维护冗余的状态快照。
  • 在 AI 安全治理层面,会话级隔离天然支持租户/模型维度的审计日志与访问控制,LiteTopic 的元数据(如 user+model)可作为策略执行点,为细粒度的合规性检查提供基础设施支撑。

立即咨询 JOTO

JOTO 提供覆盖企业智能体规划与搭建、AI 平台私有化部署、RAG 知识工程、AI 安全治理、FDE 驻场共创及持续运营优化的全周期 AI 落地服务,帮助企业把验证中的 AI 能力转化为安全、可控、可持续迭代的生产力。 联系 JOTO 获取 AI 落地咨询

想把这些做法用到你的业务里?

留下你的场景和痛点,我们帮你判断从哪一步开始。

联系我们
Contact Us

Start your enterprise AI rollout

Tell us your industry, team, and current pain points. We'll get back to you within one business day to help you decide what to tackle first, what data to prepare, and which platform fits.

WeChat
Scan to add us for a 1:1 chat
JOTO WeChat consultation QR code

Tell us what you need

Once we receive your details, we'll be in touch within one business day.