RocketMQ LiteTopic:面向 AI 长任务的轻量消息
hi,我是阿昌。学习记录一下 RocketMQ 的 LiteTopic,它不是简单地把普通 Topic 做轻,而是从 AI 应用的长任务、长会话、多 Agent 协作这些问题,重新设计对应的消息通信模型。
传统微服务大多处理毫秒级请求;AI 应用常常要跑几分钟甚至几十分钟。
通信模型不变,超时、断线、级联阻塞和资源浪费就会一起出现。
一、是什么
LiteTopic是 RocketMQ for AI 提供的一种轻量消息能力。
它不是只有一层Topic,而是采用两级模型:
1 | 父 Topic:chatAgentBotTopic |
可以这样理解:
- 父 Topic 是一个业务大类,比如
chatbot、agent-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 | Supervisor |
如果它们通过同步 HTTP 调用协作,就会有三个问题:
- Supervisor 要一直等待最慢的 Agent;
- 任意一个 Agent 超时,整条任务链都可能失败;
- 其他已经完成的 Agent 结果,也可能随着整体失败一起作废。
3. AI 会话是长时间、有状态的
用户和 AI 聊半小时,网络断了一次。
如果会话状态只在当前连接或应用内存中:
1 | 断线 → 连接丢失 → 进度丢失 → 用户重试 → 再烧一遍GPU算力 |
但如果每个会话都有自己的 LiteTopic,服务端继续把消息写进去;用户重连后带上上次消费的 offset,就能从断点续传。
后台任务可以继续跑,用户回来后继续读,不需要从头再来。
三、怎么做
1. 先创建一个“轻量消息”父 Topic
控制台中先创建一个类型为轻量消息的 Topic,比如:
1 | chatbot |
这是 LiteTopic 的父 Topic。生产和消费时都先绑定它。
然后业务再通过 LiteTopic 名称,区分具体会话或任务:
1 | chatbot/session_001 |
2. 发送时指定 LiteTopic 名称
发送消息时,除了指定父 Topic,还要设置 LiteTopic 名称:
1 | Message message = provider.newMessageBuilder() |
首次发送时,服务端会自动创建这个 LiteTopic。
所以,应用不需要提前手工创建百万个会话 Topic;根据会话 ID 或任务 ID 写入即可。
3. 消费时订阅具体 LiteTopic
消费者先绑定父 Topic,再订阅要处理的 LiteTopic:
1 | LitePushConsumer consumer = provider.newLitePushConsumerBuilder() |
一条会话就是一条有序消息流。
前端断线时,网关可以重建连接并从已有消费位置继续推送。
4. 多 Agent 改成“任务和结果分离”
任务和结果不要走一条同步 HTTP 链路,而是拆成两个方向。
1 | Supervisor |
这样 Supervisor 发布完任务后不必阻塞等 HTTP 返回,而是在收到全部结果后再做汇总,就解耦了线程,提升了吞吐量,最后通过 SSE 推给用户。
5. 用 Event-Driven Pull 承载大量订阅
一个消费者订阅很多 LiteTopic 时,如果还逐个长轮询:
1 | topic-001 有消息吗? |
网络请求会随着订阅数量增长,造成网络压力大;
LiteTopic 使用Event-Driven Pull:
Broker 维护消费者的订阅集合,再把有新消息的 LiteTopic 聚合成就绪集合,消费者一次拉取多个就绪消息。
1 | 订阅集合 Subscription Set |
这样是为了让大量会话或任务订阅时,网络开销不会跟着线性膨胀。
6. 响应结果通过选择性消费回到对应 Supervisor
多 Agent 协作时,任务可以由多个 Worker 集群消费,但结果不能被任意一个 Supervisor 随机抢走。
配置思路是:
1 | 任务 Topic:普通消息 + CLUSTERING 集群消费 |
这样,任务由多个 Worker 分摊执行;
响应则只交给发起这次协作的 Supervisor 消费。
1 | Supervisor A 发起 session_001 |
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 | if (user.requestCount > user.rateLimit) { |
一个用户对应一个 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 | 创建会话/任务通道 → 异步持续写消息 → 按序消费 → 断线后继续读取 |
它比较适合解决四类问题:
