









从生产端到消费端全链路保证 RabbitMQ/Kafka 消息不丢失的完整方案,核心是通过 “生产端确认、中间件持久化、消费端确认” 三大核心环节,堵住每个环节的丢失风险,我会分别梳理 RabbitMQ 和 Kafka 的全流程保障方案,同时标注两者的关键差异。
消息丢失的风险点主要在 3 个环节:生产端(发送失败)、中间件(宕机 / 重启)、消费端(消费异常未处理),全链路保障的核心是:
// Spring Boot 配置示例
spring:
rabbitmq:
publisher-confirm-type: correlated # 开启异步确认
publisher-returns: true # 开启消息路由失败返回(如队列不存在)
// 生产者代码示例(异步确认)
@Autowired
private RabbitTemplate rabbitTemplate;
@PostConstruct
public void init() {
// 确认回调:MQ收到消息触发
rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
if (!ack) {
// 未确认,重试发送(需做幂等,避免重复)
log.error("消息发送失败:{}", cause);
retrySend(correlationData.getId());
}
});
// 退回回调:消息路由失败触发(如队列不存在)
rabbitTemplate.setReturnsCallback(returned -> {
log.error("消息路由失败:{}", returned.getMessage());
retrySend(returned.getMessage().getMessageProperties().getCorrelationId());
});
}
mandatory=false(默认):否则路由失败的消息会直接丢弃,需设为true并通过 ReturnsCallback 处理。durable=true,否则 MQ 重启后队列 / 交换机消失,消息丢失;
// 队列声明示例(持久化)
@Bean
public Queue durableQueue() {
return QueueBuilder.durable("durable_queue") // 持久化队列
.build();
}
deliveryMode=2(持久化),否则消息只存在内存,MQ 宕机丢失;
// 发送消息时指定持久化
rabbitTemplate.convertAndSend("exchange", "routingKey", "msg", msg -> {
msg.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
return msg;
});
ha-mode=all(所有节点同步)、ha-sync-mode=automatic(自动同步)。autoAck=true,消费端刚收到消息就确认,若消费过程崩溃,消息丢失;需设为autoAck=false,消费完成后手动 ack。
// 消费端配置(手动确认)
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual # 手动确认
retry:
enabled: true # 开启消费重试
max-attempts: 3 # 最大重试次数
initial-interval: 1000ms # 重试间隔
// 消费端代码示例
@RabbitListener(queues = "durable_queue")
public void consume(Message message, Channel channel) throws IOException {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
// 执行业务逻辑
processMsg(message);
// 消费完成,手动确认(单条确认)
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
log.error("消费失败", e);
// 重试后仍失败,拒绝并死信(避免无限重试)
channel.basicNack(deliveryTag, false, false);
}
}
Kafka 的核心设计(分区副本)比 RabbitMQ 更侧重高可用,保障方案略有差异,但核心逻辑一致:
acks参数控制生产者收到的确认级别,是生产端防丢失的核心:
acks=0:生产者不等待确认,消息可能丢失(禁用);acks=1:只等 leader 节点确认,leader 宕机且副本未同步则丢失(不推荐);acks=-1/all:等待 leader + 所有 ISR 副本确认,最高可靠性(推荐)。# 生产者配置
bootstrap.servers=kafka1:9092,kafka2:9092,kafka3:9092
acks=all
retries=3 # 发送失败重试次数
retry.backoff.ms=1000 # 重试间隔
enable.idempotence=true # 开启幂等性,避免重试导致重复消息
// Kafka生产者代码示例
Producer<String, String> producer = new KafkaProducer<>(props);
ProducerRecord<String, String> record = new ProducerRecord<>("topic", "key", "msg");
producer.send(record, (metadata, exception) -> {
if (exception != null) {
log.error("发送失败", exception);
// 本地消息表重试(同RabbitMQ逻辑)
retrySend(record);
}
});
replication.factor≥3(副本数≥3),确保 leader 宕机后,follower 能成为新 leader,消息不丢失;
# 创建topic(3副本,1分区)
kafka-topics.sh --create --topic durable_topic --bootstrap-server kafka1:9092 --replication-factor 3 --partitions 1
false,避免非 ISR 副本成为 leader(可能丢失未同步的消息);log.flush.interval.messages=5000(累计 5000 条刷盘)、log.flush.interval.ms=1000(1 秒刷盘),可根据业务调整(刷盘越频繁,性能越低,可靠性越高);retention.ms(消息保留时间,默认 7 天),避免消息未消费就被删除。Kafka 通过 offset 记录消费位置,若 offset 提交早于消费完成,会导致消息丢失,核心方案:
# 消费者配置
enable.auto.commit=false # 关闭自动提交
auto.offset.reset=earliest # 消费失败时从最早offset开始(避免漏消息)
// 消费端代码示例(手动提交)
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList("durable_topic"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
try {
// 执行业务逻辑
processMsg(record);
} catch (Exception e) {
log.error("消费失败", e);
// 重试失败则记录,人工处理
continue;
}
}
// 批量提交offset(消费完成后提交)
consumer.commitSync();
}
consumer.commitSync()(同步提交),而非commitAsync()(异步可能提交失败);| 环节 | RabbitMQ 核心保障手段 | Kafka 核心保障手段 |
|---|---|---|
| 生产端 | Publisher Confirm + 本地消息表 | acks=all + 幂等性 + 重试 |
| 服务端 | 队列 / 交换机 / 消息持久化 + 镜像队列 | 副本数≥3 + 禁用 unclean leader 选举 |
| 消费端 | 手动 ack + 死信队列 | 手动提交 offset + 死信 topic |
全链路保证消息不丢失的核心要点:
额外建议:
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。