在 Java 开发中,消息队列(MQ)保证消息不丢的核心思路是:覆盖 “生产端 -> 存储端 -> 消费端” 全链路,每个环节都要有确认、持久化或重试机制。


一、全链路拆解:哪里会丢消息?

要保证不丢,先搞清楚哪几个环节可能丢:

  1. 生产端(Producer):消息没发送到 MQ Broker(网络闪断、Broker 挂了)。
  2. 存储端(Broker):消息到了 Broker,但还没持久化就宕机了。
  3. 消费端(Consumer):消息到了 Consumer,但还没处理完就挂了(或程序报错),MQ 以为消费完了。


二、方案详解:RabbitMQ 如何保证不丢

1. 生产端:确认机制(Publisher Confirm & Return)

确保消息一定从 Producer 发送到了 Broker

  • Confirm 机制:消息到达 Broker 后,Broker 会给 Producer 一个确认(ACK)。
  • Return 机制:如果消息无法路由到队列(比如 Exchange 或 RoutingKey 错了),Broker 会把消息返回给 Producer。

Java 代码示例:

@Configuration
public class RabbitConfig {

    @Bean
    public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
        RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);

        // 1. 开启 Publisher Confirm 确认机制
        rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
            if (ack) {
                System.out.println("消息成功发送到 Broker");
            } else {
                System.out.println("消息发送失败: " + cause);
                // 这里可以做重试或存入数据库后续补偿
            }
        });

        // 2. 开启 Return 回退机制(消息无法路由时触发)
        rabbitTemplate.setReturnsCallback(returned -> {
            System.out.println("消息无法路由,被退回: " + returned.getMessage());
        });

        return rabbitTemplate;
    }
}

2. 存储端:持久化(Persistence)

确保消息在 Broker 里存盘了,即使 Broker 重启消息也不丢。需要同时开启三个持久化:

  1. Exchange 持久化:声明 Exchange 时设置 durable=true
  2. Queue 持久化:声明 Queue 时设置 durable=true
  3. Message 持久化:发送消息时设置 deliveryMode=2(持久化模式)。

Java 代码示例:

@Service
public class MqProducer {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    public void sendPersistentMessage() {
        // 1. 定义持久化的 Exchange (durable=true)
        DirectExchange exchange = ExchangeBuilder.directExchange("my.exchange").durable(true).build();
        
        // 2. 定义持久化的 Queue (durable=true)
        Queue queue = QueueBuilder.durable("my.queue").build();

        // 3. 发送持久化消息 (MessageProperties.PERSISTENT_TEXT_PLAIN)
        rabbitTemplate.convertAndSend("my.exchange", "my.routing.key", "Hello World", message -> {
            message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
            return message;
        });
    }
}

3. 消费端:手动确认(Manual ACK)

确保消息被 Consumer 真正处理完了,才告诉 Broker 删除消息。

  • 默认是 autoAck=true(自动确认):消息只要被 Consumer 接收到,就立即确认,不管业务代码有没有执行完。
  • 必须改为 manualAck(手动确认):业务代码执行成功,调用 basicAck;执行失败,调用 basicNack(重回队列或丢弃)。

Java 代码示例:

@RabbitListener(queues = "my.queue")
public void handleMessage(Message message, Channel channel) throws IOException {
    long deliveryTag = message.getMessageProperties().getDeliveryTag();
    try {
        // 1. 执行业务逻辑
        System.out.println("收到消息: " + new String(message.getBody()));
        
        // 2. 业务成功,手动确认 (multiple=false 表示只确认当前这一条)
        channel.basicAck(deliveryTag, false);
        
    } catch (Exception e) {
        // 3. 业务失败,处理异常
        // requeue=true: 消息重回队列头,重新消费 (可能会无限重试,需注意)
        // requeue=false: 消息丢弃或进入死信队列 (推荐配合死信队列使用)
        channel.basicNack(deliveryTag, false, false);
    }
}

三、方案详解:Kafka 如何保证不丢

Kafka 的设计思路和 RabbitMQ 略有不同,它通过 副本(Replication)偏移量(Offset) 来保证。

1. 生产端:ACK 配置

Kafka Producer 发送消息时,通过 acks 参数控制确认级别:

  • acks=0:Producer 发出去就不管了,极易丢消息(不推荐)。
  • acks=1:Leader Partition 收到消息并写入本地日志就确认(默认,可能丢)。
  • acks=all (或 -1):Leader 等待 ISR(In-Sync Replicas,同步副本列表)中所有 Follower 都同步成功才确认(最安全,不丢消息)。

Java 配置示例:

Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

// 核心配置:acks=all 保证不丢
props.put("acks", "all");
// 重试次数:如果发送失败,自动重试
props.put("retries", 3);

Producer<String, String> producer = new KafkaProducer<>(props);

2. 存储端:副本机制 & 持久化

  • 多副本(Replication Factor):Topic 创建时设置 replication-factor >= 2(通常设为 3),Leader 挂了,Follower 可以顶上。
  • ISR 机制:只有同步跟上的 Follower 才在 ISR 列表里,配合 acks=all 使用。
  • 持久化刷盘:Kafka 是顺序写磁盘,性能很高,虽然是异步刷盘,但因为有多副本,所以一般不担心单机宕机。

3. 消费端:手动提交 Offset

和 RabbitMQ 类似,Kafka 也需要关闭自动提交,改为手动提交消费位移(Offset)

  • enable.auto.commit=false:关闭自动提交。
  • 业务逻辑执行成功后,调用 commitSync()(同步提交)或 commitAsync()(异步提交)。

Java 代码示例:

Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "my-group");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");

// 1. 关闭自动提交 Offset
props.put("enable.auto.commit", "false");

Consumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("my-topic"));

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> record : records) {
        try {
            // 2. 执行业务逻辑
            System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value());
            
        } catch (Exception e) {
            // 处理异常,记录日志,不要提交 Offset,下次 poll 还会拉到这条消息
            continue;
        }
    }
    // 3. 业务处理完,手动同步提交 Offset (确保数据真的处理完了)
    consumer.commitSync();
}

四、总结

  1. 生产端要有确认回调(Confirm/ACK),发送失败要重试。
  2. 存储端要持久化 / 多副本(RabbitMQ 开 durable,Kafka 开 replication-factor)。
  3. 消费端要手动确认(Manual ACK / Commit Offset),业务没做完不确认。
Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐