消息队列防丢全链路方案
·
在 Java 开发中,消息队列(MQ)保证消息不丢的核心思路是:覆盖 “生产端 -> 存储端 -> 消费端” 全链路,每个环节都要有确认、持久化或重试机制。
一、全链路拆解:哪里会丢消息?
要保证不丢,先搞清楚哪几个环节可能丢:
- 生产端(Producer):消息没发送到 MQ Broker(网络闪断、Broker 挂了)。
- 存储端(Broker):消息到了 Broker,但还没持久化就宕机了。
- 消费端(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 重启消息也不丢。需要同时开启三个持久化:
- Exchange 持久化:声明 Exchange 时设置
durable=true。 - Queue 持久化:声明 Queue 时设置
durable=true。 - 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();
}
四、总结
- 生产端:要有确认回调(Confirm/ACK),发送失败要重试。
- 存储端:要持久化 / 多副本(RabbitMQ 开 durable,Kafka 开 replication-factor)。
- 消费端:要手动确认(Manual ACK / Commit Offset),业务没做完不确认。
更多推荐
所有评论(0)