前面两篇我们介绍了RocketMQ的延迟消息和事务消息实现原理,今天我们就介绍最后一个高级消息类型——顺序消息。

之前在RocketMQ第4讲——生产者一篇中大致聊了一些顺序消息的相关东西,这篇就从发送、存储和消费阶段具体介绍下。

一、消息发送

首先,我们需要保证单个生产者串行(同步)来发送顺序消息。

现在生产环境基本上都是多节点部署的,如果有多个生产者分布在不同节点上,都往同一个Topic发送顺序消息,那么根本保证不了消息的顺序性:

即使它们发送的时间顺序是对的,但最终到达Broker的顺序也是无法保证的(网络延迟),所以要保证单个生产者发送顺序消息。

在保证单个消费者发送顺序消息后,还需要保证是串行发送,RocketMQ的生产者是支持多线程发送消息的,因此在使用上如果我们利用多线程来提高发送顺序消息的并发,理论上也就无法保证绝对的顺序,比如现在发送消息-1和消息-2:

按道理说消息-1会先到达Broker,因为它发送的比较早,但是由于是多线程,所以无法保证消息-1一定比消息-2更先到达,因为线程会被调度,可能线程A执行一半就停了,而线程B还在一直执行,这就是多线程的不确定性。

因此,生产者要保证两点:单个生产者和串行发送。

二、消息存储

在确定完发送的顺序性后,就来到了Broker存储阶段。

我们知道消息是按照时间顺序追加到commitlog中的,然后会被定时分发到cosumerQueue中:

同一个消费组内,一个consumerQueue只会被一个消费者消费,而且还是按照顺序来消费的,所以我们只需要让相关的顺序消息发送到同一个consumerQueue中即可:

SendResult sendResult=producer.send(message, new MessageQueueSelector() {
    /**
     * @param mqs 该Topic下所有可选的messageQueue
     * @param message 待发送的消息
     * @param arg 发送消息时传递的参数
     * @return
     */
    @Override
    public MessageQueue select(List<MessageQueue> mqs, Message message, Object arg) {
        //订单号
        Integer id= (Integer) arg;
        //根据参数计算出一个要接收消息的MessageQueue的下标
        int index=id %mqs.size();
        return mqs.get(index);
    }
},orderId);

当我们作为MQ生产者需要发送顺序消息的时候,需要在send方法中传入一个MessageQueueSelector,其中需要重写一个select方法,这个方法就是用来定义要把消息发送到哪个MessageQueue的。

这样就实现了分区顺序消息,仅保证同一笔订单的顺序性,如果我们需要保证所有订单的顺序性,那么就直接写死一个队列,让所有的消息都发往这一个队列中即可,这就是全局顺序消息(一般不用)。

三、消费消息

现在发送和存储都保证的顺序,那么消费者只需要老老实实的按照拉取的消息一条一条消费就行了,但是还是有很多细节的。

首先消费者要保证单线程消费,RocketMQ的MessageListener回调函数提供了两种消费模式,有序模式MessageListenerOrderly和并发模式MessageListenerConcurrently,所以我们要选择MessageListenerOrderly模式接收消息:

// 注册顺序消息监听器
consumer.registerMessageListener(new MessageListenerOrderly() {
    @Override
    public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, ConsumeOrderlyContext context) {
        for (MessageExt msg : msgs) {
            System.out.printf("Receive Order msg:"+ new String(msg.getBody()));
        }
        return ConsumeOrderlyStatus.SUCCESS;
    }
});

如果是多线程,那么道理同生产者多线程发送一样,是无法保证消费的顺序性的,而且还要考虑异常的场景,也就是消费失败了怎么办?

我们之前也介绍过,正常消息如果消费失败,会进行重试,重试达到16次后会进入死信队列,后续人工介入,而对于顺序消息来说,前置的消息消费失败了,那后续的消息是没法正常消费的,如果没在逻辑上做处理很容易造成脏数据,还是拿订单的例子来说,如果创建订单失败了,那么支付订单、发货订单是无法消费成功的。

所以对于顺序消息来说,失败的场景还是比较难处理的,而且RocketMQ对顺序消息默认的重试次数是Integer的最大值次,这一旦失败后续的消息不就全部阻塞了么。

所以针对顺序消息的消费失败,就需要在业务上提前做处理,也就是让相关联的消息都失败,不过不抛错(正常提交消费点位),先记着(持久化到数据库),然后人工介入处理。

四、顺序消费实现原理

那么RocketMQ是如何保证消息消费的顺序性的呢?我们应该也或多或少听说过,没错就是“三把锁”:分布式锁、Synchronized、ReentrantLock。

4.1 分布式锁

RocketMQ用ConsumMessageOrderlyService来保证顺序消费,该类在初始化的时候会启动一个定时任务:

该定时任务会向Broker申请当前消费者负责的队列锁,将自己的消费组、客户端ID和负责的队列发往Broker,Broker就会把对应队列与这个消费者绑定,并把这个关系存储在本地:

这样一来,别的消费者如果想消费对应的队列也得来加分布式锁,因此这个分布式锁保证了同一个消费组内,一个队列只会被分配给一个消费者。

4.2 Synchronized锁

接下来消息拉取过程中,消费者会一次性拉取多条消息,并且会将拉取到的消息放入到processQueue,同时将消息提交到消费线程池进行并发消费。

所以Synchronized这把锁的目的就是为了保证同一时刻只有一个线程去消费这个队列

看到这里可能会有疑问:为什么要丢给线程池并发消费呢?用一个线程不就行了,这样还不用加锁。这是因为一个消费者可能同时负责多个队列,所以需要并发消费来提高消费速度,所以要用多线程。

4.3 ReentrantLock锁

线程获取到Synchorized后,就开始处理消费消息了,但在真正开始消费前,会先获取processQueue的consumLock,这个lock是一把ReentrantLock锁:

但是,已经有了sychronized来保证一个队列只会被一个线程消费,为什么还要加一个ReentrantLock呢?

这里其实考虑的是重平衡的问题,当我们消费者集群新增了一些消费者,发生重平衡的时候,某个队列原来可能是被消费者A消费的,经过重平衡后现在由消费者B来消费了。

这时候消费者A就需要把自己加在Broker上的锁给释放掉,而这个释放的过程中就需要保证消息不能在消费过程中就被中断了,比如消费者A正在处理一部分消息,但是消费点位还没提交,此时消费者B立马去消费队列中的消息,那就会造成重复消费。

那么如何判断消息是否正在消费中呢,就是通过ReentrantLock这把锁,也就是说释放锁的线程也需要尝试对ProcessQueue进行加锁,加锁成功才能释放,失败就说明正在消费,就等下一次重平衡。

所以ReentrantLock这把锁的目的就是保证在重平衡过程中不会出现重复消费。

这三把锁就能绝对保证消息消费的顺序性吗?

没有银弹,还是保证不了,比如现有一个Broker集群:Broker-1、Broker-2,有一组顺序消息:msg-1、msg-2、msg-3,按照我们上面的取余策略,假设这3条消息都发往Broker-1的某队列中,msg-1和msg-2发送成功后Broker-1就挂掉了,msg-3就只能发往Broker-2了,那么这一组顺序消息就会被两个消费者消费,这也就没顺序性可言了。

ps:通过上面的介绍我们知道RocketMQ是通过三把锁来实现顺序消息的,这种方式带来的问题就是会降低吞吐量,如果前面有消息阻塞,那么会导致更多的消息阻塞,所以顺序消息要慎用!

End:希望对大家有所帮助,如果有纰漏或者更好的想法,请您一定不要吝啬你的赐教🙋。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐