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

RabbitMQ Java客户端消息未正常确认问题排查求助

RabbitMQ Push模式消费者消息确认异常问题

Java客户端依赖

<dependency>
    <groupId>com.rabbitmq</groupId>
    <artifactId>amqp-client</artifactId>
    <version>5.17.0</version>
</dependency>

生产者代码

package com.dsabyte.rabbitmqdemo;

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

public class Send {
    public static void main(String[] args) {
        ConnectionFactory factory = new ConnectionFactory();

        try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel();) {

            channel.queueDeclare("java_queue", true, false, false, null);
            channel.confirmSelect();
            for (int counter = 1; counter <= 10; counter++) {
                channel.basicPublish("", "java_queue", null, String.valueOf(counter).getBytes());
            }
            System.out.println("Messages sent!");
        } catch (Exception e) {
            System.out.println(e);
        }
    }
}

Push模式消费者代码

package com.dsabyte.rabbitmqdemo;

import java.io.IOException;
import java.util.concurrent.TimeoutException;

import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;

public class SubScribe {

    public static void main(String[] args) throws IOException, TimeoutException {
        ConnectionFactory factory = new ConnectionFactory();

        try (Connection conn = factory.newConnection(); Channel channel = conn.createChannel();) {

            channel.basicConsume("java_queue", false, "consumer_tag", new DefaultConsumer(channel) {
                @Override
                public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,
                        byte[] body) throws IOException {

                    System.out.println("Msg: " + new String(body) + ", redeliver: " + envelope.isRedeliver());
                    channel.basicAck(envelope.getDeliveryTag(), false);
                }
            });
        } catch (Exception ex) {
            ex.printStackTrace();
        }
    }

}

说明:使用Spring Tools 4.11.0编辑器运行程序

运行情况

  1. 生产者运行输出:
SLF4J: Failed to load class "org.slf4j.impl.StaticLoggerBinder".
SLF4J: Defaulting to no-operation (NOP) logger implementation
SLF4J: See http://www.slf4j.org/codes.html#StaticLoggerBinder for further details.
Messages sent!
  1. 首次启动消费者输出:
SLF4J: Failed to load class "org.slf4j.impl.StaticLoggerBinder".
SLF4J: Defaulting to no-operation (NOP) logger implementation
SLF4J: See http://www.slf4j.org/codes.html#StaticLoggerBinder for further details.
Msg: 1, redeliver: false
Msg: 2, redeliver: false
Msg: 3, redeliver: false
Msg: 4, redeliver: false
Msg: 5, redeliver: false
Msg: 6, redeliver: false
Msg: 7, redeliver: false
Msg: 8, redeliver: false
Msg: 9, redeliver: false
Msg: 10, redeliver: false
com.rabbitmq.client.ShutdownSignalException: clean channel shutdown; protocol method: #method<channel.close>(reply-code=200, reply-text=Closed due to exception from Consumer (consumer_tag) method handleDelivery for channel AMQChannel(amqp://guest@127.0.0.1:5672/,1), class-id=0, method-id=0)
    at com.rabbitmq.utility.ValueOrException.getValue(ValueOrException.java:66)
    at com.rabbitmq.utility.BlockingValueOrException.uninterruptibleGetValue(BlockingValueOrException.java:36)
    at com.rabbitmq.client.impl.AMQChannel$BlockingRpcContinuation.getReply(AMQChannel.java:502)
    at com.rabbitmq.client.impl.ChannelN.close(ChannelN.java:617)
    at com.rabbitmq.client.impl.ChannelN.close(ChannelN.java:542)
    at com.rabbitmq.client.impl.ChannelN.close(ChannelN.java:535)
    at com.rabbitmq.client.impl.recovery.AutorecoveringChannel.lambda$close$0(AutorecoveringChannel.java:74)
    at com.rabbitmq.client.impl.recovery.AutorecoveringChannel.executeAndClean(AutorecoveringChannel.java:102)
    at com.rabbitmq.client.impl.recovery.AutorecoveringChannel.close(AutorecoveringChannel.java:74)
    at com.dsabyte.rabbitmqdemo.SubScribe.main(SubScribe.java:29)
Caused by: com.rabbitmq.client.ShutdownSignalException: clean channel shutdown; protocol method: #method<channel.close>(reply-code=200, reply-text=Closed due to exception from Consumer (consumer_tag) method handleDelivery for channel AMQChannel(amqp://guest@127.0.0.1:5672/,1), class-id=0, method-id=0)
    at com.rabbitmq.client.impl.ChannelN.close(ChannelN.java:589)
    at com.rabbitmq.client.impl.ChannelN.close(ChannelN.java:542)
    at com.rabbitmq.client.impl.StrictExceptionHandler.handleChannelKiller(StrictExceptionHandler.java:72)
    at com.rabbitmq.client.impl.StrictExceptionHandler.handleConsumerException(StrictExceptionHandler.java:61)
    at com.rabbitmq.client.impl.ConsumerDispatcher$5.run(ConsumerDispatcher.java:154)
    at com.rabbitmq.client.impl.ConsumerWorkService$WorkPoolRunnable.run(ConsumerWorkService.java:111)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)
  1. 再次运行消费者输出:
SLF4J: Failed to load class "org.slf4j.impl.StaticLoggerBinder".
SLF4J: Defaulting to no-operation (NOP) logger implementation
SLF4J: See http://www.slf4j.org/codes.html#StaticLoggerBinder for further details.
Msg: 7, redeliver: true
Msg: 8, redeliver: true
Msg: 9, redeliver: true
Msg: 10, redeliver: true
  1. 第三次运行消费者无任何输出。

从redeliver: true可以看出消息未被正常确认,但报错信息未明确指出问题所在。

该问题仅出现在Push模式消费者中,Pull模式消费者无此问题。

Pull模式消费者代码

package com.dsabyte.rabbitmqdemo;

import java.io.IOException;
import java.util.concurrent.TimeoutException;

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.GetResponse;

public class Recv {
    public static void main(String[] args) throws InterruptedException {
        ConnectionFactory factory = new ConnectionFactory();

        DeliverCallback callback = (consumerTag, delivery) -> {
            String msg = new String(delivery.getBody());
            System.out.println("Received message: " + msg);
        };

        try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel();) {
            channel.queueDeclare("java_queue", true, false, false, null);
            System.out.println("Waiting for the message");

            while (true) {
                GetResponse msg = channel.basicGet("java_queue", false);

                if (msg != null) {
                    System.out.println("Msg: " + new String(msg.getBody()) + ", deliveryTag: "
                            + msg.getEnvelope().getDeliveryTag());
                    channel.basicAck(msg.getEnvelope().getDeliveryTag(), false);
                } else {
                    break;
                }
            }

        } catch (IOException e) {
            e.printStackTrace();
        } catch (TimeoutException e) {
            e.printStackTrace();
        }
    }
}

问题原因及解决方案

问题根源

Push模式的basicConsume是异步操作,消息投递逻辑在后台线程中处理。而代码中使用的try-with-resources会在main方法执行完basicConsume后立即关闭Channel和Connection,此时后台线程的handleDelivery方法还在处理消息,调用basicAck时通道已关闭,抛出异常导致消息未被成功确认,RabbitMQ会将这些消息标记为未确认并重新投递。

解决方案

需要让main线程保持活跃,避免连接和通道被过早关闭,以下是几种实现方式:

方案1:使用CountDownLatch等待(推荐)

package com.dsabyte.rabbitmqdemo;

import java.io.IOException;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeoutException;

import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;

public class SubScribe {

    public static void main(String[] args) throws IOException, TimeoutException, InterruptedException {
        ConnectionFactory factory = new ConnectionFactory();
        CountDownLatch latch = new CountDownLatch(1);

        Connection conn = factory.newConnection();
        Channel channel = conn.createChannel();

        try {
            channel.basicConsume("java_queue", false, "consumer_tag", new DefaultConsumer(channel) {
                @Override
                public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,
                        byte[] body) throws IOException {

                    System.out.println("Msg: " + new String(body) + ", redeliver: " + envelope.isRedeliver());
                    channel.basicAck(envelope.getDeliveryTag(), false);
                    // 若需处理完所有消息后退出,可在此判断并调用latch.countDown()
                }
            });

            latch.await(); // 阻塞main线程,直到调用countDown()触发关闭
        } catch (Exception ex) {
            ex.printStackTrace();
        } finally {
            channel.close();
            conn.close();
        }
    }
}

方案2:添加线程休眠(适合测试场景)

package com.dsabyte.rabbitmqdemo;

import java.io.IOException;
import java.util.concurrent.TimeoutException;

import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;

public class SubScribe {

    public static void main(String[] args) throws IOException, TimeoutException, InterruptedException {
        ConnectionFactory factory = new ConnectionFactory();

        Connection conn = factory.newConnection();
        Channel channel = conn.createChannel();

        try {
            channel.basicConsume("java_queue", false, "consumer_tag", new DefaultConsumer(channel) {
                @Override
                public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,
                        byte[] body) throws IOException {

                    System.out.println("Msg: " + new String(body) + ", redeliver: " + envelope.isRedeliver());
                    channel.basicAck(envelope.getDeliveryTag(), false);
                }
            });

            Thread.sleep(10000); // 休眠10秒,确保消息处理完成
        } catch (Exception ex) {
            ex.printStackTrace();
        } finally {
            channel.close();
            conn.close();
        }
    }
}

方案3:监听关闭信号(生产环境推荐)

通过注册ShutdownListener或使用信号量实现优雅关闭,确保所有消息处理完成后再关闭连接。

验证效果

修改后运行消费者,消息会被正常确认,不会出现redeliver: true的情况,也不会抛出通道关闭异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 19:17:01