Java 客户端 API 指南
概述
本指南涵盖了 RabbitMQ Java 客户端及其公共 API。假设读者正在使用最新主版本客户端,并且熟悉基础知识。
本指南的主要章节包括
- 许可
- 支持的 JDK 和 Android 版本
- 支持时间表
- 连接到 RabbitMQ
- 连接与通道的生命周期
- 客户端提供的连接名称
- 使用交换机 (Exchanges) 和队列 (Queues)
- 发布消息
- 使用订阅进行消费
- 并发注意事项与安全性
- 网络故障自动恢复
- 消费者中的未处理异常
- 指标与监控
- 使用 AddressResolver 接口进行端点解析
- 请求/响应模式 ("RPC")
- TLS 支持
- OAuth 2 支持
API 参考手册 (JavaDoc) 可单独获取。
支持时间表
请参阅 RabbitMQ Java 库支持页面了解支持时间表。
JDK 和 Android 版本支持
该库的 5.x 版本系列需要 JDK 8,无论是在编译还是运行时。在 Android 上,这意味着仅支持 Android 7.0 或更高版本。
4.x 版本系列支持 JDK 6 以及 7.0 之前的 Android 版本。
许可证
该库是开源的,在 GitHub 上开发,并采用三重授权许可
这意味着用户可以将该库视为受上述列表中的任何一种许可证所许可。例如,用户可以选择 Apache 公共许可证 2.0 并将此客户端包含在商业产品中。受 GPLv2 许可的代码库可以选择 GPLv2,依此类推。
概述
客户端 API 在 AMQP 0-9-1 协议模型中暴露了关键实体,并增加了额外的抽象以方便使用。
RabbitMQ Java 客户端使用 com.rabbitmq.client 作为其顶层包。关键的类和接口包括
- Channel: 代表一个 AMQP 0-9-1 通道,并提供大部分操作(协议方法)。
- Connection: 代表一个 AMQP 0-9-1 连接
- ConnectionFactory: 构建
Connection实例 - Consumer: 代表一个消息消费者
- DefaultConsumer: 消费者常用的基类
- BasicProperties: 消息属性(元数据)
- BasicProperties.Builder:
BasicProperties的构建器
协议操作可通过 Channel 接口进行。Connection 用于打开通道、注册连接生命周期事件处理器以及关闭不再需要的连接。Connection 通过 ConnectionFactory 实例化,这也是配置各种连接设置(如 vhost 或用户名)的方式。
连接与通道
核心 API 类是 Connection 和 Channel,分别代表 AMQP 0-9-1 连接和通道。它们通常在使用前进行导入
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.Channel;
连接到 RabbitMQ
以下代码使用给定参数(主机名、端口号等)连接到 RabbitMQ 节点
ConnectionFactory factory = new ConnectionFactory();
// "guest"/"guest" by default, limited to localhost connections
factory.setUsername(userName);
factory.setPassword(password);
factory.setVirtualHost(virtualHost);
factory.setHost(hostName);
factory.setPort(portNumber);
Connection conn = factory.newConnection();
对于在本地运行的 RabbitMQ 节点,所有这些参数都有合理的默认值。
如果在创建连接之前未指定属性,则将使用该属性的默认值
| 属性 | 默认值 |
| 用户名 | "guest" |
| 密码 | "guest" |
| 虚拟主机 | "/" |
| 主机名 | "localhost" |
| 端口 | 常规连接使用 |
使用 URI 连接
或者,可以使用 URI
ConnectionFactory factory = new ConnectionFactory();
factory.setUri("amqp://userName:password@hostName:portNumber/virtualHost");
Connection conn = factory.newConnection();
对于在本地运行的常规 RabbitMQ 服务器,所有这些参数都有合理的默认值。
成功和失败的客户端连接事件可以在 服务器节点日志中观察到。
请注意,默认情况下,用户 guest 只能从 localhost 连接。这是为了限制生产系统中广为人知的凭据使用。
应用程序开发者可以为连接分配自定义名称。如果设置,该名称将在 RabbitMQ 节点日志以及 管理界面中显示。
然后,Connection 接口可用于打开通道
Channel channel = conn.createChannel();
通道现在可以用于发送和接收消息,正如后续章节所描述的那样。
使用端点列表
可以在连接时指定一个端点列表。将使用第一个可达的端点。在连接失败的情况下,使用端点列表使得应用程序可以在原节点宕机时连接到其他节点。
要使用多个端点,请向 ConnectionFactory#newConnection 提供一个 Address 列表。Address 代表一个主机名和端口对。
Address[] addrArr = new Address[]{ new Address(hostname1, portnumber1)
, new Address(hostname2, portnumber2)};
Connection conn = factory.newConnection(addrArr);
将尝试连接到 hostname1:portnumber1,如果失败,则连接到 hostname2:portnumber2。返回的连接是数组中第一个成功的连接(不抛出 IOException)。这完全等同于在工厂上重复设置主机和端口,每次都调用 factory.newConnection(),直到其中一个成功。
如果同时提供了 ExecutorService(使用 factory.newConnection(es, addrArr) 形式),则该线程池将与(第一个)成功的连接相关联。
如果您想更灵活地控制要连接的主机,请参阅服务发现支持。
与 RabbitMQ 断开连接
要断开连接,只需关闭通道和连接即可
channel.close();
conn.close();
请注意,关闭通道通常被认为是良好的实践,但在此处并非绝对必要——当底层的连接关闭时,它会自动完成。
客户端断开连接事件可以在 服务器节点日志中观察到。
连接与通道的生命周期
客户端连接旨在长期存活。底层协议是为长期运行的连接设计和优化的。这意味着为每个操作(例如发布一条消息)打开一个新连接是不必要的,并且极力不推荐,因为它会引入大量的网络往返和开销。
通道也旨在长期存活,但由于许多可恢复的协议错误会导致通道关闭,通道的寿命可能比其所属的连接短。为每个操作关闭和打开新通道通常是不必要的,但在某些情况下是恰当的。如果不确定,请优先考虑重用通道。
通道级异常(例如尝试从不存在的队列中消费)会导致通道关闭。关闭后的通道不能再使用,也不会从服务器接收任何事件(如消息投递)。通道级异常将由 RabbitMQ 记录,并启动通道的关闭序列(见下文)。
客户端提供的连接名称
RabbitMQ 节点关于其客户端的信息有限:
- 它们的 TCP 端点(源 IP 地址和端口)
- 使用的凭据
仅凭这些信息就可能在识别应用程序和实例时遇到问题,尤其是在凭据可以共享,客户端通过负载均衡器连接,但 Proxy Protocol 无法启用时。
为了更容易在服务器日志和管理界面中识别客户端,AMQP 0-9-1 客户端连接(包括 RabbitMQ Java 客户端)可以提供一个自定义标识符。如果设置了该标识符,它将在日志条目和管理界面中提及。该标识符被称为客户端提供的连接名称。该名称可用于标识应用程序或应用程序中的特定组件。名称是可选的;但强烈建议开发者提供一个名称,因为这将极大地简化某些运维任务。
RabbitMQ Java 客户端的 ConnectionFactory#newConnection 方法重载 接受客户端提供的连接名称。以下是上面连接示例的修改版,它提供了这样一个名称
ConnectionFactory factory = new ConnectionFactory();
factory.setUri("amqp://userName:password@hostName:portNumber/virtualHost");
// provides a custom connection name
Connection conn = factory.newConnection("app:audit component:event-consumer");
使用交换机和队列
客户端应用程序与交换机和队列一起工作,这些是协议的高级构建块。它们在使用前必须进行声明。声明任何一种类型的对象只需确保存在该名称的对象,必要时创建它。
延续上一个示例,以下代码声明了一个交换机和一个服务器命名的队列,然后将它们绑定在一起。
channel.exchangeDeclare(exchangeName, "direct", true);
String queueName = channel.queueDeclare().getQueue();
channel.queueBind(queueName, exchangeName, routingKey);
这将主动声明以下对象,这两个对象都可以通过使用附加参数进行自定义。在这里它们没有任何特殊参数。
- 一个持久的、非自动删除的 "direct" 类型交换机
- 一个非持久的、排他的、自动删除的、具有生成名称的队列
上述函数调用然后将队列绑定到具有给定路由键的交换机上。
请注意,当只有一个客户端想要使用队列时,这是一种典型的声明队列的方式:它不需要众所周知的名称,没有其他客户端可以使用它(排他),并且会自动清理(自动删除)。如果多个客户端想要共享一个具有众所周知名称的队列,则此代码将是合适的
channel.exchangeDeclare(exchangeName, "direct", true);
channel.queueDeclare(queueName, true, false, false, null);
channel.queueBind(queueName, exchangeName, routingKey);
这将主动声明
- 一个持久的、非自动删除的 "direct" 类型交换机
- 一个持久的、非排他的、非自动删除的、具有众所周知名称的队列
许多 Channel API 方法都被重载。这些 exchangeDeclare、queueDeclare 和 queueBind 的便捷短形式使用合理的默认值。还有带有更多参数的长形式,以便在需要时覆盖这些默认值,并在必要时提供完全控制。
这种“短形式、长形式”模式在整个客户端 API 中都有使用。
被动声明
队列和交换机可以“被动”声明。被动声明只是检查具有提供名称的实体是否存在。如果存在,该操作则为空操作。对于队列,成功的被动声明将返回与非被动声明相同的信息,即队列中就绪状态的消费者和消息数量。
如果实体不存在,则该操作以通道级异常失败。该通道在此之后不能使用。应该打开一个新通道。通常使用一次性(临时)通道进行被动声明。
Channel#queueDeclarePassive 和 Channel#exchangeDeclarePassive 是用于被动声明的方法。以下示例演示了 Channel#queueDeclarePassive
Queue.DeclareOk response = channel.queueDeclarePassive("queue-name");
// returns the number of messages in Ready state in the queue
response.getMessageCount();
// returns the number of consumers the queue has
response.getConsumerCount();
Channel#exchangeDeclarePassive 的返回值不包含有用的信息。因此,如果该方法返回且没有发生通道异常,则意味着交换机确实存在。
带有可选响应的操作
一些常见操作还有一个不会等待服务器响应的“no wait”版本。例如,要声明一个队列并指示服务器不要发送任何响应,请使用
channel.queueDeclareNoWait(queueName, true, false, false, null);
“no wait”版本效率更高,但提供的安全保证较低,例如它们更依赖心跳机制来检测失败的操作。如果有疑问,请从标准版本开始。“no wait”版本仅在具有高拓扑(队列、绑定)变更的场景中才需要。
删除实体和清除消息
队列或交换机可以被显式删除
channel.queueDelete("queue-name")
仅当队列为空时才可能删除它
channel.queueDelete("queue-name", false, true)
或者如果它未被使用(没有任何消费者)
channel.queueDelete("queue-name", true, false)
队列可以被清除(其所有消息被删除)
channel.queuePurge("queue-name")
发布消息
要将消息发布到交换机,请使用 Channel.basicPublish,如下所示
byte[] messageBodyBytes = "Hello, world!".getBytes();
channel.basicPublish(exchangeName, routingKey, null, messageBodyBytes);
为了精细控制,使用重载变体来指定 mandatory 标志,或发送带有预设消息属性的消息(详细信息请参阅发布者指南)
channel.basicPublish(exchangeName, routingKey, mandatory,
MessageProperties.PERSISTENT_TEXT_PLAIN,
messageBodyBytes);
这将发送一条传递模式为 2(持久化)、优先级为 1 且内容类型为 "text/plain" 的消息。使用 Builder 类构建一个包含所需属性的消息属性对象,例如
channel.basicPublish(exchangeName, routingKey,
new AMQP.BasicProperties.Builder()
.contentType("text/plain")
.deliveryMode(2)
.priority(1)
.userId("bob")
.build(),
messageBodyBytes);
此示例发布一条带有自定义标头的消息
Map<String, Object> headers = new HashMap<String, Object>();
headers.put("latitude", 51.5252949);
headers.put("longitude", -0.0905493);
channel.basicPublish(exchangeName, routingKey,
new AMQP.BasicProperties.Builder()
.headers(headers)
.build(),
messageBodyBytes);
此示例发布一条带有过期时间的消息
channel.basicPublish(exchangeName, routingKey,
new AMQP.BasicProperties.Builder()
.expiration("60000")
.build(),
messageBodyBytes);
这只是一组简短的示例,并未演示每个支持的属性。
请注意 BasicProperties 是外部类 AMQP 的内部类。
如果资源驱动的警报生效,Channel#basicPublish 的调用最终会阻塞。
通道与并发注意事项(线程安全)
应避免在线程之间共享 Channel 实例。应用程序应该为每个线程使用一个 Channel,而不是在多个线程之间共享同一个 Channel。
虽然某些通道操作可以安全地并发调用,但有些则不行,这会导致线路上的帧交错不正确、重复确认等问题。
在共享通道上并发发布可能会导致线路上的帧交错不正确,触发连接级协议异常,并导致代理立即关闭连接。因此,它需要在应用程序代码中进行显式同步(Channel#basicPublish 必须在临界区中调用)。在线程之间共享通道也会干扰发布者确认 (Publisher Confirms)。最好完全避免在线程之间共享通道,例如通过为每个线程使用一个通道。
可以使用通道池来避免在共享通道上并发发布:一旦线程使用完通道,它就会将其返回到池中,从而使通道可供另一个线程使用。通道池可以被视为一种特定的同步解决方案。建议使用现有的池化库,而不是自己动手开发解决方案。例如,Spring AMQP 自带了开箱即用的通道池功能。
通道消耗资源,在大多数情况下,应用程序在同一个 JVM 进程中很少需要超过几百个打开的通道。如果我们假设应用程序为每个通道拥有一个线程(因为通道不应该并发使用),那么单个 JVM 的数千个线程已经是不小的开销,这很可能是可以避免的。此外,少数几个快速发布者很容易使网络接口和代理节点饱和:发布涉及的工作量少于路由、存储和投递消息。
要避免的一个经典反模式是为每条发布的消息打开一个通道。通道应该是相当长期的,打开一个新的通道是一个网络往返,这使得这种模式效率极低。
在一个线程中消费,并在另一个线程中在共享通道上发布是可以安全的。
服务器推送的投递(见下文)是并发分发的,并保证保留每个通道的顺序。分发机制使用一个 java.util.concurrent.ExecutorService,每个连接一个。通过 ConnectionFactory#setSharedExecutor 设置器,可以提供一个由单个 ConnectionFactory 产生的所有连接共享的自定义执行器。
当使用手动确认时,考虑哪个线程执行确认很重要。如果它与接收投递的线程不同(例如 Consumer#handleDelivery 将投递处理委托给了另一个线程),则将 multiple 参数设置为 true 进行确认是不安全的,并且会导致重复确认,从而导致通道级协议异常并关闭通道。一次确认一条消息是可以安全的。
通过订阅接收消息(“推送 API”)
import com.rabbitmq.client.Consumer;
import com.rabbitmq.client.DefaultConsumer;
接收消息的最有效方式是使用 Consumer 接口设置订阅。这样消息会在到达时自动投递,而不是必须显式请求。
在调用与 Consumer 相关的 API 方法时,单个订阅总是通过其消费者标签来引用。消费者标签是一个消费者标识符,可以是客户端生成的,也可以是服务器生成的。要让 RabbitMQ 生成节点范围内唯一的标签,请使用不带消费者标签参数的 Channel#basicConsume 重载,或者为消费者标签传递一个空字符串,并使用 Channel#basicConsume 返回的值。消费者标签用于取消消费者。
不同的 Consumer 实例必须具有不同的消费者标签。强烈不建议在连接上使用重复的消费者标签,这可能导致自动连接恢复出现问题,并在监控消费者时产生令人困惑的数据。
实现 Consumer 的最简单方法是继承便捷类 DefaultConsumer。此子类的对象可以在 basicConsume 调用上传递以设置订阅
boolean autoAck = false;
channel.basicConsume(queueName, autoAck, "myConsumerTag",
new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag,
Envelope envelope,
AMQP.BasicProperties properties,
byte[] body)
throws IOException
{
String routingKey = envelope.getRoutingKey();
String contentType = properties.getContentType();
long deliveryTag = envelope.getDeliveryTag();
<i>// (process the message components here ...)</i>
channel.basicAck(deliveryTag, false);
}
});
在这里,由于我们指定了 autoAck = false,因此必须确认投递给 Consumer 的消息,最方便的做法是在 handleDelivery 方法中完成,如图所示。
更复杂的 Consumer 将需要覆盖更多方法。特别是,当通道和连接关闭时会调用 handleShutdownSignal,而在对该 Consumer 的任何其他回调被调用之前,会将消费者标签传递给 handleConsumeOk。
Consumer 还可以实现 handleCancelOk 和 handleCancel 方法,以分别接收显式和隐式取消的通知。
您可以使用 Channel.basicCancel 显式取消特定的 Consumer
channel.basicCancel(consumerTag);
传递消费者标签。
就像发布者一样,考虑消费者的并发危害安全性非常重要。
对 Consumer 的回调是在与其实例化其 Channel 的线程不同的线程池中分发的。这意味着 Consumer 可以安全地调用 Connection 或 Channel 上的阻塞方法,例如 Channel#queueDeclare 或 Channel#basicCancel。
每个 Channel 都会按照 RabbitMQ 发送它们的顺序,将所有投递分发到其 Consumer 处理程序方法中。通道之间的投递顺序无法保证:这些投递可以并行分发。
对于每个 Channel 一个 Consumer 的最常见用例,这意味着 Consumer 不会阻碍其他 Consumer。对于每个 Channel 有多个 Consumer 的情况,请注意长期运行的 Consumer 可能会阻碍该 Channel 上对其他 Consumer 的回调分发。
请参阅并发注意事项(线程安全)部分,了解与并发和并发危害安全性相关的其他主题。
检索单个消息(“拉取 API”)
还可以按需检索单个消息(“拉取 API”,即轮询)。这种消费方法效率极低,因为它实际上是轮询,应用程序必须反复请求结果,即使绝大多数请求都没有结果。因此,极力不推荐使用这种方法。
要“拉取”一条消息,请使用 Channel.basicGet 方法。返回的值是 GetResponse 的实例,从中可以提取标题信息(属性)和消息正文
boolean autoAck = false;
GetResponse response = channel.basicGet(queueName, autoAck);
if (response == null) {
// No message retrieved.
} else {
AMQP.BasicProperties props = response.getProps();
byte[] body = response.getBody();
long deliveryTag = response.getEnvelope().getDeliveryTag();
// ...
由于此示例使用了手动确认(上面的 autoAck = false),您还必须调用 Channel.basicAck 来确认您已成功接收到该消息
// ...
channel.basicAck(method.deliveryTag, false); // acknowledge receipt of the message
}
处理不可路由的消息
如果发布消息时设置了 "mandatory" 标志,但无法路由,代理将把它返回给发送客户端(通过 AMQP.Basic.Return 命令)。
要接收此类返回的通知,客户端可以实现 ReturnListener 接口并调用 Channel.addReturnListener。如果客户端尚未为特定通道配置返回侦听器,则相关的返回消息将被静默丢弃。
channel.addReturnListener(new ReturnListener() {
public void handleReturn(int replyCode,
String replyText,
String exchange,
String routingKey,
AMQP.BasicProperties properties,
byte[] body)
throws IOException {
...
}
});
例如,如果客户端将带有 "mandatory" 标志的消息发布到未绑定到队列的 "direct" 类型交换机,则将调用返回侦听器。
关闭协议
客户端关闭过程概述
AMQP 0-9-1 连接和通道在管理网络故障、内部故障和显式本地关闭时使用相同的通用方法。
AMQP 0-9-1 连接和通道具有以下生命周期状态
-
open: 对象可以使用 -
closing: 对象已被显式通知在本地关闭,已向任何支持的底层对象发出关闭请求,并正在等待其关闭程序完成 -
closed: 对象已收到来自任何底层对象的所有关闭完成通知,因此已自行关闭
无论关闭的原因是什么(如应用程序请求、内部客户端库故障、远程网络请求或网络故障),这些对象最终总是处于关闭状态。
连接和通道对象具有以下与关闭相关的方法
-
addShutdownListener(ShutdownListener listener)和 -
removeShutdownListener(ShutdownListener listener),用于管理任何侦听器,当对象转换为closed状态时将触发这些侦听器。注意,向已关闭的对象添加 ShutdownListener 将立即触发该侦听器 -
getCloseReason(),允许调查对象关闭的原因 -
isOpen(),用于测试对象是否处于打开状态 -
close(int closeCode, String closeMessage),显式通知对象关闭
侦听器的简单用法如下所示
import com.rabbitmq.client.ShutdownSignalException;
import com.rabbitmq.client.ShutdownListener;
connection.addShutdownListener(new ShutdownListener() {
public void shutdownCompleted(ShutdownSignalException cause)
{
...
}
});
有关关闭情况的信息
可以通过显式调用 getCloseReason() 方法或在 ShutdownListener 类的 service(ShutdownSignalException cause) 方法中使用 cause 参数,检索包含有关关闭原因的所有信息的 ShutdownSignalException。
ShutdownSignalException 类提供了分析关闭原因的方法。通过调用 isHardError() 方法,我们可以获取它是连接错误还是通道错误的信息;getReason() 返回有关原因的信息,形式为 AMQP 方法(AMQP.Channel.Close 或 AMQP.Connection.Close,如果原因是库中的某些异常,如网络通信失败,则可以通过 getCause() 检索该异常)。
public void shutdownCompleted(ShutdownSignalException cause)
{
if (cause.isHardError())
{
Connection conn = (Connection)cause.getReference();
if (!cause.isInitiatedByApplication())
{
Method reason = cause.getReason();
...
}
...
} else {
Channel ch = (Channel)cause.getReference();
...
}
}
原子性与 isOpen() 方法的使用
不建议在生产代码中使用通道和连接对象的 isOpen() 方法,因为该方法返回的值取决于关闭原因是否存在。以下代码说明了竞争条件的可能性
public void brokenMethod(Channel channel)
{
if (channel.isOpen())
{
// The following code depends on the channel being in open state.
// However there is a possibility of the change in the channel state
// between isOpen() and basicQos(1) call
...
channel.basicQos(1);
}
}
相反,我们通常应该忽略此类检查,直接尝试所需的操作。如果执行代码时连接的通道已关闭,则会抛出 ShutdownSignalException,表明对象处于无效状态。我们还应该捕获由 SocketException(当代理意外关闭连接时)或 ShutdownSignalException(当代理发起干净关闭时)引起的 IOException。
public void validMethod(Channel channel)
{
try {
...
channel.basicQos(1);
} catch (ShutdownSignalException sse) {
// possibly check if channel was closed
// by the time we started action and reasons for
// closing it
...
} catch (IOException ioe) {
// check why connection was closed
...
}
}
高级连接选项
消费者操作线程池
Consumer 线程(参见下面的接收)默认情况下会自动分配在一个新的 ExecutorService 线程池中。如果需要更大的控制权,请在 newConnection() 方法上提供一个 ExecutorService,以便改用此线程池。这是一个分配了比通常分配的线程池更大的线程池的示例
ExecutorService es = Executors.newFixedThreadPool(20);
Connection conn = factory.newConnection(es);
Executors 和 ExecutorService 类都在 java.util.concurrent 包中。
当连接关闭时,默认的 ExecutorService 将被 shutdown(),但用户提供的 ExecutorService(如上面的 es)将不会被 shutdown()。提供自定义 ExecutorService 的客户端必须确保它最终被关闭(通过调用其 shutdown() 方法),否则池的线程可能会阻止 JVM 终止。
同一个执行器服务可以在多个连接之间共享,或者在重新连接时串行重用,但它不能在被 shutdown() 之后使用。
只有在有证据表明处理 Consumer 回调存在严重瓶颈时,才应考虑使用此功能。如果没有执行任何 Consumer 回调,或者很少,那么默认分配就绰绰有余了。开销最初很小,并且分配的总线程资源是有界的,即使偶尔发生突发的消费者活动也是如此。
使用 AddressResolver 接口进行服务发现
可以使用 AddressResolver 的实现来更改连接时使用的端点解析算法
Connection conn = factory.newConnection(addressResolver);
AddressResolver 接口如下所示
public interface AddressResolver {
List<Address> getAddresses() throws IOException;
}
就像端点列表一样,返回的第一个 Address 将被首先尝试,如果客户端连接到第一个失败,则尝试第二个,依此类推。
如果同时提供了 ExecutorService(使用 factory.newConnection(es, addressResolver) 形式),则线程池与(第一个)成功的连接相关联。
AddressResolver 是实现自定义服务发现逻辑的绝佳位置,这在动态基础设施中特别有用。结合自动恢复,客户端可以自动连接到即使在最初启动时没有启动的节点。亲和性和负载平衡是自定义 AddressResolver 可能有用的其他场景。
Java 客户端附带了以下实现(详细信息请参阅 javadoc)
-
DnsRecordIpAddressResolver: 给定主机名,返回其 IP 地址(根据平台 DNS 服务器解析)。这对于简单的基于 DNS 的负载平衡或故障转移很有用。 -
DnsSrvRecordAddressResolver: 给定服务名称,返回主机名/端口对。搜索实现为 DNS SRV 请求。这在使用像 HashiCorp Consul 这样的服务注册表时很有用。
心跳超时
有关心跳以及如何在 Java 客户端中配置它们的更多信息,请参阅心跳指南。
自定义线程工厂
Google App Engine (GAE) 等环境可能会限制直接实例化线程。要在这些环境中使用 RabbitMQ Java 客户端,必须配置一个自定义 ThreadFactory,它使用适当的方法来实例化线程,例如 GAE 的 ThreadManager。
下面是一个 Google App Engine 的示例。
import com.google.appengine.api.ThreadManager;
ConnectionFactory cf = new ConnectionFactory();
cf.setThreadFactory(ThreadManager.backgroundThreadFactory());
使用 Netty 进行网络 I/O
Java 客户端的 5.27.0 版本带来了对 Netty 网络 I/O 的支持。Netty 的速度并不一定比阻塞 I/O 快,但它提供了对资源(如线程)的更多控制,并提供了高级网络选项,如 带 OpenSSL 的 TLS 和 原生传输 (epoll, io_uring, kqueue)。
使用默认的阻塞 I/O 模式,每个连接都使用一个线程从网络套接字读取。使用 Netty,您可以控制读取和写入网络数据的线程数。
如果您的 Java 进程使用许多连接(几十个或几百个),请使用 Netty。您应该使用比默认阻塞模式更少的线程。设置适当数量的线程后,您不应该经历任何性能下降,特别是如果连接不是很繁忙的话。
Netty 通过 ConnectionFactory#netty() 助手激活和配置。Netty 的 EventLoopGroup 是对线程数挑剔的应用程序最重要的设置。以下是如何将其设置为 4 个线程的示例
int nbThreads = 4;
IoHandlerFactory ioHandlerFactory = NioIoHandler.newFactory();
EventLoopGroup eventLoopGroup = new MultiThreadIoEventLoopGroup(
nbThreads, ioHandlerFactory
);
connectionFactory.netty().eventLoopGroup(eventLoopGroup);
// ...
// dispose the event loop group after closing all connections
eventLoopGroup.shutdownGracefully();
请注意,事件循环组必须在连接关闭其连接后进行处置。如果不设置事件循环组,每个连接将使用其自己的 1 线程事件循环组(并负责将其关闭)。这远非最佳,这就是为什么使用 Netty 时强烈建议设置 EventLoopGroup 的原因。
Netty 使用其自己的 SslContext API 进行 TLS 配置(不是 JDK 的 SSLContext),因此当激活 Netty 时,ConnectionFactory#useSslProtocol() 方法无效。请改用 ConnectionFactory.netty().sslContext(SslContext) 以及 Netty 的 SslContextBuilder 类。这是一个示例
X509Certificate caCertificate = ...;
connectionFactory.netty()
.sslContext(SslContextBuilder
.forClient() // mandatory, do not forget to call
.trustManager(caCertificate) // pass in certificate directly
.build());
网络故障自动恢复
连接恢复
客户端和 RabbitMQ 节点之间的网络连接可能会失败。RabbitMQ Java 客户端支持连接和拓扑(队列、交换机、绑定和消费者)的自动恢复。
许多应用程序的自动恢复过程遵循以下步骤
- 重新连接
- 恢复连接侦听器
- 重新打开通道
- 恢复通道侦听器
- 恢复通道
basic.qos设置、发布者确认和事务设置
拓扑恢复包括为每个通道执行的以下操作
- 重新声明交换器(预定义交换器除外)
- 重新声明队列
- 恢复所有绑定
- 恢复所有消费者
从 Java 客户端的 4.0.0 版本开始,自动恢复默认启用(因此拓扑恢复也默认启用)。
拓扑恢复依赖于每连接的实体缓存(队列、交换机、绑定、消费者)。当在连接上声明队列时,它将被添加到缓存中。当它被删除或被安排删除(例如因为它被自动删除)时,它将被移除。此模型具有下述的一些局限性。
要禁用或启用自动连接恢复,请使用 factory.setAutomaticRecoveryEnabled(boolean) 方法。以下代码片段显示了如何显式启用自动恢复(例如对于 4.0.0 之前的 Java 客户端)
ConnectionFactory factory = new ConnectionFactory();
factory.setUsername(userName);
factory.setPassword(password);
factory.setVirtualHost(virtualHost);
factory.setHost(hostName);
factory.setPort(portNumber);
factory.setAutomaticRecoveryEnabled(true);
// connection that will recover automatically
Connection conn = factory.newConnection();
如果恢复因异常而失败(例如仍无法到达 RabbitMQ 节点),它将在固定时间间隔(默认为 5 秒)后重试。时间间隔可以配置
ConnectionFactory factory = new ConnectionFactory();
// attempt recovery every 10 seconds
factory.setNetworkRecoveryInterval(10000);
当提供地址列表时,列表会被混洗,并依次尝试所有地址
ConnectionFactory factory = new ConnectionFactory();
Address[] addresses = {new Address("192.168.1.4"), new Address("192.168.1.5")};
factory.newConnection(addresses);
何时触发连接恢复?
如果启用了自动连接恢复,它将由以下事件触发
- 连接的 I/O 循环中抛出 I/O 异常
- 套接字读取操作超时
- 检测到丢失的服务器心跳
- 连接的 I/O 循环中抛出任何其他意外异常
以先发生者为准。
如果最初连接到 RabbitMQ 节点的客户端连接失败,则不会触发自动连接恢复。应用程序开发人员负责重试此类连接、记录失败尝试、限制重试次数等。这是一个非常基础的示例
ConnectionFactory factory = new ConnectionFactory();
// configure various connection settings
try {
Connection conn = factory.newConnection();
} catch (java.net.ConnectException e) {
Thread.sleep(5000);
// apply retry logic
}
当应用程序通过 Connection.Close 方法关闭连接时,不会发起连接恢复。
通道级异常不会触发任何形式的恢复,因为它们通常表明应用程序中存在语义问题(例如尝试从不存在的队列中消费)。
恢复侦听器
可以在可恢复的连接和通道上注册一个或多个恢复侦听器。启用连接恢复后,ConnectionFactory#newConnection 和 Connection#createChannel 返回的连接将实现 com.rabbitmq.client.Recoverable,并提供两个名称非常直观的方法
addRecoveryListenerremoveRecoveryListener
请注意,目前您需要将连接和通道强制转换为 Recoverable 才能使用这些方法。
对发布的影响
连接断开时使用 Channel.basicPublish 发布的消息将会丢失。客户端不会在连接恢复后将它们排队等待投递。为确保发布的消息到达 RabbitMQ,应用程序需要使用发布者确认并考虑连接故障。
拓扑恢复
拓扑恢复涉及交换机、队列、绑定和消费者的恢复。在启用自动恢复时,它默认启用。在现代版本的客户端中,拓扑恢复默认启用。
如果需要,可以显式禁用拓扑恢复
ConnectionFactory factory = new ConnectionFactory();
Connection conn = factory.newConnection();
// enable automatic recovery (e.g. Java client prior 4.0.0)
factory.setAutomaticRecoveryEnabled(true);
// disable topology recovery
factory.setTopologyRecoveryEnabled(false);
故障检测与恢复局限性
自动连接恢复有许多限制和有意为之的设计决策,应用程序开发人员需要了解这些。
拓扑恢复依赖于每连接的实体缓存(队列、交换机、绑定、消费者)。当在连接上声明队列时,它将被添加到缓存中。当它被删除或被安排删除(例如因为它被自动删除)时,它将被移除。这使得在不同的通道上声明和删除实体而不会产生意外结果成为可能。这也意味着消费者标签(一种特定于通道的标识符)在所有使用自动连接恢复的连接上的所有通道中必须是唯一的。
当连接断开或丢失时,需要时间来检测。因此,有一段时间窗口内库和应用程序都不知道连接已经有效失败。在此期间发布的任何消息都会像往常一样序列化并写入 TCP 套接字。它们投递到代理只能通过发布者确认来保证:AMQP 0-9-1 中的发布在设计上完全是异步的。
当启用了自动恢复的连接检测到套接字或 I/O 操作错误时,恢复会在可配置的延迟后开始,默认为 5 秒。此设计假设尽管许多网络故障是暂时的且通常持续时间很短,但它们不会瞬间消失。延迟还可以避免服务器端资源清理(如排他或自动删除队列删除)与在同一资源上新打开的连接上执行的操作之间的固有竞争条件。
连接恢复尝试默认将以相同的时间间隔继续,直到成功打开新连接。可以通过向 ConnectionFactory#setRecoveryDelayHandler 提供 RecoveryDelayHandler 实现实例来使恢复延迟动态化。使用动态计算延迟间隔的实现应避免过小的值(例如低于 2 秒的值)。
当连接处于恢复状态时,在其通道上尝试的任何发布都将被异常拒绝。客户端目前不对这些传出消息执行任何内部缓冲。应用程序开发人员有责任跟踪此类消息,并在恢复成功后重新发布它们。发布者确认是一种协议扩展,应由无法承受消息丢失的发布者使用。
当通道因通道级异常而关闭时,连接恢复不会启动。此类异常通常表明应用程序级存在问题。库无法在此时做出知情的决定。
即使连接恢复启动,关闭的通道也不会恢复。这包括显式关闭的通道和上面的通道级异常情况。
手动确认与自动恢复
使用手动确认时,可能会在消息投递和确认之间发生到 RabbitMQ 节点的网络连接失败。连接恢复后,RabbitMQ 将重置所有通道上的投递标签。
这意味着带有旧投递标签的 basic.ack、basic.nack 和 basic.reject 将导致通道异常。为了避免这种情况,RabbitMQ Java 客户端会跟踪并更新投递标签,使其在恢复之间单调增长。
Channel.basicAck、Channel.basicNack 和 Channel.basicReject 然后将调整后的投递标签转换为 RabbitMQ 使用的标签。
带有过期投递标签的确认将不会被发送。使用手动确认和自动恢复的应用程序必须能够处理重复投递。
通道生命周期与拓扑恢复
自动连接恢复旨在对应用程序开发人员尽可能透明,这就是为什么即使在后台有几个连接失败和恢复,Channel 实例仍然保持不变。技术上,当自动恢复打开时,Channel 实例充当代理或装饰器:它们将 AMQP 业务委托给实际的 AMQP 通道实现,并在其周围实现一些恢复机制。这就是为什么在通道创建了一些资源(队列、交换机、绑定)后,您不应该关闭该通道,否则这些资源的拓扑恢复稍后会失败,因为该通道已被关闭。相反,请在应用程序的整个生命周期内保持创建的通道打开。
未处理异常
与连接、通道、恢复和消费者生命周期相关的未处理异常被委托给异常处理器。异常处理器是任何实现 ExceptionHandler 接口的对象。默认情况下,使用 DefaultExceptionHandler 的实例。它将异常详细信息打印到标准输出。
可以使用 ConnectionFactory#setExceptionHandler 覆盖该处理器。它将用于工厂创建的所有连接
ConnectionFactory factory = new ConnectionFactory();
cf.setExceptionHandler(customHandler);
异常处理器应仅用于异常记录。
指标与监控
客户端为活动连接收集运行时指标(例如已发布消息的数量)。指标收集是一个可选功能,应在 ConnectionFactory 级别设置,使用 setMetricsCollector(metricsCollector) 方法。此方法需要一个 MetricsCollector 实例,该实例在客户端代码的多个位置被调用。
客户端开箱即用支持 Micrometer、Dropwizard Metrics 和 OpenTelemetry。
以下是收集的指标
- 打开连接的数量
- 打开通道的数量
- 已发布消息的数量
- 已确认消息的数量
- 已负面确认 (nack-ed) 的出站消息数量
- 已返回的不可路由的出站消息数量
- 出站消息的失败数量
- 已消费消息的数量
- 已确认消息的数量
- 已拒绝消息的数量
Micrometer 和 Dropwizard Metrics 都提供计数,以及消息相关指标的平均速率、过去五分钟的速率等。它们还支持用于监控和报告的常见工具(JMX、Graphite、Ganglia、Datadog 等)。有关更多详细信息,请参阅下面的专用章节。
开发人员在启用指标收集时应记住几件事。
- 在使用 Micrometer 或 Dropwizard Metrics 时,不要忘记将适当的依赖项(在 Maven、Gradle 中,甚至作为 JAR 文件)添加到 JVM 类路径中。这些是可选依赖项,不会随 Java 客户端自动拉取。根据所使用的报告后端,您可能还需要添加其他依赖项。
- 指标收集是可扩展的。鼓励为特定需求实现自定义的
MetricsCollector。 MetricsCollector是在ConnectionFactory级别设置的,但可以在不同实例之间共享。- 指标收集不支持事务。例如,如果确认是在事务中发送的,然后事务回滚,则确认会被计入客户端指标(但显然代理不会计算)。请注意,确认实际上是发送给代理的,然后被事务回滚取消,因此客户端指标在发送的确认方面是正确的。简而言之,不要将客户端指标用于关键业务逻辑,它们不能保证完美准确。它们旨在用于简化对运行系统的理解并提高运维效率。
Micrometer 支持
必须先启用指标收集
以以下方式 Micrometer
ConnectionFactory connectionFactory = new ConnectionFactory();
MicrometerMetricsCollector metrics = new MicrometerMetricsCollector();
connectionFactory.setMetricsCollector(metrics);
...
metrics.getPublishedMessages(); // get Micrometer's Counter object
Micrometer 支持 几个报告后端:Netflix Atlas、Prometheus、Datadog、Influx、JMX 等。
您通常会将 MeterRegistry 的实例传递给 MicrometerMetricsCollector。这是一个使用 JMX 的示例
JmxMeterRegistry registry = new JmxMeterRegistry();
MicrometerMetricsCollector metrics = new MicrometerMetricsCollector(registry);
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setMetricsCollector(metrics);
Dropwizard Metrics 支持
使用 Dropwizard 启用指标收集,如下所示
ConnectionFactory connectionFactory = new ConnectionFactory();
StandardMetricsCollector metrics = new StandardMetricsCollector();
connectionFactory.setMetricsCollector(metrics);
...
metrics.getPublishedMessages(); // get Metrics' Meter object
Dropwizard Metrics 支持 几个报告后端:console、JMX、HTTP、Graphite、Ganglia 等。
您通常会将 MetricsRegistry 的实例传递给 StandardMetricsCollector。这是一个使用 JMX 的示例
MetricRegistry registry = new MetricRegistry();
StandardMetricsCollector metrics = new StandardMetricsCollector(registry);
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setMetricsCollector(metrics);
JmxReporter reporter = JmxReporter
.forRegistry(registry)
.inDomain("com.rabbitmq.client.jmx")
.build();
reporter.start();
Google App Engine 上的 RabbitMQ Java 客户端
在 Google App Engine (GAE) 上使用 RabbitMQ Java 客户端需要使用自定义线程工厂,该工厂使用 GAE 的 ThreadManager 实例化线程(见上文)。此外,必须设置较低的心跳间隔(4-5 秒),以避免在 GAE 上遇到较低的 InputStream 读取超时
ConnectionFactory factory = new ConnectionFactory();
cf.setRequestedHeartbeat(5);
警告和局限性
为了使拓扑恢复成为可能,RabbitMQ Java 客户端维护了一个声明的队列、交换机和绑定的缓存。缓存是每连接的。某些 RabbitMQ 功能使得客户端无法观察到某些拓扑变化,例如当队列由于 TTL 而被删除时。RabbitMQ Java 客户端尝试在最常见的情况下使缓存条目失效
- 当删除队列时。
- 当删除交换机时。
- 当删除绑定时。
- 当在自动删除的队列上取消消费者时。
- 当队列或交换机与自动删除的交换机解绑时。
然而,客户端无法跟踪单个连接之外的这些拓扑变化。依赖自动删除队列或交换机以及队列 TTL(注意:不是消息 TTL!)并使用自动连接恢复的应用程序应显式删除已知未使用或已删除的实体,以清除客户端拓扑缓存。通过在 RabbitMQ 3.3.x 中使 Channel#queueDelete、Channel#exchangeDelete、Channel#queueUnbind 和 Channel#exchangeUnbind 幂等(删除不存在的东西不会导致异常),这得到了促进。
RPC(请求/回复)模式:示例
作为编程上的方便,Java 客户端 API 提供了一个类 RpcClient,它使用临时回复队列通过 AMQP 0-9-1 提供简单的 RPC 风格通信设施。
该类不对 RPC 参数和返回值强加任何特定的格式。它只是提供了一种机制,用于将消息发送到具有特定路由键的给定交换机,并等待回复队列上的响应。
import com.rabbitmq.client.RpcClient;
RpcClient rpc = new RpcClient(channel, exchangeName, routingKey);
(该类如何使用 AMQP 0-9-1 的实现细节如下:请求消息发送时,basic.correlation_id 字段设置为该 RpcClient 实例唯一的值,并且 basic.reply_to 设置为回复队列的名称。)
一旦创建了该类的实例,您可以使用以下任何方法发送 RPC 请求
byte[] primitiveCall(byte[] message);
String stringCall(String message)
Map mapCall(Map message)
Map mapCall(Object[] keyValuePairs)
primitiveCall 方法传输原始字节数组作为请求和响应正文。方法 stringCall 是 primitiveCall 周围的一个精简便捷包装器,将消息正文视为默认字符编码中的 String 实例。
mapCall 变体稍微复杂一些:它们将包含普通 Java 值的 java.util.Map 编码为 AMQP 0-9-1 二进制表表示,并以相同的方式解码响应。(请注意,这里可以使用的值类型有一些限制 - 有关详细信息,请参阅 javadoc。)
所有编组/解组便捷方法都使用 primitiveCall 作为传输机制,并且只是在它上面提供一个包装层。
TLS 支持
可以使用 TLS 对客户端和代理之间的通信进行加密。还支持客户端和服务器身份验证(即对等验证)。这是 Java 客户端使用加密的最简单、最天真的方法
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setPort(5671);
// Only suitable for development.
// This code will not perform peer certificate chain verification and prone
// to man-in-the-middle attacks.
// See the main TLS guide to learn about peer verification and how to enable it.
factory.useSslProtocol();
请注意,客户端默认不在上述示例中强制执行任何服务器身份验证(对等证书链验证),因为使用了默认的“信任所有证书”的 TrustManager。这对于本地开发很方便,但易受中间人攻击,因此不推荐用于生产环境。
要了解有关 RabbitMQ 中 TLS 支持的更多信息,请参阅 TLS 指南。如果您只想配置 Java 客户端(特别是对等验证和信任管理器部分),请阅读 TLS 指南的相应部分。
请注意,当 Netty 用于网络 I/O 时,需要使用其自己的 SslContext API。有关更多详细信息,请参阅 Netty 部分。
OAuth 2 支持
客户端可以针对像 UAA 这样的 OAuth 2 服务器进行身份验证。必须在服务器端启用 OAuth 2 插件,并将其配置为使用与客户端相同的 OAuth 2 服务器。
获取 OAuth 2 令牌
Java 客户端提供 OAuth2ClientCredentialsGrantCredentialsProvider 类,以使用 OAuth 2 客户端凭据流程获取 JWT 令牌。客户端在打开连接时将在密码字段中发送 JWT 令牌。代理随后将在授权连接并授予对所请求虚拟主机的访问权限之前验证 JWT 令牌签名、有效性和权限。
优先使用 OAuth2ClientCredentialsGrantCredentialsProviderBuilder 创建 OAuth2ClientCredentialsGrantCredentialsProvider 实例,然后在 ConnectionFactory 上进行设置。以下代码片段显示了如何为 OAuth 2 插件的示例设置配置和创建 OAuth 2 凭据提供程序的实例
import com.rabbitmq.client.impl.OAuth2ClientCredentialsGrantCredentialsProvider.
OAuth2ClientCredentialsGrantCredentialsProviderBuilder;
...
CredentialsProvider credentialsProvider =
new OAuth2ClientCredentialsGrantCredentialsProviderBuilder()
.tokenEndpointUri("https://:8080/uaa/oauth/token/")
.clientId("rabbit_client").clientSecret("rabbit_secret")
.grantType("password")
.parameter("username", "rabbit_super")
.parameter("password", "rabbit_super")
.build();
connectionFactory.setCredentialsProvider(credentialsProvider);
在生产环境中,请确保令牌端点 URI 使用 HTTPS,并在必要时为 HTTPS 请求配置 SSLContext(以验证并信任 OAuth 2 服务器的身份)。以下代码片段通过使用 OAuth2ClientCredentialsGrantCredentialsProviderBuilder 中的 tls().sslContext() 方法来实现此目的
SSLContext sslContext = ... // create and initialise SSLContext
CredentialsProvider credentialsProvider =
new OAuth2ClientCredentialsGrantCredentialsProviderBuilder()
.tokenEndpointUri("https://:8080/uaa/oauth/token/")
.clientId("rabbit_client").clientSecret("rabbit_secret")
.grantType("password")
.parameter("username", "rabbit_super")
.parameter("password", "rabbit_super")
.tls() // configure TLS
.sslContext(sslContext) // set SSLContext
.builder() // back to main configuration
.build();
请参阅 Javadoc 以查看所有可用选项。
刷新令牌
令牌会过期,代理将拒绝处理使用过期令牌的连接。为了避免这种情况,可以在令牌过期前调用 CredentialsProvider#refresh() 并将新令牌发送到服务器。从应用程序的角度来看,这样做比较繁琐,因此 Java 客户端提供了 DefaultCredentialsRefreshService 来提供帮助。该工具会跟踪已使用的令牌,在过期前对其进行刷新,并为它所负责的连接发送新令牌。
以下代码片段展示了如何创建 DefaultCredentialsRefreshService 实例并将其设置到 ConnectionFactory 上
import com.rabbitmq.client.impl.DefaultCredentialsRefreshService.
DefaultCredentialsRefreshServiceBuilder;
...
CredentialsRefreshService refreshService =
new DefaultCredentialsRefreshServiceBuilder().build();
cf.setCredentialsRefreshService(refreshService);
DefaultCredentialsRefreshService 会在令牌有效时间的 80% 之后安排刷新,例如,如果令牌在 60 分钟后过期,它将在 48 分钟后刷新。这是默认行为,请查阅 Javadoc 以获取更多信息。