一、线上事故从一顿午饭说起

那天中午我正在吃饭,突然手机连环轰炸,群里一条接一条告警。我打开监控一看,RocketMQ 消费组的五个节点在疯狂刷日志,一个消息被打印了好几遍,数据库里出现大量重复数据。我刚开始以为是正常的重复投递,但仔细看消费位点发现不对劲,消费位点像过山车一样一会儿前进一会儿后退,分区分派图更奇葩,同一台机器上的队列在不停地换主人。

这种现象在专业术语里面叫分区分派震荡,通俗点说就是消费组里的成员一直在变动,每次变动都会引发 Rebalance,而每次 Rebalance 都会让一部分正在处理的消息被重新分派,这些消息被新接手的消费者当作新消息再处理一遍,于是重复消费就出现了。

那为什么消费组的成员会不停地变动呢?查了一圈日志发现心跳超时了。RocketMQ 的客户端和 Broker 之间靠心跳维持关系,一旦 Broker 认为某个客户端挂了,就会把这个客户端名下的消费队列分给其他人。而这个"认为挂了"的判断标准,就是你心跳超时了,没在指定时间内汇报。

问题来了,为什么心跳会超时?答案往往藏在消费耗时里。如果你的消费者线程在一条消息上卡了太长时间,它就没有时间去发送心跳信号,就像人迷路的时候忘了跟家里人报平安一样。家里的长辈等不到消息就以为你出事了,匆匆把你的房间分给别人住,结果你回来发现自己东西都没了,又要重新折腾一遍。

下面我要把整条链路拆开用大白话讲清楚,然后手把手告诉你怎么设置心跳间隔和消费超时,让这个后台的"物流分拣系统"不再给你惹麻烦。

二、先搞清楚 Rebalance 到底是个啥

2.1 分区分派的逻辑类似于分房子

假设你有一个消费组,里面住了三个消费者实例,相当于三户人家。Topic 里面有六个消息队列,相当于六间房。职责均摊的原则下,每户人家分两间房。这六间房就是六个队列,每间房里面的消息就是快递包裹。

RocketMQ 的 Rebalance 就是在做一件事:每隔一段时间,算一算现在有几户人家、有几间房,然后重新分配,确保尽量公平。这里分配的规则就是你指定的负载均衡策略,比如平均分布。

问题在于,房间是不能同时被两家住的。当一个消费者从一间房挪走,另一个消费者接手之前,如果原来的消费者还在处理房间里的东西,就会出岔子。就像你以为邻居搬走了,把本来放在他房间的快递拿回了自己家,结果邻居还没走,还在继续拆他的快递,这不就重复了吗?

2.2 心跳超时触发的连锁反应

RocketMQ 客户端和 Broker 之间的心跳机制特别简单,客户端每隔几秒发一个"我还活着"的信号,Broker 记录最近一次收到信号的时间。如果超过某个阈值没有收到新信号,Broker 就启动剔除程序,把那个客户端标记为不可用。

从代码层面看,这个阈值的长短取决于一个叫 heartbeatTimeout 的参数,而 heartbeatInterval 则是客户端发送心跳的间隔。假如你设置的心跳间隔是 5 秒,Broker 端的超时容忍度却是 10 秒,那么你只要哪次处理消息的时间超过了 5 秒,紧接着的下一次心跳就会迟到。迟到一次问题不大,怕的是持久性迟到,一次性处理消息就要 15 秒以上,那你的心跳信号就只能断断续续,最终被判死刑。

2.3 震荡是怎么形成的

震荡这个词很形象,就像荡秋千一样来回摆。具体的过程是这样的:

消费者 A 正在处理一条耗时 20 秒的业务消息,这 20 秒里 A 的心跳线程被业务线程拖累了,没有及时发出心跳。Broker 到了超时点等不到 A 的心跳,就把 A 标记为离线,触发一次 Rebalance。Rebalance 的结果是,原本属于 A 的队列被分给了 B 和 C。

A 处理完消息之后再尝试发心跳,Broker 发现 A 又能联系上了,又把它加回来,再次触发 Rebalance。此时队列来回倒腾,A 拿到了一些之前是 B 处理的队列。这些队列里面可能存在 B 已经消费但还没提交位点的消息,A 会从旧位点开始消费,于是重复。

重复之后 A 又被"关禁闭",再次超时,再次触发 Rebalance,循环不止,直到你想砸电脑。

三、核心参数到底怎么设置才合理

3.1 心跳间隔不能太密也不能太疏

很多新手怕心跳超时,就把心跳间隔调得非常短,比如 1 秒甚至 500 毫秒。这样看似保险,实则会增加 Broker 端处理心跳请求的压力,同时也会让你的客户端不停地发送网络请求,得不偿失。

合理的做法是设置成 5 秒到 10 秒,根据 RocketMQ 的默认值来说,5 秒是比较稳妥的。需要注意的是,这个参数在客户端生效,底层对应的是 heartbeatInterval。在启动消费者之前,你可以手动设置。

如果你用的是 RocketMQ 的 Spring Boot Starter,配置起来也很方便,在 application.properties 或者 application.yml 里指定就行。

3.2 消费超时是保护机制

消费超时指的是消息处理的最大允许时间,超过这个时间 RocketMQ 会认为当前消费失败,然后把消息重新投递。这个参数在不同版本的 RocketMQ 里有不同的叫法,原来的客户端里面叫 consumeTimeout,默认是 15 分钟。

注意,这里 15 分钟是一个非常宽裕的值,一般业务根本达不到。但恰恰因为够宽裕,很多人在消费逻辑里面写了长时间阻塞的代码,比如远程调用没有设置超时时间、数据库连接池被占满、死循环没有退出条件,等等。

如果你的消费超时时间设置得太短,会带来一个严重后果:正常的慢消息也会被判定为失败,然后被重新投递,引发重复。反之,设置得太长又会让异常消息卡住很久,阻塞队列消费进度。

3.3 一个容易被忽略的参数:消费线程数量

除了上面两个参数,consumeThreadMinconsumeThreadMax 也会影响心跳是否超时。消费线程数是处理消息的核心线程池的线程数量,如果线程数量不够,消息会堆积在阻塞队列里等待执行,造成处理延迟,进而影响心跳。

那么这些参数之间是怎样的配合关系?我用一句话概括:消费线程多了,处理快,心跳就不会被拖累;消费超时设得合理,就不会因为业务逻辑慢而误杀消息;心跳间隔设得合适,就不会因为偶发抖动造成误判。三者相辅相成,缺一不可。

四、一个完整的 Java 示例重现问题

4.1 技术栈说明

下面的示例全部基于 Java 技术栈,使用的客户端依赖是 rocketmq-client 版本 4.9.x。为了简洁,我直接写一个独立的可运行示例,模拟一条需要 20 秒才能处理完的业务消息,配合 5 秒的心跳间隔,让你看到心跳超时是怎么被触发的。

先引入 Maven 依赖,假设你已经有了一个 Maven 工程,在 pom.xml 里加入:


<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-client</artifactId>
    <version>4.9.7</version>
</dependency>

4.2 定义一个非常慢的消费者

这里模拟一个消费组,里面有多个消费者实例在运行同一个消费组,名字叫 order-consumer-group。消息是订单事件,正常的业务场景里可能在回调第三方接口,但这个第三方接口超时了 20 秒才返回。

代码里面有一个关键点,我故意把 RocketMQ 的心跳间隔调整到了 3 秒,这样问题暴露得更快。在实际项目中,你通常不直接修改这个值,因为默认值已经够用了,但了解它是怎么工作的很重要。

下面是一个消费者的代码:


import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;

import java.util.List;

/**
 * 模拟一个消费超时的消费者
 * 心跳间隔设置为3秒  故意让超时更容易被发现
 */
public class SlowConsumer {

    public static void main(String[] args) throws Exception {
        // 创建消费者实例  指定消费组名称
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("order-consumer-group");

        // 指定 NameServer 地址  换成你自己的环境地址
        consumer.setNamesrvAddr("127.0.0.1:9876");

        // 设置心跳间隔为3000毫秒  仅用于演示  日常不建议改这么短
        consumer.setHeartbeatBrokerInterval(3000);

        // 从队列最开始的消息开始消费  方便观察重复现象
        consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);

        // 消费线程池的最小线程数为5  最大线程数为10
        consumer.setConsumeThreadMin(5);
        consumer.setConsumeThreadMax(10);

        // 订阅主题  使用Tag过滤的方式只接收order开头的消息
        consumer.subscribe("OrderTopic", "order_*");

        // 注册消息监听器  每来一条消息都会回调这个方法
        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
                                                            ConsumeConcurrentlyContext context) {
                for (MessageExt msg : msgs) {
                    // 打印消息内容和当前线程名  便于观察是谁在处理
                    String body = new String(msg.getBody());
                    System.out.printf("线程 %s 收到消息 内容 %s  队列 = %d%n",
                            Thread.currentThread().getName(),
                            body,
                            msg.getQueueId());

                    // 模拟处理一条耗时很长的业务  比如调第三方的接口20秒
                    try {
                        Thread.sleep(20000);
                    } catch (InterruptedException e) {
                        Thread.currentThread().interrupt();
                    }
                }
                // 返回消费成功
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });

        // 启动消费者
        consumer.start();
        System.out.println("消费者已启动  注意观察心跳是否会超时");
    }
}

这段代码里最要命的就是 Thread.sleep(20000),它阻塞了当前消费线程 20 秒。在这个时间里,如果消费线程已经是最后一条线了,那么心跳线程也会受影响。实际上 RocketMQ 的心跳发送是单独的线程池,不会因为一条消息的阻塞而被拖死。但消费线程的阻塞会导致消息消费进度缓了下来,进而引发另一个问题:本地消息队列的积压。极端情况下,Broker 到客户端的拉取请求迟迟得不到响应,最终 Broker 会认为这个客户端已经失去了消费能力。

4.3 生产者的代码

为了补全示例,我把生产者代码也贴出来,方便你自己跑通整个链路,看到一条消息被重复消费的效果:


import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;

/**
 * 模拟发送订单消息
 * 你可以启动两个消费者实例同时消费同一个分组  观察队列重新分配
 */
public class OrderProducer {

    public static void main(String[] args) throws Exception {
        // 创建生产者  指定生产者组名 生产者组和消费者组分开
        DefaultMQProducer producer = new DefaultMQProducer("order-producer-group");

        // 指定 NameServer 地址
        producer.setNamesrvAddr("127.0.0.1:9876");

        // 启动生产者
        producer.start();

        // 连续发送十条消息  每条消息带一个唯一的订单号
        for (int i = 0; i < 10; i++) {
            String content = "订单号 ORDER202405" + i;
            Message message = new Message(
                    "OrderTopic",                  // 主题名称
                    "order_" + i,                  // 标签  用于过滤
                    content.getBytes("UTF-8")      // 消息体
            );

            // 发送消息并等待返回结果
            SendResult result = producer.send(message);
            System.out.printf("发送成功 消息ID = %s 队列 = %d%n",
                    result.getMsgId(),
                    result.getMessageQueue().getQueueId());
        }

        // 发送完成后关闭生产者  释放资源
        producer.shutdown();
    }
}

把上面的消费者代码复制一份,改一下启动类和实例名,比如 SlowConsumer2,两个消费者都订阅同一个组 order-consumer-group,然后同时启动。再启动生产者,你会看到两个消费者之间竟然出现了同一个消息被同时消费的情况,因为队列被频繁重分配了。

4.4 问题验证与日志分析

当你运行上面的代码时,可以从日志里发现几个端倪:

先看 Rebalance 的日志。你用 grep 搜索 doRebalance 之类的关键字,会看到消费组的负载均衡被反复执行。间隔非常短,有时候几秒钟就执行一次。

再看心跳日志。搜索 heartbeat 字样,你会发现有报错信息,接下来的内容我直接用命令展示:


# 这条命令在你的应用日志目录下执行  用来过滤rebalance相关日志
grep -n "doRebalance" ~/logs/rocketmq/consumer.log | tail -n 50

# 这条命令用来过滤心跳报错
grep -n "heartbeat" ~/logs/rocketmq/consumer.log | grep -i "error" | tail -n 50

如果心跳已经超时,日志中对应的地方会打出警告或者错误,并且你会观察到消费者的队列分配结果在不停变化。第一次有 3 个队列连到了这个消费者,过了一会儿变成 2 个,再过了一会儿又变成 4 个,这种异常波动就是重复消费的直接原因。

五、合理配置参数的详细方案

5.1 正确设置心跳间隔

RocketMQ 客户端中,设置心跳间隔的方法我在代码里已经用到了,就是 consumer.setHeartbeatBrokerInterval(3000)。但在生产环境中,不建议设成 3 秒这么短。原因很简单,一切正常的集群中,5 秒心跳绰绰有余,还能有效降低网络开销。

如果你用的是 Spring Boot 的 RocketMQ 集成,那么配置写在 yml 里面,像下面这样:


rocketmq:
  name-server: 127.0.0.1:9876
  consumer:
    group: order-consumer-group
    # SpringBoot集成时的心跳间隔配置项
    heartbeat-interval: 5000

注意,不同版本的 Spring Boot Starter 配置项名称可能不一样,要以你当前用的版本对应的文档为准。总体来说思路不变,间隔保持在 5 秒这个档位。

5.2 正确设置消费超时

消费超时的设置需要看你使用的客户端 API 类型。如果你是原生的 DefaultMQPushConsumer,可以直接调用下面的方法:


// 将默认的15分钟修改为更符合业务场景的5分钟
consumer.setConsumeTimeout(5 * 60 * 1000); // 单位毫秒

如果你用的是 Concurrently 模式的监听器,超时控制会比较宽松。如果你用的是顺序消息的监听器,那么超时控制会更严格一点,因为顺序消息不允许并行消费。

在 Spring Boot 的配置中,同样有对应的配置项:


rocketmq:
  consumer:
    consume-timeout: 300000 # 5分钟  单位毫秒

这里要非常注意一个细节:consumeTimeout 可以设置,但如果你在消费线程里使用了阻塞队列的 take() 方法,或者使用 Future.get() 并且没有设置超时时间,那么即使客户端层面的超时到了,底层的线程也可能被卡住不动。这个时候心跳发送虽然正常,但消费进度已经停滞,也会导致位点不更新,进而引发重复消费。

5.3 消费线程数的合理取值

消费线程数对并发能力的影响很大。假设你有 8 个队列,那么最小的合理线程数也应该是 8。不然只有 4 个线程去处理 8 个队列的消息,就算线程不阻塞,后面的队列也会排队。

推荐的做法是按队列数来设定。如果你的 Topic 有 16 个队列,则消费线程数设置为 16 到 32 之间。我的基准建议如下:

  • 消费线程最小数量 = 队列数量
  • 消费线程最大数量 = 队列数量乘以 2

示例配置如下:


consumer.setConsumeThreadMin(16);
consumer.setConsumeThreadMax(32);

5.4 批量消费大小的影响

RocketMQ 的消费监听器一次可能拿到一批消息,如果你没有设置 consumeMessageBatchMaxSize,默认是 1,也就是一次只消费一条。如果设置过大,比如一次拿 50 条,每条消息业务处理 1 秒,那么一批就是 50 秒,极其容易造成消费超时。在处理慢业务时,把批量大小调小是一个很好的策略。举例:


// 一次最多消费一条消息  避免批量处理时间过长
consumer.setConsumeMessageBatchMaxSize(1);

5.5 综合配置示例

现在把上面讲到的所有参数组合到一起,给你一个可以直接用于生产环境的基础版配置代码。我把每个参数的解释都写在注释里,方便你参考。


import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;

import java.util.List;

/**
 * 推荐配置示例
 * 该示例展示一个稳健的消费者配置  主要解决心跳超时引发的重复消费
 * 所有关键参数都有注释说明  你可以直接参考或复制
 */
public class SafeConsumer {

    public static void main(String[] args) throws Exception {
        // 创建消费组
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("order-consumer-group");

        // 指定NameServer地址
        consumer.setNamesrvAddr("127.0.0.1:9876");

        // 心跳间隔设置为5秒  这是官方推荐值  默认本来就是5秒
        consumer.setHeartbeatBrokerInterval(5000);

        // 消费超时设置为5分钟
        // 这里设置的超时时间应该大于你单条消息最大的业务处理耗时
        // 假设你调用的第三方接口最长10秒  这里可以设置为30秒或60秒
        consumer.setConsumeTimeout(300000);

        // 消费线程数量设置
        // 假如你的Topic有16个队列  最小16  最大32
        consumer.setConsumeThreadMin(16);
        consumer.setConsumeThreadMax(32);

        // 批量消费大小保持为1
        // 这个设置能让单条消息的处理时间变短  减少阻塞
        consumer.setConsumeMessageBatchMaxSize(1);

        // 从最新消费 也就是只消费新产生的消息
        consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET);

        // 订阅主题
        consumer.subscribe("OrderTopic", "order_*");

        // 注册监听器
        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
                                                            ConsumeConcurrentlyContext context) {
                for (MessageExt msg : msgs) {
                    try {
                        // 这里写真正的业务逻辑
                        System.out.printf("业务处理开始 消息内容 = %s%n",
                                new String(msg.getBody()));

                        // 模拟业务耗时10毫秒
                        Thread.sleep(10);

                    } catch (Exception e) {
                        // 业务异常时可返回重试
                        System.out.println("业务处理发生异常 稍后重试");
                        return ConsumeConcurrentlyStatus.RECONSUME_LATER;
                    }
                }
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });

        // 启动消费者
        consumer.start();
        System.out.println("安全消费者已启动");
    }
}

你看,把关键参数调整好之后,重复消费的土壤就没有了。即使偶发一次消息处理超时,RocketMQ 也会正确地重试,而不是引发分区分派震荡。

六、更深入的诊断技巧:关联技术的穿插介绍

在排查这种问题的时候,单纯靠 RocketMQ 的日志还不够,我通常会联动操作系统的网络命令和 JVM 工具一起用。这里穿插一个非常有用的技巧,你可以通过 jstack 命令去查看消费线程在干嘛。

比如你怀疑消费线程卡住了,先用 jps 找到 Java 进程 ID,再用 jstack 导出线程堆栈,然后搜索 ConsumeMessageThread 关键名称,就能看到线程的具体状态。下面给出具体命令:


# 查看当前运行的Java进程 找到你的Consumer进程ID
jps -l

# 假设进程ID是 12345  导出线程快照到文件
jstack 12345 > /tmp/rocketmq_stack.txt

# 查看消费线程的状态  从堆栈信息中分析瓶颈
grep -n "ConsumeMessageThread" /tmp/rocketmq_stack.txt

执行上面命令之后,你可能会看到下面的线程信息:


"ConsumeMessageThread_1" #20 daemon prio=5 os_prio=0 tid=0x00007f8b1c3c2800 nid=0x3b3b waiting on condition [0x00007f8b1a1d9000]
   java.lang.Thread.State: TIMED_WAITING (sleeping)

这里显示 TIMED_WAITING 状态,且处于 sleeping,说明这个线程确实在睡觉。如果所有的消费线程都在睡觉,那你就能直观地定位到问题。

此外,RocketMQ 的监控指标可以配合 Prometheus 和 Grafana 来抓取,主要关注下面几个指标:


# 查看消费者组的消费位点落后情况  通过mqadmin工具
mqadmin consumerProgress -n 127.0.0.1:9876 -g order-consumer-group

如果某一条消息的积压数量一直在增长,说明消费能力跟不上,这时候要优先考虑增加消费者实例或线程数,而不是调整心跳参数。

七、应用场景与注意事项

7.1 什么样的业务场景最容易踩坑

消息处理耗时长且波动大的业务,最容易被心跳超时和 Rebalance 震荡折磨。比如订单系统中调用第三方支付接口,如果支付回调的响应时间不可控,偶发性的 20 秒甚至更长的等待就会成为常态。再比如库存服务里面做分布式锁的操作,一个锁的等待时间可能就会超过几秒钟。

还有一类场景是大批量数据处理。比如你一条消息里面包含了 1000 个明细,逐条解析入库,总共需要 30 秒以上。如果你没有把 consumeMessageBatchMaxSize 调小,也没有设置足够长的消费超时,心跳超时几乎不可避免。

比如下面这段代码就是典型的反面教材:


// 反面教材  不要这样做
consumer.registerMessageListener(new MessageListenerConcurrently() {
    @Override
    public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
                                                    ConsumeConcurrentlyContext context) {
        // 这里默认一次拉取10条消息
        for (MessageExt msg : msgs) {
            // 每条消息处理时间2秒  10条就是20秒
            processMessage(msg);
        }
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    }
});

7.2 完整参数调整步骤

假如你现在已经出现了分区分派震荡的问题,建议按照以下步骤操作:

第一步先暂停消费者的自动启动流程,确保不再继续产生堆积。 第二步查看当前使用的 RocketMQ 版本,不同版本的参数位置不同。 第三步根据业务耗时的最大值来设置 consumeTimeout,比如最大耗时 30 秒,超时设成 60 秒。 第四步把批量大小调整为 1,确保不会因为多条消息累计耗时导致超时。 第五步重启消费者,观察心跳日志和 Rebalance 次数。 第六步如果依然震荡,检查网络波动和 GC 停顿。

最近我排查过一个真实的线上事故,最后发现引发震荡的不是消费超时,而是 Full GC 停顿时间过长。JVM 在 Full GC 的时候会暂停所有线程,包括心跳线程,如果停顿超过 10 秒,Broker 自然就会把客户端踢下线。这种情况更需要关注 JVM 的堆内存和 GC 参数。

可以顺手在启动脚本里加上下面的 JVM 参数来打印 GC 日志:


# 开启GC日志  方便排查是否是GC导致的长时间停顿
-XX:+PrintGCDetails -XX:+PrintGCDateStamps -Xloggc:/tmp/gc.log

7.3 技术优缺点分析

Heartbeat 间隔设置太短,优势是能快速感知消费者宕机,减少 Rebalance 的延迟,但劣势是加重了 Broker 网络负担,客户端频繁发送心跳包,在消费者数量多的场景下会影响 NameServer 的存储性能。

心跳间隔设置太长,优点是网络开销小,系统更稳定,但缺点是消费者挂掉之后要等很久才被感知,这期间该消费者名下的队列一直在空转,无法被其他消费者接手,造成消息积压。

消费超时时间设置得过短,优点是能快速淘汰有问题的消费请求,但缺点是一些正常的慢消费会被误杀,加入重复投递,增加下游系统的处理压力。

消费超时设置得过长,优点是允许业务处理时间更宽裕,但缺点是一旦消费逻辑卡死,消息在队列里迟迟得不到最终结果,会导致位点无法提交,消息积压越来越严重。

线程数设置过多,会消耗更多内存和 CPU,并且如果消费逻辑中存在共享资源的竞争,太多线程反而会让竞争更激烈。线程数过少,即使你有再合理的心跳间隔,消息处理不过来,一样会有反压问题。

7.4 杀招:优雅处理慢消息的思路

如果你的业务中确实存在无法缩短耗时的操作,我推荐使用事务消息的思路配合异步化处理。核心思路是把消息体中的核心信息拆出来,先快速把消费提交,然后把耗时的业务逻辑放到一个单独的线程池里面去慢慢做。只是这样需要你自己维护一个补偿机制,确保异步任务最终执行完成并发出结果。

下面给出一个简单的伪代码框架:


import java.util.concurrent.*;

/**
 * 异步处理消息的示例
 * 先快速提交耗时敏感的消息  再丢给业务线程池慢慢处理
 */
public class AsyncConsumeDemo {

    // 创建一个业务线程池  核心线程数10  最大线程数20
    private static final ExecutorService BUSINESS_POOL =
            new ThreadPoolExecutor(
                    10,
                    20,
                    60L,
                    TimeUnit.SECONDS,
                    new ArrayBlockingQueue<>(1000),
                    new ThreadPoolExecutor.CallerRunsPolicy()
            );

    public static void main(String[] args) {
        System.out.println("异步消费启动,实现快速提交避免心跳超时");
        // 业务中可以这样提交任务
        BUSINESS_POOL.submit(() -> {
            // 这里执行慢任务  例如调用第三方接口30秒
            callThirdPartyApi();
        });
    }

    private static void callThirdPartyApi() {
        // 模拟耗时操作
        try {
            Thread.sleep(30000);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

这样的好处是消息监听器很快就返回成功,RocketMQ 的位点也会快速推进,不会引发心跳超时。缺点是你需要自己处理异步任务失败的补偿逻辑,否则任务丢了没人知道。要记得给异步线程池加监控和告警。

八、文章总结

整篇文章从一次线上重复消费的事故讲起,剖析了心跳超时到分区分派震荡这条因果链。完整解释了一套通俗的流程:消费耗时过长导致心跳信号发送不及时,Broker 误判消费者下线,触发 Rebalance,Revisited 队列在多个消费者之间频繁交换,导致消息从旧位点被重复读取。

要解决这个问题,核心思路不是简单粗暴地调大某一个值,而是要让消费端的整个运行节奏保持健康。心跳间隔建议固定在 5 秒左右,不受业务复杂度影响;消费超时建议根据你实际业务的最大耗时来定,预留一定余量即可。消费线程数要与队列数和业务并发度匹配,批量消费大小尽量保持为 1,给每条消息独立的时间预算。

为了让参数设置的思路更稳固,文章提供了完整的 Java 示例代码,包含了一个慢消费者的复现程序、一个生产者的发送程序、一个安全消费者的推荐配置,以及一些诊断用的 Linux 命令和配置片段,大家在本地搭建一套最小环境就能直观看到问题。

最后一定要养成系统性排查的习惯。当你看到重复消费激增、Rebalance 日志刷屏时,不要只盯着消费超时一个参数,先抓线程栈、看 GC 日志、看 Broker 端的网络连接状态,把问题定位置准确。同时也要明白,参数设置再合理也不能替代优雅的业务代码。把消息处理做快做稳,才是真正的一劳永逸。

希望这篇文章能帮你在面对 RocketMQ 的奇怪故障时淡定不少。哪怕只有一次真正帮到你查出了问题,那这些文字也算没有白费。