You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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();
    }
}

排查建议

  1. 消费者初始化问题:移除@Async("taskExecutor")注解,异步初始化可能导致消费者未正确启动。若必须异步,需确认taskExecutor线程池配置无阻塞、核心线程数足够。
  2. 消费模式调整:当前使用MessageListenerOrderly顺序消费,并发度受队列数量限制。可临时改为MessageListenerConcurrently并发消费,测试是否能正常消费。
  3. 消费线程配置:手动设置消费者线程数,避免默认值过低:
    consumer.setConsumeThreadMin(20);
    consumer.setConsumeThreadMax(64);
    
  4. 版本兼容性:Spring Boot 2.5.12搭配rocketmq-spring-boot-starter 2.2.1存在版本适配风险,建议升级到2.2.3版本。
  5. Broker状态检查:
    • 查看Broker磁盘使用率,超过85%会触发流控,导致消息投递延迟。
    • 检查Broker的waitTimeMillsInSendQueue配置,确保未被修改为非0值。
  6. 日志增强:在消费者start()后添加日志,确认消费者是否成功启动并订阅topic:
    consumer.start();
    System.out.println("消费者启动完成,订阅topic:" + topic + ",消费者组:" + consumerGroup);
    

内容的提问来源于stack exchange,提问作者JessieJ

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 22:26:04