Apache RocketMQ 5.5.0 正式开源 LiteTopic 消息模型,阿里云 5.x 实例同步上线 RocketMQ for AI,官方定义为“AI 原生异步通信引擎”。
传统共享队列为什么扛不住 AI 会话
AI 推理时间在分钟级别之内,一个消息只对应一个会话,慢是任务本身的特性,并不是出现了问题。放入共享队列中:队首阻塞,慢会话堵塞整个队伍,增加消费者的也无济于事;积压斜率剧变,单条消费从毫秒变成分钟,流量抖动要成倍计算才能被处理掉;进程重启丢失状态、GPU 重新计算的成本会翻倍
传统的模型也做不好“同一个会话对应一个 Topic”的事情——Topic 需要事先被创建出来,并带有元数据以及路由开销,而且同一个 ConsumerGroup 必须共享订阅关系。两个模型的区别是什么,一目了然:
// 传统模型:所有会话挤在一条队列
Topic("chat") ─── [A][B][C][A][C][B]... 一条慢,全体堵
LiteTopic:一个会话一条独立通道
Topic("chat"), 父主题
├── chat/sessionA ─── A的消息流(独立保序)
├── chat/sessionB ─── B 的消息流
└── chat/sessionC ─── C 的消息流(互不阻塞)
LiteTopic:两级模型,用代码说话
Parent Topic 预创建,只挂命名空间、配额、路由;LiteTopic 悬挂于此处,发送的时候会自动创建一个,TTL 过期之后就会被自动回收掉,因此可以达到上百万的效果。发送方只需要再做一步工作,给消息加上一个会话通道:
Producerproducer= provider.newProducerBuilder();
.setTopic("chat")
.setClientConfiguration(clientConfig);
.build();
Messagemessage= provider.newMessageBuilder();
.setTopic(chat)
.setLiteTopic("chat/" + sessionId) // 通道不存在时自动创建
.setBody(payload)
.build();
producer.send(message);消费端使用 LitePushConsumer 来接管会话,在接收到消息后订阅该消息,在发送消息之前取消对该消息的订阅:
LitePushConsumerconsumer= provider.newLitePushConsumerBuilder();
.setConsumerGroup("gateway")
bindTopic("chat")
.setClientConfiguration(clientConfig);
.setMessageListener(msg -> ConsumeResult.SUCCESS)
.build();
consumer.subscribeLite("chat/" + sessionId));接管会话之后,消费进度就会自动续上
consumer.unsubscribeLite("chat/" + sessionId));结束必退订,否则占配额
应用层完全无状态:消费位点由 broker 持久化,客户端重连从上一次位点继续消费,即断点续传、推理任务不会重新运行
可以承受百万级别的共存,依靠三个硬改造:
1、索引存储:消费队列索引由 ConsumeQueue 文件改为 RocksDBKV。底层 CommitLog 按照顺序进行写入,并且没有做零拷贝、主从复制等操作,只是在上面新增了一个新的索引层
2. 订阅模式:订阅关系由 ConsumerGroup 向下推送到客户端层面,同一个 group 在不同客户端可以各自订阅自己想要的 LiteTopic,并且运行期会根据需要自动增加或减少 LiteTopic 的数量。在排他消费的情况下,broker 会选择最近的一个客户端来接替,并且通知原来的客户端退出
3、投递语义:消息写入建立索引之后 broker 发出 ready 事件,客户端从 ready 集合中合并拉取,代替扫全量队列的长轮询,避免百万通道的无效扫描
差异与硬坑
| 维度 | 普通 Topic | LiteTopic |
|---|---|---|
| 隔离粒度 | Topic 级,group 共享 | 会话级,排他消费 |
| 创建方式 | 预先手动创建 | 自动创建,TTL 自动回收 |
| 顺序性 | 需分区顺序 Topic | 单队列天然保序 |
| 并发模型 | 多队列加消费者横向扩 | 单通道单队列,吞吐靠通道数堆 |
单 LiteTopic 只有一个队列,单通道消费 TPS 有限(官方默认 200,可以调整)。单个消费者的订阅上限为 2000(可以调整)。实例级别的建/订阅量是有配额的,并且当通道被收回之后如果没有取消订阅,则该订阅会被计入总量中——如果不调用 unsubscribeLite 就会泄露配额直到打满。服务端为 5.5.0 以上版本,客户端使用的是 gRPC 5.1.0 以上的版本
我的判断
价值不在于“AI”这个标签上,在于隔离粒度之下探到会话、状态之外置到基础设施中去——正是多 Agent 编排、SSE 流式会话、RAG 增量同步等痛点所在。普通的业务消息不要变动,原来的模型已经足够使用了;被长会话断连、重复推理、慢会话互堵所困扰的人们,可以进行实测 5.5,并重点验证配额规划和消费吞吐模型



文章评论