《RocketMQ 官网》阅读笔记 RocketMQ 消费者组(Consumer Group)消费者(Consumer)

📅 2026/7/27 16:49:37 👁️ 阅读次数 📝 编程学习
《RocketMQ 官网》阅读笔记 RocketMQ 消费者组(Consumer Group)消费者(Consumer)

《RocketMQ 官网》阅读笔记 RocketMQ 消费者组(Consumer Group)

消费者组(Consumer Group)

定义

消费者组是 Apache RocketMQ 中的一个负载均衡分组,包含使用相同消费行为的消费者。
与作为运行实体的消费者不同,消费者组是逻辑资源。Apache RocketMQ 在一个消费者组中初始化多个消费者,以实现消费性能的扩展和高可用灾难恢复。
在消费者组中,消费者根据组内定义的消费行为和负载均衡策略消费消息。以下部分描述了已定义的消费行为。

  • 订阅:Apache RocketMQ 基于消费者组管理和追踪订阅。
  • 投递顺序:Apache RocketMQ Broker 通过顺序投递或并发投递方式向消费者发送消息。您可以在消费者组中配置投递方式。
  • 消费重试策略:当消费者消费消息失败时使用的重试策略。该策略包括重试次数和死信队列的设置。

模型关系

内部属性

消费者组名称

消费者组名称用于区分不同的消费者组。消费者组名称在集群内全局唯一。由用户创建和配置。

投递顺序

Apache RocketMQ 向消费者客户端投递消息的顺序。Apache RocketMQ 支持根据不同的消费场景采用顺序投递和并发投递。默认投递方式为并发投递。

消费重试策略

当消费者消费消息失败时使用的重试策略。如果消费者消费消息失败,系统将根据该策略将失败的消息重新投递给消费者进行再次消费。
消费重试策略包含以下项:

  • 最大重试次数:消息可以被重新投递的最大次数。如果消息消费失败且超过了最大重试次数,该消息将被投递到死信队列或被丢弃。
  • 重试间隔:Apache RocketMQ Broker 重新投递失败消息的间隔时间。重试间隔仅对 Push 消费者有效。

订阅 (Subscription)

与当前消费者组关联的订阅关系集合。订阅包括消费者订阅的主题以及消费者使用的消息过滤规则。消费者动态注册消费者组的订阅。Apache RocketMQ Broker 持久化订阅信息,并将订阅与消息的消费进度进行匹配。

行为约束

在 Apache RocketMQ 领域模型中,消费者管理是通过消费者分组实现的,同一组内的消费者共享消息进行消费。因此,为确保组内消息的负载均衡和正常消费,Apache RocketMQ 要求同一组内的所有消费者保持以下消费行为一致:

  • 投递顺序
  • 消费重试策略

版本兼容性

如行为约束中所述,同一组内所有消费者的投递顺序和消费重试策略需要保持一致。

  • Apache RocketMQ 服务端 5.x 版本:上述消费行为从关联的消费者组中获取。因此,同一组内所有消费者的消费行为必须保持一致,客户端无需额外关注。
  • Apache RocketMQ 服务端 3.x/4.x 历史版本:上述消费逻辑由消费者客户端接口定义。因此,在设置消费者客户端时,必须确保同一组内消费者的消费行为保持一致。

如果您使用 Apache RocketMQ 服务端 5.x 版本,但客户端使用旧版本 SDK,则消费者的消费逻辑遵循消费者客户端接口的设置。

使用说明

根据业务需求创建消费者组

在 Apache RocketMQ 中,消费者和主题之间存在多对多映射关系。我们建议您在创建消费者组前注意以下规则:

  • 保持消息投递顺序一致:消费者组内所有消费者的消息投递顺序必须一致。投递方式要么是顺序投递,要么是并发投递。我们建议不要将同一个消费者组用于不同的业务场景。
  • 业务类型一致:一个消费者组对应一个业务逻辑。不同的业务域对消息消费有不同的要求,例如消息过滤规则和消费重试策略。我们建议在不同的业务域使用不同的消费者组。我们还建议每个消费者组中包含的主题数量不超过 10 个。

避免使用自动化机制管理消费者组

在 Apache RocketMQ 架构中,消费者组是用于管理消费者状态的逻辑资源。每个消费者组都关联着各种 数据,例如消费状态、堆积消息、可观测指标和监控数据。我们建议您严格管理您的消费者组。在添加、删除、修改或查询消费者组时请务必谨慎。
Apache RocketMQ 提供自动创建消费者组的功能。但是,如果您在生产环境中启用此功能,可能会创建大量的消费者组。过多的消费者组难以管理和回收,并会导致系统资源的浪费。因此,我们建议仅在测试环境中使用此功能。

消费者(Consumer)

定义

消费者是 Apache RocketMQ 中接收并处理消息的实体。
消费者通常集成在业务系统中。它们从 Apache RocketMQ Broker 获取消息,并将消息转换为业务逻辑可感知和处理的信息。
以下因素决定了消费者的行为:

  • 消费者身份:消费者必须关联一个消费者组,以获取行为设置和消费状态。
  • 消费者类型:Apache RocketMQ 针对不同的开发场景提供了多种消费者类型,包括推模式消费者(Push Consumer)、简单消费者(Simple Consumer)和拉模式消费者(Pull Consumer)。更多信息,请参见 消费者类型。
  • 消费者本地设置:这些设置指定了消费者客户端如何根据消费者类型进行运行。例如,您可以配置消费者上的线程数和并发设置,以实现不同的传输效果。

模型关系

内部属性

消费者组名称

当前消费者关联的消费者组名称。消费者从消费者组继承其行为。消费者组是 Apache RocketMQ{#product-name} 的逻辑资源。您必须提前使用控制台或调用 API 操作来创建消费者组。

客户端 ID(Client ID)

消费者客户端的标识。此属性用于区分不同的消费者。该值在集群内必须是唯一的。客户端 ID 由 Apache RocketMQ SDK 自动生成。它主要用于日志查看和问题定位等运维目的。客户端 ID 不可修改。

通信参数

  • 接入点 (Endpoints) (必选):用于连接服务器的接入点。此接入点用于标识集群。接入点必须按格式进行配置。建议使用域名,避免使用 IP 地址,以防节点变更时无法进行热点迁移。
  • 凭证 (Credential) (可选):客户端用于身份验证的凭证。仅当服务器启用了身份识别和认证时,才需要传输此项。
  • 请求超时时间 (Request Timeout) (可选):网络请求的超时时间。

预绑定订阅列表

  • 指定消费者的订阅列表。Apache RocketMQ Broker 可以利用预绑定订阅列表在消费者初始化时(而非应用程序启动后)验证所订阅 Topic 的权限和有效性。
  • 建议在消费者初始化时指定订阅或已订阅 Topic 的列表。如果未指定订阅或变更了已订阅的 Topic,Apache RocketMQ 会动态验证这些 Topic。

消息监听器 (Message Listener)

  • 消费者在 Apache RocketMQ Broker 将消息推送给消费者后,用于调用消息消费逻辑的监听器。
  • 消息监听器的值在消费者客户端上进行配置。
  • 当您以推模式消费者身份消费消息时,必须在消费者客户端配置消息监听器。

行为约束

在 Apache RocketMQ 领域模型中,消费者管理通过消费者分组实现,同一组内的消费者共享消息进行消费。因此,为确保组内消息的正常负载和消费,Apache RocketMQ 要求同一组内的所有 消费者保持以下消费行为一致:

  • 投递顺序
  • 消费重试策略

版本兼容性

如“行为约束”所述,同一组内所有消费者的投递顺序和消费重试策略需要保持一致。

  • Apache RocketMQ 服务端 5.x 版本:上述消费行为均从关联的消费者组获取。因此,同一组内所有消费者的消费行为必须保持一致,客户端无需特别关注。
  • Apache RocketMQ 服务端 3.x/4.x 旧版本:上述消费逻辑由消费者客户端接口定义。因此,在设置消费者客户端时,必须确保同一组内消费者的消费行为一致。

如果您使用的是 Apache RocketMQ 服务端 5.x 版本,但客户端使用的是旧版本 SDK,则消费者的消费逻辑需遵循消费者客户端接口的设置。

使用说明

建议限制单个进程中的消费者数量。

Apache RocketMQ 的消费者在通信协议层面支持非阻塞传输模式。该模式具有更高的通信效率,并支持多线程并发访问。因此,在大多数场景下,单个进程中仅需为一个消费者组初始化一个消费者即可。在开发阶段,请避免使用相同的配置初始化多个消费者。

建议不要频繁创建和销毁消费者。

Apache RocketMQ 的消费者是底层资源,可以像 数据库连接池一样重复使用。您无需在每次接收消息时创建消费者,也无需在消费消息后销毁它们。如果频繁创建和销毁消费者,Broker 上会产生大量的短连接请求,这将给您的系统带来沉重的负载。
正确示例

Consumerc=ConsumerBuilder.build();for(inti=0;i<n;i++){Messagem=c.receive();//process message}c.shutdown();

错误示例

for(inti=0;i<n;i++){Consumerc=ConsumerBuilder.build();Messagem=c.receive();//process messagec.shutdown();}