RocketMQ常见面试题


1.如何保证顺序消费?

在生产中经常会有一些类似报表系统这样的系统,需要做 MySQL 的 binlog 同步。比如订单系统要同步订单表的数据到大数据部门的 MySQL 库中用于报表统计分析,通常的做法是基于 Canal 这样的中间件去监听订单数据库的 binlog,然后把这些 binlog 发送到 MQ 中,再由消费者从 MQ 中获取 binlog 落地到大数据部门的 MySQL 中。

在这个过程中,可能会有对某个订单的增删改操作,比如有三条 binlog 执行顺序是增加、修改、删除。消费者愣是换了顺序给执行成删除、修改、增加,这样能行吗?肯定是不行的。

顺序消息分为全局顺序消息和部分顺序消息,全局顺序消息指某个Topic下的所有消息都要保证顺序;部分顺序消息只要保证每一组消息被顺序消费即可,比如上面订单消息的例子,只要保证同一个订单ID的三个消息能按顺序消费即可。

1.1全局顺序消息

RocketMQ在默认情况下不保证顺序,比如创建一个Topic,默认八个写队列,八个读队列。这时候一条消息可能被写入任意一个队列里;在数据的读取过程中,可能有多个Consumer,每个Consumer也可能启动多个线程并行处理,所以消息被哪个Consumer消费,被消费的顺序和写入的顺序是否一致是不确定的。

要保证全局顺序消息,需要先把Topic的读写队列数设置为一,然后Producer和Consumer的并发设置也要是一。简单来说,为了保证整个Topic的全局消息有序,只能消除所有的并发处理,各部分都设置成单线程处理。这时高并发、高吞吐量的功能完全用不上了。

在实际应用中,更多的是像订单类消息那样,只需要部分有序即可。在这种情况下,我们经过合适的配置,依然可以利用RocketMQ高并发、高吞吐量的能力。

1.2部分顺序消息

要保证部分消息有序,需要发送端和消费端配合处理。在发送端,要做到把同一业务ID的消息发送到同一个Message Queue;在消费过程中,要做到从同一个Message Queue读取的消息不被并发处理,这样才能达到部分有序。

发送端使用MessageQueueSelector类来控制把消息发往哪个Message Queue。消费端通过使用MessageListenerOrderly类来解决单Message Queue的消息被并发处理的问题。


2.消息重复问题

产生问题的原因

RocketMQ的At least Once 机制保证消息不丢失,但是可能会造成消息重复,RocketMQ 中无法避免消息重复(Exactly-Once),在互联网应用中,尤其在网络不稳定的情况下,以下几种情况会导致消息重复问题:

  1. 生产者重复发送:生产者发送消息到 Broker 后,如果没有收到响应(例如网络闪断),会触发重试机制(默认重试 2 次),导致 Broker 存储重复的消息
  2. 消费者Offset提交失败:消费者处理消息后,Offset 提交到 Broker 是异步的(默认 5 秒一次),若消费者宕机或重启,未提交的 Offset 会导致重复消费
  3. 负载均衡导致的:消费者组扩容/缩容时,队列重新分配可能导致部分消息被重复拉取

解决方案

由于 RocketMQ 无法完全避免消息重复,业务层需实现幂等性,确保同一消息多次处理结果一致;

  • 唯一标识法

    生产者生成唯一ID,每条消息携带唯一业务 ID(如订单号+时间戳),消费者通过 Redis 或数据库记录已处理 ID。

// 生产者发送消息时添加唯一ID
Message msg = new Message("Topic", "Tag", "Order-20231001-001".getBytes());

​ 消费者处理前检查唯一 ID 是否已存在:

String msgId = message.getKeys();
if (redis.exists(msgId)) {
    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; // 已处理则跳过
}
redis.setex(msgId, 3600, "processed"); // 设置过期时间
  • 数据库约束-防重表

    插入数据时依赖数据库唯一约束,重复插入会抛出异常:

    CREATE TABLE orders (
        order_id VARCHAR(64) PRIMARY KEY,
        amount DECIMAL
    );

3.如何保证消息不丢失

生产者端保证

  • 同步发送是最可靠的发送方式,它会等待broker的确认响应
  • 失败重试机制

Broker保证

  • 消息持久化
    • 同步刷盘:消息写入 PageCache 后立即刷盘(性能较低,可靠性高)
    • 异步刷盘(默认):依赖 OS 异步刷盘,性能高但宕机可能丢失未刷盘数据
  • 主从复制:配置主从架构,并设置同步复制

消费者保证

  • 消费者确认机制

    • ACK 机制:消费者处理成功后需返回 CONSUME_SUCCESS,否则 Broker 会重试(默认 16 次)
    • 同步提交 Offset:关闭定时提交(默认 5 秒异步提交),改为处理完成后立即提交 Offset
  • 失败处理

    消费失败的进行消息回退,重试次数过多的消息放入死信队列,最后人工补偿


文章作者: Fuchanglai
版权声明: 本博客所有文章除特別声明外,均采用 CC BY 4.0 许可协议。转载请注明来源 Fuchanglai !
赏
  目录