跳至主内容
版本:4.3

消费者

概述

本指南涵盖了与消费者相关的各类主题

等等。

术语

“消费者”一词在不同语境下有不同的含义。通常,在消息传递和流处理的语境中,消费者是指消费并确认消息的应用程序(或应用程序实例)。同一个应用程序也可以同时发布消息,从而充当发布者。

消息协议中也有“持久订阅”以进行消息投递的概念。“订阅(Subscription)”是描述此类实体的常用术语,“消费者(Consumer)”是另一个术语。RabbitMQ 支持的消息协议同时使用这两个术语,但 RabbitMQ 文档倾向于使用后者。

从这个意义上讲,消费者是一种消息投递的订阅,它必须在开始投递前注册,并可由应用程序取消。

基础知识

RabbitMQ 是一个消息代理。它接收来自发布者的消息,对其进行路由;如果有对应的队列,则将其存储以供消费,或者在有消费者的情况下立即投递给消费者。

消费者从队列中消费消息。为了消费消息,必须存在一个队列。当添加新消费者时,假设队列中已有就绪的消息,投递将立即开始。

在消费者注册时,目标队列可能为空。在这种情况下,首批投递将在新消息入队时发生。

尝试从不存在的队列进行消费将导致通道级异常,错误代码为 404 Not Found,并会导致进行该尝试的通道被关闭。

消费者标签(Consumer Tags)

每个消费者都有一个标识符,客户端库使用它来确定针对给定的投递应调用哪个处理程序。其名称因协议而异。消费者标签(Consumer tags)和订阅 ID(Subscription IDs)是两个最常用的术语。RabbitMQ 文档倾向于使用前者。

消费者标签也用于取消消费者。

消费者生命周期

消费者旨在长期运行:即在消费者的整个生命周期内,它会接收多次投递。为了消费单条消息而注册一个消费者并非最优做法。

消费者通常在应用程序启动期间注册。它们往往会与连接甚至应用程序运行得一样长久。

消费者也可以更具动态性,即根据系统事件进行注册,并在不再需要时取消订阅。这在通过 Web STOMPWeb MQTT 插件使用的 WebSocket 客户端、移动端客户端等场景中很常见。

连接恢复

客户端可能会与 RabbitMQ 断开连接。当检测到连接丢失时,消息投递停止。

一些客户端库提供了包含消费者恢复功能的自动连接恢复特性。例如 Java.NETBunny。虽然连接恢复无法覆盖 100% 的场景和工作负载,但它通常对消费应用程序非常有效,且被推荐使用。

对于其他客户端库,应用程序开发人员需要负责执行连接恢复。通常以下恢复顺序效果良好:

  • 恢复连接
  • 恢复通道
  • 恢复队列
  • 恢复交换器
  • 恢复绑定
  • 恢复消费者

换句话说,消费者通常在最后被恢复,即在目标队列及其绑定的队列就绪之后。

重要

注意,对于使用自动删除(auto-delete)和排他队列(exclusive queues)的自动恢复连接,应确保这些队列是服务器命名的

注册消费者(订阅,“推送 API”)

应用程序可以订阅以让 RabbitMQ 将入队的消息(投递)推送给它们。这是通过在队列上注册消费者(订阅)来完成的。订阅生效后,RabbitMQ 将开始投递消息。对于每次投递,系统将调用用户提供的处理程序。根据所使用的客户端库,这可能是一个用户提供的函数或遵循特定接口的对象。

成功的订阅操作会返回一个订阅标识符(消费者标签)。它稍后可用于取消该消费者。

Java 客户端

请参阅 Java 客户端指南以获取示例。

.NET 客户端

请参阅 .NET 客户端指南以获取示例。

消息属性与投递元数据

每次投递都结合了消息元数据和投递信息。不同的客户端库提供访问这些属性的方式略有不同。通常,投递处理程序可以访问投递数据结构。

以下属性是投递和路由详细信息;它们本身不是消息属性,而是由 RabbitMQ 在路由和投递时设置的:

属性类型描述
投递标签(Delivery tag)正整数

投递标识符,请参阅确认(Confirms)

重投(Redelivered)布尔值如果此消息之前曾被投递并重新入队,则设为 true
交换器(Exchange)字符串路由此消息的交换器
路由键字符串发布者使用的路由键
消费者标签(Consumer tag)字符串消费者(订阅)标识符

以下是消息属性。它们大多数是可选的。它们由发布者在发布时设置:

属性类型描述必须?
投递模式枚举(1 或 2)

2 表示“持久化”,1 表示“瞬态”。某些客户端库将此属性公开为布尔值或枚举。

类型字符串应用程序特定的消息类型,例如 "orders.created"
Headers映射(字符串 => 任意值)带有字符串头名称的任意头部映射
内容类型字符串内容类型,例如 "application/json"。由应用程序使用,而非 RabbitMQ 核心使用
内容编码字符串内容编码,例如 "gzip"。由应用程序使用,而非 RabbitMQ 核心使用
消息 ID字符串任意消息 ID
关联 ID(Correlation ID)字符串帮助关联请求与响应,请参阅教程 6
回复地址(Reply To)字符串携带响应队列名称,请参阅教程 6
过期时间(Expiration)字符串每条消息的 TTL
时间戳时间戳应用程序提供的时间戳
用户 ID字符串用户 ID,如果已设置则会被验证
应用 ID字符串应用程序名称

消息类型

消息上的类型属性是一个任意字符串,有助于应用程序交流该消息属于何种类型。它由发布者在发布时设置。该值可以是发布者和消费者商定的任何领域特定的字符串。

RabbitMQ 不会验证或使用此字段,它供应用程序和插件使用及解释。

在实践中,消息类型自然会划分为不同的组,点分隔命名约定很常见(虽然 RabbitMQ 或客户端不强制要求),例如 orders.createdlogs.lineprofiles.image.changed

如果消费者收到无法处理的类型的消息,强烈建议记录此类事件以便于故障排查。

内容类型与内容编码

内容(MIME 媒体)类型和内容编码字段允许发布者告知消费者应如何反序列化和解码消息载荷。

RabbitMQ 不会验证或使用这些字段,它们供应用程序和插件使用及解释。

例如,带有 JSON 载荷的消息应使用 application/json。如果载荷使用 LZ77 (GZip) 算法压缩,其内容编码应为 gzip

通过用逗号分隔,可以指定多种编码。

确认模式

在注册消费者时,应用程序可以选择两种投递模式之一:

  • 自动(投递无需确认,即“即发即忘”)
  • 手动(投递需要客户端确认)

消费者确认是单独文档指南的主题,与发布者确认(针对发布者的紧密相关概念)一起讨论。

通过预取限制同时投递的数量

在手动确认模式下,消费者有一种限制“在途”(通过网络传输中或已投递但未确认)投递数量的方法。这可以避免消费者过载。

该特性与消费者确认是单独文档指南的主题。

消费者容量指标

RabbitMQ 管理界面以及监控数据端点(例如用于 Prometheus 抓取的端点)会显示单个队列的消费者容量(consumer capacity,以前称为消费者利用率)指标。

该指标计算为队列能够立即向消费者投递消息的时间比例。它有助于操作员注意到可能需要为队列添加更多消费者(应用程序实例)的情况。

如果此数值小于 100%,队列领导者副本可能在以下情况下能够更快地投递消息:

  • 有更多的消费者;或者
  • 消费者处理投递的时间更短;或者
  • 消费者通道使用了更高的预取值

对于没有消费者的队列,消费者容量将为 0%。对于有在线消费者但没有消息流的队列,该值为 100%:其理念是任何数量的消费者都能维持这种投递速率。

注意,消费者容量仅仅是一个提示。消费者应用程序可以也应该收集关于其操作的更具体指标,以帮助进行规模调整和可能的容量变更。

取消消费者(取消订阅)

要取消消费者,必须知道其标识符(消费者标签)。

消费者被取消后,将不会有未来的投递分发给它。注意,可能仍会有之前已分发的“在途”投递。取消消费者既不会丢弃也不会重新排队这些投递。

被取消的消费者除了 RabbitMQ 处理 basic.cancel 方法时已经在途的投递外,不会观察到任何新的投递。所有之前未确认的投递不会受到任何影响。若要重新排队在途投递,应用程序必须关闭通道。

Java 客户端

请参阅 Java 客户端指南以获取示例。

.NET 客户端

请参阅 .NET 客户端指南以获取示例。

轮询单条消息(“拉取 API”)

危险

本节描述的机制是一种轮询形式。作为分布式系统中任何基于轮询的方法,它效率极低,特别是在队列可能在一段时间内为空的情况下。

除集成测试外,强烈建议不要使用此 AMQP 0-9-1 消费机制。

RabbitMQ 管理Prometheus 插件提供了几个指标,有助于检测使用轮询(basic.get)的应用程序。

提示

请使用长期运行的消费者而不是轮询。

使用 AMQP 0-9-1,可以通过 basic.get 协议方法逐条获取消息。消息以 FIFO(先进先出)顺序获取。和使用消费者(订阅)一样,可以使用自动或手动确认。

强烈不建议逐条获取消息,因为与常规长期运行的消费者相比,它效率极低。与任何基于轮询的算法一样,在消息发布零星且队列可能长时间保持为空的系统中,它将极其浪费资源

如有疑问,请优先使用常规长期运行的消费者。

Java 客户端

请参阅 Java 客户端指南以获取示例。

.NET 客户端

请参阅 .NET 客户端指南以获取示例。

投递确认超时

重要

从 RabbitMQ 4.3 开始,仅仲裁队列(quorum queues)支持投递确认超时。

RabbitMQ 对消费者投递确认强制执行超时。这是一个保护机制,用于检测消费者何时未确认消息投递。配置投递确认超时有助于防止磁盘数据压缩异常及节点磁盘空间耗尽。

工作原理

如果消费者在超时值内未确认其投递,其通道将因 PRECONDITION_FAILED 通道异常而关闭。错误消息如下所示:

Consumer 'consumer-tag-998754663370' on channel 1 and queue 'qq.1' in vhost '/' has timed out
waiting for a consumer acknowledgement of a delivery with delivery tag = 10. Timeout used: 180000 ms.
This timeout value can be configured, see consumers doc guide to learn more

该错误由消费者连接到的节点记录。随后,该通道上所有消费者的所有后续投递都将被重新排队。要解决 PRECONDITION_FAILED 通道异常,请重新评估您的消费者并考虑增加超时值。

RabbitMQ 的默认超时值为 30 分钟。是否强制执行超时会按一分钟的时间间隔周期性评估。不支持低于一分钟的值,也不建议低于五分钟的值。

节点级配置

超时值可在 rabbitmq.conf 中配置(以毫秒为单位):

# 30 minutes in milliseconds
consumer_timeout = 1800000
# one hour in milliseconds
consumer_timeout = 3600000

可以使用 advanced.config 停用超时。不建议这样做。

%% advanced.config
[
{rabbit, [
{consumer_timeout, undefined}
]}
].

与其完全禁用超时,不如考虑使用一个较大的值(例如几小时)。

队列级配置

从 RabbitMQ 3.12 开始,也可以按队列配置超时值。

使用策略进行队列级投递超时设置

设置 consumer-timeout 策略键。

该值必须以毫秒为单位。是否强制执行超时会按一分钟的时间间隔周期性评估。

# override consumer timeout for a group of queues using a policy
rabbitmqctl set_policy queue_consumer_timeout "with_delivery_timeout\.*" '{"consumer-timeout":3600000}' --apply-to quorum_queues

使用可选队列参数进行队列级投递超时设置

在声明队列时,设置 x-consumer-timeout 可选队列参数。超时以毫秒为单位指定。是否强制执行超时会按一分钟的时间间隔周期性评估。

限制每个通道的消费者数量

在某些可能发生消费者泄漏的场景中,限制每个通道上可以活跃的消费者数量是明智的。这可以在 rabbitmq.conf 中使用设置 consumer_max_per_channel 进行配置。

consumer_max_per_channel = 100

独占性

仅对于经典队列,当使用 AMQP 0-9-1 客户端注册消费者时,可以将 basic.consume 方法的 exclusive 标志设置为 true,以请求该消费者成为目标队列上的唯一消费者。仅当此时没有已经在该队列上注册的消费者时,调用才会成功。这确保了同一时间只有一个消费者从队列消费。

如果独占消费者被取消或死亡,应用程序有责任注册一个新的消费者以继续从队列消费。

如果需要独占消费需要消费连续性,请使用单活跃消费者

重要

仲裁队列将忽略 basic.consume 帧上的 exclusive 标志。对于仲裁队列,请改为使用单活跃消费者

单活跃消费者

单活跃消费者(SAC)可以实现同一时间只有一个消费者从队列消费。如果活跃消费者被取消或断开连接,另一个已注册的消费者将接替其位置。当消息必须按照到达队列的顺序进行消费和处理时,仅使用一个消费者进行消费非常有用(请参阅保持消息顺序)。

典型的事件顺序如下:

  • 声明一个队列,一些消费者大致在同一时间注册到该队列。
  • 第一个注册的消费者成为单活跃消费者:消息被分发给它,其他消费者被忽略。
  • 如果该队列是仲裁队列,并且一个优先级更高的消费者注册,则队列将停止向当前的活跃消费者投递消息。当所有消息都被确认后,新的消费者成为活跃消费者。
  • 当单活跃消费者因某种原因被取消或死亡时,另一个消费者被选为活跃消费者。换句话说,队列会自动故障转移到另一个消费者。有关新消费者如何被选择的更多详细信息,请参阅SAC 行为

注意,如果不启用单活跃消费者功能,消息将使用循环调度(round-robin)分发给所有消费者。

信息

本节涵盖了适用于经典和仲裁队列上的 AMQP 1.0 和 AMQP 0-9-1 客户端的单活跃消费者功能。它与流(streams)的单活跃消费者功能有显著不同。

尝试使用 AMQP 0-9-1 客户端在流上启用 SAC 不会生效。要在流上使用 SAC,必须使用原生的 RabbitMQ 流协议客户端

在仲裁和经典队列上启用单活跃消费者

可以在声明队列时启用 SAC,并将 x-single-active-consumer 参数设置为 true

使用 RabbitMQ AMQP 1.0 Java 客户端

connection.management().queue()
.name("my-queue")
.quorum()
.queue()
.singleActiveConsumer(true)
.declare();

与独占消费者的区别

AMQP 0-9-1 独占消费者相比,单活跃消费者减轻了应用程序端维持消费连续性的压力。消费者只需注册,故障转移会自动处理,无需检测活跃消费者故障并注册新消费者。

确定哪个消费者当前处于活跃状态

管理界面CLI 可以报告在启用了该功能的队列上哪个消费者是当前的活跃消费者。

AMQP 1.0 客户端的消费者活动状态通知

从仲裁队列消费的 AMQP 1.0 客户端会在 AMQP 1.0 flow 帧中接收到一个名为 rabbitmq:active链路状态属性。当启用 SAC 时,此属性指示消费者是活跃的(true)还是不活跃且处于等待状态(false)。当未启用 SAC 时,该值始终为 true

每个消费者都会接收此属性:

  • 在首次授予信用后(初始状态);
  • 每当其活动状态发生变化时(仅在启用 SAC 时相关)。

此功能特定于 AMQP 1.0,不适用于 AMQP 0-9-1 客户端。

以下使用 amqp10_client 库的 Erlang 示例展示了如何接收活动状态通知。

%% Attach a receiver opting in to receive notifications on link state property changes.
{ok, Receiver} = amqp10_client:attach_link(
Session,
#{name => <<"my-receiver">>,
role => {receiver, #{address => <<"/queues/my-queue">>}, self()},
snd_settle_mode => unsettled,
notify_when_state_properties_changed => true}),
receive {amqp10_event, {link, Receiver, attached}} -> ok
end,
ok = amqp10_client:flow_link_credit(Receiver, 10, never),

%% Wait for the initial activity status notification.
receive {amqp10_event,
{link, Receiver,
{state_properties, #{<<"rabbitmq:active">> := Active}}}} ->
io:format("Consumer active: ~p~n", [Active])
end.

初始 SAC 选择

当与经典队列一起使用时,初始活跃消费者是随机选取的,即使使用了消费者优先级也是如此。

如果该队列是仲裁队列,并且一个优先级更高的消费者注册,则队列将停止向当前的活跃消费者投递消息。当所有消息都被确认后,新的消费者成为活跃消费者。

宣布此功能的博客文章中了解有关此行为的更多信息。

SAC 和独占消费者是互斥的

尝试在 SAC 中注册独占消费者将导致错误。SAC 的定义假设会有多个消费者在线。

SAC 不能通过策略启用

单活跃消费者功能不能通过策略启用。由于 RabbitMQ 中的策略本质上是动态的,它们可能会随时变更,从而启用或禁用它们声明的功能。想象一下突然在队列上禁用单活跃消费者:代理将开始向不活跃的消费者发送消息,消息将被并行处理,这与单活跃消费者试图实现的目标完全相反。由于单活跃消费者的语义与策略的动态特性不兼容,此功能只能在声明队列时使用队列参数启用。

消费者活动状态

管理界面list_consumers CLI 命令报告消费者的 active 标志。该标志的值取决于几个参数:

  • 对于经典队列,当未启用单活跃消费者时,该标志始终为 true
  • 对于仲裁队列且未启用单活跃消费者时,该标志默认情况下为 true,如果消费者连接到的节点被怀疑宕机,则设为 false
  • 如果启用了单活跃消费者,该标志仅对于当前的单活跃消费者设为 true,队列上的其他消费者在等待活跃消费者离开时被提升,因此它们的 active 被设为 false

优先级

通常,连接到队列的活跃消费者以循环方式从队列接收消息。

消费者优先级允许您确保高优先级消费者在活跃时接收消息,只有在高优先级消费者因有效预取设置而阻塞时,消息才会流向低优先级消费者。

当使用消费者优先级时,如果存在多个具有相同高优先级的活跃消费者,消息将以循环方式分发。

消费者优先级在单独的指南中介绍。

异常处理

消费者应处理在处理投递或任何其他消费者操作期间出现的任何异常。此类异常应被记录、收集并忽略。

如果消费者由于依赖项不可用或类似原因无法处理投递,它应该明确记录这一点并取消自身注册,直到能够再次处理投递为止。这将使消费者的不可用性对 RabbitMQ 和监控系统可见。

并发考量

消费者并发主要取决于客户端库的实现细节和应用程序配置。对于大多数客户端库(如 Java、.NET、Go、Erlang),投递被分发到一个处理所有异步消费者操作的线程池(或类似结构)中。该池通常具有可控的并发度。

Java 和 .NET 客户端保证单个通道上的投递将按照它们接收到的相同顺序进行分发,无论并发度如何。注意,一旦分发,对投递的并发处理将在处理线程之间导致自然的竞争条件。

某些客户端(如 Bunny)和框架可能会选择将消费者分发池限制为单个线程(或类似),以避免在并发处理投递时的自然竞争条件。一些应用程序依赖于严格按顺序处理投递,因此必须使用并发系数 1 或在自己的代码中处理同步。能够并发处理投递的应用程序可以使用高达其可用核心数的并发度。

队列并行考量

单个 RabbitMQ 队列绑定到单个核心。请使用多个队列来提高节点上的 CPU 利用率。诸如分片(sharding)一致性哈希交换器(consistent hash exchange)等插件有助于提高并行性。

© . This site is unofficial and not affiliated with VMware.