
一、开场先把两个问题拆开看面试官RocketMQ 如何保证消息不丢失又如何保证消息不被重复消费这两个问题看起来是独立的实际上背后考察的是你对 RocketMQ 消息链路的整体理解。很多人一上来就背“同步刷盘、同步复制、消费重试、幂等处理”但没有讲清楚它们分别解决的是哪一个环节的问题也没有讲清楚 RocketMQ 为什么无法天然做到“既不丢失又不重复”。这篇文章不会只给你一张面试背诵清单而是沿着一条消息从生产者发出、经过 Broker 存储、再到消费者消费的完整链路把消息丢失和重复消费的根因、配置、代码实践、事务消息机制、面试回答思路一次性讲透。先记住一个核心结论消息不丢失需要同时保证生产端可靠发送、Broker 端可靠存储、消费端可靠消费三个环节。消息不重复消费本质不是靠 RocketMQ 直接提供精确一次语义而是靠业务侧做好幂等控制。RocketMQ 默认提供的是“至少一次投递”网络超时、重试、故障恢复都可能带来重复消息。所以“消息不丢失”和“消息不重复”通常要做成一套组合方案可靠投递 消费幂等。先理解这一点面试时你的回答会更像真正做过生产实践的人而不是只会背配置。二、RocketMQ 消息链路回顾在展开可靠性方案之前有必要先回顾 RocketMQ 的核心组件和消息流转路径。只有理解链路才能知道消息在哪些节点、哪些时刻可能丢失。2.1 核心角色NameServer路由注册中心维护 Topic 和 Broker 的映射关系。NameServer 本身不参与消息落盘只告诉生产者和消费者应该连接哪些 Broker。Broker消息存储和服务节点负责接收生产者发来的消息写入 CommitLog再构建 ConsumeQueue 等索引最后投递给消费者。Producer消息生产者将消息发送到 Broker 的指定 Topic 和 MessageQueue。Consumer消息消费者从 Broker 拉取消息并进行业务处理。Topic消息主题逻辑概念。一个 Topic 可以包含多个 MessageQueue用于并行和负载均衡。MessageQueue一个 Topic 下的队列可以理解为一个分区。如何选择 MessageQueue 会影响顺序性。2.2 一条消息的完整旅程Producer 向 NameServer 查询 Topic 的 Broker 地址和 MessageQueue 信息。Producer 按照负载均衡策略选择一个 MessageQueue。Producer 将消息发送到对应 Broker。Broker 将消息顺序写入 CommitLog。Broker 异步构建 ConsumeQueue 和索引文件。Consumer 向 Broker 拉取消息处理完成后返回消费结果。Consumer 维护消费位点表示当前这条消息已经被处理。消息链路中容易出现问题的三个环节分别是生产发送环节、Broker 存储环节、消费处理环节。下面分别讨论。三、消息不丢失的总纲RocketMQ 的消息可靠性设计可以从三个角度来记忆环节可能丢失的原因对应解决手段生产端网络超时、Broker 不可用、发送结果未确认、异步发送回调丢失同步发送、失败重试、失败告警落库、事务消息Broker 存储端进程崩溃后内存数据未落盘、主节点宕机后数据未同步到从节点同步刷盘、主从同步复制、DLedger 多副本提交消费端业务处理未完成就确认消费、消费异常后位点错误前进处理成功后再确认消费、消费重试、死信队列兜底这三层必须同时保障否则任意一个环节出问题都可能导致消息最终丢失。很多人只关注 Broker 的刷盘和复制但生产端发送失败没有处理或者消费端提前 ACK同样会造成消息丢失。四、生产端如何保证消息不丢失4.1 同步发送并确认结果RocketMQ 提供了三种发送方式同步发送、异步发送、单向发送。其中单向发送不等结果性能最高但风险最大通常只适合日志采集等允许少量丢失的场景。异步发送通过回调接收结果性能较好但代码相对复杂需要正确处理回调异常和失败重试。同步发送等待发送结果返回只有 Broker 明确返回成功才认为消息已经送达。核心业务建议使用同步发送。对于核心链路比如订单支付、库存扣减等必须使用同步发送并判断返回的发送状态。只有返回SEND_OK才能认为消息被 Broker 接收。4.2 发送失败自动重试网络是有可能抖动的。如果第一次发送失败不能直接放弃应该让 Producer 自动重试。RocketMQ 生产者提供了重试次数配置retryTimesWhenSendFailed同步发送失败后的重试次数。retryTimesWhenSendAsyncFailed异步发送失败后的重试次数。retryAnotherBrokerWhenNotStoreOK当 Broker 存储未成功时是否尝试发送到其他 Broker。需要特别注意重试可能造成消息重复。比如第一次发送已经成功写入 Broker但网络超时导致客户端没有收到成功响应客户端发起第二次发送就会产生两条内容相同的消息。这也是后续必须做幂等控制的原因之一。4.3 核心业务发送代码示例DefaultMQProducer producer new DefaultMQProducer(lossless-producer-group); producer.setNamesrvAddr(127.0.0.1:9876); producer.setRetryTimesWhenSendFailed(3); producer.setRetryTimesWhenSendAsyncFailed(3); producer.setRetryAnotherBrokerWhenNotStoreOK(true); producer.start(); Message msg new Message(OrderTopic, OrderPaid, order-1001, 订单已支付.getBytes(StandardCharsets.UTF_8)); SendResult result producer.send(msg, new MessageQueueSelector() { Override public MessageQueue select(ListMessageQueue mqs, Message msg, Object arg) { Long orderId (Long) arg; int index (int) (orderId % mqs.size()); return mqs.get(index); } }, 1001L); if (!SendStatus.SEND_OK.equals(result.getSendStatus())) { // 写本地待补偿表或直接告警并人工介入 }这里选择 MessageQueue 时使用订单号取模目的是把同一订单的消息尽量路由到同一个队列为后续可能需要的顺序消费创造条件。4.4 发送后仍需做补偿自动重试并不能覆盖所有失败场景。例如 Broker 长时间不可用、NameServer 路由异常、客户端进程崩溃等。为了让消息最终不丢失生产侧还需要有一套本地消息表或可靠事件表机制业务操作和本地消息记录放在同一个数据库事务中。业务事务提交后通过定时任务扫描未成功发送的消息。对未发送成功的消息持续投递到 RocketMQ。收到发送成功结果后更新本地事件状态。这种方案本质是通过“先落地本地事务再异步可靠投递”来弥补单纯网络重试的不足。五、Broker 端如何保证消息不丢失当消息到达 Broker 后风险就转到了存储侧。RocketMQ 的丢失风险主要集中在两个地方刷盘策略消息是先写内存缓存还是必须写磁盘后才返回成功。主从复制策略主节点写完消息后是否需要等待从节点同步完成。5.1 同步刷盘与异步刷盘Broker 收到消息后通常先写入 PageCache再根据刷盘策略决定何时写入磁盘。刷盘方式有两种方式行为优点缺点异步刷盘消息写入 PageCache 后即返回成功后台线程异步写磁盘吞吐量高Broker 宕机时尚未落盘的消息可能丢失同步刷盘每条消息都必须写入磁盘后才返回成功数据可靠性高吞吐量相对较低在消息不允许丢失的核心业务中应该配置同步刷盘。虽然同步刷盘会带来一定性能损耗但这是数据可靠性的基础底线。5.2 主从同步复制与异步复制单台 Broker 做同步刷盘仍然面临机器宕机、磁盘损坏、机房故障等风险。因此 RocketMQ 支持主从模式主 Broker 接收读写请求从 Broker 同步数据主节点故障时可以切换。主从复制也有两种模式异步复制主节点写成功后即可返回从节点后台异步复制。主节点宕机但数据未同步到从节点时消息可能丢失。同步复制主节点必须等待从节点确认同步完成后才向生产者返回成功。可靠性更强。真正要求消息不丢失需要同时开启同步刷盘和同步复制。Broker 角色推荐配置为SYNC_MASTER。# 同步刷盘消息写入磁盘后才返回成功 flushDiskTypeSYNC_FLUSH 同步复制主节点需要等待从节点确认同步 brokerRoleSYNC_MASTER5.3 DLedger 多副本机制传统主从模式在主从切换时仍然可能存在选主延迟、数据一致性窗口等问题。RocketMQ 引入 DLedger 后可以将多个 Broker 组成 Raft 协议的多副本集群。消息只有被多数派节点确认后才认为提交成功。这样即使单节点或少数节点故障已经确认提交的消息也不会丢失。在高要求的生产环境中DLedger 模式比单纯主从同步复制更有保障。面试中如果能主动提到 DLedger 和多数派提交会是一个明显的加分点。六、消费端如何保证消息不丢失6.1 消费位点的含义Consumer 处理完一条消息后会记录消费位点。下次重启或重新负载均衡时从上次记录的位点之后继续消费。如果业务还没处理成功位点却已经前进了后续就不会再拿到这条消息于是表现为消息丢失。因此消费端最关键的原则是必须在业务逻辑真正处理成功之后才返回消费成功才允许更新消费位点。在 RocketMQ 的DefaultMQPushConsumer中使用MessageListenerConcurrently并发消费时通过返回状态决定是否确认消费返回CONSUME_SUCCESS表示消费成功Broker 可以推进位点。返回RECONSUME_LATER表示消费失败RocketMQ 会按策略重新投递这条消息。6.2 消费端示例DefaultMQPushConsumer consumer new DefaultMQPushConsumer(lossless-consumer-group); consumer.setNamesrvAddr(127.0.0.1:9876); consumer.subscribe(OrderTopic, OrderPaid); consumer.registerMessageListener(new MessageListenerConcurrently() { Override public ConsumeConcurrentlyStatus consumeMessage(ListMessageExt msgs, ConsumeConcurrentlyContext context) { try { for (MessageExt msg : msgs) { String body new String(msg.getBody(), StandardCharsets.UTF_8); orderService.payOrder(body); } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } catch (Exception e) { log.error(消息消费失败, e); return ConsumeConcurrentlyStatus.RECONSUME_LATER; } } }); consumer.start();上面的代码中orderService.payOrder(body)只有真正执行业务成功后才返回随后方法才返回CONSUME_SUCCESS。任何异常都会返回RECONSUME_LATER触发重试。6.3 重试与死信队列RocketMQ 提供消费重试机制。消费失败的消息会进入重试队列经过若干次重试后如果仍然失败会被转入死信队列避免一直占用重试资源。死信队列不是“垃圾桶”而是一个兜底机制。运维或业务系统可以订阅死信队列对最终消费失败的消息进行告警、人工处理或补偿。消费端如果不设置死信队列策略异常消息反复重试反而会拖垮系统。七、为什么消息会被重复消费很多面试者会误以为“消息不丢失”做好之后重复消费就不会发生。其实不然。RocketMQ 无法在分布式环境下天然做到完全不重复原因是它默认提供的是至少一次投递语义。重复消费的常见来源包括生产端重试Broker 已经收到消息但生产者因为网络超时没有收到确认于是再次发送。消费端位点回退消费者处理完成后位点提交失败下一次从旧位点开始拉取导致同一批消息被再次消费。网络分区或故障恢复Broker 故障恢复后部分消息可能被重复投递。消费超时同一消费组内发生负载均衡未处理完的消息可能重新分配给其他消费者。所以重复消费不是单独去关闭重试就可以解决的问题。关闭重试反而可能导致消息丢失。正确做法是允许 RocketMQ 重试以保证不丢失同时在业务侧做幂等以保证重复不产生副作用。八、如何保证消息不被重复消费业务幂等8.1 幂等的本质幂等的意思是同一个请求执行一次和执行多次产生的业务结果相同。对于消息消费来说就是要保证同一条消息被处理多次时不会导致订单重复支付、库存重复扣减、积分重复发放等问题。实现幂等需要业务先定义一个唯一业务键然后围绕这个唯一键做去重控制。常见的唯一键可以来自消息 ID、订单号、流水号、业务主键等。8.2 基于数据库唯一约束在数据库中创建一张消费记录表并通过唯一索引去重。消费前先尝试插入一条记录如果唯一键冲突说明消息已经处理过直接忽略。CREATE TABLE order_consume_log ( id BIGINT PRIMARY KEY AUTO_INCREMENT, msg_id VARCHAR(64) NOT NULL, order_id BIGINT NOT NULL, consume_status TINYINT NOT NULL DEFAULT 0, create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE KEY uk_msg_id_order_id (msg_id, order_id) );Transactional public void handlePayMessage(String msgId, Long orderId) { int inserted consumeLogMapper.insert(msgId, orderId, 0); if (inserted ! 1) { return; // 唯一键冲突说明已经处理过 } orderService.pay(orderId); consumeLogMapper.updateStatus(msgId, orderId, 1); }8.3 基于业务状态机实现幂等数据库唯一约束适合有明确唯一键的场景。如果业务本身已经有状态流转也可以直接借助状态机的条件更新来实现幂等不必额外维护一张消费记录表。核心思路是更新时带上前置状态作为条件只允许符合条件的数据发生状态变化并通过受影响行数判断本次操作是否真正生效。UPDATE t_order SET status PAID, paid_time NOW() WHERE order_id #{orderId} AND status PENDING;如果返回受影响行数为 1说明这是第一次把订单从“待支付”推进到“已支付”可以继续执行后续逻辑如果返回 0说明订单已经不是待支付状态直接跳过即可。状态机幂等的优点是它和业务表天然绑定不用额外建表也不容易漏掉去重逻辑缺点是它只适合状态清晰、可定义前置条件的业务。8.4 基于 Redis 去重的幂等方案对于一些高并发、允许最终一致的场景也可以使用 Redis 的SET NX EX特性做幂等控制。核心思想是在消息处理前把唯一业务键写入 Redis如果写入成功就继续处理如果写入失败就认为这条消息已经被处理过。String key consume:lock: orderId : msgId; Boolean first redisTemplate.opsForValue() .setIfAbsent(key, 1, Duration.ofMinutes(30)); if (Boolean.FALSE.equals(first)) { return; // 已经被处理忽略重复消息 } try { orderService.pay(orderId); } finally { // 不要无脑删除 key否则并发场景下可能导致重复处理 }使用 Redis 方案时要特别注意两点必须设置过期时间如果没有过期时间一旦业务执行到一半进程崩溃key 会长期存在后续消息会被错误地拒掉。不要盲目删除 key如果提前删除 key其他并发消费者可能又在处理中写入成功导致重复执行。合理做法是让 key 自然过期或者只在业务状态已确定提交后再做清理。Redis 方案性能高但可靠性不如数据库唯一约束通常用于可容忍极低概率异常的旁路去重核心资金类业务仍然建议以数据库为准。8.5 幂等方案对比与选型方案可靠性性能实现复杂度适用场景数据库唯一约束高中低核心交易、订单、资金等强一致场景业务状态机高高中状态流转清晰、可由业务字段判断的场景Redis 去重中高中高并发可容忍少量异常的营销、积分等旁路场景选型建议核心交易、资金、库存等强一致场景优先选择“数据库唯一约束 业务状态机”组合高并发可容忍少量异常的营销、积分、风控旁路场景可以使用 Redis 去重降低数据库压力。九、事务消息让“业务操作”和“消息发送”尽量一致9.1 为什么还需要事务消息前面讲到生产端发送补偿时提到“本地消息表”它的目的是解决“业务成功但消息没发出”或“消息发出但业务回滚”的不一致问题。RocketMQ 提供的事务消息可以把这个过程标准化。事务消息解决的不是消息重复或消息丢失本身而是业务动作与消息发送之间的不一致要么两者都成功要么都不生效避免业务已经成功但消息没有发出或消息已经发出但业务回滚。9.2 事务消息的执行流程Producer 发送一条 Half 消息到 BrokerBroker 存储后返回成功但消费者暂时读取不到这条消息。Producer 执行本地事务根据结果返回COMMIT_MESSAGE或ROLLBACK_MESSAGE。如果本地事务提交Broker 将 Half 消息变为可投递消息如果回滚Broker 删除该消息。如果 Broker 长时间没有收到本地事务结果会向 Producer 发起事务状态回查Producer 再确认本地事务真实状态。9.3 事务消息代码示例TransactionMQProducer producer new TransactionMQProducer(tx-producer-group); producer.setNamesrvAddr(127.0.0.1:9876); producer.setTransactionListener(new TransactionListener() { Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { try { orderService.createOrderAndPay((Long) arg); return LocalTransactionState.COMMIT_MESSAGE; } catch (Exception e) { log.error(本地事务执行失败, e); return LocalTransactionState.ROLLBACK_MESSAGE; } } Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { Long orderId Long.valueOf(msg.getKeys()); boolean success orderService.isOrderPaid(orderId); return success ? LocalTransactionState.COMMIT_MESSAGE : LocalTransactionState.ROLLBACK_MESSAGE; } }); producer.start(); Message msg new Message(OrderTopic, OrderPaid, order-2001, 订单支付事务消息.getBytes(StandardCharsets.UTF_8)); TransactionSendResult result producer.sendMessageInTransaction(msg, 2001L);9.4 事务消息仍然需要幂等即使使用了事务消息消息投递阶段依然可能出现重复。例如事务回查期间网络抖动、消费端确认超时等Consumer 仍可能重复拉取同一条消息。因此事务消息不能替代业务幂等两者要叠加使用。十、RocketMQ 5.x 对本主题的补充在 RocketMQ 5.x 中消费模型和重试机制发生了一些调整。对可靠性相关问题的思考方式也需要更新POP 消费模式更接近“按主题拉取”消费位点管理由 Broker 统一维护客户端可以更灵活处理单条消息的 ACK但可靠性原则不变仍然要业务成功后再确认。消费分组与重试5.x 对重试队列和死信队列进行了调整配置方式变化但“重试是可靠性手段、幂等是重复性手段”的结论不改变。Serverless 形态托管集群由平台维护刷盘和副本开发人员可以把更多精力放在生产端确认和消费幂等上。面试时如果能结合 4.x 和 5.x 的差异简单聊两句会给面试官留下不错的印象但不要为了显示新特性而偏离本题主线。十一、面试回答模板如果面试官直接问“RocketMQ 怎么保证消息不丢失、不重复消费”可以直接用下面的框架回答RocketMQ 默认是“至少一次”投递本身不提供端到端精确一次语义所以这个问题要从两层拆开。消息不丢失要覆盖三个环节生产端用同步发送、失败重试、本地消息表或事务消息Broker 端用同步刷盘和同步复制必要时上 DLedger 多副本消费端必须业务成功后才返回成功失败交给重试和死信队列。消息不重复则主要靠业务幂等可以用数据库唯一约束、业务状态机或 Redis 去重最终做到可靠投递加消费幂等的组合。回答到这一段通常已经能覆盖面试官的核心考察点。如果对方继续追问某个环节再沿着生产端、Broker 端、消费端或幂等实现细节展开。十二、总结RocketMQ 的“消息不丢失”和“消息不被重复消费”看起来是两个问题实际是一个组合方案可靠投递解决消息不丢失生产端同步发送与补偿、Broker 端同步刷盘与同步复制、消费端成功后确认。幂等控制解决消息不重复数据库唯一约束、业务状态机、Redis 去重等方案根据业务强度选择。事务消息解决业务与消息发送之间的一致性但无法替代消费幂等。任何只背配置、不区分环节的回答都很难真正应对反问。最后记住一句话先用可靠机制保证消息至少到达一次再用幂等机制保证多次到达也只有一次业务效果。