RocketMQ生产者可发消息但消费者无法正常消费问题求助
RocketMQ消费者消费延迟/不消费问题排查求助
生产者可正常向队列发送消息,但消费者无法及时消费,偶尔延迟约10秒,多数时候完全不消费。查看RocketMQ控制面板,存在未消费消息堆积。
application.yml配置
rocketmq: name-server: 127.0.0.1:9876 producer: group: ProducerGroup consumer: group: ConsumerGroup topic: MyTopic
Maven依赖
<properties> <java.version>11</java.version> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding> </properties> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.5.12</version> </parent> .... <dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-spring-boot-starter</artifactId> <version>2.2.1</version> </dependency>
生产者代码
import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.client.producer.DefaultMQProducer; import org.apache.rocketmq.common.message.Message; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; @Component public class RocketMQProducer { private DefaultMQProducer producer; @Autowired public RocketMQProducer(@Value("${rocketmq.producer.group}") String producerGroup, @Value("${rocketmq.name-server}") String nameServerAddress) { producer = new DefaultMQProducer(producerGroup); producer.setNamesrvAddr(nameServerAddress); try { producer.start(); } catch (MQClientException e) { System.out.println("RocketMQProducer is wrong: "); e.printStackTrace(); } } public void sendMessages(String topic, String tags, String message) { try { Message mqMessage = new Message(topic, tags, message.getBytes()); producer.send(mqMessage); System.out.println("Message sent to queue: " + message); } catch (Exception e) { System.out.println("RocketMQProducer sendMessages() is wrong: "); e.printStackTrace(); } } public void shutdown() { producer.shutdown(); } }
Netty客户端调用生产者代码(冗余代码已省略)
@Component public class NettyClientHandler extends ChannelInboundHandlerAdapter { @Autowired private RocketMQProducer rocketMQProducer; @Value("${rocketmq.consumer.topic}") private String topic; @Override public void channelActive(ChannelHandlerContext ctx) throws Exception { LogUtil.log("Client,channelActive"); } @Override public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception { LogUtil.log("Client,Received a message from the server"); if (msg instanceof ByteBuf) { ByteBuf byteBuf = (ByteBuf) msg; String message = byteBuf.toString(StandardCharsets.UTF_8); rocketMQProducer.sendMessages(topic, "", message); System.out.println("Received message: " + message); } } }
消费者代码
@Service public class RocketMQCommonConsumerListener implements CommandLineRunner { @Autowired private subway.service.Service service; @Value("${rocketmq.consumer.group}") private String consumerGroup; @Value("${rocketmq.name-server}") private String nameServerAddress; @Value("${rocketmq.consumer.topic}") private String topic; public void consumeMessages() { DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(consumerGroup); consumer.setNamesrvAddr(nameServerAddress); consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); System.out.println("consume_step1"); try { consumer.subscribe(topic, "*"); System.out.println("consume_step2"); consumer.registerMessageListener(new MessageListenerOrderly() { @Override public ConsumeOrderlyStatus consumeMessage(List<MessageExt> messages, ConsumeOrderlyContext context) { StringBuilder sb = new StringBuilder(); boolean needMerge = true; System.out.print("consume_step3, "); long threadId = Thread.currentThread().getId(); System.out.println("Current Thread ID: " + threadId); for (MessageExt message : messages) { String str = new String(message.getBody()); System.out.println("RocketMQ received message:" + str); } return ConsumeOrderlyStatus.SUCCESS; } }); consumer.start(); } catch (MQClientException e) { e.printStackTrace(); } } private void handleJson(String json) { System.out.println("json data is :"); System.out.println(json); System.out.println("\t\t\t============\t\t\t\n"); } @Async("taskExecutor") @Override public void run(String... args){ consumeMessages(); } }
排查建议
- 消费者初始化问题:移除
@Async("taskExecutor")注解,异步初始化可能导致消费者未正确启动。若必须异步,需确认taskExecutor线程池配置无阻塞、核心线程数足够。 - 消费模式调整:当前使用
MessageListenerOrderly顺序消费,并发度受队列数量限制。可临时改为MessageListenerConcurrently并发消费,测试是否能正常消费。 - 消费线程配置:手动设置消费者线程数,避免默认值过低:
consumer.setConsumeThreadMin(20); consumer.setConsumeThreadMax(64); - 版本兼容性:Spring Boot 2.5.12搭配rocketmq-spring-boot-starter 2.2.1存在版本适配风险,建议升级到2.2.3版本。
- Broker状态检查:
- 查看Broker磁盘使用率,超过85%会触发流控,导致消息投递延迟。
- 检查Broker的
waitTimeMillsInSendQueue配置,确保未被修改为非0值。
- 日志增强:在消费者
start()后添加日志,确认消费者是否成功启动并订阅topic:consumer.start(); System.out.println("消费者启动完成,订阅topic:" + topic + ",消费者组:" + consumerGroup);
内容的提问来源于stack exchange,提问作者JessieJ
相关产品推荐
相关产品推荐

