消费负载均衡
当消费者组中的消费者从 Apache RocketMQ 主题中拉取消息时,会使用负载均衡策略来决定如何将消息分配给消费者。负载均衡策略可以提高服务并发能力和应用的可扩展性。本主题介绍了 Apache RocketMQ 为消费者提供的负载均衡策略。
背景信息
熟悉 Apache RocketMQ 提供的负载均衡策略,有助于您在面对以下场景时确定应采取的相应措施:
灾难恢复:您可以确定在本地节点发生故障时,消息如何进行重试和切换。
消息顺序:您可以更好地理解 Apache RocketMQ 如何确保严格的先进先出(FIFO)消息顺序。
水平分区:您可以根据消息分配方式来规划流量迁移和水平扩容操作。
广播消费与集群消费
Apache RocketMQ 允许订阅同一消息的多个消费者组,并且每个消费者组可以初始化多个消费者。消费者组和消费者可以在以下场景中配置以消费消息:
消费者组间的广播消费:该场景如图左侧所示。每个消费者组初始化自己的消费者,该消费者消费所有消息。消息以一对多的关系从主题分发给多个订阅者。
此模式通常用于网关推送和配置推送等场景。
消费者组内的集群消费:该场景如图右侧所示。每个消费者组初始化多个消费者,消息被发送给组内的所有消费者。当您希望在组内实现水平流量分区和负载均衡时,此模式非常有用。
此模式适用于微服务解耦。
消费者负载均衡策略简介
在广播消费场景中,无需进行负载均衡,因为每个消费者组中只有一个消费者。
然而,在集群消费场景中,每个消费者组包含多个消费者。负载均衡策略有助于决定如何分配消息。
根据消费者类型,负载均衡策略可分为以下两种类型:
基于消息的负载均衡
使用范围
基于消息的负载均衡是 PushConsumer 和 SimpleConsumer 的唯一且默认的策略。
工作原理
基于消息的负载均衡将主题中的消息均匀地分配给消费者组中的多个消费者。
如上图所示,消费者组 A 由三个消费者组成:A1、A2 和 A3。这三个消费者共同消费主题中 Queue1 的消息。
基于消息的负载均衡确保了队列中的消息可以被多个消费者并发处理。但是,消息是随机发送给消费者的,这意味着您无法指定消息如何分配给消费者。
基于消息的负载均衡基于主题中单条消息的确认语义。消费者接收到消息后,Broker 会锁定该消息,以确保其在被消费或超时之前对其他消费者不可见。这防止了同一队列的消息被不同的消费者多次消费。
顺序消息的负载策略
对于顺序消息,消息的顺序是指消息组中多条消息的序列。这些消息必须按照它们在 Broker 上存储的完全相同的顺序进行处理。因此,基于消息的负载均衡需要确保消息组中的消息按照服务器上的存储顺序进行消费。当不同的消费者处理同一组中的消息时,系统会严格按照消息顺序锁定消息,以确保消息按顺序消费。
在上图中,Queue1 的消息组 G1 中有四条顺序消息。它们的存储顺序由 M1 到 M4 表示。在消费过程中,当消息 M1 和 M2 被消费者 A1 处理时,如果 M1 和 M2 的消费状态未提交,消费者 A2 无法并行消费消息 M3 和 M4。消费者只有在成功提交前序消息的消费状态后,才能消费后续消息。
特性
与基于队列的负载均衡相比,基于消息的负载均衡具有以下特性:
- 更均衡的消费分配。在传统的基于队列的负载均衡中,队列数量和消费者数量可能无法很好地平衡。这导致系统出现部分消费者空闲而部分消费者负载过重的情况。相比之下,基于消息的负载均衡确保了消费者之间的均衡负载,无需管理队列和消费者的数量。
- 对网络能力差异更具包容性。在在线生产环境中,由于实际网络条件或硬件规格不一致,消费者的处理能力可能存在差异。如果基于队列分配消息,可能会出现某些消费者消息积压而另一些消费者空闲的情况。相反,基于消息的负载均衡按需分配消息,从而在消费者之间实现更均衡的负载分布。
- 更简单的队列分配运维。在使用传统的基于队列的负载均衡场景中,必须确保队列数量大于或等于消费者数量,以确保没有消费者处于空闲状态。基于消息的负载均衡不存在此问题。
适用场景
由于队列中的消息是离散分配给消费者的,因此基于消息的负载均衡适用于大多数在线事件处理场景。在这些场景中,消费者仅需基础处理能力,而无需进行消息的批量聚合。至于需要聚合和批量处理消息的流处理和聚合计算等场景,基于队列的负载均衡是更好的选择。示例
消费者无需为基于消息的负载均衡执行额外配置。默认情况下,PushConsumer 和 SimpleConsumer 已启用此策略。
SimpleConsumer simpleConsumer = null;
// Consumption example 1: When push consumers consume normal messages, they need only to process messages on a message listener and do not need to consider load balancing.
MessageListener messageListener = new MessageListener() {
@Override
public ConsumeResult consume(MessageView messageView) {
System.out.println(messageView);
// Return the status based on the consumption result.
return ConsumeResult.SUCCESS;
}
};
// Consumption example 2: When simple consumers consume normal messages, they obtain and submit messages. The consumers obtain messages based on the subscribed topic and do not need to consider load balancing.
List<MessageView> messageViewList = null;
try {
messageViewList = simpleConsumer.receive(10, Duration.ofSeconds(30));
messageViewList.forEach(messageView -> {
System.out.println(messageView);
// After consumption is complete, consumers must invoke ACK to submit the consumption result.
try {
simpleConsumer.ack(messageView);
} catch (ClientException e) {
e.printStackTrace();
}
});
} catch (ClientException e) {
// If the pull fails due to system traffic throttling or other reasons, consumers must re-initiate the request to obtain the message.
e.printStackTrace();
}
基于队列的负载均衡
使用范围
对于 Broker 版本 4.x 和 3.x 的消费者,包括 PullConsumer、DefaultPushConsumer、DefaultPullConsumer 和 DefaultLitePullConsumer,只能使用基于队列的负载均衡。
工作原理
在基于队列的负载均衡策略中,同一消费者组中的消费者消费分配给它们的队列中的消息。每个队列由一个消费者消费。
如上图所示,主题中的三个队列(Queue1、Queue2 和 Queue3)被分配给消费者组中的两个消费者。由于每个队列只能分配给一个消费者,消费者 A2 被分配了两个队列。如果队列数量少于消费者数量,则部分消费者将没有队列可分配。
基于队列的负载均衡根据队列数量和消费者数量等运行数据来分配消息。每个队列绑定到一个特定的消费者。然后,每个消费者根据获取消息的消费语义进行处理:>提交位点(Offsets)>持久化位点。当消费者获取消息时,消费状态不会返回给队列。因此,为了避免多个消费者重复消费消息,每个队列只能由一个消费者消费。
基于队列的负载均衡保证了一个队列仅由一个消费者处理。但是,此策略的实现依赖于消费者和 Broker 之间的信息协商机制。
Apache RocketMQ 不保证队列中的消息一定仅由一个消费者处理。因此,当消费者数量和队列数量发生变化时,可能会出现队列分配的临时不一致,从而导致少量消息被多次处理。
特性
与基于消息的负载均衡相比,基于队列的负载均衡粒度更粗,灵活性较差。但是,基于队列的负载均衡非常适合流处理场景。它确保队列中的消息由一个消费者处理。因此,基于队列的负载均衡更适合需要处理聚合消息或批量消息的场景。
适用场景
基于队列的负载均衡适用于处理聚合消息或批量消息的场景。这些是流计算和数据聚合应用中的常见场景。
示例
消费者无需为基于队列的负载均衡执行额外配置。默认情况下,Broker 4.x 和 3.x 版本的 PullConsumer 已启用此策略。
有关示例代码的更多信息,请访问 Apache RocketMQ 代码库。
版本兼容性
基于消息的负载均衡策略在 Apache RocketMQ Broker 5.0 版本中提供。对于 Broker 4.x 和 3.x 版本,仅支持基于队列的负载均衡策略。
对于 Apache RocketMQ 5.x 版本,基于消息和基于队列的负载均衡策略均可用。哪种策略生效取决于客户端版本和消费者类型。
使用说明
在消费逻辑中实现消息幂等性。
基于消息和基于队列的负载均衡策略在添加消费者、移除消费者和 Broker 扩缩容等场景下都会触发临时的重平衡。这可能导致临时的负载不一致,并导致少量消息被多次消费。因此,必须进行去重以确保消息消费的幂等性。