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编辑器运行程序
运行情况
- 生产者运行输出:
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!
- 首次启动消费者输出:
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)
- 再次运行消费者输出:
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
- 第三次运行消费者无任何输出。
从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
相关产品推荐
相关产品推荐

