首先要明确一个关键点:RocketMQ 保证的是一种局部顺序(Partially Ordered)​,而非全局顺序(Globally Ordered)。这意味着消息的顺序性只在某个特定维度(比如同一个订单ID)下保证,而不是整个 Topic 下的所有消息都严格有序。

保证顺序的核心在于:​将需要保证顺序的一批消息,通过相同的策略,路由到同一个消息队列(MessageQueue)中,然后由同一个消费者单线程地消费这个队列。​

保证顺序消费的三大步骤

整个过程涉及三个环节:​顺序发送顺序存储​ 和 ​顺序消费。任何一个环节出错都会导致顺序紊乱。

步骤

角色

关键动作

目的

1. 顺序发送

Producer

使用 ​MessageQueueSelector,将同一业务ID(如OrderId)的消息发送到同一个MessageQueue

从源头上保证具有顺序性的消息被集中到同一队列

2. 顺序存储

Broker

自然保证,Broker 会顺序地将消息写入到其持有的 MessageQueue 中

存储层面维持消息的写入顺序

3. 顺序消费

Consumer

以 ​MessageListenerOrderly​ 方式消费,对每个 MessageQueue ​加锁并进行单线程消费

消费层面确保一个队列只被一个线程顺序处理

1. 顺序发送 (Producer)

生产者必须将具有顺序性的一批消息(例如,同一个订单的创建、付款、发货)发送到同一个 MessageQueue。

// 示例:发送顺序消息
public class OrderProducer {
    public static void main(String[] args) throws Exception {
        DefaultMQProducer producer = new DefaultMQProducer("please_rename_unique_group_name");
        producer.start();

        for (int i = 0; i < 10; i++) {
            // 假设这是订单ID,同一订单的操作会有相同的orderId
            int orderId = i % 3; // 共有3个订单,订单ID为0, 1, 2
            String body = "Order " + orderId + " - Step " + i;
            Message msg = new Message("OrderTopic", "TagA", body.getBytes(RemotingHelper.DEFAULT_CHARSET));

            // 关键:使用MessageQueueSelector,根据orderId选择发送到哪个队列
            SendResult sendResult = producer.send(msg, new MessageQueueSelector() {
                @Override
                public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
                    // arg 就是传入的orderId
                    Integer id = (Integer) arg;
                    int index = id % mqs.size(); // 通过订单ID取模,固定选择某个队列
                    return mqs.get(index);
                }
            }, orderId); // 将orderId作为参数传入

            System.out.printf("SendResult status:%s, queue:%s, body:%s%n",
                    sendResult.getSendStatus(),
                    sendResult.getMessageQueue().getQueueId(),
                    body);
        }
        producer.shutdown();
    }
}
2. 顺序存储 (Broker)

Broker 接收到消息后,会将其顺序追加(Append)到指定的 MessageQueue 对应的日志文件上。这个过程是天然顺序的,先到的消息先存储。

3. 顺序消费 (Consumer)

这是最复杂的一环。消费者必须使用 MessageListenerOrderly来监听消息。

// 示例:顺序消费消息
public class OrderConsumer {
    public static void main(String[] args) throws Exception {
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("please_rename_unique_group_name");
        consumer.setNamesrvAddr("localhost:9876");

        // 关键:注册顺序消息监听器 MessageListenerOrderly
        consumer.registerMessageListener(new MessageListenerOrderly() {
            @Override
            public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, ConsumeOrderlyContext context) {
                context.setAutoCommit(true); // 通常设置为自动提交
                for (MessageExt msg : msgs) {
                    System.out.println("Received msg: " + new String(msg.getBody()) + ", QueueId: " + msg.getQueueId());
                    // 模拟业务处理
                    // 如果处理失败,返回 ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT;
                    // RocketMQ 会暂停当前队列片刻,然后重试这批消息,而不是跳过。
                }
                return ConsumeOrderlyStatus.SUCCESS;
            }
        });

        consumer.subscribe("OrderTopic", "*");
        consumer.start();
        System.out.println("Consumer Started.");
    }
}

总结

保证 RocketMQ 顺序消费的秘诀:

  1. 发送端​:使用 MessageQueueSelector,将同一业务ID的消息发送到同一个队列

  2. 消费端​:使用 MessageListenerOrderly进行消费,依靠 Broker 的队列锁机制实现单线程顺序处理失败暂停重试

    注意事项与局限性

  3. 性能牺牲​:顺序消费意味着并发度的下降。​一个队列只能被一个线程消费。如果要提高吞吐量,通常需要增加队列数(Topic 的 queueNum),但这会使得顺序的“局部”范围变小(例如,能保证的订单顺序从1000个变成了100个)。

  4. 失败重试​:如果某条消息消费失败,它会阻塞后续同一队列的所有消息,直到它自己被成功消费或超过最大重试次数进入死信队列。​业务逻辑需要保证幂等性,因为同一条消息可能会被重复消费。

  5. 局部顺序​:只能保证同一个 MessageQueue 内的消息顺序,​不同 MessageQueue 之间的消息消费顺序是无法保证的。这就是“局部顺序”的含义。

  6. 广播模式不支持​:顺序消费只在集群模式(CLUSTERING)​​ 下有效。在广播模式(BROADCASTING)下,每个消费者都会消费所有消息,无法保证顺序。

更多推荐