跳至主内容
版本:4.3

消费者确认和发布者确认 (Publisher Confirms)

概述

本指南涵盖了与数据安全相关的两个功能:消费者确认(Consumer Acknowledgements)和发布者确认(Publisher Confirms)。

以及更多内容。在消息传递应用中,消费者端和发布者端的确认机制对于数据安全至关重要。

更多相关主题请参阅发布者消费者指南。

基础知识

使用 RabbitMQ 等消息代理的系统本质上是分布式的。由于发送的协议方法(消息)无法保证到达对端或被其成功处理,因此发布者和消费者都需要一种交付和处理确认机制。RabbitMQ 支持的多种消息协议都提供了此类功能。本指南涵盖 AMQP 0-9-1 中的功能,但其理念在其他受支持的协议中基本相同。

消费者到 RabbitMQ 的交付处理确认在消息协议中称为“确认(acknowledgements)”;从代理到发布者的代理确认是一种名为发布者确认(publisher confirms)的协议扩展。这两个功能基于相同的理念,并受 TCP 的启发。

它们对于实现从发布者到 RabbitMQ 节点以及从 RabbitMQ 节点到消费者的可靠交付至关重要。换句话说,它们对于数据安全是必不可少的,应用程序与 RabbitMQ 节点一样,都需要为此负责。

发布者确认与消费者交付确认有关联吗?

发布者确认消费者交付确认是非常相似的功能,它们在不同的上下文中解决相似的问题。

  1. 消费者确认(顾名思义)涵盖了 RabbitMQ 与消费者之间的通信。
  2. 发布者确认涵盖了发布者与 RabbitMQ 之间的通信。

然而,这两个功能完全正交且互不感知。

发布者确认不感知消费者:它们仅涵盖发布者与其连接的节点以及队列(或)领导者副本之间的交互。

消费者确认不感知发布者:其目标是向 RabbitMQ 节点确认给定的交付已被成功接收和处理,以便标记已交付的消息以供日后删除。

有时发布和消费应用程序需要通过请求和响应进行通信,这需要来自对端的显式确认。RabbitMQ 教程 #6 演示了实现此功能的基础知识,而直接回复(Direct Reply-to)提供了一种无需创建响应队列即可实现的方法。

然而,本指南不涉及此类通信,仅将其与本指南所述的更专注的消息协议功能进行对比。

(消费者)交付确认

当 RabbitMQ 将消息交付给消费者时,它需要知道何时认为消息已成功发送。什么样的逻辑是最佳的取决于系统本身,因此这主要是一个应用程序的决策。在 AMQP 0-9-1 中,这是在使用 basic.consume 方法注册消费者或使用 basic.get 方法按需获取消息时决定的。

如果您偏好面向示例和循序渐进的材料,消费者确认也在RabbitMQ 教程 #2 中有详细介绍。

交付标识符:交付标签(Delivery Tags)

在讨论其他主题之前,有必要解释交付是如何被标识的(以及确认如何指明它们各自的交付)。当消费者(订阅)注册后,消息将通过 RabbitMQ 使用 basic.deliver 方法被交付(推送)。该方法携带一个“交付标签(delivery tag)”,它唯一标识通道上的交付。因此,交付标签的作用域仅限于通道。

交付标签是单调递增的正整数,客户端库会按此形式呈现。用于确认交付的客户端库方法将交付标签作为参数。

由于交付标签的作用域是通道级的,因此必须在接收到交付的同一通道上进行确认。在不同的通道上进行确认将导致“未知交付标签(unknown delivery tag)”协议异常,并关闭通道。

消费者确认模式与数据安全考量

当节点将消息交付给消费者时,它必须决定是否应认为消息已由消费者处理(或至少是已接收)。由于多种原因(客户端连接、消费者应用程序等)都可能失败,这一决定涉及数据安全。消息协议通常提供一种确认机制,允许消费者向与其连接的节点确认交付。是否使用该机制是在消费者订阅时决定的。

根据所使用的确认模式,RabbitMQ 可以在消息发送后(写入 TCP 套接字)立即认为其已成功交付,或者在收到显式(“手动”)客户端确认时认为其已成功交付。手动发送的确认可以是肯定的或否定的,并使用以下协议方法之一:

  • basic.ack 用于肯定确认
  • basic.nack 用于否定确认(注意:这是 RabbitMQ 对 AMQP 0-9-1 的扩展
  • basic.reject 用于否定确认,但与 basic.nack 相比有一个局限性

这些方法如何在客户端库 API 中呈现将在下面讨论。

肯定确认只是指示 RabbitMQ 将消息记录为已交付,随后可以将其丢弃。使用 basic.reject 进行的否定确认具有相同的效果。区别主要在于语义:肯定确认假设消息已成功处理,而否定确认则暗示交付未被处理但仍应被删除。

在自动确认模式下,消息在发送后立即被视为已成功交付。这种模式以降低交付和消费者处理的安全性为代价换取了更高的吞吐量(只要消费者能跟上处理速度)。此模式通常被称为“即发即弃(fire-and-forget)”。与手动确认模型不同,如果在成功交付之前消费者的 TCP 连接或通道关闭,服务器发送的消息将会丢失。因此,自动消息确认应被视为不安全的,并不适合所有工作负载。

使用自动确认模式时另一个需要考虑的重要因素是消费者过载。手动确认模式通常与有界的通道预取(prefetch)配合使用,该预取限制了通道上未完成(“正在进行”)交付的数量。然而,使用自动确认时,定义上没有任何限制。因此,消费者可能会被交付速率压垮,导致内存中积累大量积压数据,耗尽堆内存或被操作系统终止进程。一些客户端库会应用 TCP 反压(停止从套接字读取,直到未处理交付的积压下降到一定限制之下)。因此,自动确认模式仅推荐用于能够高效且以稳定速率处理交付的消费者。

肯定确认交付

用于交付确认的 API 方法通常在客户端库中作为通道上的操作公开。Java 客户端用户将使用 Channel#basicAckChannel#basicNack 来分别执行 basic.ackbasic.nack。以下是一个展示肯定确认的 Java 客户端示例:

// this example assumes an existing channel instance

boolean autoAck = false;
channel.basicConsume(queueName, autoAck, "a-consumer-tag",
new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag,
Envelope envelope,
AMQP.BasicProperties properties,
byte[] body)
throws IOException
{
long deliveryTag = envelope.getDeliveryTag();
// positively acknowledge a single delivery, the message will
// be discarded
channel.basicAck(deliveryTag, false);
}
});

在 .NET 客户端中,这些方法分别是 IModel#BasicAckIModel#BasicNack。以下是一个展示该客户端肯定确认的示例:

// this example assumes an existing channel (IModel) instance

var consumer = new EventingBasicConsumer(channel);
consumer.Received += (ch, ea) =>
{
var body = ea.Body.ToArray();
// positively acknowledge a single delivery, the message will
// be discarded
channel.BasicAck(ea.DeliveryTag, false);
};
String consumerTag = channel.BasicConsume(queueName, false, consumer);

批量确认交付

手动确认可以进行批量处理以减少网络流量。这可以通过将确认方法(见上文)的 multiple 字段设置为 true 来实现。请注意,basic.reject 历史上没有该字段,这也是 RabbitMQ 引入 basic.nack 作为协议扩展的原因。

multiple 字段设置为 true 时,RabbitMQ 将确认所有未完成的交付标签,直到并包括确认中指定的标签。与确认相关的所有其他事项一样,这也仅限于通道范围内。例如,假设通道 Ch 上有未确认的交付标签 5、6、7 和 8,当在该通道上收到一个 delivery_tag 设置为 8multiple 设置为 true 的确认帧时,从 5 到 8 的所有标签都将被确认。如果 multiple 设置为 false,交付 5、6 和 7 则仍未确认。

若要使用 RabbitMQ Java 客户端确认多个交付,请为 Channel#basicAckmultiple 参数传递 true

// this example assumes an existing channel instance

boolean autoAck = false;
channel.basicConsume(queueName, autoAck, "a-consumer-tag",
new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag,
Envelope envelope,
AMQP.BasicProperties properties,
byte[] body)
throws IOException
{
long deliveryTag = envelope.getDeliveryTag();
// positively acknowledge all deliveries up to
// this delivery tag
channel.basicAck(deliveryTag, true);
}
});

.NET 客户端的理念也非常相似。

// this example assumes an existing channel (IModel) instance

var consumer = new EventingBasicConsumer(channel);
consumer.Received += (ch, ea) =>
{
var body = ea.Body.ToArray();
// positively acknowledge all deliveries up to
// this delivery tag
channel.BasicAck(ea.DeliveryTag, true);
};
String consumerTag = channel.BasicConsume(queueName, false, consumer);

否定确认和重新入队

有时消费者无法立即处理交付,但其他实例可能可以。在这种情况下,可能希望将其重新入队并让另一个消费者接收和处理它。basic.rejectbasic.nack 是用于此目的的两个协议方法。

这些方法通常用于否定确认交付。此类交付可以由代理丢弃、发送到死信队列(Dead-lettered)或重新入队。此行为由 requeue 字段控制。当该字段设置为 true 时,代理将重新入队带有指定交付标签的交付(或多个交付,稍后说明)。或者,当此字段设置为 false 时,如果配置了死信交换机,消息将被路由到该交换机,否则将被丢弃。

这两个方法通常在客户端库中作为通道上的操作公开。Java 客户端用户将使用 Channel#basicRejectChannel#basicNack 来分别执行 basic.rejectbasic.nack

// this example assumes an existing channel instance

boolean autoAck = false;
channel.basicConsume(queueName, autoAck, "a-consumer-tag",
new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag,
Envelope envelope,
AMQP.BasicProperties properties,
byte[] body)
throws IOException
{
long deliveryTag = envelope.getDeliveryTag();
// negatively acknowledge, the message will
// be discarded
channel.basicReject(deliveryTag, false);
}
});
// this example assumes an existing channel instance

boolean autoAck = false;
channel.basicConsume(queueName, autoAck, "a-consumer-tag",
new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag,
Envelope envelope,
AMQP.BasicProperties properties,
byte[] body)
throws IOException
{
long deliveryTag = envelope.getDeliveryTag();
// requeue the delivery
channel.basicReject(deliveryTag, true);
}
});

在 .NET 客户端中,这些方法分别是 IModel#BasicRejectIModel#BasicNack

// this example assumes an existing channel (IModel) instance

var consumer = new EventingBasicConsumer(channel);
consumer.Received += (ch, ea) =>
{
var body = ea.Body.ToArray();
// negatively acknowledge, the message will
// be discarded
channel.BasicReject(ea.DeliveryTag, false);
};
String consumerTag = channel.BasicConsume(queueName, false, consumer);
// this example assumes an existing channel (IModel) instance

var consumer = new EventingBasicConsumer(channel);
consumer.Received += (ch, ea) =>
{
var body = ea.Body.ToArray();
// requeue the delivery
channel.BasicReject(ea.DeliveryTag, true);
};
String consumerTag = channel.BasicConsume(queueName, false, consumer);

当消息被重新入队时,如果可能,它将被放回队列中的原始位置。如果不行(由于多个消费者共享队列时存在并发交付和确认),消息将被重新入队到更靠近队列头部的位置。

重新入队的消息可能会根据它们在队列中的位置以及活跃消费者通道使用的预取值而立即准备好进行重新交付。这意味着,如果所有消费者因为无法处理交付(由于临时条件)而进行重新入队,它们将创建重新入队/重新交付循环。此类循环在网络带宽和 CPU 资源方面代价高昂。消费者实现可以跟踪重新交付的次数并彻底拒绝消息(丢弃它们)或在延迟后安排重新入队。

可以使用 basic.nack 方法一次性拒绝或重新入队多条消息。这就是它与 basic.reject 的区别所在。它接受一个额外的参数 multiple。以下是一个 Java 客户端示例:

// this example assumes an existing channel instance

boolean autoAck = false;
channel.basicConsume(queueName, autoAck, "a-consumer-tag",
new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag,
Envelope envelope,
AMQP.BasicProperties properties,
byte[] body)
throws IOException
{
long deliveryTag = envelope.getDeliveryTag();
// requeue all unacknowledged deliveries up to
// this delivery tag
channel.basicNack(deliveryTag, true, true);
}
});

.NET 客户端的工作方式非常相似。

// this example assumes an existing channel (IModel) instance

var consumer = new EventingBasicConsumer(channel);
consumer.Received += (ch, ea) =>
{
var body = ea.Body.ToArray();
// requeue all unacknowledged deliveries up to
// this delivery tag
channel.BasicNack(ea.DeliveryTag, true, true);
};
String consumerTag = channel.BasicConsume(queueName, false, consumer);

通道预取设置(QoS)

消息异步交付(发送)给客户端,在任何给定时刻,通道上可能会有多条消息处于“传输中(in flight)”。来自客户端的手动确认本质上也是异步的,但流向相反。

这意味着存在一个未确认交付的滑动窗口。

对于大多数消费者而言,限制该窗口的大小以避免消费者端的无界缓冲区(堆)增长问题是有意义的。这可以通过使用 basic.qos 方法设置“预取计数(prefetch count)”值来实现。该值定义了通道上允许的最大未确认交付数量。当数量达到配置的计数时,RabbitMQ 将停止在该通道上交付更多消息,直到至少有一条未完成的消息被确认。

值为 0 表示“无限制”,允许任意数量的未确认消息。

例如,假设通道 Ch 上有四个交付标签为 5、6、7 和 8 的交付未确认,且 Ch 的预取计数设置为 4,除非至少有一个未完成的交付被确认,否则 RabbitMQ 不会再向 Ch 推送更多交付。

当在该通道上收到一个 delivery_tag 设置为 5(或 678)的确认帧时,RabbitMQ 会注意到并再交付一条消息。确认多条消息将使多于一条的消息可供交付。

值得重申的是,交付和手动客户端确认的流程是完全异步的。因此,如果在已经有交付在途的情况下更改预取值,就会出现自然的竞争条件,并且通道上的未确认消息可能会暂时超过预取计数。

“无限制”预取

在 AMQP 0-9-1 协议中,值 0 表示“无限制”,允许通道上有任意数量的未确认消息。

重要

消息源(如队列)仍然可以引入有限的限制,使得在没有消费者交付确认的情况下,通道上的“无限制预取”必须遵守特定限制。

例如,仲裁队列(quorum queues)将消费者预取最大值限制为 2,000,以限制潜在的失控 Raft 日志增长,这类似于仲裁队列强制执行消费者交付确认超时的方式。

正如下文所述的消费者确认模式、预取和吞吐量所述,在消费者吞吐量方面,2,000 的预取值与 300 的预取值并没有显著区别。

在最大吞吐量是最高优先级的场景中,请使用带有 RabbitMQ 流协议客户端的流和分区流,而不是增加通道预取。

通道级、消费者级和全局预取

QoS 设置可以针对特定通道或特定消费者进行配置。消费者预取指南解释了此作用域的效果。

预取与轮询消费者

QoS 预取设置对使用 basic.get(“拉取 API”)获取的消息没有影响,即使在手动确认模式下也是如此。

消费者确认模式、预取和吞吐量

确认模式和 QoS 预取值对消费者吞吐量有显著影响。通常情况下,增加预取将提高向消费者交付消息的速率。自动确认模式产生最佳的交付速率。然而,在这两种情况下,已交付但尚未处理的消息数量也会增加,从而增加消费者的 RAM 消耗。

应谨慎使用自动确认模式或具有无限制预取的手动确认模式。在不确认的情况下消费大量消息的消费者,会导致它们连接到的节点上的内存消耗增长。找到合适的预取值是一个反复试验的过程,且会因工作负载而异。100 到 300 范围内的值通常能提供最佳吞吐量,且不会有压垮消费者的重大风险。更高的值通常会遇到边际收益递减规律

1 的预取值是最保守的。它会显著降低吞吐量,尤其是在消费者连接延迟较高的环境中。对于许多应用程序,更高的值将是适当且最佳的。

当消费者失败或连接丢失时:自动重新入队

使用手动确认时,任何未被确认的交付(消息)在交付发生的通道(或连接)关闭时都会自动重新入队。这包括客户端的 TCP 连接丢失、消费者应用程序(进程)失败以及通道级协议异常(下文介绍)。

请注意,检测不可用的客户端需要一段时间。

由于此行为,消费者必须准备好处理重新交付的消息,并在实现时考虑幂等性。重新交付的消息将具有一个特殊的布尔属性 redeliver,RabbitMQ 会将其设置为 true。对于第一次交付,它将被设置为 false。请注意,消费者可能会收到之前交付给另一个消费者的消息。

客户端错误:重复确认和未知标签

如果客户端对同一个交付标签进行了多次确认,RabbitMQ 将导致通道错误,例如 PRECONDITION_FAILED - unknown delivery tag 100。如果使用了未知的交付标签,也会抛出相同的通道异常。

代理抱怨“未知交付标签”的另一种情况是,当尝试在与接收到交付的通道不同的通道上进行确认(无论是肯定的还是否定的)时。交付必须在同一通道上确认。

发布者确认

网络的失败方式可能并不明显,且检测某些失败需要时间。因此,向套接字写入协议帧或一组帧(例如已发布的消息)的客户端不能假设消息已到达服务器并已成功处理。它可能在途中丢失,或者其交付可能被显著延迟。

使用标准 AMQP 0-9-1,保证消息不丢失的唯一方法是使用事务——使通道事务化,然后为每条消息或每组发布的消息执行提交。在这种情况下,事务过于重量级,会将吞吐量降低 250 倍。为了解决这个问题,引入了确认机制。它模仿了协议中已经存在的消费者确认机制。

要启用确认,客户端发送 confirm.select 方法。根据是否设置了 no-wait,代理可能会响应 confirm.select-ok。一旦在通道上使用了 confirm.select 方法,该通道就被称为处于确认模式。事务化通道无法进入确认模式,且一旦通道处于确认模式,它就无法被事务化。

一旦通道处于确认模式,代理和客户端都会对消息进行计数(在第一次 confirm.select 时从 1 开始计数)。然后,代理在处理消息时通过在同一通道上发送 basic.ack 来确认消息。delivery-tag 字段包含已确认消息的序列号。代理也可以在 basic.ack 中设置 multiple 字段,以指示已处理序列号在该数字及之前的所有消息。

发布的否定确认

在特殊情况下,当代理无法成功处理消息时,代理将发送 basic.nack 而不是 basic.ack。在此上下文中,basic.nack 的字段与 basic.ack 中的对应字段含义相同,且 requeue 字段应被忽略。通过对一条或多条消息执行 nack,代理表明它无法处理这些消息并拒绝承担责任;此时,客户端可以选择重新发布这些消息。

在通道进入确认模式后,所有随后发布的消息都将被确认或 nack 一次。对于消息何时被确认不做任何保证。没有消息会同时被确认和 nack。

仅当负责队列的 Erlang 进程中发生内部错误时,才会发送 basic.nack

代理何时确认已发布的消息?

对于不可路由的消息,代理将在交换机验证消息不会路由到任何队列(返回一个空队列列表)后发出确认。如果消息也被发布为 mandatory,则 basic.return 会在 basic.ack 之前发送给客户端。否定确认(basic.nack)的情况相同。

对于可路由的消息,当消息已被所有队列接受时,会发送 basic.ack。对于路由到持久化队列的持久消息,这意味着持久化到磁盘。对于仲裁队列,这意味着法定人数的副本已接受并向选定的领导者确认了该消息。

持久消息的确认延迟

对于路由到持久队列的持久消息,basic.ack 将在将消息持久化到磁盘后发送。RabbitMQ 消息存储会在一段时间(几百毫秒)后批量将消息持久化到磁盘,以尽量减少 fsync(2) 调用的次数,或者在队列空闲时执行此操作。

这意味着在持续负载下,basic.ack 的延迟可能会达到几百毫秒。为了提高吞吐量,强烈建议应用程序异步处理确认(作为流),或者批量发布消息并等待未完成的确认。具体 API 因客户端库而异。

发布者确认的顺序考虑

在大多数情况下,RabbitMQ 会按照消息发布的顺序向发布者确认消息(这适用于在单个通道上发布的消息)。但是,发布者确认是异步发出的,可以确认单条消息或一组消息。确认发出的确切时刻取决于消息的交付模式(持久化与瞬态)以及消息路由到的队列的属性(见上文)。换句话说,不同的消息可能在不同的时间被认为已准备好确认。这意味着确认到达的顺序可能与其各自消息的顺序不同。应用程序应尽可能不依赖确认的顺序。

发布者确认与保证交付

如果 RabbitMQ 节点在持久消息写入磁盘之前失败,则可能会丢失这些消息。例如,考虑以下场景:

  1. 客户端向持久队列发布持久消息
  2. 客户端从队列消费消息(注意消息是持久化的,队列也是持久化的),但确认功能未激活
  3. 代理节点失败并重启,以及
  4. 客户端重新连接并开始消费消息

此时,客户端可能会合理地假设消息会被再次交付。事实并非如此:重启导致代理丢失了该消息。为了保证持久性,客户端应该使用确认功能。如果发布者的通道处于确认模式,发布者将不会收到丢失消息的 ack(因为该消息尚未写入磁盘)。

限制

最大交付标签

交付标签是一个 64 位长值,因此其最大值为 9223372036854775807。由于交付标签的作用域仅限于通道,实际上发布者或消费者几乎不可能超过此值。

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