RocketMQ 4.8消费者仅消费部分队列的原因及修复方案
问题描述
向RocketMQ 4.8发送消息后,消费者仅消费了部分队列的消息,所用的RocketMQ消费者代码如下:
public void appConsumer(Long appId, List<String> topics) throws MQClientException { DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(appId.toString()); RocketMQConfigDTO rocketMQConfigDTO = MqConfigHandler.MQConfig(); consumer.setNamesrvAddr(rocketMQConfigDTO.getHost()); consumer.setMessageModel(MessageModel.CLUSTERING); consumer.setMaxReconsumeTimes(rocketMQConfigDTO.getMaxReconsumeTimes()); consumer.setConsumeThreadMin(rocketMQConfigDTO.getConsumeThreadMin()); consumer.setConsumeThreadMax(rocketMQConfigDTO.getConsumeThreadMax()); consumer.setConsumeMessageBatchMaxSize(rocketMQConfigDTO.getConsumeMessageBatchMaxSize()); consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET); for (String topic : topics) { consumer.subscribe(topic, "*"); } consumer.registerMessageListener(new MessageListenerConcurrently() { @Override public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) { Callable<ConsumeConcurrentlyStatus> task = new Callable<ConsumeConcurrentlyStatus>() { @Override public ConsumeConcurrentlyStatus call() throws Exception { return handleMessage(msgs, context,appId); } }; asyncService.submit(task); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); consumer.start(); eventConsumers.put(appId, consumer); }
登录RocketMQ控制台查看时,发现同一Broker下消费者仅消费了部分队列的消息:
队列0和队列1的消息未被消费者消费,需要排查问题产生原因、修复方案,以及消费者仅消费队列3和队列4消息的原因。
问题根因
- 核心bug出在消息监听器的实现逻辑:
consumeMessage方法中把实际消费逻辑提交到自定义asyncService异步线程池后,直接返回了CONSUME_SUCCESS。RocketMQ PushConsumer本身内置了完整的消费线程池和消费进度管控机制,这种额外加一层异步提交、提前返回消费成功的写法,会直接触发客户端的队列流控机制,导致部分队列停止拉取消息。 - RocketMQ客户端默认配置
pullThresholdForQueue=1000,即单个队列在客户端本地缓存的消息数超过1000时,客户端会自动停止从该队列拉取新消息。提前返回消费成功后,客户端会误判这批消息已经处理完成,立刻继续从对应队列拉取新消息填充本地缓存;如果自定义异步线程池的实际处理速度跟不上拉取速度,部分队列的本地缓存会快速堆积到阈值,触发流控停止拉取。截图里只有队列3、4能被消费,本质是这两个队列的拉取速度和异步线程池的处理速度刚好匹配,没有触发流控阈值,剩下的0、1队列都因为缓存堆积被暂停拉取了。 - 存在小概率场景:集群消费模式下,同一个消费组的队列会平均分配给所有在线的消费者实例,如果消费组下还有其他运行中的消费者实例,队列0、1可能已经被分配给其他实例,当前实例自然不会消费这两个队列的消息。但结合给出的代码,90%以上的概率是异步提前返回ACK导致的流控问题。
修复方案
- 优先删掉自定义的异步提交逻辑:不要在
consumeMessage里用asyncService.submit提交任务,直接在方法内同步调用handleMessage,等业务逻辑实际执行完成后再返回对应的消费状态。代码里已经配置了consumeThreadMin、consumeThreadMax参数,PushConsumer自带的线程池完全可以满足异步消费的需求,不需要额外套一层自定义线程池。 - 如果业务场景必须用自定义异步线程池处理,绝对不能提前返回
CONSUME_SUCCESS,必须等异步任务实际执行完成后,再返回对应消费状态,同时要根据异步线程池的实际处理能力,调整pullThresholdForQueue等流控参数,避免本地缓存消息堆积触发流控。 - 登录RocketMQ控制台查看消费组的连接信息,确认当前消费组下的在线实例数和队列分配结果:如果是集群模式下队列被分配给其他正常运行的实例,属于正常负载均衡逻辑,不需要额外处理;如果所有实例都存在部分队列消费卡住的情况,优先修复异步提前返回ACK的代码问题。
- 代码修复后重启消费实例,观察控制台各队列的消费进度是否正常推进,客户端本地缓存的消息数是否稳定在流控阈值以下即可。
内容的提问来源于stack exchange,提问作者Dolphin
相关产品推荐
相关产品推荐

