RabbitMQ 手动 ACK 消息确认机制详解
基本概念
ACK (Acknowledgment):消费者处理完消息后,通知 RabbitMQ 可以删除消息的确认信号。

手动 ACK:由开发者在代码中显式调用确认方法,而非自动确认。
为什么需要手动 ACK?
| 自动 ACK | 手动 ACK |
|---|---|
| 消息投递后立即删除 | 消费者处理后才会删除 |
| 可能丢失消息(消费者崩溃时) | 保证消息不丢失 |
| 无法重试处理失败的消息 | 支持失败重试 |
基础实现
1 生产者(不变)
// 生产者无需改动,只需确保消息持久化
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.deliveryMode(2) // 持久化消息
.build();
channel.basicPublish("", queueName, props, message.getBytes());
2 消费者手动 ACK
Channel channel = connection.createChannel();
channel.basicQos(1); // 每次只处理一个消息
// 关键参数:autoAck = false
channel.basicConsume(queueName, false, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope,
AMQP.BasicProperties properties, byte[] body) throws IOException {
try {
String message = new String(body, "UTF-8");
System.out.println("Received: " + message);
// 处理业务逻辑
processMessage(message);
// 手动确认:处理成功
channel.basicAck(envelope.getDeliveryTag(), false);
} catch (Exception e) {
// 处理失败:拒绝消息(可决定是否重新入队)
channel.basicNack(envelope.getDeliveryTag(), false, true);
// 或者使用 basicReject
// channel.basicReject(envelope.getDeliveryTag(), true);
}
}
});
Spring AMQP 实现
@Component
public class MessageConsumer {
@RabbitListener(queues = "myQueue")
public void handleMessage(String message, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) {
try {
System.out.println("Received: " + message);
processMessage(message);
// 手动确认
channel.basicAck(tag, false); // false = 只确认当前消息
} catch (Exception e) {
try {
// 拒绝消息,重新入队
channel.basicNack(tag, false, true);
// 或 basicReject
// channel.basicReject(tag, true);
} catch (IOException ex) {
ex.printStackTrace();
}
}
}
}
配置 Spring Boot
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual # 关键配置:手动模式
prefetch: 1 # 每次预取1条
retry:
enabled: true # 启用重试
max-attempts: 3
@Configuration
public class RabbitConfig {
@Bean
public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory(
ConnectionFactory connectionFactory) {
SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
factory.setConnectionFactory(connectionFactory);
factory.setAcknowledgeMode(AcknowledgeMode.MANUAL); // 手动确认
factory.setPrefetchCount(1);
return factory;
}
}
核心 API
| 方法 | 参数 | 说明 |
|---|---|---|
basicAck(tag, multiple) |
tag: 消息标识 multiple: 是否批量确认 |
确认消息已处理 |
basicNack(tag, multiple, requeue) |
requeue: 是否重新入队 | 拒绝消息 |
basicReject(tag, requeue) |
同Nack但不可批量 | 拒绝消息(旧API) |
常见场景处理
1 消息处理失败重试
try {
// 业务处理
processMessage(message);
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
// 判断重试次数(需要自己维护)
if (retryCount < MAX_RETRIES) {
channel.basicNack(deliveryTag, false, true); // 重新入队
} else {
// 超过重试次数:记录日志、发送到死信队列
channel.basicNack(deliveryTag, false, false); // 不入队
// 或发送到死信交换机
sendToDeadLetter(message);
}
}
2 批量确认
// 批量确认:确认所有小于等于 tag 的消息 channel.basicAck(deliveryTag, true); // 批量拒绝 channel.basicNack(deliveryTag, true, false);
注意事项
-
消息持久化:确保消息不会丢失
// 生产者设置 channel.confirmSelect(); // 发布者确认
-
QoS 设置:控制预取数量
channel.basicQos(1); // 每次处理1条
-
异常处理:
// 确保任何时候都能执行 ACK finally { if (!acked) { channel.basicNack(tag, false, true); } } -
死信队列:配置失败消息处理
Map<String, Object> args = new HashMap<>(); args.put("x-dead-letter-exchange", "dlx-exchange"); args.put("x-dead-letter-routing-key", "dlx-routing-key");
完整示例:可靠消费
public class ReliableConsumer {
private final int MAX_RETRIES = 3;
private final Map<Long, Integer> retryCounts = new ConcurrentHashMap<>();
public void consume() throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
channel.basicQos(1);
channel.basicConsume("myQueue", false, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope,
AMQP.BasicProperties properties, byte[] body) {
long deliveryTag = envelope.getDeliveryTag();
try {
String message = new String(body, "UTF-8");
System.out.println("Processing: " + message);
// 业务处理
processMessage(message);
// 成功确认
channel.basicAck(deliveryTag, false);
retryCounts.remove(deliveryTag);
} catch (Exception e) {
// 获取重试次数
int retryCount = retryCounts.getOrDefault(deliveryTag, 0);
if (retryCount < MAX_RETRIES) {
retryCounts.put(deliveryTag, retryCount + 1);
System.out.println("Retry " + (retryCount + 1) + "/" + MAX_RETRIES);
try {
channel.basicNack(deliveryTag, false, true);
} catch (IOException ex) {
ex.printStackTrace();
}
} else {
System.out.println("Max retries reached, sending to DLQ");
try {
// 不重新入队,消息会进入死信队列
channel.basicNack(deliveryTag, false, false);
retryCounts.remove(deliveryTag);
} catch (IOException ex) {
ex.printStackTrace();
}
}
}
}
});
// 保持运行
Thread.sleep(Long.MAX_VALUE);
}
}
}
调试技巧
# 查看未确认消息数 rabbitmqctl list_queues name messages_unacknowledged # 查看消费者状态 rabbitmqctl list_consumers # 查看死信队列 rabbitmqctl list_queues name messages messages_unacknowledged
手动 ACK 是保证消息可靠性的关键机制,核心思想是业务处理成功后才确认消息已消费,确保消息不会在消费者处理过程中丢失。