AMQP 1.0 流控制的十个好处
这篇博客文章概述了 AMQP 1.0 流控制相比 AMQP 0.9.1 的十个优势,并通过两个基准测试证明了显著的性能提升。此外,我们深入探讨了强大的 AMQP 1.0 流控制原语以及它们在 RabbitMQ 中的使用方法。
流量控制是管理两个节点之间数据传输速率的过程,旨在防止快速发送者淹没慢速接收者。
AMQP 1.0 协议在两个不同层面定义了流量控制
链路流量控制
在 AMQP 1.0 中,消息通过链路 (link)发送。链路要么将发送方客户端应用程序连接到 RabbitMQ 中的交换机,要么将 RabbitMQ 中的队列连接到消费方客户端应用程序。
虽然 AMQP 1.0 规范使用“发送者 (senders)”和“接收者 (receivers)”这些术语,但 RabbitMQ 文档通常指代为“发布者 (publishers)”(或“生产者”)和“消费者 (consumers)”。在讨论客户端应用程序时,这些术语可以互换使用。因此,一个客户端应用程序实例如果:
- 向 RabbitMQ 发送消息,则它是 发送者 / 发布者 / 生产者(此时 RabbitMQ 作为 接收者)。
- 从 RabbitMQ 接收消息,则它是 接收者 / 消费者(此时 RabbitMQ 作为 发送者)。
链路信用值 (Link-Credit)
AMQP 1.0 链路流量控制背后的核心思想很简单:为了接收消息,消费者必须向发送队列授予信用值 (credits)。
一个信用值对应一条消息。例如,当消费者授予 10 个信用值时,RabbitMQ 允许发送 10 条消息。这种接收方向发送方提供反馈的直观原则,确保了发送方永远不会淹没接收方。
接收方和发送方都维护自己的“链路状态”。该状态的一部分是当前的链路信用值。每传输一条消息,链路信用值就会减少 1。具体而言,发送方在发送消息时将链路信用值减 1,接收方在接收消息时也将链路信用值减 1。当发送方的链路信用值达到 0 时,它必须停止发送消息。
消息在 传输 (transfer) 帧中发送。
信用值在 流控 (flow) 帧中授予。
<field name="link-credit" type="uint"/>
正如您可能猜到的,它们被称为“流控”帧是因为这些帧承载了流量控制信息。类型 uint 代表无符号整数,其值为 0 到一个很大的数字 (2^32 - 1) 之间。
即使链路已成功建立(在 AMQP 1.0 术语中称为“挂载 (attached)”),在消费者发送第一个 flow 帧并向发送队列授予链路信用值之前,RabbitMQ 也不允许开始向消费者发送消息。
在最简单的形式中,当客户端(接收方)向队列(发送方)授予单个信用值时,队列将发送一条消息,如 AMQP 1.0 规范的 图 2.43:同步获取 所示。
Receiver Sender
=================================================================
...
flow(link-credit=1) ---------->
+---- transfer(...)
*block until transfer arrives* /
<---+
...
-----------------------------------------------------------------
一次同步获取一条消息会导致吞吐量较低。因此,客户端通常会向队列授予多个信用值,如 AMQP 1.0 规范的 图 2.45:异步通知 所示。
Receiver Sender
=====================================================================
...
<---------- transfer(...)
<---------- transfer(...)
flow(link-credit=delta) ---+ +--- transfer(...)
\ /
x
/ \
<--+ +-->
<---------- transfer(...)
<---------- transfer(...)
flow(link-credit=delta) ---+ +--- transfer(...)
\ /
x
/ \
<--+ +-->
...
---------------------------------------------------------------------
如果接收方授予 N 个信用值,并等待 所有 N 条消息到达后再授予下一次 N 个信用值,则吞吐量将高于 N=1 的图 2.43。然而,如果您仔细观察图 2.45,您会发现接收方在接收完之前所有消息之前就授予了更多信用值。这种方法能获得最高的吞吐量。例如,在图 2.45 中,接收方最初可能授予了 6 个信用值,然后每当收到 3 条消息时,就向 RabbitMQ 发送另一个 link-credit = 6 的 flow 帧。
授予链路信用值不是累加的。
当接收方发送一个 link-credit = N 的 flow 帧时,接收方会将当前信用值设置为 N,而不是增加 N 个信用值。例如,如果接收方在没有传输任何消息的情况下连续发送了两个 link-credit = 50 的 flow 帧,接收方将拥有 50 个信用值,而不是 100 个。
接收方了解其当前的处处理能力,因此始终由接收方(而非发送方)决定当前的链路信用值。发送方仅通过发送更多消息来“消耗”接收方授予的链路信用值。
接收方可以根据其当前的处处理能力动态增加或减少链路信用值的数量。
消费客户端应用程序可以动态调整希望从特定队列接收的消息数量。
这是 AMQP 1.0 链路流量控制相较于 AMQP 0.9.1 中消费者预取 (prefetch) 的一大优势。在 AMQP 0.9.1 中,basic.qos 方法适用于给定 AMQP 0.9.1 信道上的所有消费者。此外,动态更新消费者预取并不可能或不方便,正如在 #10174 中讨论的那样。
消费客户端应用程序可以动态地在同一会话的多个队列中,优先选择从哪些队列接收消息。
这是 AMQP 1.0 链路流量控制优于 AMQP 0.9.1 消费者预取的另一个优势。一旦 AMQP 0.9.1 客户端在多个队列上调用 basic.consume,它将持续从所有这些队列接收消息,直到调用 basic.cancel。
您可能会问:link-credit 的合适值是多少,客户端应该多久补充一次?通常情况下,答案是您需要使用不同的值对您的特定工作负载进行基准测试以找出答案。
与其实现复杂的算法,我建议从简单的开始:例如,客户端可以最初授予 200 个链路信用值,并且每当剩余链路信用值低于 100 时,发送一个 link-credit = 200 的流控帧。
事实上,RabbitMQ 反向操作也是如此:RabbitMQ AMQP 1.0 会话进程最初授予发布者 170 个链路信用值,当剩余信用值低于一半(即 85)且未确认消息数量少于 170 时,再次授予 170 个链路信用值。(在代理内部,AMQP 1.0 会话进程与目标队列之间始终启用发布者确认,即使没有向发布客户端发送确认也是如此。这意味着如果目标队列确认不够快,RabbitMQ 会停止向发送应用程序授予链路信用值。)请注意,这些 RabbitMQ 实现细节随时可能更改。
170 这个值可以通过 advanced.config 中的 rabbit.max_link_credit 设置进行配置。
当一个目标队列过载时,发布者可以继续高速向所有其他目标队列发布消息。
应用程序可以在同一个 AMQP 1.0 连接或会话上向多个队列发送消息(通过挂载多个链路)。假设一个简单的场景,客户端打开了两个链路:
- 链路 1 发送到一个经典队列 (classic queue)。
- 链路 2 发送到一个 5 副本的仲裁队列 (quorum queue)。
在确认消息之前,仲裁队列必须已将消息复制到大多数副本,并且每个副本都已将消息 fsync 到其本地磁盘。
相比之下,经典队列不复制消息。此外,当消息被足够快地消费和确认时,经典队列可以(随后)向发布者确认消息,而无需将其写入磁盘。因此,在这种情况下,经典队列的吞吐量将远高于仲裁队列。
AMQP 1.0 链路流量控制的美妙之处在于,RabbitMQ 可以减慢链路 2 的信用值授予,同时继续高频率地授予链路 1 的信用值。因此,即使 5 副本仲裁队列处理消息的速度不如(单副本)经典队列快,客户端也可以继续全速向经典队列发送消息。
下图复制自 之前 的 AMQP 0.9.1 博客文章

此图中的“信用值 (credit)”一词是指 RabbitMQ 对 AMQP 0.9.1 连接的内部流量控制,与 AMQP 1.0 中的链路信用值无关。
上图中的 reader 是从套接字读取 AMQP 0.9.1 帧的 Erlang 进程。该图显示,对于 AMQP 0.9.1 连接,RabbitMQ 会阻塞 reader,导致 TCP 背压被施加到客户端。因此,当单个目标队列变得过载时,RabbitMQ 会限制整个 AMQP 0.9.1 连接,从而影响到向所有其他目标队列的发布。
以下基准测试显示了当连接向多个目标队列发送消息时,AMQP 1.0 如何提供比 AMQP 0.9.1 高出数倍的吞吐量。
基准测试:两个发送者
为了将我们刚才讨论的理论付诸实践,two_senders 程序模拟了一个类似的场景。
该程序打开了一个 AMQP 0.9.1 连接和信道,以及一个 AMQP 1.0 连接和会话。
在 AMQP 0.9.1 信道和 AMQP 1.0 会话上,都有两个 goroutine 尽可能快地发布消息到一个经典队列和一个仲裁队列。这总共导致了四个目标队列:
main.go RabbitMQ
+-------------+ +----------------------------------+
| | AMQP 0.9.1 connection | |
| |#####################################| |
| +---+ |-------------------------------------| +------------------------+ |
| | P | | classic-queue-amqp-091 | |
| +---+ +------------------------+ |
| AMQP 0.9.1 channel |
| +---+ +------------------------+ |
| | P | | quorum-queue-amqp-091 | |
| +---+ |-------------------------------------| +------------------------+ |
| |#####################################| |
| | | |
| | | |
| | | |
| |#####################################| |
| +---+ |-------------------------------------| +-----------------------+ |
| | P |O============================================>O| classic-queue-amqp-10 | |
| +---+ +-----------------------+ |
| AMQP 1.0 session |
| +---+ +-----------------------+ |
| | P |O======================================+=====>O| quorum-queue-amqp-10 | |
| +-+-+ |----------------------------------|--| +-----------------------+ |
| | |##################################|##| |
| | | AMQP 1.0 connection | | |
+------|------+ | +----------------------------------+
| |
| |
Publisher AMQP 1.0 link
goroutine
按如下方式运行基准测试:
- 在 Ubuntu 机器上使用
make run-broker启动 RabbitMQ 服务器 v4.0.0-beta.6。(在 macOS 上,fsync不会物理写入盘片)。 - 使用
go run two_senders/main.go执行 Go 程序。10 秒后,Go 程序将运行完成。 - 列出每个队列中的消息数量:
./sbin/rabbitmqctl --silent list_queues name type messages --formatter=pretty_table
┌────────────────────────┬─────────┬──────────┐
│ name │ type │ messages │
├────────────────────────┼─────────┼──────────┤
│ classic-queue-amqp-091 │ classic │ 159077 │
├────────────────────────┼─────────┼──────────┤
│ quorum-queue-amqp-091 │ quorum │ 155782 │
├────────────────────────┼─────────┼──────────┤
│ classic-queue-amqp-10 │ classic │ 1089075 │
├────────────────────────┼─────────┼──────────┤
│ quorum-queue-amqp-10 │ quorum │ 148580 │
└────────────────────────┴─────────┴──────────┘
正如 AMQP 1.0 基准测试 博客文章中所解释的,仲裁队列会进行 fsync(准确地说是 fdatasync),而经典队列则不会。因此,即使没有复制,仲裁队列也会比经典队列慢得多,因为我使用的是消费级磁盘,每次 fsync 至少需要 5 毫秒。对于生产集群,建议使用 fsync 更快的高端企业级磁盘。
结果显示,单个 AMQP 0.9.1 连接向目标经典队列和目标仲裁队列发送的消息数量大致相同。这是因为 quorum-queue-amqp-091 导致整个 AMQP 0.9.1 连接在我的基准测试中每秒被阻塞(和解除阻塞)约 80 次。因此,在单个 AMQP 0.9.1 连接上向多个目标队列(classic-queue-amqp-091 和 quorum-queue-amqp-091)的发布速率受限于最慢的目标队列(quorum-queue-amqp-091)。AMQP 0.9.1 连接总共发送了 159,077 + 155,782 = 314,859 条消息。
相比之下,由于链路流量控制,RabbitMQ 仅限制了通往 quorum-queue-amqp-10 目标的链路,从而允许 AMQP 1.0 客户端继续全速向 classic-queue-amqp-10 目标发布。AMQP 1.0 连接总共发送了 1,089,075 + 148,580 = 1,237,655 条消息。
因此,在我们的简单基准测试中,AMQP 1.0 的总发送吞吐量是 AMQP 0.9.1 的四倍(!)。
当一个目标队列过载时,客户端可以继续高速从所有源队列消费。因此,AMQP 1.0 客户端可以使用单个连接同时进行高吞吐量的发布和消费。
优势 #3 描述了单个过载的目标队列如何导致 RabbitMQ 阻塞读取进程,使其无法读取任何 AMQP 0.9.1 帧。这意味着客户端不仅无法发布消息,还无法消费消息。这是因为客户端的消息确认 (acknowledgements) 不再被 RabbitMQ 处理,一旦消费者的预取 (prefetch) 达到限制,就无法向其投递新消息。
虽然这种消费限制是暂时的(AMQP 0.9.1 reader 进程每秒被阻塞和解除阻塞多次),但它会显著降低消费速率。
RabbitMQ AMQP 0.9.1 文档建议:
因此建议在可能的情况下,发布者和消费者使用单独的连接,以便消费者与可能应用于发布连接的流量控制隔离开来,避免流量控制影响到手动的消费者确认。
这导致整个 AMQP 0.9.1 客户端库生态系统都采用了这种“最佳实践”,即使用单独的连接进行发布和消费。例如 RabbitMQ AMQP 0.9.1 C++ 库 指出:
发布和消费在不同的连接上进行
一个常见的应用程序陷阱是在同一个连接上进行消费和生产。这可能导致消费速率下降,因为 RabbitMQ 会向快速发布者施加背压——根据被消费/发布的确切队列,这可能会导致恶性循环。
相比之下,AMQP 1.0 链路流量控制允许仅减慢客户端应用程序中的链路发送方。所有其他链路(无论是发送还是消费)都可以继续全速运行。
因此,在 AMQP 1.0 中,客户端可以使用单个连接进行发布和消费。
基准测试:一个发送者,一个接收者
one_sender_one_receiver 程序模拟了客户端打开两个链路的场景:
- 链路 1 从经典队列接收。
- 链路 2 向仲裁队列发送。
该程序打开了一个 AMQP 0.9.1 连接和信道,以及一个 AMQP 1.0 连接和会话。
为了准备基准测试,程序向每个经典队列写入一百万条消息。
在 AMQP 0.9.1 信道和 AMQP 1.0 会话上,都有两个 goroutine:
- 一个 goroutine(链路 1)以 200 的预取值从经典队列接收消息并确认每一条。
- 一个 goroutine(链路 2)以每批 10,000 条消息的速度发布到仲裁队列。(收到所有 10,000 条确认后,再发布下一批。)
main.go RabbitMQ
+-------------+ +----------------------------------+
| | AMQP 0.9.1 connection | |
| |#####################################| |
| +---+ |-------------------------------------| +------------------------+ |
| | C | | classic-queue-amqp-091 | |
| +---+ +------------------------+ |
| AMQP 0.9.1 channel |
| +---+ +------------------------+ |
| | P | | quorum-queue-amqp-091 | |
| +---+ |-------------------------------------| +------------------------+ |
| |#####################################| |
| | | |
| | | |
| | | |
| |#####################################| |
| +---+ |-------------------------------------| +-----------------------+ |
| | C |O<============================================O| classic-queue-amqp-10 | |
| +---+ +-----------------------+ |
| AMQP 1.0 session |
| +---+ +-----------------------+ |
| | P |O======================================+=====>O| quorum-queue-amqp-10 | |
| +-+-+ |----------------------------------|--| +-----------------------+ |
| | |##################################|##| |
| | | AMQP 1.0 connection | | |
+------|------+ | +----------------------------------+
| |
| |
Publisher or Consumer AMQP 1.0 link
goroutine
按如下方式运行基准测试:
- 在 Ubuntu 机器上使用
make run-broker启动 RabbitMQ 服务器 v4.0.0-beta.6。 - 使用
go run one_sender_one_receiver/main.go执行 Go 程序。 - 程序运行完成后,列出每个队列中的消息数量:
./sbin/rabbitmqctl --silent list_queues name type messages --formatter=pretty_table
┌────────────────────────┬─────────┬──────────┐
│ name │ type │ messages │
├────────────────────────┼─────────┼──────────┤
│ classic-queue-amqp-091 │ classic │ 990932 │
├────────────────────────┼─────────┼──────────┤
│ quorum-queue-amqp-091 │ quorum │ 172800 │
├────────────────────────┼─────────┼──────────┤
│ classic-queue-amqp-10 │ classic │ 336229 │
├────────────────────────┼─────────┼──────────┤
│ quorum-queue-amqp-10 │ quorum │ 130000 │
└────────────────────────┴─────────┴──────────┘
虽然 AMQP 0.9.1 客户端仅消费了 1,000,000 - 990,932 = 9,068 条消息,但 AMQP 1.0 客户端消费了 1,000,000 - 336,229 = 663,771 条消息。
因此,在此基准测试中,AMQP 1.0 客户端接收的消息比 AMQP 0.9.1 客户端多 73 倍(!)。
在 AMQP 0.9.1 中,消费者预取限制了未确认消息的数量。当消费者通过发送 basic.ack 帧确认消息时,RabbitMQ 会交付更多消息。
在 AMQP 1.0 中,消息确认独立于链路流量控制。消费者通过发送 处理状态 (disposition) 帧来确认消息,但这不会提示 RabbitMQ 交付更多消息。相反,客户端必须通过发送 flow 帧来补充链路信用值,以便 RabbitMQ 继续发送消息。为了方便起见,一些 AMQP 1.0 客户端库在您的应用程序确认消息时会自动发送 disposition 和 flow 帧。
投递计数 (Delivery-Count)
到目前为止,我们只了解了 流控 (flow) 帧中的一个字段:link-credit。
在以下场景中会发生什么?
Receiver Sender
=======================================================================
...
link state: link state:
link-credit = 3 link-credit = 3
flow(link-credit = 6) ---+ +--- transfer(...)
\ /
x
/ \
<--+ +-->
link state: link state:
link-credit = 5 link-credit = 6
最初,link-credit 为 3。接收方决定发送一个 flow 帧,将新的 link-credit 设置为 6。与此同时,发送方发送了一个 transfer 帧。
由于接收方在发送 flow 帧后收到了 transfer 帧,它将计算其新的链路信用值为 6 - 1 = 5。但是,因为发送方在发送 transfer 帧后收到了 flow 帧,它会将信用值设置为 6。结果,状态——以及对链路还剩多少信用值的看法——变得不一致。这是有问题的,因为发送方可能会淹没接收方。
为了防止这种不一致,链路状态和 flow 帧中都需要第二个字段:
<field name="delivery-count" type="sequence-no"/>
每传输一条消息,投递计数就会增加 1。具体而言,发送方每发送一条消息就增加投递计数,接收方每接收一条消息也就增加投递计数。
当发送方收到 flow 帧(包含链路信用值和投递计数)时,发送方根据以下 公式 设置其链路信用值:
link-credit(snd) := delivery-count(rcv) + link-credit(rcv) - delivery-count(snd).
(snd) 指发送方的链路状态,而接收方的链路状态 (rcv) 则在 flow 帧内发送给发送方。
在发送方,这个公式的意思是:“将新的链路信用值设置为我在流控帧中收到的链路信用值减去任何在途 (in-flight) 的投递。”
投递计数的作用是在事件序列中建立顺序,这些事件包括:
- 发送方发送消息
- 接收方接收消息
- 接收方授予链路信用值
- 发送方计算接收方的链路信用值
使用投递计数解决了我们之前讨论的不一致问题。
假设投递计数最初设置为 20
Receiver Sender
========================================================================================
...
link state: link state:
delivery-count = 20 delivery-count = 20
link-credit = 3 link-credit = 3
flow(delivery-count = 20,
link-credit = 6) ---+ +--- transfer(...)
\ /
x
/ \
<--+ +-->
link state: link state:
delivery-count = 21 delivery-count = 21
link-credit = 5 link-credit = 20+6-21 = 5 (above formula)
某些 AMQP 1.0 字段(包括投递计数)的类型为 序列号 (sequence-no)。这些是 32 位 RFC-1982 序列号,范围从 [0 .. 4,294,967,295] 并循环:4,294,967,295 加 1 结果为 0。
投递计数由发送方初始化,发送方在 挂载 (attach) 帧的 initial-delivery-count 字段中发送其选择的值。
发送方可以将投递计数初始化为它选择的任何值,例如 0、10 或 4,294,967,295。该值除了我们之前讨论的目的外没有内在含义:即通过比较接收方的投递计数与发送方的投递计数来确定有多少条消息正在传输中。
据我所知,AMQP 1.0 链路流量控制似乎是基于 1995 年的论文 Credit-Based Flow Control for ATM Networks。
公式:
link-credit(snd) := delivery-count(rcv) + link-credit(rcv) - delivery-count(snd).
出自 AMQP 1.0 规范第 2.6.7 节,与论文中的公式匹配:
Crd_Bal = Buf_Alloc - (Tx_Cnt - Fwd_Cnt)(1)
(见论文中)。
AMQP 1.0 规范甚至采用了论文中类似的术语,如“链路 (link)”和“节点 (node)”。此外,规范之所以称之为“delivery-count(投递计数)”,可能是因为论文中也称之为“count(计数)”(尽管规范澄清了在 AMQP 中它实际上不是计数,而是一个序列号)。
论文中的图 2 很好地说明了这个概念:

在此图中:
Tx_Cnt对应delivery-count(snd)Fwd_Cnt对应delivery-count(rcv)Crd_Bal对应link-credit(snd)Buf_Alloc对应link-credit(rcv)
论文详细解释了接收方补充链路信用值的频率和数量。如果您想成为 AMQP 1.0 流量控制方面的专家,我建议阅读 AMQP 1.0 规范和这篇论文。
消耗 (Drain)
理解了 流控 (flow) 帧中两个最重要的链路控制字段(link-credit 和 delivery-count)后,让我们看看 drain 字段:
<field name="drain" type="boolean" default="false"/>
默认情况下,drain 字段设置为 false。接收方决定是否启用消耗模式,发送方则将 drain 设置为来自接收方的最后一个已知值。
消耗 (Draining) 意味着发送方应通过发送可用消息来使用接收方的所有链路信用值。如果没有足够的消息发送,发送方仍必须按如下方式耗尽所有链路信用值:
- 将投递计数推进剩余的链路信用值数量。
- 将链路信用值设置为 0。
- 向接收方发送一个
flow帧。
通过将 drain 字段设置为 true,消费者向 RabbitMQ 请求“要么发送一个 传输 (transfer) 帧,要么发送一个 流控 (flow) 帧”。如果源队列为空,RabbitMQ 将立即仅回复一个 flow 帧。
因此,drain 字段允许消费者在同步获取消息时设置超时,如 AMQP 1.0 规范的 图 2.44:带超时的同步获取 所示。
Receiver Sender
=================================================================
...
flow(link-credit=1) ---------->
*wait for link-credit <= 0*
flow(drain=True) ---+ +--- transfer(...)
\ /
x
/ \
(1) <--+ +-->
(2) <---------- flow(...)
...
-----------------------------------------------------------------
(1) If a message is available within the timeout, it will
arrive at this point.
(2) If a message is not available within the timeout, the
drain flag will ensure that the sender promptly advances the
delivery-count until link-credit is consumed.
由于链路信用值被快速消耗,消费者可以明确判断是收到了一条消息,还是操作已超时。
回显 (Echo)
与 drain 字段类似,echo 字段:
- 默认设置为
false - 由消费者决定,因为 RabbitMQ 在当前的实现中不设置此字段
<field name="echo" type="boolean" default="false"/>
消费者可以设置 echo 字段来请求 RabbitMQ 回复一个 flow 帧。一个用例在 AMQP 1.0 规范的 图 2.46:停止传入消息 中有描述。
Receiver Sender
================================================================
...
<---------- transfer(...)
flow(..., ---+ +--- transfer(...)
link-credit=0, \ /
echo=True) x
/ \
(1) <--+ +-->
(2) <---------- flow(...)
...
----------------------------------------------------------------
(1) In-flight transfers can still arrive until the flow state
is updated at the sender.
(2) At this point no further transfers will arrive.
在 AMQP 1.0 中,消费者可以被停止/暂停,并在稍后恢复。
在 AMQP 0.9.1 中,消费者无法暂停和恢复。相反,在通过 basic.consume 注册新消费者之前,必须使用 basic.cancel 方法取消该消费者。
AMQP 1.0 允许从一个单一活跃消费者 (Single Active Consumer) 到下一个消费者的优雅切换,同时保持消息顺序。
在 AMQP 1.0 中,消费者可以先停止链路并在卸载 (detach) 链路之前确认消息(首选方式),或者直接卸载链路。直接卸载链路将导致任何未确认的消息重新入队。无论哪种方式,卸载链路都会导致下一个消费者被激活,并按原始顺序接收消息。
相比之下,AMQP 0.9.1 中的单一活跃消费者无法优雅且安全地移交给下一个消费者。当 AMQP 0.9.1 消费者通过 basic.cancel 取消消费但仍有未确认消息时,队列将激活下一个消费者。如果 AMQP 0.9.1 客户端不久后崩溃,检出给旧消费者的消息将重新入队,这可能会破坏消息顺序。为了保持消息顺序,AMQP 0.9.1 客户端必须关闭整个信道(而不先调用 basic.cancel),以便在下一个消费者激活之前消息已重新入队。
可用量 (Available)
RabbitMQ 在 流控 (flow) 帧中设置 available 字段,以通知消费者有多少条消息可用:
<field name="available" type="uint"/>
available 值只是一个近似值,因为从队列发出此信息到 flow 帧到达消费者期间,此信息可能已经过时,例如当其他客户端向此队列发布消息或从中消费消息时。
对于经典队列和仲裁队列,available 表示准备好交付的消息数量(即队列长度减去已检出给消费者的消息数量)。
对于流 (streams),available 表示已提交偏移量与最后消费偏移量之间的差值。大致上,已提交偏移量是 RabbitMQ 集群中不同流副本同意的流末尾。最后消费的偏移量可能与已提交偏移量相同,也可能滞后。available 值仅是一个估算值,因为流偏移量不一定代表一条消息。
如果启用了单一活跃消费者功能,仲裁队列将为所有非活跃(等待中)的消费者返回 available = 0。这是合理的,因为对于非活跃消费者而言,无论他们补充了多少信用值,都没有可用的消息。
在上一节中,我们了解到消费者设置 echo 字段的一个用例是停止链路。另一个用例是当消费者想了解其消费队列中可用的消息数量时。
AMQP 1.0 会告知消费者队列中可用的消息数量。
在从 RabbitMQ 发送到消费者的每个 flow 帧中包含此信息是优于 AMQP 0.9.1 的一个优势。在 AMQP 0.9.1 中,查询可用消息的方法更繁琐且效率较低:
- 使用
passive=true进行queue.declare:queue.declare_ok响应将包含message_count字段。 basic.get:basic.get_ok响应将包含message_count字段。
available 字段也可以由发布者设置,以告诉 RabbitMQ 该发布者有多少可用消息待发送。目前,RabbitMQ 忽略此信息。
属性 (Properties)
流控 (flow) 帧的 properties 字段可以携带应用程序特定的链路状态属性:
<field name="properties" type="fields"/>
目前,RabbitMQ 尚未使用此字段。
AMQP 1.0 链路流量控制具有可扩展性。
想象一下,RabbitMQ 希望发布者能够直接向仲裁队列主副本 (leader) 发送消息。由于包括入队消息在内的所有 Ra 命令都必须先通过主副本,因此让客户端直接连接到托管仲裁队列主副本的 RabbitMQ 节点是有意义的。这种“队列局部性 (queue locality)”将减少 RabbitMQ 集群内部流量,从而提高延迟和吞吐量。当主副本发生变化时,RabbitMQ 可以通过 properties 字段向发布者推送主副本变更通知(包括托管仲裁队列新主副本的新 RabbitMQ 节点),而不是导致 RabbitMQ 卸载 (detach) 链路(这可能会干扰发布应用程序)。然后,应用程序可以决定何时方便卸载链路,并在不同的连接上挂载新链路以继续“本地”发布。
或者,RabbitMQ 可以在每个 flow 帧的 properties 字段中发送一个名为 local 的布尔键值,以指示该链路当前是否在本地发布或消费。如果 local 字段从 true 变为 false,客户端可以通过 HTTP over AMQP 1.0 查询新的队列拓扑和主副本。
这些只是说明 RabbitMQ 未来如何利用链路流量控制可扩展性的假设示例,这相比 AMQP 0.9.1 是一个优势。
总结
我们了解到,AMQP 1.0 链路流量控制保护了单个消费者或队列不被消息淹没。接收方会定期向发送方提供关于其当前处理能力的反馈。
每个链路端点(发送方和接收方)维护完整的链路流量控制状态,并通过 流控 (flow) 帧与另一侧交换此状态。
以下链路流量控制字段由接收方独立决定:
link-credit (链路信用值)drain (消耗状态)
以下链路流量控制字段由发送方独立决定:
delivery-count (投递计数)available (可用量)
会话流量控制 (Session Flow Control)
在深入探讨 会话流量控制 之前,让我们先回顾一下什么是会话 (session)。
会话 (Session)
每个 AMQP 连接都会由客户端库创建一个 TCP 连接。AMQP 1.0 连接使用 开启 (open) 帧建立。
在 AMQP 1.0 连接中,客户端可以启动多个 AMQP 1.0 会话。这类似于 AMQP 0.9.1 客户端在 AMQP 0.9.1 连接内打开多个 AMQP 0.9.1 信道 (channels)。AMQP 1.0 会话使用 开始 (begin) 帧启动。
在会话内,客户端随后可以通过发送 挂载 (attach) 帧来创建链路 (links)。
Client App RabbitMQ
+-------------+ +-------------+
| |################################| |
| +---+ |--------------------------------| +---+ |
| | C |O<===============================+========O| Q | |
| +-+-+ \ |-------------------+-------|----| |+-+-+ |
| | \ |#######+###########|#######|####| | | |
+-----|-----\-+ | | | +---|--|------+
| \ | | | | |
| Target | | | Source |
| | | | |
Consumer | | Link Queue
| Session
Connection
客户端应用程序可以创建,例如:
- 具有单个会话的单个连接,或者
- 具有多个会话的单个连接,或者
- 具有一个或多个会话的多个连接
通常情况下,具有单个会话的单个连接就足够了。
打开 AMQP 连接会有一些开销:必须建立 TCP 连接,在客户端和服务器端分配套接字和 TCP 缓冲区等操作系统资源,并产生 TCP 或 TLS 握手的延迟。此外,在 RabbitMQ 节点上,会为每个传入的 AMQP 连接创建一个管理树 (supervision tree)。
会话可以被认为是“轻量级”连接。每个 AMQP 1.0 会话目前都实现为独立的 Erlang 进程。因此,如果连接已经建立,创建一个会话的成本非常低。
在 AMQP 连接中创建第二个会话在以下场景中可能很有用:
- 快速且廉价地创建另一个“虚拟连接”,而无需承担上述 AMQP 连接设置开销。
- 确保高优先级(需要低延迟)的传输 (transfer) 帧不被其他传输帧阻塞。
- 当 RabbitMQ 会话进程变得非常繁忙时(例如在将消息路由到队列时),增加并行度。
流量字段 (Flow Fields)
会话基于传输帧的数量提供流量控制方案。由于给定连接的帧有最大尺寸限制,这实际上提供了基于传输字节数的流量控制。
请记住,大消息会被分割成多个传输帧。
AMQP 1.0 中的会话流量控制工作在比链路流量控制更高的层级。链路流量控制保护的是单个消费者和队列,而会话流量控制旨在保护整个客户端应用程序和 RabbitMQ 整体。
正如每个链路端点维护链路流量控制状态并通过 flow 帧交换该状态一样,每个会话端点维护会话流量控制状态并在同一个 flow 帧内交换该状态。流控 (flow) 帧中的其余字段因此用于会话流量控制:
<field name="next-incoming-id" type="transfer-number"/>
<field name="incoming-window" type="uint" mandatory="true"/>
<field name="next-outgoing-id" type="transfer-number" mandatory="true"/>
<field name="outgoing-window" type="uint" mandatory="true"/>
incoming-window (接收窗口) 与 link-credit 字段类似,接收侧通过它通知发送侧可以容忍接收多少个“单位”。区别在于,对于 link-credit,单位是潜在很大的应用程序消息,而对于 incoming-window,单位是 transfer 帧。
为了更好地理解,举一个极端的例子:如果在连接上协商的 max-frame-size (最大帧大小) 非常小,而单条消息非常大,即便接收方向发送方授予了一个链路信用值,发送方仍可能因为会话流量控制而无法完整发送该消息。
next-incoming-id 和 next-outgoing-id 是序列号,其作用与链路流量控制中的 投递计数 相同:它们确保当 transfer 帧与 flow 帧并行发送时,窗口在另一侧能被正确计算。会话流量控制需要两个序列号,而链路流量控制只需要一个,这是因为会话是双向的,而链路是单向的。
最初,RabbitMQ 允许发布者发送 400 个传输帧。每当 RabbitMQ 会话进程处理完其中的一半(200 个传输帧)时,RabbitMQ 就会通过向发布者发送一个包含 incoming_window = 400 的 flow 帧来扩大此窗口。
400 这个值可以通过 advanced.config 中的 rabbit.max_incoming_window 设置进行配置。
RabbitMQ 告警 (Alarms)
这一规则唯一的例外是当触发 内存或磁盘告警 时。为了保护 RabbitMQ 不会耗尽内存或磁盘空间,每个会话都会通过向发布者发送一个 incoming_window = 0 的 flow 帧来关闭其接收窗口,从而有效地阻止发布者发送任何进一步的 transfer 帧。
在发生全集群范围的内存或磁盘告警时,RabbitMQ 只会阻止 AMQP 1.0 客户端发布 transfer 帧。其他操作,例如 AMQP 1.0 客户端在同一会话上进行消费,或创建旨在排空 RabbitMQ 队列(从而减少内存和磁盘使用)的新消费连接,仍将被允许。
在 AMQP 0.9.1 中,RabbitMQ 将完全阻塞从连接套接字的读取,并阻止打开新连接,直到告警清除。AMQP 0.9.1 客户端必须使用单独的发布和消费连接,才能在告警期间继续消费。
如优势 #4 和 #9 所述,因为客户端可以安全且高效地使用单个 AMQP 1.0 连接进行发布和消费,从而降低了 RabbitMQ 中的内存消耗。
AMQP 0.9.1 实际上需要的连接数量是 AMQP 1.0 的两倍。AMQP 1.0 基准测试 提供了关于可以节省多少内存的见解。
接收窗口 (Incoming-Window)
AMQP 1.0 流量控制的缺点在于其复杂性。虽然链路流量控制和会话流量控制背后的想法和出发点是合理的,但在每一层之上再实现另一层(会话流量控制之于链路流量控制)要做到万无一失且在所有场景下都高效是具有挑战性的。考虑到客户端可以:
- 独立修改会话流量控制和链路流量控制(在同一个
flow帧或不同的flow帧中)。 - 随时以不同的频率发送
flow帧。 - 动态增加或减少会话窗口或链路信用值。
- 授予链路信用值为 0(导致队列停止发送)或巨量的链路信用值(高达 40 亿)。
- 关闭其会话
incoming-window(导致服务器会话停止发送)或将其敞开(高达 40 亿)。 - 停止从其套接字读取,向服务器施加 TCP 背压。
- 加入其他特殊逻辑,如
drain=true。
在 AMQP 1.0 程序中,并发性是必不可少的。取决于编程语言,客户端或代理的不同部分由不同的线程或进程实现。
RabbitMQ 通过在不同的 Erlang 进程中运行 连接读取进程(解析帧)、连接写入进程(序列化帧)、会话进程(如路由消息)和 队列进程(存储和转发消息)来实现并发。
在一种场景下,如果客户端向发送队列授予了巨量的 link-credit,但维持了一个非常小的会话 incoming-window,此时队列进程被允许交付消息,但会话进程却不被允许将其转发给客户端。在这种情况下,消息会被缓冲在会话进程中,直到客户端的 incoming-window 允许它们被转发。
为了防范这种大量消息可能积压在会话进程中的场景,RabbitMQ 在其当前的实现中,内部从会话进程向队列进程授予链路信用值时,每批最多只授予 256 个。即使客户端授予了大量的链路信用值,对于给定的消费者,队列也只会看到最多 256 个信用值。一旦这 256 条消息被服务器会话进程发送出去,服务器会话进程将向队列授予下一批 256 个信用值。
256 这个值可以通过 advanced.config 中的 rabbit.max_queue_credit 设置进行配置。
总结
这篇博客文章解释了 AMQP 1.0 链路流量控制和会话流量控制的工作原理。
我们了解了 RabbitMQ 内部的各个进程如何保护自己免受过载:
- 队列进程受链路流量控制保护。
- 会话进程和 RabbitMQ 整体受会话流量控制保护。
- 当连接读取进程无法足够快地读取帧时,通过施加 TCP 背压来受到保护(这种情况应该很少见)。
这篇博客文章还强调了 AMQP 1.0 流量控制相较于 AMQP 0.9.1 流量控制的十大优势。其主要优点包括:
- 对消费者进行细粒度控制,允许他们随时决定希望从特定队列消费的确切消息数量。
- 当目标队列达到其限制时,在单个 AMQP 连接上进行发布和消费的吞吐量更高。我们运行了两个基准测试:
- 向两个不同的目标队列发送消息时,AMQP 1.0 的总发送吞吐量比 AMQP 0.9.1 高出四倍。
- 在极端情况下,当向一个队列发送过快而从另一个队列接收时,AMQP 1.0 的消费速率比 AMQP 0.9.1 高出 73 倍。
- 安全且高效地使用单个连接同时进行发布和消费。
- 消费者能够随时停止和恢复。
- 针对未来用例的可扩展性。
