Skip to content
Go back

RocketMQ 消息积压处理:监控、流控、扩容与 DLQ 的全链路方案

RocketMQ 消息积压处理:从发现到消除的完整流程

一句话结论(30s)

消息积压本质不是 RocketMQ 的 Bug,而是「生产速度 > 消费速度」的容量信号,因为积压持续增长最终会打满磁盘、拖垮服务。处理的关键设计是「先监控发现、再流控止血、后扩容加速、最后 DLQ 兜底」的全链路分层,而最核心的判断是区分一过性高峰还是持续性消费能力不足。核心权衡:扩容 Consumer 前必须先确认 Queue 数量是否够,因为 Consumer 实例数超过 Queue 数时多出来的实例只会空转。

核心原理(2min)

主流程分四步:监控阶段用 Diff/Lag、Delay、Consume/Produce TPS 定位积压,关键是看 Diff 的增长趋势而非绝对值;止血阶段先对生产者限流防止磁盘打满,再按需扩 Broker 节点、扩 Topic 队列;加速阶段靠增加 Consumer 实例(上限 = Queue 数)、调大消费线程池与批量拉取、优化消费逻辑(异步化 / 批量入库)与消息批次处理;兜底阶段靠 DLQ——毒消息重试 16 次失败后移入 %DLQ%Group,防止单条卡死全队并保留现场。关键机制:顺序消息积压最棘手,因为不能加线程也不能跳过,只能设超时 + SUSPEND,或手动把毒消息发 DLQ 后返回 SUCCESS。预防靠监控 + 限流 + 压测 + 基于 Lag 的 K8s HPA 自动扩缩容。

底层深入(5-10min)

积压不是 Bug,是信号

消息积压本身不是 RocketMQ 出了问题——它是系统处理能力不足的信号。生产速度 > 消费速度 → 积压持续增长 → 最终磁盘打满、服务不可用。

处理积压的关键是快速发现、精准定位、分层处理

🤔 思考穿插:为什么说积压是「信号」而不是「Bug」?—— 因为积压的根本原因是「生产速度 > 消费速度」,是容量问题而非代码缺陷;把积压当 Bug 会去瞎改代码,把它当信号才能聚焦到「扩消费能力」或「限生产速度」上。所以处理积压的第一步是正确归因。

一、监控:先看见,才能处理

核心指标

指标含义告警阈值建议
Diff / Lag未消费的消息数量(Producer offset - Consumer offset)> 10 万条(超过消费速率对应的时间量)
Delay最旧一条未消费消息的延迟时间> 1 分钟
Consume TPS消费速率(每分钟/每秒钟消费条数)低于基准线的 50%
Produce TPS生产速率用于对比消费速率

RocketMQ Console 监控

Topic 详情页:
┌──────────────────────────────────────────────────┐
│ 生产情况                    消费情况              │
│ Produce TPS: 5000 msg/s    Consume TPS: 1000/s   │ ← 积压速率 = 4000/s
│ Max Offset: 1,200,000      Consumer Offset: 800,000 │
│                            Diff: 400,000          │ ← 积压 40 万条
│                            Delay: 400s            │ ← 延迟 400 秒
└──────────────────────────────────────────────────┘

关键判断:Diff 的增长趋势比绝对值更重要。 如果 Diff 稳定在 10 万,说明生产 = 消费,只是之前遗留的积压;如果 Diff 以每秒 100 的速度增长,说明消费严重跟不上。

🤔 思考穿插:为什么 Diff 的「增长趋势」比「绝对值」更重要?—— 因为绝对值大但稳定,说明只是历史遗留、生产 = 消费,等它慢慢消化即可;只有 Diff 持续增长,才说明消费能力真的跟不上、需要立即干预。所以看趋势能区分「一过性高峰」和「持续性能力不足」。

自定义监控

// 通过 RocketMQ Admin API 获取消费进度
DefaultMQAdminExt admin = new DefaultMQAdminExt();
admin.start();

// 消费者 offset
long consumerOffset = admin.queryConsumerOffset(
    "DefaultGroup", new MessageQueue("TopicTest", "broker-a", 0));

// 队列最大 offset
long maxOffset = admin.maxOffset(new MessageQueue("TopicTest", "broker-a", 0));

long lag = maxOffset - consumerOffset;

// 上报到 Prometheus / 公司监控平台
gauge.labels("TopicTest", "DefaultGroup").set(lag);

二、应急止损:先止血

止损步骤

Step 1:生产者流控(最快生效)

如果积压已经导致 Broker 磁盘趋近打满,必须立刻限流:

// 生产者端限流——RocketMQ 提供的 Hook
producer.setSendMsgTimeout(5000);  // 超时 5 秒直接失败,不堆积在客户端

// 更暴力的方案:业务层拦截
if (cacheManager.get("mq:overload").equals("true")) {
    log.warn("MQ overload, message dropped: {}", msg);
    return;  // 丢弃或降级
}

Step 2:扩容 Broker(如果需要)

如果瓶颈在 Broker(磁盘 IO、CPU、网络),横向增加 Broker 节点:

# 新增 Broker 节点,加入集群
# Broker 的 Master-Slave 对可以动态加入
sh mqbroker -c broker-new.properties -n namesrv:9876

Step 3:Topic 队列扩容

增加 Queue 数量,允许更多消费者线程并行处理:

mqadmin updateTopic -n namesrv:9876 -t OrderTopic -c DefaultCluster -r 16 -w 16
# -r: readQueueNums, -w: writeQueueNums(通常保持一致)

三、消费端加速:核心手段

3.1 增加消费者实例数量

当前:3 个 Consumer,每个消费 4 个 Queue → 12 个线程
扩容:6 个 Consumer,每个消费 2 个 Queue → 12 个线程?❌ 不对

Queue 数量 = 16
3 Consumer → 分配不均(16/3,有 Consumer 拿 6 个 Queue 有拿 5 个)
4 Consumer → 刚好 4 个线程/Consumer → 并行度 = 16

规则:Consumer 实例数 <= Queue 数。 多出来的 Consumer 拿不到 Queue,空转。所以扩容 Consumer 前,先确认 Queue 数是否足够。

🤔 思考穿插:为什么 Consumer 实例数不能超过 Queue 数,多了反而空转?—— 因为 Queue 是最小负载均衡单位,Rebalance 时每个 Queue 最多分配给一个实例;实例数超过 Queue 数时,多出来的实例分不到任何 Queue、只能干等。所以扩容前必须先确认 Queue 数量是否足够。

3.2 增加单个 Consumer 的并行度

consumer.setConsumeThreadMin(20);  // 最小消费线程
consumer.setConsumeThreadMax(64);  // 最大消费线程
consumer.setConsumeMessageBatchMaxSize(32);  // 一次拉取 32 条,减少网络 RTT

注意:ConsumeThreadMax全局的线程池大小,不是每个 Queue 的。每个 Queue 可以分配到线程池中的任意线程。

3.3 优化消费逻辑

很多时候积压不是 MQ 的问题,是消费逻辑慢:

// 慢的原因排查
// 1. 是否有同步 IO?(数据库查询、HTTP 调用)
//    → 改成异步 / 批量(batch insert、pipeline)
// 2. 是否有锁竞争?
//    → 无锁化 / 分段锁
// 3. 是否有大对象创建?
//    → 对象池 / 复用
// 4. 是否每条消息都打日志?
//    → 采样日志或去掉不必要的日志

3.4 消息批次处理

单条 DB 插入 → 批量 DB 插入:

consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
    List<Order> orders = new ArrayList<>(msgs.size());
    for (MessageExt msg : msgs) {
        orders.add(parseOrder(msg.getBody()));
    }
    // 批量入库,1 次 IO 完成 N 条消息
    orderDao.batchInsert(orders);
    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});

四、死信队列(DLQ):为失败的消息兜底

什么是 DLQ?

%DLQ% + ConsumerGroup——当一条消息重试消费 maxReconsumeTimes(默认 16 次)仍然失败时,RocketMQ 不会无限重试,而是将它移入死信队列。

正常消费流程:
消息 → 消费失败 → 1s 后重试 → 10s 后重试 → ... → 16 次后 → 进入 DLQ

进入 DLQ = 人工介入的信号

DLQ 的作用

  1. 防止积压雪崩:一条”毒消息”(如格式错误导致永远解析失败)会让该 Queue 的消费卡住。DLQ 把毒消息移走,让后续正常消息能继续消费。
  2. 保留现场:进入 DLQ 的消息不会丢失,人工可以查询、修复、重新发送。

DLQ 管理

# 查看某个 ConsumerGroup 的 DLQ
mqadmin queryMsgByKey -n namesrv:9876 -t %DLQ%OrderConsumerGroup -k ""

# 重发 DLQ 中的消息(修复后)
mqadmin sendMessage -n namesrv:9876 -t OrderTopic -p "fixed body"

# 删除 DLQ 中的消息(确认不需要的)
mqadmin deleteExpiredCommitLog -n namesrv:9876

五、特殊情况:顺序消息积压

顺序消息(MessageListenerOrderly)的积压处理更加棘手,因为:

  1. 不能加线程:一个 Queue 同一时刻只能一个线程消费。
  2. 不能跳过:一条失败必须等它重试成功,后面的消息全部阻塞。

🤔 思考穿插:为什么顺序消息积压最棘手,既不能加线程也不能跳过?—— 因为顺序消费的前提是「同一 Queue 单线程串行」,加线程会破坏顺序、跳过会让后面的消息失去它依赖的前序状态;所以只能靠「超时挂起 + 把毒消息手动发 DLQ」来缓解。顺序和吞吐,本就是天平的两端。

处理策略

consumer.registerMessageListener(new MessageListenerOrderly() {
    @Override
    public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs,
                                                ConsumeOrderlyContext context) {
        for (MessageExt msg : msgs) {
            try {
                // 每条消息设置超时
                Future<?> future = executor.submit(() -> process(msg));
                future.get(5000, TimeUnit.MILLISECONDS);  // 5 秒超时
            } catch (TimeoutException e) {
                // 超时 → 挂起,稍后重试(不阻塞整个队列)
                context.setSuspendCurrentQueueTimeMillis(3000);
                return ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT;
            }
        }
        return ConsumeOrderlyStatus.SUCCESS;
    }
});

如果序列化失败(毒消息),直接发往 DLQ 并返回 SUCCESS(手动跳过):

try {
    Order order = JSON.parseObject(msg.getBody(), Order.class);
    process(order);
} catch (Exception e) {
    // 无法解析 → 发往 DLQ 并跳过
    log.error("Poison pill detected, skip: {}", msg.getMsgId());
    return ConsumeOrderlyStatus.SUCCESS;  // 不卡住后面的消息
}

六、预防体系

层级措施工具
监控Diff/Lag/Delay 实时大盘,趋势告警Prometheus + Grafana
限流生产者端限流,触发阈值后拒绝新消息Sentinel / 自研限流器
压测定期全链路压测,确认消费极限 QPSJMeter / 内部压测平台
弹性Consumer 自动扩缩容(K8s HPA 基于 Lag)K8s HPA + 自定义 Metrics
降级非核心消息可丢弃/延迟处理业务标识 + 分级 Topic

基于 Lag 的自动扩容

# Kubernetes HPA 配置(需要在 Prometheus Adapter 中暴露 Lag 指标)
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
spec:
  metrics:
  - type: Pods
    pods:
      metric:
        name: rocketmq_consumer_lag
      target:
        type: AverageValue
        averageValue: "10000"  # Lag > 1 万 → 扩容

总结

消息积压是系统性问题,不是 RocketMQ 自身的问题。处理流程:

1. 发现 → 监控告警(Lag > 阈值)
2. 止血 → 生产者限流,防止磁盘打满
3. 扩容 → 增加 Broker / Queue / Consumer 实例
4. 加速 → 消费逻辑优化、批量处理、异步化
5. 兜底 → DLQ 接收毒消息,防止单条卡死全队
6. 预防 → 监控 + 限流 + 压测 + 弹性伸缩

最重要的判断:积压是一过性(瞬时高峰)还是持续性(消费能力不足)?一过性等它消掉就行;持续性必须扩容或优化消费逻辑。

章末提问

追问 1:发现消息积压后,你的处理顺序是什么?

结论先行:先监控定位(看 Diff 趋势),再止血(生产者限流防磁盘打满),再扩容加速(Broker/Queue/Consumer),最后 DLQ 兜底毒消息。

因为:积压一旦打满磁盘就全站不可用,所以要先用限流止血;判断是「一过性高峰」还是「持续性能力不足」决定要不要扩容;毒消息会卡死单队,要靠 DLQ 隔离并保留现场。四步是「发现 → 止血 → 加速 → 兜底」的递进关系。

追问 2:为什么扩容 Consumer 前要先确认 Queue 数量?

结论先行:因为 Queue 是最小负载均衡单位,Consumer 实例数超过 Queue 数时,多出来的实例分不到 Queue、只会空转。

因为:Rebalance 时每个 Queue 最多分配给一个实例,实例数 > Queue 数意味着必然有实例拿不到队列;所以必须先扩 Queue(或确认 Queue 足够),再扩 Consumer 实例,扩容才有意义。

追问 3:顺序消息积压为什么比并发消息更难处理?

结论先行:因为顺序消费要保证单 Queue 串行,既不能加线程也不能跳过失败消息,一条卡住会阻塞全队。

因为:并发消费可以多线程并行、失败可跳过或快速重试;顺序消费里加线程会乱序、跳过会丢失前序状态,恢复手段只剩「超时挂起」或「把毒消息手动发 DLQ 再返回 SUCCESS」。可用的优化空间更小,所以更难。


Share this post on:

Previous Post
RocketMQ 消息过滤:Tag 过滤、SQL92 表达式与过滤原理
Next Post
RocketMQ 消息幂等性:为什么 MessageId 不可靠,以及如何真正去重