跳至主内容
版本:4.3

发布者

概述

本指南涵盖了与发布者相关的各类主题

等等。

本指南重点介绍 AMQP 0-9-1,并提及 RabbitMQ 支持的其他协议(AMQP 1.0、MQTT 和 STOMP)中关键的协议特定差异。

术语

“发布者”(publisher)在不同语境下含义不同。通常在消息传递中,发布者(也称为“生产者”producer)是一个发布(生产)消息的应用程序(或应用实例)。同一个应用程序也可以消费消息,从而同时成为消费者

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

基础知识

RabbitMQ 是一个消息代理。它接收来自发布者的消息,对其进行路由,并将其存储以供消费(如果有可路由的队列),或者在有消费者存在时立即将其投递给消费者。

发布者发布消息的目标因协议而异。

在 AMQP 0-9-1 中,发布者发布到交换器。在 AMQP 1.0 中,发布发生在链路(link)上。在 MQTT 中,发布者发布到主题(topics)。最后,STOMP 支持多种目标类型:主题、队列、AMQP 0-9-1 交换器。

有关此内容的详细信息,请参阅协议特定差异部分。

发布的消息必须路由到队列(主题等)。队列(主题)可能有在线的消费者。当消息成功路由到队列且有在线消费者可以接收更多投递时,消息将被发送给消费者。

尝试发布到不存在的队列(主题)将导致通道级异常,错误代码为 404 Not Found,并导致该操作所处的通道被关闭。

发布者生命周期

发布者通常生命周期较长:即在发布者的整个生命周期内会发布多条消息。为发布单条消息而开启一个连接或通道(会话)是不优化的做法。

发布者通常在应用程序启动时开启连接。它们通常与连接甚至应用程序运行的时间一样长。

发布者可以更加动态,根据系统事件开始发布,并在不再需要时停止。这在通过 Web STOMPWeb MQTT 插件使用的 WebSocket 客户端、移动客户端等场景中很常见。

协议差异

RabbitMQ 支持的所有协议中,发布消息的过程非常相似。所有四种协议都允许用户发布带有负载(主体)和一条或多条消息属性(头信息)的消息。

所有四种协议也都支持发布者确认机制,这允许发布应用程序跟踪哪些消息已被代理成功接收,并继续发布下一批或重试发布当前批次。

差异通常更多在于所使用的术语,而非语义。消息属性也因协议而异。

AMQP 0-9-1

在 AMQP 0-9-1 中,发布发生在通道上的交换器中。交换器使用路由拓扑,通过定义一个或多个队列与交换器之间的绑定,或者源交换器和目标交换器之间的绑定来完成。成功路由的消息被存储在队列中。

每个实体的角色在 AMQP 0-9-1 概念指南中有详细说明。

发布者确认是发布者确认机制。

有几种常见的发布者错误类型,可通过不同的协议功能进行处理

  • 发布到不存在的交换器会导致通道错误,该错误会关闭通道,因此在此通道上不允许进行后续的发布(或其他任何操作)。
  • 当发布的消息无法路由到任何队列(例如,因为目标交换器没有定义绑定)且发布者将 mandatory 消息属性设置为 false(这是默认值)时,消息会被丢弃或重新发布到备用交换器(如果有的话)。
  • 当发布的消息无法路由到任何队列,且发布者将 mandatory 消息属性设置为 true 时,消息将被退回给发布者。发布者必须设置退回消息处理器才能处理此退回(例如记录错误或尝试使用其他交换器)。

AMQP 1.0

在 AMQP 1.0 中,发布发生在链路(link)的上下文中。

MQTT

在 MQTT 中,消息是在连接上发布到主题的。服务端 MQTT 连接进程通过主题交换器将消息路由到队列

当发布者选择使用 QoS 1 时,RabbitMQ 会使用 PUBACK 数据包对已发布的消息进行确认。

发布者可以向服务器提示,发布到该主题的消息必须是保留的(存储以供未来投递给新订阅者)。每个主题只能保留最新发布的一条消息。

MQTT 5.0 的 PUBACK 数据包包含一个原因代码,用于告知发布者发布是否成功。RabbitMQ 返回的原因代码包括:

  • 0 - 成功:消息被成功路由到的所有队列均已接收该消息。
  • 16 - 无匹配的订阅者:RabbitMQ 无法将消息路由到任何队列(因为主题交换器未定义绑定)。
  • 131 - 实现特定的错误:RabbitMQ 拒绝了该消息(例如,当目标经典队列不可用时)。

在 MQTT 3.1 和 3.1.1 中,除了关闭连接外,服务器没有向客户端传达发布错误的机制。

请参阅 MQTTMQTT-over-WebSockets 指南以了解更多信息。

STOMP

STOMP 客户端在连接上发布到一个或多个目的地,这些目的地在 RabbitMQ 的情况下可能具有不同的语义。

STOMP 提供了一种让服务器向发布者传达消息处理错误的方法。其发布者确认的变体被称为回执(receipts),这是客户端在发布时可以启用的功能。

请参阅 STOMP 指南STOMP-over-WebSocketsSTOMP 1.2 规范以了解更多信息。

路由

AMQP 0-9-1

AMQP 0-9-1 中的路由由交换器执行。交换器是具名的路由表。表项称为绑定。这在 AMQP 0-9-1 概念指南中有详细说明。

有几种内置的交换器类型

  • Topic(主题)
  • Fanout(扇出)
  • Direct(直连,包括默认交换器)
  • Headers

前三种类型在教程中都有示例说明。

更多交换器类型可通过插件提供。一致性哈希交换器随机路由交换器内部事件交换器都是随 RabbitMQ 提供的交换器插件。与所有插件一样,它们必须先启用才能使用。

不可路由消息的处理

客户端可能会尝试将消息发布到不存在的目标(交换器、主题、队列)。本节涵盖不同协议在处理此类情况时的差异。

RabbitMQ 收集并公开指标,可用于检测发布不可路由消息的发布者。

AMQP 0-9-1

当发布的消息无法路由到任何队列(例如,因为目标交换器没有定义绑定)且发布者将 mandatory 消息属性设置为 false(这是默认值)时,消息会被丢弃或重新发布到备用交换器(如果有的话)。

当发布的消息无法路由到任何队列,且发布者将 mandatory 消息属性设置为 true 时,消息将被退回给发布者。发布者必须设置退回消息处理器才能处理此退回(例如记录错误或尝试使用其他交换器)。

备用交换器(Alternate Exchanges)是一种 AMQP 0-9-1 交换器功能,允许客户端处理交换器无法路由的消息(即因为没有绑定的队列或没有匹配的绑定)。典型示例包括检测客户端意外或恶意发布的不可路由消息,或者实现“否则”路由语义,即将某些消息进行特殊处理,其余消息交由通用处理器处理。

MQTT

发布到新主题会为其建立一个队列。不同的主题/QoS 级别组合将使用具有不同属性的队列。因此,发布者和消费者必须使用相同的 QoS 级别。

STOMP

STOMP 支持多种不同的目的地,包括那些假设预先存在拓扑的目的地。

  • /topic:发布到没有消费者的主题将导致消息被丢弃。该主题上的第一个订阅者将为其声明一个队列。
  • /exchange:目标交换器必须存在,否则服务器将报告错误。
  • /amq/queue:目标队列必须存在,否则服务器将报告错误。
  • /queue:发布到不存在的队列将自动创建该队列。
  • /temp-queue:发布到不存在的临时队列将自动创建该队列。

指标

存在一个用于记录“不可路由且被丢弃”消息的指标。

Unroutable message metrics

在上述示例中,所有发布的消息都被当作不可路由(且非 mandatory)丢弃了。

消息属性

AMQP 0-9-1

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

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

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

投递标识符,请参阅确认

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

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

属性类型描述是否必需?
Delivery mode(投递模式)枚举 (1 或 2)

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

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

消息类型

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

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

在实践中,消息类型自然会分为几组,通常使用点分隔的命名约定(尽管 RabbitMQ 或客户端并不强制要求),例如 orders.createdlogs.lineprofiles.image.changed

如果消费者收到了未知类型的投递,强烈建议记录此类事件以便于故障排查。

内容类型与编码

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

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

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

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

发布者确认(Confirms)与数据安全

确保数据安全是应用程序、客户端库和 RabbitMQ 集群节点的共同责任。本节涵盖了许多与数据安全相关的主题。

网络可能会以不那么明显的方式失败,并且检测某些故障需要时间。因此,向 socket 写入了协议帧或帧集合(例如已发布的消息)的客户端,不能假设消息已经到达服务器并被成功处理。消息可能在传输过程中丢失,或者其投递可能被显著延迟。

为解决此问题,开发了发布者端确认机制。它模仿了协议中已存在的消费者确认机制

使用发布者确认的策略

发布者确认为应用程序开发者提供了一种跟踪 RabbitMQ 已成功接收哪些消息的机制。有几种常用的发布者确认使用策略:

  • 单独发布消息并使用流式确认(异步 API 元素:确认事件处理器、Future/Promise 等)
  • 发布一批消息并等待所有挂起的确认
  • 单独发布消息,并在继续发布之前等待其被确认。强烈不建议使用此选项,因为它对发布者吞吐量有严重的负面影响。

它们在吞吐量影响和易用性方面各不相同。

流式确认

大多数客户端库通常提供一种方式,让开发者在收到服务器发来的确认时处理它们。确认将异步到达。由于在 AMQP 0-9-1 中发布本身也是异步的,此选项允许以极小的开销进行安全发布。算法通常类似于:

  • 在通道上启用发布者确认
  • 对于每条发布的消息,添加一个映射当前序列号到该消息的条目
  • 收到正向确认 (ack) 时,删除该条目
  • 收到负向确认 (nack) 时,删除该条目并调度该消息以进行重新发布(或执行其他合适的操作)

在 RabbitMQ Java 客户端中,确认处理器通过 ConfirmCallbackConfirmListener 接口公开。必须向通道添加一个或多个监听器

批量发布

该策略涉及发布成批的消息,并等待整个批次被确认。重试也是针对批次进行的。

  • 在通道上启用发布者确认
  • 对于每个已发布的消息批次,等待所有挂起的确认
  • 当所有确认均为正向时,发布下一批次
  • 如果出现负向确认或超时,请重新发布整个批次或仅重新发布相关消息

某些客户端提供了等待所有挂起确认的便捷 API 元素。例如,在 Java 客户端中,有 Channel#waitForConfirms(timeout)

由于此方法涉及等待确认,它会对发布者吞吐量产生负面影响。批次越大,影响越小。

发布并等待

该策略可被视为反模式,记录它主要是为了完整性。它涉及发布一条消息并立即等待挂起的确认到达。可以将其视为上述批量发布策略中批大小等于 1 的情况。

此方法对吞吐量有非常显著的负面影响,不建议使用。

从连接故障中恢复

客户端与 RabbitMQ 节点之间的网络连接可能会失败。应用程序如何处理此类故障直接关系到整个系统的数据安全性。

几个 RabbitMQ 客户端支持自动恢复连接和拓扑(队列、交换器、绑定和消费者):Java、.NET、Bunny 等。

其他客户端不提供自动恢复功能,但会提供应用程序开发者如何实现恢复的示例。

许多应用程序的自动恢复过程遵循以下步骤:

  1. 重新连接到可访问的节点
  2. 恢复连接侦听器
  3. 重新打开通道
  4. 恢复通道侦听器
  5. 恢复通道 basic.qos 设置、发布者确认和事务设置

连接和通道恢复后,拓扑恢复开始。拓扑恢复包括为每个通道执行以下操作:

  1. 重新声明交换器(预定义交换器除外)
  2. 重新声明队列
  3. 恢复所有绑定
  4. 恢复所有消费者

异常处理

发布者通常会遇到两类异常:

  • 由于写入失败或超时引起的网络 I/O 异常
  • 确认投递超时

请注意,此处的“异常”是指通用意义上的错误;某些编程语言根本没有异常,因此客户端会以不同方式传达错误。本节的讨论和建议应同样适用于大多数客户端库和编程语言。

前者在写入期间可能立即发生,也可能有一定的延迟。这是因为某些类型的 I/O 故障(例如由于网络拥塞或丢包率高)需要时间来检测。在连接恢复后可以继续发布,但如果连接因警报而阻塞,除非警报解除,否则所有后续尝试都将失败。这在下面的资源警报的影响部分有更详细的介绍。

后者只有在应用程序开发者提供超时设置时才会发生。给定的应用程序合理的超时值由开发者决定。它不应低于有效的心跳超时

资源警报的影响

当集群节点出现资源警报时,集群中所有尝试发布消息的连接都将被阻塞,直到集群内所有警报解除。

当连接被阻塞时,该连接发送的任何数据都不会被读取、解析或处理。当连接解除阻塞后,所有客户端流量处理恢复。

兼容的 AMQP 0-9-1 客户端在被阻塞和解除阻塞时会收到通知

阻塞连接上的写入将超时或因 I/O 写入异常而失败。

指标

指标收集与监控对发布者而言,与应用程序中的任何其他部分一样重要。当涉及到发布者时,RabbitMQ 收集的以下几个指标特别值得关注:

发布和确认率通常不言自明。流失率之所以重要,是因为它们有助于检测未最佳化使用连接或通道的应用程序,从而导致次优的发布率并浪费资源。

不可路由消息率有助于检测发布无法路由到任何队列的消息的应用程序。例如,这可能暗示了配置错误。

客户端库也可以收集指标。RabbitMQ Java 客户端就是一个例子。这些指标可以洞察应用程序特定的架构(例如,哪个发布组件发布了不可路由的消息),而这是 RabbitMQ 节点无法推断的。

并发注意事项

并发主题完全取决于客户端库的实现细节,但可以提供一些通用建议。通常应避免在共享的“发布上下文”(AMQP 0-9-1 中的通道、STOMP 中的连接、AMQP 1.0 中的会话等)上进行发布,这被认为是不安全的。

这样做会导致线路上的数据帧帧结构不正确。这将导致连接关闭。

如果单个应用程序中只有少量并发发布者,每个发布者使用一个线程(或类似结构)是最佳解决方案。如果有大量发布者(例如数百或数千个),请使用线程池。

暂时阻塞发布

通过将内存高水位线设置为 0,从而使资源警报立即触发,可以有效地阻塞集群中的所有发布操作。

rabbitmqctl set_vm_memory_high_watermark 0

排查发布者问题

本节涵盖了一些发布者的常见问题,以及如何识别和解决它们。分布式系统中的故障有多种形态和形式,因此此列表绝非详尽。

连接失败

像任何客户端一样,发布者必须首先成功连接并成功通过认证。

潜在的连接问题范围相当广泛,有一个专门的指南

认证与授权

像任何客户端一样,发布者可能无法认证,或者没有权限访问其目标虚拟主机(vhost),或者无法发布到目标交换器。

此类失败由 RabbitMQ 记录为错误。

请参阅访问控制指南中关于认证授权故障排查的部分。

连接流失(Connection Churn)

某些应用程序为每条发布的消息开启一个新连接。这是非常低效的,也不是消息协议设计的使用方式。可以通过连接指标检测此类状况。

尽可能优先使用长连接。

连接中断

网络连接可能会失败。一些客户端库支持自动连接和拓扑恢复,另一些则使得在应用程序代码中实现连接恢复变得容易。

当连接中断时,发布将无法进行,或不会被客户端内部排队(延迟)。此外,之前序列化并写入 socket 的消息不能保证到达目标节点。因此,对于需要可靠发布和数据安全的发布者来说,至关重要的是使用发布者确认来跟踪哪些发布已被 RabbitMQ 确认。未被确认的消息在一段时间后应被视为未投递。如果是应用程序安全的,这些消息可以被重新发布。这在教程 7 和本指南的数据安全部分有介绍。

详情请参阅从网络连接故障中恢复

路由问题

发布者可以成功连接、认证并获得发布到交换器(主题、目的地)的权限。但是,这些消息可能无法路由到任何队列或消费者。这可能是由于:

  • 应用程序之间的配置不匹配,例如发布者和消费者使用的主题不匹配
  • 发布者配置错误(交换器、主题、路由键不正确)
  • 对于 AMQP 0-9-1,目标交换器上缺少绑定
  • 资源警报生效:请参阅下文
  • 网络连接已失败且客户端未能恢复:请参阅上文

检查拓扑和指标通常有助于快速缩小问题范围。例如,管理 UI 中各个交换器的页面可用于确认是否存在入站消息活动(入口速率大于零)以及绑定情况是什么。

在以下示例中,该交换器没有绑定,因此没有任何消息会被路由到任何地方。

An exchange without bindings

绑定也可以使用 rabbitmq-diagnostics 列出。

# note that the implicit default exchange bindings won't
# be listed as of RabbitMQ 3.8
rabbitmq-diagnostics list_bindings --vhost "/"
=> Listing bindings for vhost /...

在上述示例中,该命令未产生任何结果。

从 RabbitMQ 3.8 开始,有一个用于记录不可路由消息丢弃的新指标。

Unroutable message metrics

在上述示例中,所有发布的消息都被当作不可路由(且非 mandatory)丢弃了。请参阅本指南中的不可路由消息的处理部分。

集群范围和连接指标以及服务器日志将有助于发现生效的资源警报。

资源警报

当资源警报生效时,所有发布消息的连接都将被阻塞,直到警报解除。客户端可以选择接收阻塞通知。请在资源警报指南中了解更多信息。

协议异常

使用某些协议(例如 AMQP 0-9-1 和 STOMP)时,发布者可能会遇到协议错误(异常)。例如,发布到不存在的交换器或将交换器绑定到不存在的交换器将导致通道异常并使通道关闭。关闭的通道无法进行发布。此类事件会被发布者所连接的 RabbitMQ 节点记录。根据所使用的客户端库,失败的发布尝试也会导致客户端返回异常或错误。

在共享通道上并发发布

客户端库不支持在共享通道上进行并发发布。请在并发注意事项部分了解更多信息。

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