运维知识
悠悠
2026年9月10日

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

输出里重点看这几列:

字段含义
DiffbrokerOffset - consumerOffset,就是堆积量
brokerOffsetBroker 端当前最新的消息偏移量
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),小内存机器慎用。
  • flushDiskTypeSYNC_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 消除消费链路中的阻塞

那天排查消费逻辑的时候发现了两个"瘤子":

  1. 消费逻辑里同步调了一个外部 HTTP 接口,超时设了 5 秒,偶尔会卡住
  2. 每条消息都单独写一次数据库

改法:

// 改之前:同步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_ratioBroker 磁盘使用率> 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 分钟。核心操作其实就三板斧:加线程、加实例、优化代码。但每一步的前提是你得知道瓶颈在哪儿——盲目调参反而可能把事情搞更糟。


避坑清单

  1. 别上来就扩 Broker,先确认是不是消费端自己的问题。80% 的堆积都是消费端瓶颈。
  2. 消费者实例数不要超过队列数,超了也分不到队列,白白浪费资源。
  3. consumeThreadMinconsumeThreadMax 都要改,只改 max 不生效。
  4. 扩队列数之前评估顺序消费的影响,如果业务用了 MessageQueueSelector,扩队列会打乱分区。
  5. 异步刷盘虽然快,但要做好主从同步。用 SYNC_MASTER + ASYNC_FLUSH 是个不错的折中方案。
  6. 批量消费要处理好幂等,一批 32 条如果写到一半失败了,RocketMQ 会整批重试,别写出重复数据。
  7. 别忘了消息压缩。消息体大的话,生产端开启压缩能显著降低 Broker 的磁盘 IO 和网络开销:
// 生产端开启消息压缩(消息体>4KB自动压缩)
producer.setCompressMsgBodyOverHowmuch(4096);

一张图总结调优路径

消息堆积 → 定位瓶颈在哪一层
  │
  ├─ 消费端瓶颈
  │   ├─ 增加 consumeThread(64~128)
  │   ├─ 增加消费者实例(≤ 队列数)
  │   ├─ 优化消费代码(异步化/批量/去阻塞)
  │   └─ 消息过滤前置(Tag/SQL92)
  │
  ├─ Broker 端瓶颈
  │   ├─ 增加队列数(readQueueNums/writeQueueNums)
  │   ├─ 调整 sendMessageThreadPoolNums
  │   ├─ 异步刷盘 ASYNC_FLUSH
  │   ├─ 开启 transientStorePool
  │   └─ 横向扩展 Broker
  │
  └─ 硬件瓶颈
      ├─ 升级 SSD(CommitLog 性能瓶颈)
      ├─ 增加内存(PageCache 命中率)
      └─ 升级网卡(万兆网络)

消息堆积这事,说复杂也复杂,说简单也简单。关键是平时把监控做好,别等到堆积了 2000 万条才急匆匆地翻文档。


我是「运维躬行录」,专注分享云计算和运维的生产实践经验。如果这篇文章对你有帮助,帮忙点个赞、转发一下。 关注公众号:耕云躬行录 个人博客:躬行笔记

文章目录

博主介绍

热爱技术的云计算运维工程师,Python全栈工程师,分享开发经验与生活感悟。
欢迎关注我的微信公众号@运维躬行录,领取海量学习资料

微信二维码