消息队列速查手册 - RocketMQ 5.x 面试版

最后更新:2026-08-10


1. 消息队列基础

1.1 为什么要用 MQ

一句话总结:MQ 就是系统之间的"传话筒",让系统不用面对面说话,而是通过中间人传话。

三大核心作用

1. 解耦(最常见)

想象一下供热平台的场景:用户缴费成功后,需要通知计费系统更新账单、通知发票系统开票、通知客服系统更新状态。如果不用 MQ,缴费系统就得直接调用这三个系统的接口,代码耦合度极高,任何一个系统挂了都会影响缴费。

用 MQ 之后,缴费系统只管发一条"缴费成功"的消息出去,谁关心谁来订阅。这就是解耦的精髓。

2. 异步

还是缴费场景,如果同步调用三个系统,每个系统响应 200ms,总耗时 600ms。用 MQ 之后,缴费系统只需要 50ms 把消息扔进去就返回了,用户体验直接起飞。

3. 削峰

供热缴费高峰期,比如供暖季开始前那几天,瞬间可能有上万笔缴费请求。如果直接打到数据库,数据库直接躺平。用 MQ 把请求先存起来,消费者按自己的节奏慢慢处理,数据库就不会被打崩了。

面试怎么说: "在我们的供热平台里,MQ 主要解决三个问题。第一是解耦,缴费系统通过消息通知下游系统,不用直接调用接口;第二是异步,用户缴费后不用等所有系统都处理完才能看到结果;第三是削峰,供暖季开始前的高峰期,消息队列能扛住瞬时流量,保护下游系统。"

1.2 MQ 引入带来的问题

别以为引入 MQ 就万事大吉了,它也会带来新的问题:

问题具体表现解决方案
系统可用性降低MQ 挂了,整个链路就断了集群部署、主从切换、多活部署
消息丢失消息从生产到消费任一环节出问题生产者确认机制、持久化、手动ACK
重复消费网络抖动导致消息重发幂等性设计(后面详细讲)
消息顺序性消息乱序到达分区顺序消息
消息积压消费速度跟不上生产速度扩容消费者、优化消费逻辑

面试怎么说: "引入 MQ 确实解决了耦合和性能问题,但也带来了新的挑战。比如消息丢失问题,我们在生产者端用同步发送加确认机制,Broker 端用同步刷盘保证持久化,消费者端用手动 ACK 确认。对于重复消费,通过业务层面的幂等性设计来解决,比如用唯一业务 ID 加去重表。"

1.3 RocketMQ vs Kafka vs RabbitMQ 对比

特性RocketMQ 5.xKafkaRabbitMQ
开发语言JavaScala/JavaErlang
协议自定义协议自定义协议AMQP
单队列吞吐十万级/秒百万级/秒万级/秒
延迟ms 级别ms 级别μs 级别
消息可靠性支持事务消息、延迟消息支持副本机制支持确认机制
消息回溯支持按时间戳回溯支持按 offset 回溯不支持
适用场景业务场景(订单、支付)日志收集、大数据企业应用、即时通讯
生态阿里生态大数据生态企业级应用生态
运维难度中等较低较高

面试怎么说: "这三款 MQ 各有特点。RocketMQ 是阿里开源的,最适合国内的业务场景,支持事务消息和延迟消息,对 Java 生态友好。Kafka 吞吐量最高,适合日志收集和大数据场景。RabbitMQ 延迟最低,但吞吐量相对较小,适合企业级应用。我们供热平台选 RocketMQ,主要是因为需要事务消息保证支付场景的一致性,而且团队都是 Java 技术栈。"


2. RocketMQ 5.x 架构

2.1 整体架构

RocketMQ 5.x 的架构由四个核心组件组成:

1. NameServer(名字服务器)

  • 轻量级的服务注册中心

  • 每个 NameServer 独立工作,节点之间不通信

  • 存储 Broker 的路由信息

  • 生产者、消费者定期从 NameServer 拉取路由信息

2. Broker(消息服务器)

  • 核心组件,负责消息存储、转发

  • 分为 Master 和 Slave,Master 可读写,Slave 只能读

  • 支持多副本机制保证高可用

3. Producer(生产者)

  • 从 NameServer 获取路由信息

  • 向 Broker 发送消息

  • 支持多种发送方式:同步、异步、单向

4. Consumer(消费者)

  • 从 NameServer 获取路由信息

  • 从 Broker 拉取消息(Pull 模式)

  • 支持集群消费和广播消费两种模式

2.2 5.x 新架构变化

RocketMQ 5.x 最大的变化是计算存储分离,引入了 Controller 模式

4.x 架构的问题

  • Broker 既负责消息存储,又负责消息处理

  • 扩容需要整个 Broker 扩,资源利用率低

  • 主从切换依赖 ZooKeeper 或 DLedger,复杂度高

5.x 架构的改进

  • 存储层:专门负责消息存储,可以独立扩容

  • 计算层:专门负责消息处理,可以独立扩容

  • Controller:统一管理元数据,简化主从切换

面试怎么说: "RocketMQ 5.x 最大的改进是计算存储分离。之前的版本,Broker 既要存消息又要处理消息,扩容只能整个扩。5.x 把存储和计算拆开,存储层可以单独扩容磁盘,计算层可以单独扩容 CPU 和内存。这样资源利用率更高,弹性伸缩更灵活。"

2.3 和 4.x 的关键区别

特性RocketMQ 4.xRocketMQ 5.x
架构模式Broker 一体化计算存储分离
主从管理依赖 DLedger/ZooKeeper内置 Controller
消息类型固定级别延迟消息任意时间延迟消息
弹性伸缩整体扩容存储和计算独立扩容
云原生一般原生支持 K8s
事务消息支持性能优化,更稳定

2.4 消息存储模型

RocketMQ 的消息存储采用三个文件配合:

1. CommitLog(提交日志)

  • 存储所有消息的实际内容

  • 顺序写入,性能极高

  • 所有 Topic 的消息都写同一个文件

2. ConsumeQueue(消费队列)

  • 相当于消息的索引

  • 存储消息在 CommitLog 中的偏移量(offset)

  • 每个 Topic 的每个 Queue 对应一个 ConsumeQueue

3. IndexFile(索引文件)

  • 支持按消息 Key 查询

  • 可以快速找到某个 Key 对应的所有消息

面试怎么说: "RocketMQ 的消息存储分三层。CommitLog 存实际消息,所有 Topic 的消息顺序写入同一个文件,这样顺序写磁盘性能最高。ConsumeQueue 是索引,记录消息在 CommitLog 的位置。IndexFile 支持按业务 Key 快速查询。这种设计既保证了写入性能,又支持快速查询。"


3. 消息类型

3.1 普通消息

最简单的消息类型,发出去就不管了,适合对可靠性要求不高的场景。

使用场景:日志收集、数据同步、非核心业务通知

3.2 顺序消息

两种类型

  • 全局顺序:整个 Topic 的消息都按顺序,性能差

  • 分区顺序:同一个 Queue 的消息按顺序,性能好(推荐)

实现原理: 发送时,通过 MessageQueueSelector 选择同一个 Queue,保证同一个业务的消息发到同一个 Queue。消费时,单线程消费这个 Queue。

使用场景:订单状态变更、库存变更

// 生产者:选择同一个 Queue
producer.send(msg, new MessageQueueSelector() {
    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
        Integer orderId = (Integer) arg;
        int index = orderId % mqs.size();
        return mqs.get(index);
    }
}, orderId);

// 消费者:顺序消费
consumer.registerMessageListener(new MessageListenerOrderly() {
    public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, 
                                                ConsumeOrderlyContext context) {
        // 处理消息
        return ConsumeOrderlyStatus.SUCCESS;
    }
});

面试怎么说: "顺序消息分全局顺序和分区顺序。全局顺序性能差,实际用得少。我们用的是分区顺序,通过 MessageQueueSelector 把同一个订单的消息发到同一个 Queue,然后单线程消费这个 Queue,这样就能保证同一个订单的消息按顺序处理。"

3.3 延迟消息

4.x 版本:只支持固定级别(1s 5s 10s 30s 1m ... 2h),共 18 个级别

5.x 版本:支持任意时间延迟,精确到毫秒

使用场景:订单超时取消、延迟任务、定时通知

// 5.x 任意时间延迟
Message msg = new Message("Topic", "Tag", body);
msg.setDeliveryTimestamp(System.currentTimeMillis() + 3600000); // 1小时后
producer.send(msg);

3.4 事务消息(重点!)

事务消息是 RocketMQ 的杀手锏,解决分布式事务问题。

流程

  1. 生产者发送 half 消息(半消息)到 Broker

  2. Broker 存储 half 消息,但对消费者不可见

  3. 生产者执行本地事务

  4. 根据本地事务结果,提交或回滚消息

    • 提交:消费者可以消费

    • 回滚:删除消息

  5. 如果生产者没有提交或回滚,Broker 定时回查

为什么需要事务消息?

假设供热平台缴费场景:缴费成功后,需要更新缴费记录和生成发票。如果这两个操作在不同系统,用普通消息可能出现:

  • 本地事务成功,消息发送失败 → 发票没生成

  • 消息发送成功,本地事务失败 → 发票生成了但缴费没成功

事务消息保证这两个操作要么都成功,要么都失败。

面试怎么说: "事务消息的核心是 half 消息加回查机制。生产者先发送一个 half 消息,这个消息对消费者不可见。然后执行本地事务,成功后提交消息,失败就回滚。如果生产者挂了,Broker 会定时回查本地事务状态,保证最终一致性。这样就能解决分布式事务问题,保证本地事务和消息发送的原子性。"


4. 消息可靠性保证

消息丢失是面试必问的问题。消息从生产到消费经过三个阶段,每个阶段都可能丢失:

4.1 生产者端

问题:消息发送失败怎么办?

解决方案

  1. 同步发送:等待 Broker 确认,最可靠

  2. 重试机制:发送失败自动重试,默认 2 次

  3. 发送确认:通过 SendResult 判断是否成功

// 同步发送 + 重试
SendResult result = producer.send(msg);
if (result.getSendStatus() != SendStatus.SEND_OK) {
    // 发送失败,记录日志,补偿处理
    log.error("消息发送失败: {}", msg);
}

4.2 Broker 端

问题:消息存储到 Broker 后,Broker 挂了怎么办?

解决方案

1. 刷盘策略

  • 同步刷盘:消息写入磁盘后才返回成功,最可靠但性能差

  • 异步刷盘:消息写入内存就返回,定期刷盘,性能好但可能丢消息

2. 主从同步

  • 同步复制:Master 和 Slave 都写入成功才返回,最可靠

  • 异步复制:Master 写入成功就返回,Slave 异步同步,性能好但可能丢消息

最佳实践:同步刷盘 + 同步复制(金融级可靠性要求)

面试怎么说: "Broker 端保证可靠性主要靠两点。第一是刷盘策略,同步刷盘保证消息写入磁盘才返回,不会因为宕机丢消息。第二是主从同步,同步复制保证 Master 和 Slave 都有消息,Master 挂了 Slave 可以接管。对于支付这种核心场景,我们用同步刷盘加同步复制,虽然性能低一点,但数据不会丢。"

4.3 消费者端

问题:消息消费失败怎么办?

解决方案

  1. 手动 ACK:消费成功后才确认,失败不确认

  2. 重试队列:消费失败的消息进入重试队列,延迟重发

  3. 死信队列:重试多次还失败的消息进入死信队列,人工处理

// 手动 ACK
consumer.registerMessageListener(new MessageListenerConcurrently() {
    public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, 
                                                     ConsumeConcurrentlyContext context) {
        try {
            // 业务处理
            processMessage(msg);
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; // 成功
        } catch (Exception e) {
            return ConsumeConcurrentlyStatus.RECONSUME_LATER; // 失败,稍后重试
        }
    }
});

4.4 端到端可靠性方案

完整的可靠性保证需要三个环节配合:

  1. 生产者:同步发送 + 失败重试 + 发送确认

  2. Broker:同步刷盘 + 同步复制 + 多副本

  3. 消费者:手动 ACK + 重试机制 + 死信队列

面试怎么说: "端到端的可靠性需要三个环节配合。生产者用同步发送加失败重试,Broker 用同步刷盘加同步复制,消费者用手动 ACK 加重试机制。对于支付这种核心场景,我们还会加死信队列,重试多次还失败的消息进入死信队列,人工介入处理。这样就能保证消息从生产到消费不丢失。"


5. 消息幂等消费

5.1 为什么需要幂等

MQ 为了保证消息不丢失,会重发消息。比如:

  • 网络抖动导致消息重发

  • 消费者处理超时,Broker 认为没消费成功重发

  • 消费者扩容导致 rebalance,消息被重新分配

如果不做幂等,同一条消息可能被消费多次,导致业务逻辑出错。比如:

  • 缴费消息被消费两次 → 用户被扣两次钱

  • 订单状态变更消息被消费两次 → 订单状态异常

5.2 实现方案

方案一:唯一 ID + 数据库去重表

原理:每条消息带一个唯一业务 ID(如订单号),消费前先查去重表,存在就跳过。

// 消费消息
String orderId = msg.getKeys(); // 订单号
int count = deduplicationMapper.selectByOrderId(orderId);
if (count > 0) {
    log.info("消息已消费: {}", orderId);
    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}

// 业务处理
processOrder(orderId);

// 插入去重表
deduplicationMapper.insert(orderId);

优点:实现简单,可靠性高 缺点:需要额外的数据库表,有性能开销

方案二:Redis SETNX

原理:用 Redis 的 SETNX 命令,如果 Key 不存在就设置成功,存在就设置失败。

String msgId = msg.getMsgId();
boolean success = redis.setnx(msgId, "1", 3600); // 1小时过期
if (!success) {
    log.info("消息已消费: {}", msgId);
    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}

// 业务处理
processMessage(msg);

优点:性能好,适合高并发场景 缺点:Redis 挂了可能重复消费,需要补偿机制

方案三:状态机

原理:业务状态只能单向变更,重复消费不会影响状态。

// 订单状态:待支付 -> 已支付 -> 已发货
String orderId = msg.getKeys();
int affected = orderMapper.updateStatus(orderId, "已支付", "待支付");
if (affected == 0) {
    log.info("订单状态已变更,跳过: {}", orderId);
    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}

// 业务处理
processOrder(orderId);

优点:不需要额外的表或 Redis,性能好 缺点:只适用于状态变更场景

5.3 各方案对比

方案可靠性性能实现复杂度适用场景
去重表通用场景
Redis SETNX高并发场景
状态机状态变更场景

5.4 实战案例

供热平台缴费场景

用户缴费成功后,MQ 发送缴费成功消息。消费者需要:

  1. 更新缴费记录状态为"已缴费"

  2. 生成发票

幂等方案

@Transactional
public ConsumeConcurrentlyStatus consumeMessage(MessageExt msg) {
    String paymentId = msg.getKeys();
    
    // 1. 查询缴费记录
    Payment payment = paymentMapper.selectById(paymentId);
    if (payment.getStatus() == PaymentStatus.PAID) {
        log.info("缴费消息已处理: {}", paymentId);
        return CONSUME_SUCCESS;
    }
    
    // 2. 更新缴费状态(状态机)
    int affected = paymentMapper.updateStatus(paymentId, PAID, UNPAID);
    if (affected == 0) {
        return CONSUME_SUCCESS; // 状态已变更,跳过
    }
    
    // 3. 生成发票
    invoiceService.createInvoice(payment);
    
    return CONSUME_SUCCESS;
}

关键点

  • 用状态机保证"更新状态"操作的幂等性

  • 整个方法加事务,保证原子性

  • 如果生成发票失败,事务回滚,下次重试

面试怎么说: "幂等消费的核心思路是'先检查,再处理'。我们的缴费场景用的是状态机方案,因为缴费记录有明确的状态流转(待缴费→已缴费),重复消费不会影响状态。对于没有明确状态的场景,可以用去重表或 Redis SETNX。关键是要根据业务场景选择最合适的方案。"


6. 消息积压处理

消息积压是生产环境最常见的问题之一。

6.1 排查步骤

第一步:确认积压情况

# 查看消费者组积压情况
mqadmin consumerProgress -g consumerGroup

第二步:定位原因

  1. 消费者慢?

    • 查看消费者日志,是否有异常

    • 查看消费者 TPS,是否下降

    • 查看下游服务(数据库、Redis)是否正常

  2. 生产者突增?

    • 查看生产者 TPS,是否突增

    • 是否有定时任务或批量任务触发

    • 是否有促销活动导致流量激增

  3. 消费者数量不足?

    • 查看消费者实例数

    • 查看每个实例的消费 TPS

    • 是否达到消费者瓶颈

6.2 应急方案

方案一:临时扩容消费者

如果消费者实例数 < Queue 数,可以增加消费者实例。

// 临时启动新的消费者实例
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("consumerGroup");
consumer.subscribe("Topic", "*");
consumer.registerMessageListener(listener);
consumer.start();

方案二:临时增加 Queue 数量

如果消费者实例数已经等于 Queue 数,需要增加 Queue 数量。

# 更新 Topic 的 Queue 数量
mqadmin updateTopic -n nameserver:9876 -b broker:10911 -t Topic -w 16 -r 16

方案三:跳过非关键消息

对于非核心业务,可以跳过积压消息,快速处理最新消息。

// 重置消费位点,跳过积压消息
mqadmin resetOffsetByTimestamp -n nameserver:9876 -g consumerGroup -t Topic -s timestamp

6.3 长期优化方案

  1. 优化消费逻辑

    • 批量处理:一次处理多条消息

    • 异步处理:耗时操作异步化

    • 缓存优化:减少数据库查询

  2. 优化消费者配置

    consumer.setConsumeThreadMin(20); // 最小消费线程
    consumer.setConsumeThreadMax(64); // 最大消费线程
    consumer.setConsumeMessageBatchMaxSize(32); // 批量消费数量
    

  3. 监控告警

    • 监控消息积压量

    • 设置告警阈值

    • 及时发现问题

6.4 生产案例

供热平台缴费高峰场景

供暖季开始前一周,缴费量暴增 10 倍,消息积压达到 100 万条。

排查过程

  1. 查看消费者 TPS:从 1000 降到 200

  2. 查看消费者日志:数据库连接超时

  3. 查看数据库:CPU 100%,大量慢查询

解决方案

  1. 应急:临时扩容消费者实例(4 → 16)

  2. 应急:增加 Topic 的 Queue 数量(8 → 32)

  3. 优化:批量处理消息(每次处理 1 条 → 每次处理 32 条)

  4. 优化:优化数据库查询(加索引、减少查询字段)

  5. 长期:引入 Redis 缓存,减少数据库压力

结果:积压消息在 2 小时内处理完毕,后续高峰期再未出现积压。

面试怎么说: "处理消息积压首先要定位原因,是消费者慢还是生产者突增。我们的供暖季高峰案例,排查发现是数据库慢查询导致消费者变慢。应急方案是扩容消费者和 Queue 数量,长期优化是批量处理、优化 SQL、引入缓存。关键是建立监控告警,及时发现积压。"


7. RocketMQ vs Kafka 深度对比

7.1 架构差异

特性RocketMQ 5.xKafka
核心组件NameServer + BrokerZooKeeper + Broker
元数据管理NameServer(轻量级)ZooKeeper(重量级)
消息存储CommitLog + ConsumeQueuePartition + Segment
副本机制主从同步ISR 机制
主从切换Controller 自动切换Controller 自动切换
扩展性计算存储分离,独立扩容整体扩容
协议自定义协议自定义协议

7.2 吞吐量对比

Kafka 吞吐量更高的原因

  1. 批量压缩:Kafka 支持批量发送和压缩,减少网络传输

  2. 零拷贝:Kafka 使用 sendfile 系统调用,减少数据拷贝

  3. 顺序写盘:两者都支持,但 Kafka 优化更极致

  4. PageCache:Kafka 充分利用操作系统的 PageCache

实际数据

  • Kafka:单 Broker 百万级 TPS

  • RocketMQ:单 Broker 十万级 TPS

7.3 延迟对比

RocketMQ 延迟更低的原因

  1. 内存映射:RocketMQ 使用 MMAP,Kafka 使用 sendfile

  2. 长轮询:RocketMQ 使用长轮询,Kafka 使用短轮询

  3. 消息存储:RocketMQ 的 CommitLog 更小,查询更快

实际数据

  • RocketMQ:ms 级别延迟

  • Kafka:ms 级别延迟(略高于 RocketMQ)

7.4 消息可靠性对比

特性RocketMQKafka
事务消息原生支持不支持
延迟消息原生支持(5.x 任意时间)不支持
消息回溯支持按时间戳支持按 offset
消息过滤支持 Tag 和 SQL92不支持
死信队列原生支持需要自己实现
消息轨迹原生支持不支持

7.5 适用场景对比

RocketMQ 适合

  • 业务场景:订单、支付、库存

  • 需要事务消息

  • 需要延迟消息

  • 需要消息过滤

  • Java 技术栈

Kafka 适合

  • 日志收集

  • 大数据处理(流计算、ETL)

  • 高吞吐场景

  • 多语言环境

7.6 选型建议

场景推荐原因
业务系统(订单、支付)RocketMQ支持事务消息、延迟消息
日志收集Kafka吞吐量高,成本低
大数据处理Kafka生态完善,支持流计算
企业应用RabbitMQ协议标准,兼容性好
高并发场景Kafka吞吐量最高
低延迟场景RocketMQ延迟更低

面试怎么说: "RocketMQ 和 Kafka 各有优势。RocketMQ 在业务场景更有优势,支持事务消息、延迟消息、消息过滤,对 Java 生态友好。Kafka 在吞吐量和大数据生态更有优势,适合日志收集和流处理。我们供热平台选 RocketMQ,主要是因为需要事务消息保证支付一致性,而且团队都是 Java 技术栈。如果是做日志收集,我会选 Kafka。"


8. 生产环境 RocketMQ 调优

8.1 Broker 调优参数

1. 刷盘策略

# 同步刷盘(高可靠)
flushDiskType=SYNC_FLUSH

# 异步刷盘(高性能)
flushDiskType=ASYNC_FLUSH

2. 主从同步

# 同步复制(高可靠)
brokerRole=SYNC_MASTER

# 异步复制(高性能)
brokerRole=ASYNC_MASTER

3. 内存优化

# 堆内存
-Xms8g -Xmx8g

# 堆外内存(用于 MMAP)
-XX:MaxDirectMemorySize=16g

4. 线程池优化

# 发送消息线程池
sendMessageThreadPoolNums=16

# 拉取消息线程池
pullMessageThreadPoolNums=16

# 消费消息线程池
consumeMessageThreadPoolNums=32

8.2 生产者调优

1. 批量发送

// 批量发送(不超过 4MB)
List<Message> msgs = new ArrayList<>();
for (int i = 0; i < 100; i++) {
    msgs.add(new Message("Topic", "Tag", ("Message " + i).getBytes()));
}
producer.send(msgs);

2. 压缩

// 开启压缩(默认开启)
producer.setCompressMsgBodyOverHowmuch(4096); // 超过 4KB 压缩

3. 异步发送

// 异步发送(不阻塞)
producer.send(msg, new SendCallback() {
    public void onSuccess(SendResult sendResult) {
        // 发送成功
    }
    public void onException(Throwable e) {
        // 发送失败
    }
});

4. 超时设置

producer.setSendMsgTimeout(3000); // 发送超时 3 秒

8.3 消费者调优

1. 并发消费

// 设置消费线程数
consumer.setConsumeThreadMin(20);
consumer.setConsumeThreadMax(64);

2. 批量消费

// 一次拉取多条消息
consumer.setConsumeMessageBatchMaxSize(32);

3. 拉取间隔

// 拉取间隔(减少空轮询)
consumer.setPullInterval(0);

4. 消费超时

// 消费超时时间
consumer.setConsumeTimeout(15); // 15 分钟

8.4 监控告警方案

1. 监控指标

  • 消息积压量(最重要)

  • 生产者 TPS

  • 消费者 TPS

  • 消息延迟

  • Broker 磁盘使用率

  • Broker CPU 使用率

2. 告警规则

# Prometheus 告警规则
groups:
  - name: rocketmq
    rules:
      - alert: MessageAccumulation
        expr: rocketmq_consumer_offset_diff > 100000
        for: 5m
        labels:
          severity: critical
        annotations:
          summary: "消息积压超过 10 万条"
          
      - alert: BrokerDiskUsage
        expr: rocketmq_broker_disk_usage > 0.8
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "Broker 磁盘使用率超过 80%"

3. 监控工具

  • RocketMQ Dashboard:官方监控面板

  • Prometheus + Grafana:指标采集和可视化

  • SkyWalking:分布式链路追踪

面试怎么说: "生产环境调优主要从四个方面。第一是 Broker 调优,根据业务场景选择刷盘策略和主从同步方式。第二是生产者调优,用批量发送和压缩提高吞吐量。第三是消费者调优,增加消费线程数和批量消费数量。第四是监控告警,重点监控消息积压量、TPS、延迟等指标,设置合理的告警阈值。"


9. 面试高频问答

Q1:为什么选择 RocketMQ 而不是 Kafka?

回答: "我们供热平台是典型的业务系统,需要处理订单、支付等场景。选择 RocketMQ 有三个原因:第一,RocketMQ 支持事务消息,能保证支付场景的分布式事务一致性;第二,支持延迟消息,订单超时取消等业务用得很方便;第三,团队都是 Java 技术栈,RocketMQ 对 Java 生态更友好。如果是做日志收集,我会选 Kafka,因为吞吐量更高。"

Q2:如何保证消息不丢失?

回答: "消息从生产到消费经过三个阶段,每个阶段都要保证不丢失。生产者用同步发送加失败重试,确保消息发到 Broker。Broker 用同步刷盘加同步复制,确保消息持久化到磁盘且有副本。消费者用手动 ACK,消费成功后才确认,失败进入重试队列。对于核心业务,还会加死信队列,重试多次还失败的消息进入死信队列,人工介入处理。"

Q3:如何处理消息重复消费?

回答: "重复消费是 MQ 的常见问题,因为网络抖动、消费者扩容等原因,同一条消息可能被消费多次。我们的解决方案是幂等性设计。对于有明确状态的场景,用状态机方案,比如订单状态只能单向变更。对于没有状态的场景,用唯一业务 ID 加去重表,消费前先查去重表。高并发场景可以用 Redis SETNX,但要注意 Redis 挂了的补偿机制。"

Q4:消息积压怎么处理?

回答: "首先定位原因,是消费者慢还是生产者突增。如果是消费者慢,查看下游服务是否正常,优化消费逻辑。如果是生产者突增,查看是否有批量任务或促销活动。应急方案有三个:第一,临时扩容消费者实例;第二,增加 Topic 的 Queue 数量;第三,对于非核心业务,可以跳过积压消息。长期优化要批量处理、异步化、引入缓存。关键是要建立监控告警,及时发现积压。"

Q5:RocketMQ 的事务消息是怎么实现的?

回答: "事务消息的核心是 half 消息加回查机制。生产者先发送一个 half 消息到 Broker,这个消息对消费者不可见。然后生产者执行本地事务,成功后提交消息,失败就回滚。如果生产者没有提交或回滚(比如生产者挂了),Broker 会定时(默认 60 秒)回查生产者的本地事务状态,根据回查结果决定提交还是回滚。这样就能保证本地事务和消息发送的原子性,解决分布式事务问题。"

Q6:顺序消息怎么实现?

回答: "顺序消息分全局顺序和分区顺序。全局顺序是整个 Topic 的消息都按顺序,性能差,实际用得少。我们用分区顺序,通过 MessageQueueSelector 把同一个业务的消息发到同一个 Queue,然后单线程消费这个 Queue。比如订单场景,通过订单号取模选择 Queue,保证同一个订单的消息按顺序处理。注意消费者要用 MessageListenerOrderly,不能用 MessageListenerConcurrently。"

Q7:RocketMQ 的存储模型是什么?

回答: "RocketMQ 的消息存储分三层。第一层是 CommitLog,存储所有 Topic 的实际消息,顺序写入磁盘,性能最高。第二层是 ConsumeQueue,相当于消息的索引,记录消息在 CommitLog 的偏移量,每个 Topic 的每个 Queue 对应一个 ConsumeQueue。第三层是 IndexFile,支持按业务 Key 快速查询消息。消费者消费时,先读 ConsumeQueue 拿到偏移量,再从 CommitLog 读实际消息。"

Q8:RocketMQ 5.x 和 4.x 有什么区别?

回答: "最大的区别是 5.x 实现了计算存储分离。4.x 的 Broker 既要存消息又要处理消息,扩容只能整体扩。5.x 把存储和计算拆开,存储层可以单独扩容磁盘,计算层可以单独扩容 CPU 和内存。另外 5.x 的延迟消息支持任意时间,4.x 只支持固定级别。5.x 还内置了 Controller,主从切换不依赖 ZooKeeper,架构更简洁。"

Q9:如何提高消息吞吐量?

回答: "生产者端:用批量发送,一次发送多条消息;开启压缩,减少网络传输;用异步发送,不阻塞主线程。Broker 端:用异步刷盘,提高写入性能;用异步复制,提高主从同步性能。消费者端:增加消费线程数;批量消费,一次处理多条消息;优化消费逻辑,减少耗时操作。"

Q10:如何保证消息的顺序性?

回答: "顺序性分生产、存储、消费三个阶段。生产阶段,通过 MessageQueueSelector 把同一个业务的消息发到同一个 Queue。存储阶段,RocketMQ 保证同一个 Queue 的消息顺序写入 CommitLog。消费阶段,用 MessageListenerOrderly 单线程消费同一个 Queue。三个阶段配合,就能保证同一个业务的消息按顺序处理。"

Q11:消息队列有哪些缺点?

回答: "引入 MQ 会带来三个问题。第一是系统可用性降低,MQ 挂了整个链路就断了,需要集群部署保证高可用。第二是消息一致性问题,消息可能丢失或重复,需要各种机制保证可靠性。第三是系统复杂度增加,需要考虑消息积压、消息顺序、幂等消费等问题。所以不是所有场景都适合用 MQ,只有真正需要解耦、异步、削峰的场景才用。"

Q12:RocketMQ 的主从切换是怎么实现的?

回答: "4.x 版本依赖 DLedger 或 ZooKeeper 实现主从切换。5.x 版本内置了 Controller,Master 和 Slave 定期向 Controller 发送心跳,Controller 检测到 Master 挂了,会选择一个 Slave 升级为 Master。整个过程对生产者和消费者透明,他们从 NameServer 获取最新的路由信息。"

Q13:如何处理消息丢失?

回答: "消息丢失分三个阶段。生产者丢失:用同步发送加失败重试,确保消息发到 Broker。Broker 丢失:用同步刷盘保证消息写入磁盘,用同步复制保证有副本。消费者丢失:用手动 ACK,消费成功后才确认,失败进入重试队列。对于核心业务,还要加死信队列,人工介入处理。"

Q14:RocketMQ 的 NameServer 为什么不用 ZooKeeper?

回答: "NameServer 的设计哲学是简单。每个 NameServer 独立工作,节点之间不通信,通过心跳机制检测 Broker 是否存活。这种设计足够满足 RocketMQ 的需求,而且运维简单,不需要依赖 ZooKeeper 这种重量级的组件。如果 NameServer 挂了,生产者和消费者本地有路由信息的缓存,短时间内不影响消息发送和消费。"

Q15:如何设计一个高可用的消息队列系统?

回答: "高可用要从四个层面保证。第一是 NameServer 层面,至少部署 3 个节点,任意一个挂了不影响服务。第二是 Broker 层面,每个 Topic 至少 2 个副本(1 主 1 从),用同步复制保证数据不丢。第三是生产者层面,用同步发送加失败重试,发送失败可以降级到本地文件。第四是消费者层面,用集群消费模式,任意一个消费者挂了不影响整体消费。"

Q16:消息队列的吞吐量受哪些因素影响?

回答: "主要受五个因素影响。第一是磁盘 IO,消息要持久化到磁盘,顺序写比随机写快。第二是网络带宽,消息要在网络上传输。第三是 CPU,消息要序列化、压缩。第四是内存,消息要先写入内存。第五是消息大小,消息越大吞吐量越低。所以提高吞吐量的方法是批量发送、压缩、异步刷盘、顺序写盘。"

Q17:如何监控消息队列?

回答: "监控主要关注四个指标。第一是消息积压量,这是最重要的指标,积压说明消费能力不足。第二是生产者和消费者的 TPS,反映系统的吞吐量。第三是消息延迟,反映消息处理的速度。第四是 Broker 的资源使用率,包括 CPU、内存、磁盘、网络。我们用 Prometheus 采集指标,Grafana 做可视化,设置合理的告警阈值,及时发现问题。"

Q18:消息队列如何处理大消息?

回答: "RocketMQ 默认消息大小不超过 4MB。如果消息超过这个限制,有三种方案。第一是消息引用,消息体只存文件的 URL,实际数据存在 OSS 或 HDFS。第二是消息分片,把大消息拆分成多个小消息,消费者收到所有分片后再组装。第三是压缩,如果消息是文本类型,可以压缩后发送。一般推荐第一种方案,简单可靠。"

Q19:如何保证消息的可靠性?

回答: "消息可靠性要从端到端保证。生产者用同步发送加失败重试,Broker 用同步刷盘加同步复制,消费者用手动 ACK 加重试机制。对于核心业务,还要加死信队列和消息轨迹追踪。关键是要根据业务场景选择合适的可靠性级别,不是所有业务都需要最高的可靠性,那样会牺牲性能。"

Q20:RocketMQ 和 RabbitMQ 有什么区别?

回答: "主要区别在三个方面。第一是吞吐量,RocketMQ 单 Broker 十万级 TPS,RabbitMQ 单 Broker 万级 TPS。第二是延迟,RabbitMQ 延迟更低,可以达到微秒级。第三是适用场景,RocketMQ 适合业务场景,RabbitMQ 适合企业应用。RabbitMQ 支持 AMQP 协议,兼容性好,但性能不如 RocketMQ。我们选 RocketMQ 主要是因为吞吐量更高,而且支持事务消息。"


总结

消息队列是分布式系统的核心组件,面试必问。记住几个核心点:

  1. 为什么要用 MQ:解耦、异步、削峰

  2. MQ 的问题:可用性、一致性、重复消费

  3. RocketMQ 架构:NameServer + Broker + Producer + Consumer

  4. 消息类型:普通、顺序、延迟、事务

  5. 消息可靠性:生产者同步发送、Broker 同步刷盘、消费者手动 ACK

  6. 幂等消费:唯一 ID + 去重表、Redis SETNX、状态机

  7. 消息积压:定位原因、应急扩容、长期优化

  8. RocketMQ vs Kafka:业务场景选 RocketMQ,大数据选 Kafka

面试口诀

  • MQ 三大作用:解耦异步削峰

  • MQ 三大问题:可用一致重复

  • 可靠性三段:生产 Broker 消费

  • 幂等三方案:去重 Redis 状态机

  • 选型两句话:业务 RocketMQ,数据 Kafka

祝面试顺利!


本内容由 Coze AI 生成,请遵循相关法律法规及《人工智能生成合成内容标识办法》使用与传播。

技术面试笔记 / 05_消息队列速查手册_2026 0 0 cosolar
2026-08-30T02:14:22.282452079Z 2026-08-30T02:17:13.372503392Z