RocketMQ LiteTopic:面向 AI 长任务的轻量消息
阿昌 Java小菜鸡

RocketMQ LiteTopic:面向 AI 长任务的轻量消息

hi,我是阿昌。学习记录一下 RocketMQ 的 LiteTopic,它不是简单地把普通 Topic 做轻,而是从 AI 应用的长任务、长会话、多 Agent 协作这些问题,重新设计对应的消息通信模型。

传统微服务大多处理毫秒级请求;AI 应用常常要跑几分钟甚至几十分钟。

通信模型不变,超时、断线、级联阻塞和资源浪费就会一起出现。

一、是什么

LiteTopic是 RocketMQ for AI 提供的一种轻量消息能力。

它不是只有一层Topic,而是采用两级模型:

1
2
3
4
5
父 Topic:chatAgentBotTopic
├── LiteTopic:session_1
├── LiteTopic:session_2
├── LiteTopic:session_3
└── ...

可以这样理解:

  • 父 Topic 是一个业务大类,比如chatbotagent-task
  • LiteTopic 是这个大类下面的细粒度消息通道;
  • 一个 LiteTopic 可以对应一次用户会话、一个 Agent,或者一个长耗时任务;
  • LiteTopic 由生产端首次写消息时自动创建,并按 TTL 自动清理。

Session-as-Topic一个会话就是一个 LiteTopic。

以前会话状态经常放在应用内存、Redis 或数据库里;

现在,会话产生的过程消息、推理结果、进度事件可以持续写入这个会话对应的 LiteTopic。

二、为了什么

LiteTopic 主要是在解决 AI 应用和传统微服务的三个差别。

1. AI 调用太慢,HTTP 同步连接扛不住

传统微服务一次请求通常在几十到几百毫秒内结束:

1
请求 → 处理 → 响应

AI 应用不一样。一次大模型推理、代码生成、图片生成、深度检索,可能从几十秒持续到几十分钟。

如果还让客户端、网关、Supervisor Agent 一直保持 HTTP 同步等待,就容易出现:

  • 连接超时;
  • 网络抖动后上下文丢失;
  • 已经跑完大半的 LLM 任务没有结果可取;
  • GPU 仍在继续计算,但用户只能重新发起请求。

所以这里不是“接口慢一点”的问题,而是同步调用的假设已经不成立了。

2. 多 Agent 协作会出现级联阻塞

假设一个 Supervisor Agent 要依次或并行协调多个 Sub-Agent:

1
2
3
4
5
Supervisor
├── Weather Agent(30 秒)
├── Travel Agent(2 分钟)
├── Finance Agent(45 秒)
└── 汇总结果

如果它们通过同步 HTTP 调用协作,就会有三个问题:

  • Supervisor 要一直等待最慢的 Agent;
  • 任意一个 Agent 超时,整条任务链都可能失败;
  • 其他已经完成的 Agent 结果,也可能随着整体失败一起作废。

3. AI 会话是长时间、有状态的

用户和 AI 聊半小时,网络断了一次。

如果会话状态只在当前连接或应用内存中:

1
断线 → 连接丢失 → 进度丢失 → 用户重试 → 再烧一遍GPU算力

但如果每个会话都有自己的 LiteTopic,服务端继续把消息写进去;用户重连后带上上次消费的 offset,就能从断点续传。

后台任务可以继续跑,用户回来后继续读,不需要从头再来。

三、怎么做

1. 先创建一个“轻量消息”父 Topic

控制台中先创建一个类型为轻量消息的 Topic,比如:

1
chatbot

这是 LiteTopic 的父 Topic。生产和消费时都先绑定它。

然后业务再通过 LiteTopic 名称,区分具体会话或任务:

1
2
3
chatbot/session_001
chatbot/session_002
chatbot/session_003

2. 发送时指定 LiteTopic 名称

发送消息时,除了指定父 Topic,还要设置 LiteTopic 名称:

1
2
3
4
5
6
7
Message message = provider.newMessageBuilder()
.setTopic("chatbot")
.setLiteTopic("session_001")
.setBody(body)
.build();

producer.send(message);

首次发送时,服务端会自动创建这个 LiteTopic。

所以,应用不需要提前手工创建百万个会话 Topic;根据会话 ID 或任务 ID 写入即可。

3. 消费时订阅具体 LiteTopic

消费者先绑定父 Topic,再订阅要处理的 LiteTopic:

1
2
3
4
5
6
7
8
9
10
11
LitePushConsumer consumer = provider.newLitePushConsumerBuilder()
.setClientConfiguration(clientConfiguration)
.setConsumerGroup("chatbot-push-group")
.bindTopic("chatbot")
.setMessageListener(messageView -> {
// 处理本次会话的消息
return ConsumeResult.SUCCESS;
})
.build();

consumer.subscribeLite("session_001");

一条会话就是一条有序消息流。

前端断线时,网关可以重建连接并从已有消费位置继续推送。

4. 多 Agent 改成“任务和结果分离”

任务和结果不要走一条同步 HTTP 链路,而是拆成两个方向。

1
2
3
4
5
6
7
8
9
10
Supervisor
├── 发布任务到 WeatherAgentTask
├── 发布任务到 TravelAgentTask
└── 订阅自己的 Response LiteTopic

Weather Agent
└── 消费任务 → LLM 推理 → 结果写回 Response LiteTopic

Travel Agent
└── 消费任务 → LLM 推理 → 结果写回 Response LiteTopic

这样 Supervisor 发布完任务后不必阻塞等 HTTP 返回,而是在收到全部结果后再做汇总,就解耦了线程,提升了吞吐量,最后通过 SSE 推给用户。

5. 用 Event-Driven Pull 承载大量订阅

一个消费者订阅很多 LiteTopic 时,如果还逐个长轮询:

1
2
3
4
topic-001 有消息吗?
topic-002 有消息吗?
topic-003 有消息吗?
...

网络请求会随着订阅数量增长,造成网络压力大

LiteTopic 使用Event-Driven Pull

Broker 维护消费者的订阅集合,再把有新消息的 LiteTopic 聚合成就绪集合,消费者一次拉取多个就绪消息。

1
2
3
4
5
订阅集合 Subscription Set

就绪集合 Ready Set

一次拉取多个有消息的 LiteTopic

这样是为了让大量会话或任务订阅时,网络开销不会跟着线性膨胀。

6. 响应结果通过选择性消费回到对应 Supervisor

多 Agent 协作时,任务可以由多个 Worker 集群消费,但结果不能被任意一个 Supervisor 随机抢走。

配置思路是:

1
2
任务 Topic:普通消息 + CLUSTERING 集群消费
响应 Topic:轻量消息 + LITE_SELECTIVE 选择性消费 + 顺序投递

这样,任务由多个 Worker 分摊执行;

响应则只交给发起这次协作的 Supervisor 消费。

1
2
3
4
5
6
7
8
9
Supervisor A 发起 session_001

写任务到 WeatherAgentTask

任意一个 Weather Worker 消费并执行

结果写入 session_001 对应的 Response LiteTopic

只有订阅 session_001 的 Supervisor A 收到

CLUSTERING:解决 Worker 扩容和任务分摊。

LITE_SELECTIVE:解决结果不能被其他会话、其他 Supervisor 误消费。

顺序投递:保证同一次会话中的“进度 → 中间结果 → 最终结果”按正确顺序到达。

这里的关键不是“物理机器 A 必须收到”,而是负责 session_001 这个会话的逻辑消费者收到。

如果 Supervisor A 重启了,新实例接管 session_001 后,仍然可以继续从对应 LiteTopic 消费结果;

这正是消息化比 HTTP 回调更稳的地方。

四、特点有哪些

1. 两级 Topic 模型

普通 Topic 是单层模型;

LiteTopic 是父 Topic + LiteTopic的两级模型。

  • 父 Topic 用于管理业务大类
  • LiteTopic 用于隔离具体会话、具体 Agent 或具体任务。

2. 自动创建和自动回收

LiteTopic 在首次写入时自动创建;TTL 内没有新消息写入时,会自动清理。

这点特别适合会话、临时任务、短生命周期 Agent 协作这类场景,避免大量长尾 Topic 一直占着运维成本。

3. 一个 LiteTopic 内天然有序

每个 LiteTopic 底层对应一个队列,生产、存储和消费都围绕这个队列进行,因此同一个 LiteTopic 内的消息可以严格按顺序处理。

例如 RAG 知识库更新:

1
删除文档 v1 → 写入文档 v2

只要这两个事件在同一个 LiteTopic 内,就能避免后到的旧操作把新状态覆盖掉。

4. 排他消费

LiteTopic 支持Exclusive排他消费模式。

对于同一个 LiteTopic,Broker 只允许同一个消费组下最新连接的客户端消费;旧客户端会被取消订阅。

这个机制很适合“一个会话只应由一个网关连接持续推送”的场景,避免两个客户端同时推送同一条会话消息。

5. 支持更细粒度的限流

LiteTopic 引入Suspend消费状态。

它不是成功,也不是失败,而是“本次先暂停,到指定时间后自动恢复”:

1
2
3
if (user.requestCount > user.rateLimit) {
return ConsumeResult.suspend(Duration.ofMillis(500));
}

一个用户对应一个 LiteTopic 时,暂停这个 LiteTopic,就相当于只暂停这个用户;其他用户仍然可以正常消费。

6. 面向 5.x gRPC SDK

LiteTopic 只支持 RocketMQ 5.x 集群和 5.x gRPC SDK:

  • Java SDK 需要v5.1.0及以上;
  • Go SDK 需要v5.1.4及以上;
  • 功能仍处于灰度阶段,使用前需要确认集群是否已开通能力。

7. 使用 RocksDB 管理海量轻量队列索引

RocketMQ 的消息主体仍然统一写入 CommitLog,但为了支撑百万级 LiteTopic,索引层使用RocksDB管理海量轻量队列的元数据。

消息存储主干没有推倒重来,重点是把原本难以支撑海量小队列的索引管理方式升级了。

五、有什么优点

1. 会话不会因为断线直接丢掉

会话消息持续落在对应 LiteTopic 中。

用户断线后,后台任务继续执行;重连后从原来的 offset 继续读取。这个模型对 AI 流式输出特别实用。

2. 多 Agent 不再互相“卡死”

任务发布和任务执行解耦后,慢 Agent、失败 Agent 不会把 Supervisor 的 HTTP 线程一起拖住。

已经完成的结果也可以持续保留在结果 LiteTopic 中,等汇总方恢复后再继续处理。

3. 会话和租户隔离更细

一个会话、一个知识库或一次任务就是独立 LiteTopic。

因此慢租户、大批量更新、单次异常任务,更容易限制在自己的边界内,不至于把公共队列一起堵住。

4. 长尾资源可以自动释放

会话结束、任务完成、知识库长期不更新后,TTL 自动回收对应 LiteTopic。

相比大量人工创建和清理普通 Topic,这种方式更适合百万级、短生命周期的对象。

5. 从“全局限流”变成“按人限流”

传统系统一旦全局限流,可能所有用户一起受影响。

LiteTopic 和Suspend结合后,可以只暂停某个用户或任务对应的消息流,达到文“千人千面”流控。

六、缺点是什么

这里的缺点,一部分来自产品当前限制,一部分来自使用这种模型时需要接受的工程成本。

1. 不是所有 RocketMQ 环境都能直接用

文档明确说明,功能仅面向5.x集群,并且当前处于灰度阶段;集群没有入口时,还需要联系腾讯云开通或升级。

所以它不是改几行业务代码就能接入的通用能力。

2. 客户端升级成本

需要使用 5.x gRPC SDK,旧版本客户端或旧协议不能直接套用。

如果已有大量基于旧 SDK 的生产者和消费者,接入前要评估升级、联调和灰度成本。

3. Group 的使用方式和普通消息不同

消费 LiteTopic 时,Group 需要配置为Lite Topic 消费,并且当前一个 Group 只支持绑定一个父 Topic。

另外,顺序投递场景的重试使用固定间隔,默认重试间隔为 1 秒。这些都要在消费组设计时提前考虑。

4. 架构复杂度上升

把同步调用拆成任务 Topic、结果 LiteTopic、Supervisor 汇总、SSE 推送后,系统弹性更好,但也需要多设计几件事:

  • 如何关联任务和响应;
  • 多个 Agent 结果何时算齐;
  • 超时任务如何结束;
  • 消费失败和重复消息如何幂等;
  • LiteTopic 的 TTL 如何匹配业务留存时间。

它解决了同步链路的问题,但不会自动解决业务编排问题。

5. 更适合 AI 长任务,不一定适合所有业务

如果只是一个几十毫秒就能返回的普通 CRUD 接口,直接 HTTP/RPC 往往更简单。

LiteTopic 的价值主要体现在:长耗时、长会话、高并发隔离、顺序事件和大量动态任务这些场景。

七、最后对比普通 Topic 和 LiteTopic

对比项 普通 Topic LiteTopic
模型 单层 Topic 父 Topic + LiteTopic 两级模型
创建方式 通常预先创建 LiteTopic 首次写入自动创建
生命周期 永久存在需手动人工删除 可按 TTL 自动清理
单实例 Topic 规模 万级性能开始衰退 百万级共存
顺序性 取决于具体队列和消费设计 同一 LiteTopic 内顺序投递
消费方式 常规消费模型 支持 LiteTopic 消费和排他消费
限流粒度 通常是 Topic、Group 或全局维度 可结合 Suspend 按 LiteTopic 精细限流
推荐场景 常规业务事件、稳定的固定 Topic AI 会话、Agent 协作、RAG 同步、长耗时任务
重平衡 集群级全局重平衡 仅更新单Topic绑定记录

普通 Topic:适合固定、常规的业务消息通道。
LiteTopic:适合按会话、用户、任务动态出现,并且需要顺序、隔离、续传的消息通道。

一个简单的选型判断

  • 一个请求很快结束、Topic 数量固定:优先普通 Topic 或直接 RPC;
  • 一个任务要运行很久,可能断线、重试、恢复:可以评估 LiteTopic;
  • 每个用户、会话、知识库都要独立排队:LiteTopic 更贴合;
  • 需要一个会话只由一个连接推送:评估排他消费;
  • 需要控制单个用户的推理节奏:评估Suspend和按 LiteTopic 限流。

八、优缺点

优点

1. 专为AI场景设计的LiteTopic百万级轻量主题,自动创建和销毁,每个AI会话可以映射为独立Topic。这是RocketMQ for AI最核心的竞争力。

2. 解决AI长耗时阻塞问题将同步调用转为异步非阻塞,系统吞吐量大幅提升。

3. 分布式会话状态管理通过LiteTopic实现会话状态外置,应用节点无状态化,断线可续传。

4. 智能算力调度流量整形、消息优先级、定速消费三重机制,让GPU算力用在刀刃上。

5. 生态完善原生支持MCP(Model Context Protocol)和A2A(Agent-to-Agent)协议,与LangChain、CrewAI、AutoGen、Dify等主流AI框架无缝集成。

6. 万亿级消息规模验证在阿里内部经过万亿级消息规模的实战检验。

缺点

1. LiteTopic目前主要在云上版本虽然会逐步贡献到Apache RocketMQ开源社区,但目前完整的AI能力在开源版本中还在逐步落地中。

2. 学习曲线LiteTopic、Lite Mode等新概念需要一定学习成本。

3. 需要重新设计架构从同步调用改为异步消息驱动,需要对现有系统架构进行调整。

LiteTopic,它最重要的变化不是“多了一种 Topic 类型”,而是把 AI 应用中的会话、任务和消息通道连在了一起。

以前的思路常常是:

1
HTTP 调用 → 等结果 → 返回

LiteTopic 的思路则变成:

1
创建会话/任务通道 → 异步持续写消息 → 按序消费 → 断线后继续读取

它比较适合解决四类问题:

image

 请作者喝咖啡