
消息队列在现代分布式系统里已经是标配组件了解耦、削峰、异步这三板斧几乎每个业务都能用上。但消息队列一旦引入可靠性问题就跟着来了如果一条消息在传输过程中丢了到底会出多大事往小了说日志消息丢几条可能没人察觉往大了说订单消息丢一条可能就导致用户付了钱但订单没创建或者支付成功的回调没送达最后出现资损。很多团队在引入消息队列时第一反应是“消费者里面拉一下数据不就行了”结果一上线就发现消息悄悄丢了排查起来又是半天。这篇文章我要讲清楚一个核心问题消息队列的“不丢失”不是mq单方面保证的而是生产端、Broker端、消费端三个环节一起配合才能实现。任何一个环节没有做好消息就可能静默丢失。读完这篇文章你会得到一套完整的排查思路和落地配置也能从容应对面试里关于“消息不丢失”“重复消费”“消息可靠性”这些高频问题。1. 消息为什么会丢三个环节断在哪消息从产生到被消费完整链路是这样的生产者 - Broker存储 - 消费者这条链路看起来简单但每一段都可能出问题。最常见的丢失场景有三种第一生产端发送失败。网络抖动、Broker 宕机、序列化异常都可能导致生产者发出的消息根本没有到达 Broker。如果生产者代码里不管发送结果直接 fire-and-forget这些消息就等于丢了。第二Broker 端存储丢失。消息虽然到了 Broker但 Broker 只是先放在内存里还没来得及落盘就宕机了或者副本还没有同步完成主节点就挂了。这种情况在 Kafka、RocketMQ、RabbitMQ 里都有对应的配置开关用默认配置实际上是很危险的。第三消费端处理失败。这是最容易忽略的一环。消费者从 Broker 拉取到消息后先自动提交 offset然后业务逻辑再处理。如果业务代码抛异常或者宕机了offset 已经提交这个消息就不会再被投递。看起来消息被消费了实际上业务没执行成功。所以要保证“消息不丢失”实际上要做三件事生产端确认消息真的发到了 Broker。Broker 端确认消息真的落了盘、同步到了副本。消费端确认消息真的处理成功才提交消费位点。任何一个环节没确认都不能说消息“刚好被处理了”。这篇文章会用 Kafka、RocketMQ、RabbitMQ 三类主流消息队列分别讲配置和代码也会单独讲一个轻量场景用 Redis Stream 做消息队列时怎么保证消息不丢。最后再聊重复消费和幂等设计以及面试里最常见的几种追问。2. 生产端如何确保消息真的发到了 Broker生产端的核心问题是** 一句send()调用到底是“发出去就行”还是“确认 Broker 接收了”**2.1 Kafka 生产端的 acks 配置很多刚接触 Kafka 的开发者容易踩的坑是生产者代码里调用了send()就默认消息到了 Kafka。实际上send()是异步的真正保证消息到达 Broker 要靠 acks 参数。Properties props new Properties(); props.put(bootstrap.servers, node1:9092,node2:9092,node3:9092); props.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(value.serializer, org.apache.kafka.common.serialization.StringSerializer); // 关键配置等待所有 ISR 副本确认 props.put(acks, all); // 发送失败后重试次数 props.put(retries, 3); // 每个连接最多未确认的请求数 props.put(max.in.flight.requests.per.connection, 1);acks 有三个取值面试容易问到acks 取值含义可靠性0生产者不等待 Broker 确认发完就认为成功最差可能丢消息1Leader 写入本地日志后即返回成功中等Leader 宕机时可能丢数据all或 -1所有 ISR 副本都确认后才返回成功最强但延迟更高如果消息对账敏感比如交易、订单、支付场景acksall是最低要求。retries控制发送失败后的重试配合max.in.flight.requests.per.connection1可以避免重试时消息乱序。2.2 同步发送与回调处理Kafka 的send()是异步的不能只看返回结果。最稳妥的做法是通过回调检查发送结果。import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; // 发送消息并检查回调 producer.send(new ProducerRecord(order_topic, orderId, orderJson), (metadata, exception) - { if (exception ! null) { // 发送失败记录日志并考虑重试或告警 log.error(消息发送失败topic{}, key{}, order_topic, orderId, exception); // 这里可以投递到本地重试表由定时任务补偿 } else { log.info(消息发送成功offset{}, partition{}, metadata.offset(), metadata.partition()); } });注意retries只能处理 Broker 临时不可用、网络抖动这类问题处理不了消息内容序列化异常、topic 不存在这类永久性错误。对于后者生产上常见做法是加一个本地消息表先把消息存到数据库再异步发送发送成功后才标记状态。这样即使 MQ 长时间不可用也能靠定时任务把消息捞回来重发而不是直接丢在内存里。2.3 RocketMQ 和 RabbitMQ 的确认机制RocketMQ 的send()支持同步发送并返回SendResult其中sendStatus为SEND_OK才表示真正发送成功。SendResult sendResult producer.send(message); if (sendResult.getSendStatus() SendStatus.SEND_OK) { // 发送成功 } else { // 发送失败需要重试或记录 }RabbitMQ 在 Spring Boot 中配置了 publisher-confirm 后也可以通过回调感知投递结果spring: rabbitmq: publisher-confirm-type: correlated publisher-returns: truerabbitTemplate.setConfirmCallback((correlationData, ack, cause) - { if (ack) { // Broker 确认收到 } else { log.error(消息投递失败{}, cause); } });到这里可以下一个阶段性结论生产端不丢失的基础是“发送后必须确认”。无论用哪种消息队列都不要只调用发送 API 就结束必须处理确认回执或者失败后的补偿逻辑。3. Broker 端消息落盘与多副本才是真安全消息到达 Broker 之后是不是就安全了未必。如果 Broker 只把数据放在内存里一旦宕机内存数据全部丢失。所以 Broker 端要解决两个问题持久化和副本同步。3.1 Kafka 的副本与 ISR 机制Kafka 的每个分区可以有多个副本其中一个是 Leader其余是 Follower。生产者只写 LeaderFollower 从 Leader 拉取数据。副本之间通过 ISRIn-Sync Replicas同步副本集合来跟踪哪些副本跟得上主节点的进度。生产配置里有两项缺一不可# Broker 端配置 server.properties # 分区副本数生产环境至少 2推荐 3 default.replication.factor3 # 最少同步副本数必须小于等于副本数 min.insync.replicas2这一个配置组合含义是一条消息写 Leader 后必须至少还有 1 个 ISR 副本同步完成Kafka 才会向生产者返回成功。如果 ISR 里可用副本数小于 2Broker 会拒绝写入宁可写不进去也不能写一条只存在 Leader 内存里的消息。注意min.insync.replicas2并不是解决所有问题的银弹它跟生产端的acksall是配套使用的。生产者设置acksallBroker 端设置min.insync.replicas2才能达到“Leader 至少一个 Follower 都确认”的效果。只有 Broker 端配置生产端用默认的 acks1可靠性是打折的。3.2 RocketMQ 的同步刷盘与同步复制RocketMQ 的关键设计是 CommitLog所有消息顺序写入 CommitLog再异步构建消费队列。Broker 端有两个刷盘策略和两种复制方式。# Broker 端配置 broker.conf # 刷盘方式ASYNC_FLUSH 异步刷盘SYNC_FLUSH 同步刷盘 flushDiskTypeSYNC_FLUSH # 复制方式ASYNC_MASTER 异步复制SYNC_MASTER 同步复制 brokerRoleSYNC_MASTER同步刷盘SYNC_FLUSH消息写入内存后立即刷到磁盘返回成功。性能略低但不会因为断电丢消息。异步刷盘ASYNC_FLUSH消息写入内存即返回成功后台定时刷盘。性能高但断电可能丢最近一小段时间的消息。同步复制SYNC_MASTERMaster 和 Slave 都写入成功才返回生产者成功。异步复制ASYNC_MASTERMaster 写入成功就返回再异步同步到 Slave。如果消息重要建议SYNC_FLUSH和SYNC_MASTER组合但性能开销最大。生产中常见做法是核心交易消息用同步刷盘日志、埋点类消息用异步刷盘把成本花在刀刃上。3.3 RabbitMQ 的持久化与镜像队列RabbitMQ 的消息可靠性要靠交换机、队列、消息三者的持久化配合。只把消息标记为持久化不够队列和交换机也得持久化否则 Broker 重启后队列都没了消息自然也丢了。// 创建持久化队列 MapString, Object args new HashMap(); // 队列持久化存储在磁盘 args.put(x-ha-policy, all); // 镜像队列所有节点同步 channel.queueDeclare(order.queue, true, false, false, args); // 发送持久化消息 AMQP.BasicProperties properties new AMQP.BasicProperties.Builder() .deliveryMode(2) // 2 表示持久化 .build(); channel.basicPublish(order.exchange, order.routing.key, properties, messageBody);RabbitMQ 3.8 之后推荐使用 Quorum Queue仲裁队列替代镜像队列。仲裁队列基于 Raft 协议实现天然支持多副本比镜像队列更可靠也是 RabbitMQ 官方在 4.x 版本中的推荐方案。到这里可以形成一个判断Broker 端的可靠性本质上是“存储写入策略 多副本确认策略”的组合。不要简单相信消息队列“落盘了就是安全的”必须确认它是落在哪个节点、有没有同步给其他副本。4. 消费端手动提交 offset 才是真正的“消费成功”消费端是最容易丢消息的环节而丢失原因往往是开发者主动造成的——大多数消息队列默认“自动确认”也就是消费端一拉到消息就提交位点不管业务有没有处理成功。4.1 Kafka 消费者手动提交 offsetKafka 默认enable.auto.committrue每 5 秒自动提交一次 offset。这个设计在消息重要时很危险拉到消息后业务代码还没执行完进程宕机offset 已提交重启后认为这条消息已经消费过了消息就永久丢了。正确做法是关闭自动提交改为手动提交并确保业务逻辑成功后再提交。import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import java.time.Duration; import java.util.Collections; import java.util.Properties; Properties props new Properties(); props.put(bootstrap.servers, node1:9092,node2:9092,node3:9092); props.put(group.id, order-consumer-group); props.put(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(value.deserializer, org.apache.kafka.common.serialization.StringDeserializer); // 关键配置关闭自动提交 props.put(enable.auto.commit, false); KafkaConsumerString, String consumer new KafkaConsumer(props); consumer.subscribe(Collections.singletonList(order_topic)); try { while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { try { // 真正的业务逻辑入库、调用第三方、处理业务状态 processOrder(record.value()); } catch (Exception e) { // 业务处理失败记录失败消息不要提交 offset log.error(业务处理失败等待重试或人工补偿{}, record.value(), e); // 根据业务决定阻塞重试、进入重试队列或 DLT return; } } // 一批消息全部处理成功才提交 offset consumer.commitSync(); } } finally { consumer.close(); }4.2 消费失败后的处理策略手动提交 offset 只是第一步关键的难点是“业务处理失败后接下来怎么办”直接抛异常停止消费。这样消息不会被提交消费者以一定策略重置 offset 后重新拉取但如果消息本身是“脏数据”或者第三方服务持续不可用会造成消费阻塞。重试 N 次后进入死信队列DLQ。Kafka 里可以定义一个order_topic_dlt重试多次仍失败的消息发到 DLT 单独保存主流程继续消费后续消息。这样可以防止一条坏消息卡死整个队列。本地记录 定时补偿。业务处理失败后把消息体持久化到数据库定时任务扫出来重新处理。这种方式可控性最强适合核心链路。消费端不丢失的核心原则可以总结成一句提交位点要能反映“业务处理成功”而不是“消息已经拉取”。5. 完整案例一个“不丢消息”的订单系统示例前面把生产端、Broker、消费端分别讲了一遍这里用一个完整的订单支付回调场景串起来展示每一步到底怎么写。场景描述用户在商城支付成功后支付平台回调业务后端后端解析回调结果后要把订单状态更新为“已支付”并发送一条消息通知积分服务增加积分。5.1 生产端支付回调消息发送用 Kafka 作为消息中间件生产端增加回调监听和重试机制。// 文件路径src/main/java/com/example/order/mq/OrderEventPublisher.java Component public class OrderEventPublisher { private static final Logger log LoggerFactory.getLogger(OrderEventPublisher.class); private final KafkaTemplateString, String kafkaTemplate; public OrderEventPublisher(KafkaTemplateString, String kafkaTemplate) { this.kafkaTemplate kafkaTemplate; } /** * 发送订单支付成功消息 */ public void publishOrderPaidEvent(OrderPaidEvent event) { String message JSON.toJSONString(event); // 用订单号作为 key保证同一个订单的消息进入同一个分区顺序有保障 kafkaTemplate.send(order_paid_topic, event.getOrderId(), message) .addCallback(result - { log.info(订单支付消息发送成功orderId{}, offset{}, event.getOrderId(), result.getRecordMetadata().offset()); }, ex - { log.error(订单支付消息发送失败需要进入补偿表orderId{}, event.getOrderId(), ex); // 实际项目中这里应该把消息保存到本地消息表由定时任务补偿重发 saveToLocalMessageTable(event); }); } private void saveToLocalMessageTable(OrderPaidEvent event) { // 把消息插入本地消息表状态为 NEW由定时任务扫描重发 } }application.yml中 Kafka 生产者配置spring: kafka: bootstrap-servers: node1:9092,node2:9092,node3:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer # 关键配置 acks: all retries: 3 properties: max.in.flight.requests.per.connection: 15.2 Broker 端Kafka 集群配置Broker 的server.properties核心配置# 节点id broker.id1 # 分区副本数 default.replication.factor3 # 最少同步副本数 min.insync.replicas2 # 允许 topic 被自动创建时的副本数配置 offsets.topic.replication.factor3 transaction.state.log.replication.factor3 transaction.state.log.min.isr2创建 topic 时显式指定副本数kafka-topics.sh --bootstrap-server node1:9092 --create \ --topic order_paid_topic \ --partitions 6 \ --replication-factor 3 \ --config min.insync.replicas25.3 消费端手动提交 offset 失败重试// 文件路径src/main/java/com/example/order/mq/OrderPaidConsumer.java Component public class OrderPaidConsumer { private static final Logger log LoggerFactory.getLogger(OrderPaidConsumer.class); KafkaListener(topics order_paid_topic, groupId order-paid-group) public void onMessage(ConsumerRecordString, String record, Acknowledgment ack) { try { OrderPaidEvent event JSON.parseObject(record.value(), OrderPaidEvent.class); // 更新订单状态 orderService.markOrderPaid(event.getOrderId(), event.getPaidAmount()); // 增加积分 pointsService.addPoints(event.getUserId(), event.getPointAmount()); // 业务处理成功后手动提交 offset ack.acknowledge(); } catch (Exception e) { // 记录错误根据需求选择重试、死信队列、人工补偿 log.error(消费订单支付消息失败key{}, value{}, record.key(), record.value(), e); // 不调用 ack.acknowledge()会触发消费者重试 throw new RuntimeException(处理失败触发重试, e); } } }Spring Kafka 的Acknowledgment模式需要配套配置spring: kafka: consumer: group-id: order-paid-group key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer enable-auto-commit: false auto-offset-reset: earliest listener: type: batch ack-mode: manual这里的关键点是把enable-auto-commit设为 false并把 listener 的 ack-mode 改为 manual这样只有在acknowledge()被调用时offset 才会推进。业务成功才确认这是消费端不丢失的最终防线。5.4 运行与验证启动消费端后观察日志# 生产一条测试消息 kafka-console-producer.sh --bootstrap-server node1:9092 --topic order_paid_topic {orderId:2024001,userId:10001,paidAmount:99.5,pointAmount:99} # 消费端预期输出 订单支付消息消费成功orderId2024001, offset3如果想验证“消费失败不提交 offset”的效果可以在processOrder中临时抛一个异常让消费者反复重试。观察到行为是消息会重新投递offset 始终不推进。如果消费者端设置为抛错立即停止也可以看到消费者不再拉取新消息。这些表现都能证明手动提交机制生效了。6. 用 Redis Stream 实现不丢消息的轻量队列并不是所有场景都值得引入 Kafka 或 RocketMQ。中小项目如果本来就有 Redis用 Redis Stream 做轻量消息队列是合理选择。Redis Stream 也有自己的“不丢消息”机制只是和 Kafka 的思路不同需要单独掌握。Redis Stream 是 Redis 5.0 引入的数据结构支持消息持久化、消费者组、消费确认ACK。从消息可靠性角度看它有几个关键点Producer 发送使用XADD消息会持久化到 Redis 内存和 RDB/AOF 中。Consumer 读取使用XREADGROUP消息会进入 Pending 列表只有调用XACK后才会从 Pending 中移除。如果消费者崩溃没有XACK的消息会一直留在 Pending 中可以用XPENDING查看用XAUTOCLAIM将超时消息转移给其他消费者。6.1 Spring Boot 中拉取 Redis Stream 消息// 文件路径src/main/java/com/example/order/mq/RedisStreamConsumer.java Component public class RedisStreamConsumer { private static final String STREAM_KEY order:paid:stream; private static final String GROUP_NAME order-paid-group; private static final String CONSUMER_NAME consumer-1; private final StringRedisTemplate redisTemplate; public RedisStreamConsumer(StringRedisTemplate redisTemplate) { this.redisTemplate redisTemplate; } public void pullMessage() { // 创建消费者组如果已存在会报错忽略即可 try { redisTemplate.opsForStream().createGroup(STREAM_KEY, GROUP_NAME); } catch (Exception e) { // group already exists } // 读取消息最多读取 10 条阻塞 5 秒 ListMapRecordString, Object, Object records redisTemplate.opsForStream() .read(Consumer.from(GROUP_NAME, CONSUMER_NAME), StreamReadOptions.empty().count(10).block(Duration.ofSeconds(5)), StreamOffset.create(STREAM_KEY, ReadOffset.lastConsumed())); for (MapRecordString, Object, Object record : records) { try { String orderId String.valueOf(record.getValue().get(orderId)); // 处理业务 processOrder(orderId); // 处理成功后确认消息才从 Pending 列表移除 redisTemplate.opsForStream().acknowledge(STREAM_KEY, GROUP_NAME, record.getId()); } catch (Exception e) { log.error(消费 Redis Stream 消息失败orderId{}, record.getValue(), e); // 不 ack消息保留在 Pending后续使用 XAUTOCLAIM 重新消费 } } } Scheduled(fixedDelay 10000) public void schedulePull() { pullMessage(); } }6.2 异常消息的重新处理如果消费者崩溃消息留在 Pending 中需要一个回收机制。Redis 6.2 之后推荐XAUTOCLAIM# 查看 Pending 消息数量 XPENDING order:paid:stream order-paid-group # 将挂起超过 5 分钟的消息转移给 consumer-2 XAUTOCLAIM order:paid:stream order-paid-group consumer-2 300000 0Redis Stream 相较 Kafka、RocketMQ 最大的优势是部署简单不需要额外维护集群劣势是没有原生分区、副本能力不如专业 MQ不适合超大数据量和严格顺序要求的场景。用它做“不丢消息”的轻量队列需要开发者自行处理好XACK和异常回收这两个环节少一个都会造成消息丢失或积压。7. 重复消费与幂等设计不丢失的另一面消息不丢往往伴随另一个问题重复消费。生产者发送时可能重试消费者处理成功后还没来得及提交 offset 就宕机了重启后同一批消息会被再次投递。所以分布式环境下“不丢失”和“不重复”是必须同时考虑的一对问题。解决重复消费的核心思路是幂等。所谓幂等就是同一个操作执行一次和执行多次结果完全一致。7.1 方案一业务去重表在业务数据库中建一张去重表把消息的唯一业务 ID 作为主键。消费者处理消息前先尝试插入如果插入成功说明第一次处理继续执行如果插入失败说明重复消息直接跳过。CREATE TABLE msg_consume_log ( msg_id VARCHAR(64) NOT NULL COMMENT 消息唯一ID例如订单号, biz_type VARCHAR(32) NOT NULL COMMENT 业务类型, create_time DATETIME NOT NULL COMMENT 首次消费时间, PRIMARY KEY (msg_id, biz_type) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;// 处理前先插入去重表 boolean inserted msgConsumeLogDao.insertIfAbsent(event.getOrderId(), order_paid); if (!inserted) { // 重复消息直接 ack不处理业务 ack.acknowledge(); return; } try { // 真正处理业务 doBiz(event); ack.acknowledge(); } catch (Exception e) { // 注意如果业务失败去重表里的记录要回滚否则消息不会被重试 throw e; }注意去重表应该和业务操作在同一个数据库事务里否则可能出现“去重记录插入成功业务操作失败”的中间态。7.2 方案二Redis SETNX 实现短时间幂等对幂等窗口要求不高的场景可以用 Redis 的 SETNX 配合过期时间String redisKey order:paid:idempotent: event.getOrderId(); // 5 分钟内重复消息直接丢弃 Boolean success redisTemplate.opsForValue().setIfAbsent(redisKey, 1, Duration.ofMinutes(5)); if (Boolean.TRUE.equals(success)) { // 第一次消费执行业务 doBiz(event); } else { // 重复消息跳过 log.warn(重复消息直接跳过{}, event.getOrderId()); }7.3 方案三数据库乐观锁 / 状态机控制如果业务本身就是更新订单状态可以通过乐观锁控制UPDATE order SET status PAID, paid_time NOW() WHERE order_id #{orderId} AND status UNPAID;如果影响行数为 0说明订单已经不是待支付状态说明这条消息是重复消息或者业务状态已经流转到更后面的阶段。幂等设计没有银弹核心要结合业务选择。最简单的判断纬度是这个业务操作天然支持重复执行吗不支持就要做幂等去重表是最常见的选择。8. 常见问题与排查方法问题现象可能原因排查方式解决方案消费者没收到消息生产端发送失败或发送成功但 topic 写错查看生产者日志是否有异常回调检查 callback 是否捕获异常补发消息消息丢失但日志没有任何异常Broker 端默认配置不持久化或副本不足检查 Broker 端 flush 和 replication 配置acksallmin.insync.replicas2消费端业务没执行offset 却推进了自动提交 offset 导致查看enable.auto.commit是否为 true改为手动提交业务成功后 ACK消息重复消费消费成功后未来得及提交 offset或重试机制触发msg_consume_log 查看同一 msg_id 多条记录引入幂等设计最推荐数据库去重表消费阻塞一条消息反复报错业务处理失败但消费循环未跳过查看异常堆栈和重试策略失败 N 次后发到死信队列不阻塞主流程Redis Stream 消息积压在 Pending消费者崩溃或未调用 XACKXPENDING查看 Pending 数量编写XAUTOCLAIM回收策略9. 最佳实践与工程建议消息可靠性是个系统工程单一配置解决不了所有问题。下面这些建议是我在项目里比较常用、也见效比较快的一组实践。第一核心消息必须三端全查。生产端看发送回执Broker 端看复制配置消费端看提交逻辑。每端都要写清楚确认方式不要靠猜。特别是消费端自动提交建议全局关掉。第二给消费端设计死信队列。消费失败不能无限重试。一般策略是重试 3 次后把消息转存到死信队列并触发告警通知。主链路照常跑人工介入处理死信队列里的异常消息。第三消息体要带唯一业务 ID。这个 ID 最好是业务本身的单号比如订单号、用户 ID 时间戳。没有唯一 ID后面做幂等、去重、对账都很麻烦。第四生产端要有本地消息表兜底。分布式系统里消息队列短暂不可用是常态。核心消息发送前先写本地表状态为待发送发送成功后更新状态定时任务扫描待发送且超时的记录进行补发。这是保证“最终一致性”非常实用的一招。第五压测验证不丢消息。没有压测不要轻易上线。可以用 JMeter 或者自研脚本压一批消息消费端消费完成后对比生产总数和消费总数。如果对不上优先查消费端有没有漏处理再看 Kafka 的 offset 有没有提交异常。第六监控报警必须覆盖消费滞后和重试次数。消息队列不是“发完就不管了”。消费者线程挂了、消费变慢了都不能靠人肉发现。建议对消费组 lag、消费异常次数、死信队列消息数都配置监控。10. 面试题高频追问把可靠性的底层逻辑答清楚“消息队列如何保证消息不丢失”是面试高频题但很多候选人的答案止步于“acksall”和“手动 ack”这两个关键词。能拉开差距的是下面几个追问。Q1acksall 一定能保证不丢吗不能。acksall要求所有 ISR 副本都返回确认但如果 ISR 只有一个 Leader消息实际上只存在单副本上。所以生产端用acksall时Broker 端必须配上min.insync.replicas2才能保证至少两个副本同步成功。Q2消费端手动提交 offset 后消息还会重复吗会。手动提交只能保证“业务成功后提交”但无法保证“提交前宕机”。如果业务执行成功、offset 还没提交进程就挂了重启后消息会重新投递这时候就要靠幂等兜底。Q3Kafka、RocketMQ、RabbitMQ 三个组件里哪个丢消息概率最低没有绝对的结论。三个组件在不同配置体系下表现不同Kafka 和 RocketMQ 适合大数据量、高吞吐的场景RabbitMQ 在轻量级、灵活路由的场景更有优势。丢不丢消息取决于你对每个组件做了多少可靠性配置不取决于组件名字本身。Q4如果消息发了但没收到回执怎么判断到底发没发成功这是生产中很难处理的场景之一。建议的做法是生产端消息带唯一 ID发送时写入本地表消费者处理完把消息 ID 写回或进行对账。通过周期性对账可以发现那些“不知道有没有发成功”的消息这是分布式系统里的最终一致性思路。如果你在准备面试建议把本文的生产端链路、Broker 副本 ISR、消费端手动 ACK、幂等设计这四条主线用自己的话复述一遍并配合一个订单场景说明整体流程基本就能覆盖这一题的考察点。11. 总结与后续学习方向消息不丢失不是某个中间件的单一特性而是生产端、Broker、消费端三个环节一起构成的能力。生产端用 acksall 发送回执确认消息到达Broker 端用多副本 同步刷盘保证存储可靠消费端用手动 ack 幂等设计保证业务真正被处理。三个环节环环相扣任何一环的默认配置都不能作为生产环境的“安全假设”。下一步可以从这几个方向继续深入用 Kafka 自带工具排查消费组 lag理解“积压”和“丢失”的关系。在真实项目中实现本地消息表 定时补偿跑一遍生产者发送失败的流程。为消费者增加死信队列把你现在的消费失败从“卡死主流程”变成“隔离到 DLQ”。学习 RocketMQ 的事务消息、Kafka 的 Exactly-Once 语义这些是消息可靠性话题的进阶内容。最后提醒一句这套可靠性机制的核心是权衡。你付出的配置和代码越多消息越不容易丢但系统吞吐和运维复杂度也会上升。建议先梳理清楚自己业务里的核心链路把重点保障放在会出资损、影响用户关键流程的消息上普通日志消息则不必过度设计。这样既守住了底线也不至于让团队在技术上过度建设。