RocketMQ 消息堆积了2000万条,我用这6步把消费延迟从3小时干到了5分钟
周三下午两点半,告警群炸了。
RocketMQ Dashboard 上某个 Topic 的消费延迟直接飙到了 3 小时,堆积量 2000 多万条,还在涨。业务那边反馈订单状态不更新,客服电话快被打爆了。
看了眼监控面板,生产端 TPS 是平时的 4 倍——原来是运营搞了个促销活动,没提前通知我们。消息像洪水一样往 Broker 里灌,消费端完全跟不上。
这种场景其实不少见。消息堆积的本质就一句话:生产速度 > 消费速度。但排查和调优的路径远没这么简单,因为瓶颈可能在消费者、Broker、网络、磁盘任意一个环节。
下面是我那天下午完整的排查和调优过程,每一步都是实打实在生产环境上操作的。
第一步:快速定位堆积的根因
别急着改配置,先搞清楚到底卡在哪儿。
打开 RocketMQ Dashboard(或者用命令行),看两个关键指标:
# 查看消费组的堆积情况
sh mqadmin consumerProgress -g your_consumer_group -n 127.0.0.1:9876
输出里重点看这几列:
| 字段 | 含义 |
|---|---|
| Diff | brokerOffset - consumerOffset,就是堆积量 |
| brokerOffset | Broker 端当前最新的消息偏移量 |
| consumerOffset | 消费者当前消费到的偏移量 |
如果 Diff 在持续增大,说明消费速度一直追不上生产速度。如果 Diff 很大但不再增长了,说明生产端已经降下来了,消费者在慢慢追。
还有一个命令很好用:
# 查看消费者连接信息和消费线程数
sh mqadmin consumerConnection -g your_consumer_group -n 127.0.0.1:9876
那天我一看,4 个队列只连了 2 个消费者实例,而且每个实例的消费线程才 20 个(默认值)。业务逻辑里还有一个远程 HTTP 调用,单条消息处理时间 200ms 左右。
算一下:2 个实例 × 20 线程 = 40 并发,每条 200ms,理论 TPS 也就 200。而生产端这会儿 TPS 是 800+。消费速度只有生产的 1/4,不堆积才怪。
第二步:增加消费者并发数
这是最直接、见效最快的手段。
RocketMQ 的 DefaultMQPushConsumer 有两个关键参数:
consumer.setConsumeThreadMin(20); // 默认20
consumer.setConsumeThreadMax(20); // 默认20
很多人不知道,这两个值默认都是 20,而且底层用的是 ThreadPoolExecutor,核心线程数等于最大线程数,意味着线程池其实不会动态扩缩。想加并发,得两个都改。
我当时的调整:
consumer.setConsumeThreadMin(64);
consumer.setConsumeThreadMax(64);
为什么不设更大?因为这台机器是 4 核 8G,线程设太多反而上下文切换开销大。一般经验是:
- CPU 密集型消费逻辑:线程数 = CPU 核数 × 2
- IO 密集型消费逻辑(比如有数据库操作、RPC 调用):线程数 = CPU 核数 × 4 ~ 8
我这边消费逻辑有 HTTP 调用,属于 IO 密集型,4 核机器设 64 线程不算过分。
改完重启,单实例的消费 TPS 从 100 涨到了 280 左右。
⚠️ 注意:consumeThreadMax 设得太大,可能导致下游服务(比如数据库)被打爆。改之前先确认下游能扛住。第三步:横向扩容消费者实例
光调线程数还不够,2 个实例 × 280 = 560 TPS,还是追不上 800+ 的生产速度。
RocketMQ 的消费模型是:一个队列同一时间只能被同一个消费组里的一个消费者实例消费。所以消费者实例数超过队列数就没意义了。
先看看当前有多少个队列:
sh mqadmin topicRoute -t your_topic -n 127.0.0.1:9876
输出里找 "readQueueNums": 4,说明只有 4 个读队列。那消费者实例最多有效的就是 4 个。
我当时先紧急扩了 2 个消费者实例(凑成 4 个),让每个队列都有一个专属消费者。
4 实例 × 280 = 1120 TPS,终于超过了生产端的 800。堆积量开始往下掉了。
但这只是堪堪追上。如果要更快地消化堆积,得增加队列数。
第四步:提升 Broker 端配置
如果消费者已经拉满了,还是追不上,就得从 Broker 端下手。
4.1 增加写队列数和读队列数
# 动态修改Topic的队列数(不用重启Broker)
sh mqadmin updateTopic -t your_topic -n 127.0.0.1:9876 -b 127.0.0.1:10911 -r 8 -w 8
-r 是读队列数,-w 是写队列数,一般保持一致。从 4 扩到 8,然后消费端也扩到 8 个实例。
⚠️ 改队列数时注意:如果你的业务依赖消息顺序性(比如用了 MessageQueueSelector 按订单ID分区),扩队列会打乱原有的分区映射,可能导致乱序。这种情况要评估好再操作。
4.2 调整 Broker 关键配置
编辑 broker.conf:
# 发送消息的线程池大小,默认是1,太小了
sendMessageThreadPoolNums=16
# CommitLog 刷盘策略:ASYNC_FLUSH(异步)比 SYNC_FLUSH(同步)快很多
flushDiskType=ASYNC_FLUSH
# 开启 transientStorePool,利用堆外内存做写缓冲
transientStorePoolEnable=true
# 拉取消息的线程数
pullMessageThreadPoolNums=32
# 单个 Broker 上一个 Topic 的队列数
defaultTopicQueueNums=8
几个坑:
sendMessageThreadPoolNums默认值是 1,对你没看错就是 1。很多人部署完就没改过这个参数,生产上 Broker 一到高峰就报[REJECTREQUEST]system busy,大概率跟这个有关。transientStorePoolEnable=true会额外占用堆外内存(默认每个 buffer 1G),小内存机器慎用。flushDiskType从SYNC_FLUSH改成ASYNC_FLUSH能显著提升吞吐,但断电可能丢少量消息。大部分业务场景都能接受异步刷盘。
改完配置后需要重启 Broker:
# 先优雅关闭
sh mqshutdown broker
# 再启动
nohup sh mqbroker -c /opt/rocketmq/conf/broker.conf -n 127.0.0.1:9876 &
4.3 横向扩展 Broker
如果单台 Broker 的磁盘 IO 或者网络带宽已经到瓶颈了(iostat 看 %util 长期 90%+),加机器是最直接的办法。
RocketMQ 天然支持多 Broker 部署。新加一台 Broker 注册到同一个 NameServer,然后把 Topic 的队列分散到新 Broker 上就行:
# 在新Broker上创建Topic队列
sh mqadmin updateTopic -t your_topic -n 127.0.0.1:9876 \
-b new_broker_ip:10911 -r 4 -w 4
这样消息会自动负载均衡到多个 Broker。
第五步:优化消费者代码
硬件和配置都拉满了,消费速度还是上不去?那八成是消费逻辑本身的问题。
5.1 消除消费链路中的阻塞
那天排查消费逻辑的时候发现了两个"瘤子":
- 消费逻辑里同步调了一个外部 HTTP 接口,超时设了 5 秒,偶尔会卡住
- 每条消息都单独写一次数据库
改法:
// 改之前:同步HTTP调用
String result = httpClient.get("http://xxx/api/check", 5000); // 可能卡5秒
// 改之后:异步化,不阻塞消费主线程
CompletableFuture.supplyAsync(() -> {
return httpClient.get("http://xxx/api/check", 2000);
}).thenAccept(result -> {
// 异步处理结果
processResult(result);
});
5.2 批量消费
RocketMQ 支持批量消费模式,一次拉多条消息一起处理:
// 设置每次拉取的消息数量
consumer.setConsumeMessageBatchMaxSize(32); // 默认是1
配合批量写库:
public ConsumeConcurrentlyStatus consumeMessage(
List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
// 把多条消息攒一起,批量写库
List<OrderDTO> orders = msgs.stream()
.map(msg -> JSON.parseObject(msg.getBody(), OrderDTO.class))
.collect(Collectors.toList());
orderMapper.batchInsert(orders); // 一次IO搞定
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
批量消费 + 批量写库,单次 IO 处理 32 条消息,TPS 直接翻了十几倍。
5.3 消息过滤前置
如果你的 Topic 里混了多种消息类型,消费者只关心其中一部分,那就用 Tag 过滤或者 SQL92 过滤,在 Broker 端就把不需要的消息过滤掉,减少网络传输和消费端的无效处理:
// 只订阅 tag 为 "orderPaid" 的消息
consumer.subscribe("OrderTopic", "orderPaid");
// 或者用SQL92表达式做更复杂的过滤
consumer.subscribe("OrderTopic",
MessageSelector.bySql("amount > 100 AND region = 'east'"));
注意:使用 SQL92 过滤需要在 broker.conf 里开启 enablePropertyFilter=true。第六步:建立监控和告警体系
问题解了不代表以后不会再来。消息堆积这种事,靠告警比靠人盯靠谱得多。
6.1 用 rocketmq-exporter + Prometheus + Grafana
社区有现成的 Prometheus Exporter:
# 下载并启动 rocketmq-exporter
java -jar rocketmq-exporter.jar \
--rocketmq.config.namesrvAddr=127.0.0.1:9876
然后在 Prometheus 里配置采集,关键指标:
# prometheus.yml
scrape_configs:
- job_name: 'rocketmq'
static_configs:
- targets: ['localhost:5557']
Grafana 里重点看这几个指标:
| 指标 | 含义 | 告警阈值建议 |
|---|---|---|
rocketmq_consumer_diff | 消费堆积量 | > 10000 告警 |
rocketmq_consumer_tps | 消费 TPS | 持续 < 生产 TPS 的 50% 告警 |
rocketmq_broker_disk_ratio | Broker 磁盘使用率 | > 85% 告警 |
rocketmq_send_tps | 生产 TPS | 作为基线对比 |
6.2 配置告警规则
Prometheus AlertManager 规则示例:
groups:
- name: rocketmq_alerts
rules:
- alert: MessageBacklogHigh
expr: rocketmq_consumer_diff > 10000
for: 5m
labels:
severity: warning
annotations:
summary: "RocketMQ消息堆积告警"
description: "消费组 {{ $labels.group }} 堆积量 {{ $value }},持续5分钟"
- alert: MessageBacklogCritical
expr: rocketmq_consumer_diff > 100000
for: 2m
labels:
severity: critical
annotations:
summary: "RocketMQ消息堆积严重告警"
description: "消费组 {{ $labels.group }} 堆积量 {{ $value }},需要立即处理"
6.3 定时巡检脚本
写个简单的巡检脚本,每天跑一次,输出各消费组的堆积状态:
#!/bin/bash
# rocketmq_check.sh - 消费堆积巡检
NAMESRV="127.0.0.1:9876"
GROUPS=$(sh mqadmin consumerProgress -n $NAMESRV 2>/dev/null | grep -v "^#" | awk '{print $1}' | sort -u)
echo "========== RocketMQ 消费堆积巡检 $(date '+%Y-%m-%d %H:%M:%S') =========="
echo ""
for group in $GROUPS; do
total_diff=$(sh mqadmin consumerProgress -g $group -n $NAMESRV 2>/dev/null \
| grep -v "^#" | awk '{sum+=$NF} END {print sum}')
if [ "$total_diff" -gt 10000 ] 2>/dev/null; then
echo "[WARNING] $group 堆积: $total_diff"
elif [ "$total_diff" -gt 0 ] 2>/dev/null; then
echo "[OK] $group 堆积: $total_diff"
fi
done
复盘:那天下午的时间线
| 时间 | 操作 | 效果 |
|---|---|---|
| 14:30 | 收到告警,堆积 2000w | - |
| 14:35 | 定位根因:消费者并发不足 | 确认方向 |
| 14:40 | 调整 consumeThread 从 20 → 64 | 单实例 TPS 100 → 280 |
| 14:50 | 扩消费者实例 2 → 4 | 总 TPS 560 → 1120 |
| 15:00 | 扩队列数 4 → 8,消费者扩到 8 | 总 TPS → 2200 |
| 15:10 | 优化消费逻辑,批量写库 | 总 TPS → 4000+ |
| 15:30 | 堆积从 2000w 降到 500w | 持续下降中 |
| 16:15 | 堆积清零,消费延迟 < 5min | 恢复正常 |
| 17:00 | 部署监控告警 | 防患于未然 |
从告警到恢复,大概 1 小时 45 分钟。核心操作其实就三板斧:加线程、加实例、优化代码。但每一步的前提是你得知道瓶颈在哪儿——盲目调参反而可能把事情搞更糟。
避坑清单
- 别上来就扩 Broker,先确认是不是消费端自己的问题。80% 的堆积都是消费端瓶颈。
- 消费者实例数不要超过队列数,超了也分不到队列,白白浪费资源。
consumeThreadMin和consumeThreadMax都要改,只改 max 不生效。- 扩队列数之前评估顺序消费的影响,如果业务用了 MessageQueueSelector,扩队列会打乱分区。
- 异步刷盘虽然快,但要做好主从同步。用
SYNC_MASTER+ASYNC_FLUSH是个不错的折中方案。 - 批量消费要处理好幂等,一批 32 条如果写到一半失败了,RocketMQ 会整批重试,别写出重复数据。
- 别忘了消息压缩。消息体大的话,生产端开启压缩能显著降低 Broker 的磁盘 IO 和网络开销:
// 生产端开启消息压缩(消息体>4KB自动压缩)
producer.setCompressMsgBodyOverHowmuch(4096);
一张图总结调优路径
消息堆积 → 定位瓶颈在哪一层
│
├─ 消费端瓶颈
│ ├─ 增加 consumeThread(64~128)
│ ├─ 增加消费者实例(≤ 队列数)
│ ├─ 优化消费代码(异步化/批量/去阻塞)
│ └─ 消息过滤前置(Tag/SQL92)
│
├─ Broker 端瓶颈
│ ├─ 增加队列数(readQueueNums/writeQueueNums)
│ ├─ 调整 sendMessageThreadPoolNums
│ ├─ 异步刷盘 ASYNC_FLUSH
│ ├─ 开启 transientStorePool
│ └─ 横向扩展 Broker
│
└─ 硬件瓶颈
├─ 升级 SSD(CommitLog 性能瓶颈)
├─ 增加内存(PageCache 命中率)
└─ 升级网卡(万兆网络)
消息堆积这事,说复杂也复杂,说简单也简单。关键是平时把监控做好,别等到堆积了 2000 万条才急匆匆地翻文档。
我是「运维躬行录」,专注分享云计算和运维的生产实践经验。如果这篇文章对你有帮助,帮忙点个赞、转发一下。 关注公众号:耕云躬行录 个人博客:躬行笔记