RocketMQ 已正式接入 AI,可以硬扛百万级 AI 会话!

来源:AI开源无界
作者:小知
原文:https://mp.weixin.qq.com/s/A30t-8v3_RhH_JqjJE4XTA

最近手里一个项目要用消息队列,正好赶上它接入 AI 这件事,然后我把 RocketMQ 重新捡起来研究了一下。

五月底发布的 RocketMQ 5.5.0,把专门为 AI 场景设计的 LiteTopic 消息模型放进了开源版本,也就是社区提案 RIP-83 里定义的那套东西。我之前用 RocketMQ,跑的都是订单、日志这一类传统场景,这次看到官方针对 AI 负载专门做了消息模型,还顺手把 MCP、A2A 这些智能体协议的支持也做了,就花时间把官方文档和代码都翻了一遍。

先说清楚一个容易误会的地方,RocketMQ 接入 AI,它自己并没有变成大模型,干的还是消息队列的老本行。变化主要在这套新的主题模型上,它是照着 AI 应用的通信特点做出来的。另外 LangChain、CrewAI、AutoGen、Dify 这些主流 Agent 框架也都能对接。

这篇文章就说两件事:它是怎么接入 AI 的,以及对比以前的 RocketMQ,到底多解决了哪些问题。

它是怎么接入 AI 的

核心是一套叫 LiteTopic 的轻量主题模型。

以前的 RocketMQ 里,Topic 是比较重的资源,要提前创建,数量也有限,一群消费者围着 Topic 抢消息。这套模型服务订单、支付、日志这些业务很多年,很成熟,但直接搬到 AI 场景就不太合适。AI 应用的会话量太大了,一个用户一个会话,一个 Agent 一个任务,动辄几十万上百万条独立通道,传统 Topic 扛不住这个数量级。

LiteTopic 的做法是把结构分成两层,父 Topic 当命名空间用,同一业务的通道都挂在它下面,挂着的这些 LiteTopic 才是真正干活的通道。每个会话、每个任务都可以单独占一条,第一次发消息或者订阅的时候自动创建,不用提前建,用完了按 TTL 自动回收,也不用人盯着清理。

RocketMQ 已正式接入 AI,可以硬扛百万级 AI 会话!
LiteTopic 的主题结构图

对比以前,解决了什么问题

这部分我对比着说,看得清楚一些。

以前:Topic 要提前建,数量有限,百万级会话没法玩;现在:一个会话一条通道,按需创建。

LiteTopic 第一次发消息或者订阅时就自动生效,单个集群能承载百万级通道共存。能扛住这个量,主要是因为底层索引从传统的 ConsumeQueue 文件换成了 RocksDB,消息本身还是走 CommitLog 顺序追加,写入链路没动。

以前:消费者重启,进度和状态要业务自己管;现在:断点存在 Broker,重启自动续传。

做 AI 对话产品的人都怕一件事,用户的 WebSocket 一断,重连以后上下文没了,已经花掉的算力也白费。

阿里安全团队做“安全小蜜”智能助手的时候就遇到过这个问题,高并发会话下上下文丢失、任务中断,后来靠 LiteTopic 重构了会话保持机制才解决。

做法是把应用节点做成无状态的,每个会话映射成一条 LiteTopic,交互历史按顺序存在通道里,消费位点以“内存快照加增量持久化”的方式存在 Broker 端。用户重连到任何一台机器,重新订阅同一条通道就能从断点接着收,后台的推理任务不受影响。

RocketMQ 已正式接入 AI,可以硬扛百万级 AI 会话!
会话保持方案,状态托管在 LiteTopic 里

以前:消费靠长轮询,通道多了空转烧 CPU;现在:事件驱动,有消息才唤醒。

以前的做法是 Broker 要遍历主题来检查有没有新消息,通道少的时候无所谓,规模到了百万级,这么扫就不现实了。

LiteTopic 改成了事件驱动,Broker 内部维护一个就绪集合(Ready Set),有消息写入或者可读事件触发时才唤醒对应的消费者,平时不做这种全量扫描。

以前:限流一刀切,一个慢用户堵死所有人;现在:限流精确到单个会话。

GPU 资源紧张的时候,AI 平台经常要做差异化限流,付费用户优先,普通用户排队。传统限流的毛病是一个慢用户能把整个消费线程堵死。

LiteTopic 提供了一个叫 Consume Suspend 的能力,某个用户超限了,只挂起他那一条通道,释放线程去处理别人,到时间自动恢复,不算失败也不进死信。再配合消息优先级,高峰期让高价值任务先跑。

阿里云百炼网关目前在生产环境就是用这套机制做 AI 推理的流量治理。

RocketMQ 已正式接入 AI,可以硬扛百万级 AI 会话!
Consume Suspend 的流控流程

这几条放到 Multi-Agent 系统里,好处会更明显,Supervisor 把任务拆开往各个子 Agent 的通道一丢就返回,子 Agent 各自消费处理,干完活把结果写到结果通道里,Supervisor 订阅着等汇总,全程非阻塞。

A2A 协议本来就推荐这种异步通信方式,RocketMQ 这套模型刚好能把它落地。

RocketMQ 已正式接入 AI,可以硬扛百万级 AI 会话!
基于 RocketMQ 构建异步 Multi-Agent 系统

我自己跑了一遍

光看文档不过瘾,我在本地把 5.5.0 装起来跑了一遍,过程记在这里,想试的话可以参考。

先下载解压,然后往 broker.conf 里追加三行配置:

wget https://mirrors.aliyun.com/apache/rocketmq/5.5.0/rocketmq-all-5.5.0-bin-release.zip
unzip rocketmq-all-5.5.0-bin-release.zip && cd rocketmq-all-5.5.0-bin-release

cat >> conf/broker.conf << EOF
enableLmq=true
enableMultiDispatch=true
storeType=defaultRocksDB
EOF

接着把 NameServer、Broker、Proxy 依次拉起来,再创建父 Topic 和一个绑定它的消费组:

nohup sh bin/mqnamesrv &
nohup sh bin/mqbroker -n localhost:9876 -c conf/broker.conf &
nohup sh bin/mqproxy -n localhost:9876 &

# 创建父 Topic,消息类型声明为 LITE
sh bin/mqadmin updateTopic -b localhost:10911 -t AGENT_TASK_NS -a +message.type=LITE

# 创建消费组并绑定父 Topic
sh bin/mqadmin updateSubGroup -b localhost:10911 -g executor-group -o true --attributes +lite.bind.topic=AGENT_TASK_NS

发消息这一侧,和普通 Producer 的写法差不多,只是多了一句 setLiteTopic() 指定目标通道。通道不存在会自动创建,不用提前建:

// 父 Topic 作为命名空间
static
 final String PARENT_TOPIC = "AGENT_TASK_NS";

Producer
 producer = provider.newProducerBuilder()
        .setClientConfiguration(clientConfig)
        .setTopics(PARENT_TOPIC)
        .build();

// 给 agent_001 派发任务,TASK_agent_001 这条通道会自动创建
Message
 task = provider.newMessageBuilder()
        .setTopic(PARENT_TOPIC)
        .setLiteTopic("TASK_" + executorAgentId)
        .setBody(taskPayload.getBytes(StandardCharsets.UTF_8))
        .build();

producer.send(task);

消费这一侧用新的 LitePushConsumer,绑定父 Topic 以后,用 subscribeLite() 动态订阅自己的专属通道:

LitePushConsumer consumer = provider.newLitePushConsumerBuilder()
        .setClientConfiguration(clientConfig)
        .setConsumerGroup("executor-group")
        .bindTopic(PARENT_TOPIC)
        .setMessageListener(msg -> {
            String
 task = StandardCharsets.UTF_8.decode(msg.getBody()).toString();
            // 这里接自己的 AI 推理逻辑,耗时几秒也没关系,不阻塞别人
            String
 result = callLlm(task);   // 示意,换成自己的实现
            sendResult(task, result);        // 示意,把结果写回结果通道
            return
 ConsumeResult.SUCCESS;
        })
        .build();

// 订阅专属通道,首次订阅自动创建
consumer.subscribeLite("TASK_" + executorAgentId);

Worker 中途挂了重启也不用慌,Broker 那边存着消费位点,重新订阅以后会自动从断点继续投递,这段容错逻辑不用自己实现。

几个要注意的地方

LiteTopic 解决的是通信和状态管理的问题,它管不了模型本身的效果,prompt 设计、上下文窗口这些还是得自己处理。如果你的系统就是普通的订单、日志异步解耦,传统 Topic 完全够用,没必要为了新而新。

开源版和商业版的差异也要留意,核心的 LiteTopic 模型、断点续传、事件驱动、Suspend 这些开源版都有,但 Serverless 弹性伸缩、EventBridge 生态集成这些增强能力目前只在阿里云的商业版里,要上生产的团队得先评估一下。

还有就是这套东西出来时间不长,5.5.0 之后社区又发了 5.5.1,修了一批 Lite Mode 相关的问题,看得出来还在打磨期,尝鲜没问题,大规模用之前建议先压测。

最后

整体看下来,RocketMQ 接入 AI 这件事做得比较实在。AI 应用那些麻烦事,任务耗时长、会话量大、状态不好管,本来就是消息队列擅长处理的问题,LiteTopic 就是冲着这些问题去做的。

项目里如果正好有 Multi-Agent 协作或者长会话的需求,可以装一个试试,半天就能跑通,另外 Java 人上手这一套,几乎零门槛。

开源地址:https://github.com/apache/rocketmq

版权声明:本文内容转自互联网,本文观点仅代表作者本人。本站仅提供信息存储空间服务,所有权归原作者所有。如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至1393616908@qq.com 举报,一经查实,本站将立刻删除。

(0)

相关推荐