Skip to content

顺序消息

什么是顺序消息

顺序消息是支持消费者按照发送消息的先后顺序获取消息,从而实现业务场景中的顺序处理,相比其他类型消息,顺序消息在发送、存储和投递的处理过程中,更多强调多条消息间的先后顺序关系。

零基础可以先这样理解:

顺序消息不是让所有消息全局有序,而是让同一个业务键的一组消息有序。

例如订单 order-1001 的状态必须是“创建 -> 支付 -> 发货”,但订单 order-1002 可以和它并发处理。这样既能保证单个订单正确,又不会把所有订单都堵在一条队列里。

为什么顺序消息不是全局有序

如果要求所有订单事件全局有序,就只能让所有消息进入同一条队列并串行消费。这样虽然顺序最强,但吞吐量会非常差:一个订单处理慢,所有订单都要等它。RocketMQ 的顺序消息更推荐按业务键有序,也就是同一个 orderId 有序,不同 orderId 并发。

mermaid
flowchart TD
    A["全局有序"] --> B["所有消息进入一个队列"]
    B --> C["吞吐低<br/>一个慢消息阻塞所有消息"]
    D["业务键有序"] --> E["同一 orderId 进入同一队列"]
    E --> F["单个订单顺序正确<br/>不同订单可以并发"]

所以顺序消息的关键不是“所有消息排成一条线”,而是先判断业务真正需要哪一级顺序。订单状态、账户流水、同一设备上报通常需要按业务键有序;日志、通知、积分发放很多时候只需要最终处理成功,不需要严格顺序。

工作原理图

mermaid
flowchart TD
    A["订单事件\norderId=1001"] --> B["按 orderId 选择队列"]
    C["订单事件\norderId=1002"] --> B
    B --> D["Queue 0\norder-1001: 创建 -> 支付 -> 发货"]
    B --> E["Queue 1\norder-1002: 创建 -> 取消"]
    D --> F["消费者顺序拉取 Queue 0"]
    E --> G["消费者顺序拉取 Queue 1"]
    F --> H["同一订单内部顺序一致"]
    G --> H

如果两个线程同时发送同一订单的事件,发送端本身已经乱序,Broker 无法凭空恢复正确顺序。因此顺序消息要求生产端也要按业务键串行发送。

应用场景

有序事件处理、撮合交易、数据实时增量同步等场景下,异构系统间需要维持强一致的状态同步,上游的事件变更需要按照顺序传递到下游进行处理。

如何保证消息的顺序性

RocketMQ消息的顺序性分为两部分,生产顺序性和消费顺序性。

生产顺序性

RocketMQ通过生产者和服务端的协议保障单个生产者串行地发送消息,并按序存储和持久化。 保证消息生产的顺序性必须满足以下条件:

  • 单一生产者:消息生产的顺序性仅支持单一生产者,不同生产者分布在不同的系统,即时设置相同的消息组,不同生产者之间产生的消息也无法判定其先后顺序。
  • 串行发送:RocketMQ生产者客户端支持多线程访问,但如果生产者使用了多线程并行发送,则不同线程间的消息将无法判定其先后顺序。 满足以上条件的生产者,将顺序消息发送至RocketMQ后,会保证设置了同一消息组的消息,按照发送顺序存储在同一队列中。服务端顺序存储逻辑如下:
  • 相同消息组的消息按照先后顺序被存储在同一队列。
  • 不同消息组的消息可以混合在同一个队列中,且不保证延续。 顺序消息 如上图所示,消息组1和消息组4的消息混合存储在队列1中,RocketMQ保证消息组1中的消息G1-M1,G1-M2,G1-M3是按发送顺序存储,且消息组4的消息G4-M1,G4-M2也是按顺序存储,但消息组1和消息组4中的消息不涉及顺序关系。

消费顺序性

RocketMQ通过消费者和服务端的协议保障消息消费严格按照存储的先后顺序来处理。 保证消息消费的顺序性,必须满足以下条件:

  • 投递顺序:RocketMQ通过客户端SDK和服务端通信协议保障消息按照服务端存储顺序投递,但业务方消费消息时需要严格按照接收->处理->应答的语义处理消息,避免因异步处理导致消息乱序。

消费者类型为PushConsumer时,RocketMQ保证消息按照存储顺序一条一条投递给消费者,若消费者类型为SimpleConsumer,则消费者有可能一次拉取多条消息。此时,消息消费的顺序性需要由业务方自行保证。

  • 有限重试:RocketMQ顺序消息投递仅在重试次数限定范围内,即一条消息如果一直重试失败,超过最大重试次数后将不再重试,跳过这条消息消费,不会一直阻塞后续消息处理。

对于需要严格保证消费顺序的场景,请务设置合理的重试次数,避免参数不合理导致消息乱序。

生产顺序性和消费顺序性组合

如果消息需要严格按照先进先出(FIFO)的原则处理,即先发送的先消费、后发送的后消费,则必须要同时满足生产顺序性和消费顺序性。

一般业务场景下,同一个生产者可能对接多个下游消费者,不一定所有的消费者业务都需要顺序消费,您可以将生产顺序性和消费顺序性进行差异化组合,应用于不同的业务场景。例如发送顺序消息,但使用非顺序的并发消费方式来提高吞吐能力。更多组合方式如下表所示:

生产顺序消费顺序顺序性效果
设置消息组,保证消息顺序发送。顺序消费按照消息组粒度,严格保证消息顺序。 同一消息组内的消息的消费顺序和发送顺序完全一致。
设置消息组,保证消息顺序发送。并发消费并发消费,尽可能按时间顺序处理。
未设置消息组,消息乱序发送。顺序消费按队列存储粒度,严格顺序。 基于 Apache RocketMQ 本身队列的属性,消费顺序和队列存储的顺序一致,但不保证和发送顺序一致。
未设置消息组,消息乱序发送。并发消费并发消费,尽可能按照时间顺序处理。

顺序消息声明周期

生命周期

  • 初始化:消息被生产者构建并完成初始化,待发送到服务端的状态。
  • 待消费:消息被发送到服务端,对消费者可见,等待消费者消费的状态。
  • 消费中:消息被消费者获取,并按照消费者本地的业务逻辑进行处理的过程。 此时服务端会等待消费者完成消费并提交消费结果,如果一定时间后没有收到消费者的响应,Apache RocketMQ会对消息进行重试处理。具体信息,请参见消费重试。
  • 消费提交:消费者完成消费处理,并向服务端提交消费结果,服务端标记当前消息已经被处理(包括消费成功和失败)。 Apache RocketMQ 默认支持保留所有消息,此时消息数据并不会立即被删除,只是逻辑标记已消费。消息在保存时间到期或存储空间不足被删除前,消费者仍然可以回溯消息重新消费。
  • 消息删除:Apache RocketMQ按照消息保存机制滚动清理最早的消息数据,将消息从物理文件中删除。更多信息,请参见消息存储和清理机制。

消息消费失败或消费超时,会触发服务端重试逻辑,重试消息属于新的消息,原消息的生命周期已结束。

顺序消息消费失败进行消费重试时,为保障消息的顺序性,后续消息不可被消费,必须等待前面的消息消费完成后才能被处理。

使用限制

顺序消息仅支持使用MessageType为FIFO的主题,即顺序消息只能发送至类型为顺序消息的主题中,发送的消息的类型必须和主题的类型一致。

如果普通消息误发到 FIFO Topic,或 FIFO 消息发到普通 Topic,服务端可能拒绝或表现不符合预期。Topic 类型约束的目的就是让消息语义在创建 Topic 时被固定下来,避免一个 Topic 里混入普通、延迟、顺序等不同语义,后续排查时无法判断到底应该按什么规则消费。

代码 Demo:按订单号发送顺序消息

顺序消息的关键是:同一个业务键始终路由到同一个队列。下面按 orderId 选择队列,保证同一订单的创建、支付、发货事件有序。

java
DefaultMQProducer producer = new DefaultMQProducer("ORDER_FIFO_PRODUCER_GROUP");
producer.setNamesrvAddr("127.0.0.1:9876");
producer.start();

String orderId = "order-1001";
List<String> events = List.of("CREATED", "PAID", "SHIPPED");

for (String event : events) {
    Message message = new Message(
        "ORDER_FIFO_EVENT",
        event,
        orderId,
        ("{\"orderId\":\"" + orderId + "\",\"event\":\"" + event + "\"}")
            .getBytes(StandardCharsets.UTF_8)
    );

    producer.send(message, (queues, msg, arg) -> {
        String key = (String) arg;
        int index = Math.abs(key.hashCode()) % queues.size();
        return queues.get(index);
    }, orderId);
}

producer.shutdown();

消费端处理顺序消息时不要把业务逻辑再丢到异步线程池里乱序执行,否则会破坏顺序消费语义。

排查顺序错乱的 checklist

现象重点检查
同一订单事件乱序是否多个生产者或多个线程同时发送同一业务键
消费端处理乱序是否在监听器里提交到异步线程池
失败后后续消息先执行顺序消费失败重试策略是否设置合理
不同订单互相等待是否把所有消息都路由到一个队列,导致并发度过低