在线 Java 面试刷题(持续更新):https://www.quanxiaoha.com/java-interview

面试考察点

面试官提出这个问题,核心是考察:

  1. 系统性思维:你是否能从消息的完整生命周期(生产 -> 传输 -> 存储 -> 消费)来系统性分析可靠性保障,而不仅仅是孤立的某个点。

  2. 对 RocketMQ 架构与核心机制的理解深度:你是否真正理解其持久化、高可用(主从复制)、消息确认等核心机制的工作原理及配置。

  3. 工程实践与权衡能力:在保证消息不丢失的同时,往往需要牺牲一定的性能(如吞吐量、延迟)。你是否了解关键的配置项,并能根据业务场景(如金融交易 vs 日志收集)做出合理的权衡。

  4. 问题排查与设计能力:当线上出现消息丢失时,你的排查思路是什么?这反映了你的实战经验和系统设计能力。

核心答案

保证 RocketMQ 消息不丢失,需要从消息生产端、Broker 服务端、消息消费端三个环节进行全链路保障,三者缺一不可。

  1. 生产阶段:确保消息成功发送并存储到 Broker。

  • 使用同步发送,并检查发送结果(SendResult)。

  • 合理配置重试机制(retryTimesWhenSendFailed)。

  • 对发送异常进行妥善处理(如记录日志、落库、告警)。

  • Broker 存储阶段:确保消息被可靠地持久化并完成高可用复制。

    • 主从架构:配置 brokerRole 为 SYNC_MASTER 或 ASYNC_MASTER,搭配从节点 (SLAVE)。

    • 刷盘策略:对可靠性要求极高的场景,将 flushDiskType 配置为 SYNC_FLUSH(同步刷盘)

    • 复制策略:对可靠性要求极高的场景,将 brokerRole 配置为 SYNC_MASTER,启用同步双写,确保主从都写入成功后才返回生产 ACK。

  • 消费阶段:确保消息被业务逻辑成功处理。

    • 在业务逻辑成功执行完毕后,再返回 ConsumeConcurrentlyStatus.CONSUME_SUCCESS

    • 谨慎使用 “异步消费” 或手动管理 offset,避免消息未处理就确认成功。

    • 利用消费重试机制:对于消费失败的消息(返回 RECONSUME_LATER 或抛出异常),RocketMQ 会将其投递到重试队列,在延迟后再次消费。

    深度解析

    原理与机制

    • 同步刷盘 vs 异步刷盘:

      • SYNC_FLUSH:生产者消息写入 CommitLog 后,Broker 会等待数据刷入磁盘后才返回成功响应。这是最可靠的模式。

      • ASYNC_FLUSH:消息写入 PageCache 后就返回成功,由后台线程定期刷盘。性能更高,但在 Broker 宕机且未刷盘时,会丢失 PageCache 中的数据。

    • 同步复制 vs 异步复制:

      • SYNC_MASTER:生产者消息写入主节点后,主节点会同步等待从节点 (SLAVE) 写入成功,才向生产者返回 ACK。这保证了主从数据强一致。

      • ASYNC_MASTER:主节点写入成功后立即返回 ACK,数据异步复制到从节点。在主节点宕机且数据未复制时,若从节点未完成数据同步,则会丢失消息。

    • 消费确认 (ACK) :RocketMQ 采用 offset 推进机制。消费者成功消费一批消息后,会将本地的 offset 提交给 Broker。Broker 记录此进度,后续从此 offset 之后投递消息。若消费者消费失败或未提交 offset,则下次会从上次提交的 offset 开始重新消费。

    代码示例与最佳实践

    // 1. 生产者:同步发送并检查
    DefaultMQProducer producer = new DefaultMQProducer("ProducerGroupName");
    producer.setNamesrvAddr("name-server-ip:9876");
    // 设置同步发送失败重试次数(默认2次,共3次)
    producer.setRetryTimesWhenSendFailed(3);
    producer.start();
    
    try {
        Message msg = new Message("TopicTest", "TagA", "OrderId001", "Hello, RocketMQ".getBytes());
        // 关键:同步发送,并获取发送结果
        SendResult sendResult = producer.send(msg);
        System.out.printf("消息发送成功:MsgId=%s, Queue=%s%n", 
                          sendResult.getMsgId(), 
                          sendResult.getMessageQueue());
    } catch (Exception e) {
        // 关键:必须处理异常,如记录到数据库或日志,并触发告警
        log.error("消息发送失败,将进行重试或人工处理", e);
        // 业务上:可将消息存入本地数据库,由定时任务扫描重发
    }
    
    // 2. 消费者:正确处理业务与ACK
    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("ConsumerGroupName");
    consumer.setNamesrvAddr("name-server-ip:9876");
    consumer.subscribe("TopicTest", "*");
    
    consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
        for (MessageExt msg : msgs) {
            try {
                // 关键:执行核心业务逻辑
                String orderId = new String(msg.getBody());
                boolean success = processOrder(orderId); // 你的业务处理
    
                if (success) {
                    // 业务成功,才返回CONSUME_SUCCESS
                    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
                } else {
                    // 业务失败,返回RECONSUME_LATER,消息将进入重试队列
                    // 注意:重试有最大次数(默认16次),超过后会进入死信队列(%DLQ%ConsumerGroupName)
                    return ConsumeConcurrentlyStatus.RECONSUME_LATER;
                }
            } catch (Exception e) {
                log.error("消费过程发生异常,消息将重试", e, msg);
                // 发生未捕获异常,也视为消费失败,会触发重试
                return ConsumeConcurrentlyStatus.RECONSUME_LATER;
            }
        }
        // 理论上不会到达这里
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    });
    consumer.start();

    对比分析与常见误区

    • 性能与可靠性权衡同步刷盘 + 同步复制 配置最安全,但吞吐量最低(可能只有异步模式的十分之一)。通常建议异步刷盘 + 同步复制,在保证主从不丢消息的前提下,获得较高的写入性能。

    • 常见误区  :

    1. 以为异步发送 (sendOneWay 或异步回调) 不会丢消息:异步发送不等待响应,网络抖动或 Broker 异常时,生产者无法感知发送失败。

    2. 消费逻辑中随意返回 CONSUME_SUCCESS :这是导致 “消费端消息丢失” 最常见的原因。务必确保业务逻辑执行成功。

    3. 忽视重试队列和死信队列:需要监控死信队列,处理反复重试都失败的消息,这通常是业务逻辑有严重问题的信号。

    4. Broker 单点部署:未配置主从,一旦 Broker 磁盘损坏或主机宕机,所有未消费消息全部丢失。

    总结

    保证 RocketMQ 消息不丢失是一个 “三端协同” 的工程:生产者要同步发送并处理异常,Broker 需配置合理的刷盘和复制策略,消费者则必须在业务成功后确认消息,任何一端的疏忽都会导致链路断裂。

    👉 欢迎加入小哈的星球,你将获得: 专属的项目实战(多个项目) / 1v1 提问 / Java 学习路线 / 学习打卡 / 每月赠书 / 社群讨论

    • 新项目:《Spring AI 项目实战》正在更新中..., 基于 Spring AI + Spring Boot 3.x + JDK 21;

    • 《从零手撸:仿小红书(微服务架构)》 已完结,基于 Spring Cloud Alibaba + Spring Boot 3.x + JDK 17..., 点击查看项目介绍;演示地址:http://116.62.199.48:7070/

    • 《从零手撸:前后端分离博客项目(全栈开发)》 2期已完结,演示链接:http://116.62.199.48/;

    • 专栏阅读地址:https://www.quanxiaoha.com/column

    截止目前,累计输出 100w+ 字,讲解图 4013+ 张,还在持续爆肝中.. 后续还会上新更多项目,目标是将 Java 领域典型的项目都整一波,如秒杀系统, 在线商城, IM 即时通讯,Spring Cloud Alibaba 等等,戳我加入学习,解锁全部项目,已有4200+小伙伴加入

    图片

    图片

    图片

    1. 我的私密学习小圈子,从0到1手撸企业实战项目~
    2. 面试官:RocketMQ 消息堆积了怎么处理?
    3. 高并发下是先写数据库,还是先写缓存?
    4. 从夯到拉,锐评13个 Java Web 框架!
    最近面试BAT,整理一份面试资料《Java面试BATJ通关手册》,覆盖了Java核心技术、JVM、Java并发、SSM、微服务、数据库、数据结构等等。
    获取方式:点“在看”,关注公众号并回复 Java 领取,更多内容陆续奉上。
    PS:因公众号平台更改了推送规则,如果不想错过内容,记得读完点一下“在看”,加个“星标”,这样每次新文章推送才会第一时间出现在你的订阅列表里。
    点“在看”支持小哈呀,谢谢啦
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐