RabbitMQ消息确认手动ACK

wen java案例 2

RabbitMQ 手动 ACK 消息确认机制详解

基本概念

ACK (Acknowledgment):消费者处理完消息后,通知 RabbitMQ 可以删除消息的确认信号。

RabbitMQ消息确认手动ACK

手动 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);

注意事项

  1. 消息持久化:确保消息不会丢失

    // 生产者设置
    channel.confirmSelect(); // 发布者确认
  2. QoS 设置:控制预取数量

    channel.basicQos(1); // 每次处理1条
  3. 异常处理

    // 确保任何时候都能执行 ACK
    finally {
        if (!acked) {
            channel.basicNack(tag, false, true);
        }
    }
  4. 死信队列:配置失败消息处理

    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 是保证消息可靠性的关键机制,核心思想是业务处理成功后才确认消息已消费,确保消息不会在消费者处理过程中丢失。

抱歉,评论功能暂时关闭!