跳至主内容

AMQP 1.0 客户端库

本页面记录了针对 RabbitMQ 4.0 及以上版本AMQP 1.0 客户端库的使用方法。

RabbitMQ 团队提供对以下库的支持:

应用开发者可以在此找到如何将这些库用于最常见的使用场景。有关许可、下载、依赖管理、高级特性、特定用法和配置等其他信息,请参阅各库代码仓库中的 README 页面。

概述

RabbitMQ 团队维护着一套专为 RabbitMQ 设计并优化的 AMQP 1.0 客户端库。它们在 AMQP 1.0 之上提供了简单、安全且功能强大的 API。应用程序可以使用这些库来发布和消费消息,并以跨编程语言一致的方式管理服务器拓扑。这些库还提供了诸如自动连接和拓扑恢复、队列连接亲和性等高级功能。

注意

RabbitMQ 与任何兼容 AMQP 1.0 的客户端库兼容。虽然不强制要求必须使用 RabbitMQ AMQP 1.0 客户端库,但强烈建议应用程序使用它们以获得最佳体验。

安全性

RabbitMQ AMQP 1.0 客户端库默认是安全的,它们总是创建持久化实体并总是发布持久化消息。

保证

RabbitMQ AMQP 1.0 客户端库提供“至少一次”(at-least-once)交付保证。

代理(Broker)总是会确认对已发布消息的处理情况。发布者通过在创建时使用 unsettled 发送方结算模式first 接收方结算模式来实现这一点。

消费者在处理完消息后,必须始终向代理发出信号以表明处理结果。消费者在创建时使用与发布者相同的设置(first 接收方结算模式unsettled 发送方结算模式)。

客户端 API

本节涵盖如何使用 RabbitMQ AMQP 1.0 客户端库连接到集群,以及如何发布和消费消息。

连接

库为节点或节点集群提供了入口点。它的名称被称为“环境”(environment)。环境允许创建连接。它可以包含在连接之间共享的与基础设施相关的配置设置(例如 Java 的线程池)。以下是创建环境的方法:

创建环境
import com.rabbitmq.client.amqp.*;
import com.rabbitmq.client.amqp.impl.AmqpEnvironmentBuilder;

// ...

// create the environment instance
Environment environment = new AmqpEnvironmentBuilder()
.build();
// ...
// close the environment when the application stops
environment.close();

通常一个应用程序进程对应一个环境实例。当应用程序退出时,必须关闭环境以释放资源。

应用程序从环境中打开连接。它们必须指定连接到集群节点的适当设置(URI,凭据)。

打开连接
// open a connection from the environment
Connection connection = environment.connectionBuilder()
.uri("amqp://admin:admin@localhost:5672/%2f")
.build();
// ...
// close the connection when it is no longer necessary
connection.close();

库默认使用 ANONYMOUS SASL 身份验证机制。连接应该是长生命周期的对象,应用程序应避免频繁创建和销毁连接。当不再需要连接时,必须将其关闭。

发布

发布消息必须创建一个发布者。发布者发送消息的目标通常在创建时设置,但也可以按消息进行设置。

以下是在创建时设置目标来声明发布者的方法:

创建发布者
Publisher publisher = connection.publisherBuilder()
.exchange("foo").key("bar")
.build();
// ...
// close the publisher when it is no longer necessary
publisher.close();

在前面的示例中,使用该发布者发布的每条消息都将发送到 foo 交换机,路由键为 bar

信息

RabbitMQ 使用包含交换机、队列和绑定的 AMQ 0.9.1 模型

消息是从发布者实例创建的。它们遵循 AMQP 1.0 消息格式。可以定义消息体(作为字节数组)、标准属性和应用程序属性。

当消息被发布时,代理会在异步回调中指示其处理结果。客户端应用程序根据代理为该消息返回的状态(AMQP 术语中的 "outcome")采取适当措施(例如,如果消息未被 accepted,则将其存储在其他地方)。

以下代码片段展示了如何创建消息、发布消息以及处理来自代理的响应:

发布消息
// create the message
Message message = publisher
.message("hello".getBytes(StandardCharsets.UTF_8))
.messageId(1L);

// publish the message and deal with broker feedback
publisher.publish(message, context -> {
// asynchronous feedback from the broker
if (context.status() == Publisher.Status.ACCEPTED) {
// the broker accepted (confirmed) the message
} else {
// deal with possible failure
}
});

上面的发布者示例将消息发送到特定的交换机和特定的路由键,但这不是发布者支持的唯一目标。以下是发布者支持的非空目标类型:

创建具有不同目标的发布者
// publish to an exchange with a routing key
Publisher publisher1 = connection.publisherBuilder()
.exchange("foo").key("bar") // /exchanges/foo/bar
.build();

// publish to an exchange without a routing key
Publisher publisher2 = connection.publisherBuilder()
.exchange("foo") // /exchanges/foo
.build();

// publish to a queue
Publisher publisher3 = connection.publisherBuilder()
.queue("some-queue") // /queues/some-queue
.build();
信息

库将 API 调用转换为 地址格式 v2

也可以按消息设置目标。发布者必须在定义时不指定任何目标,且每条消息在属性部分的 to 字段中定义其目标。库在消息创建 API 中提供了辅助工具来定义消息目标,这避免了直接处理地址格式

以下代码片段展示了如何创建一个无目标的发布者,并定义具有不同目标类型的消息:

在消息中设置目标
// no target defined on publisher creation
Publisher publisher = connection.publisherBuilder()
.build();

// publish to an exchange with a routing key
Message message1 = publisher.message()
.toAddress().exchange("foo").key("bar")
.message();

// publish to an exchange without a routing key
Message message2 = publisher.message()
.toAddress().exchange("foo")
.message();

// publish to a queue
Message message3 = publisher.message()
.toAddress().queue("my-queue")
.message();

流支持

如果消息旨在发送到 流(stream),则可以使用 x-stream-filter-value 消息注释设置其 过滤值

在消息注释中设置流过滤值
Message message = publisher.message(body)
.annotation("x-stream-filter-value", "invoices"); // set filter value
publisher.publish(message, context -> {
// confirm callback
});

消费

消费者创建

创建消费者包括指定消费的队列以及处理消息的回调函数。

创建消费者
Consumer consumer = connection.consumerBuilder()
.queue("some-queue")
.messageHandler((context, message) -> {
byte[] body = message.body();
// ...
context.accept(); // settle the message
})
.build(); // do not forget to build the instance!

一旦应用程序处理完消息,就必须对其进行结算(settle)。这向代理指示处理结果以及应对消息采取的操作(例如删除消息)。应用程序必须结算消息,否则它们将耗尽信用额度(credits),代理将停止向其发送消息。

下一节将涵盖消息结算的语义。

消息处理结果(结果)

库允许应用程序以不同的方式结算消息。它们在消息传递应用程序的上下文中使用了尽可能明确的术语。每个术语都映射到 AMQP 规范中的特定结果(outcome)

  • accept:应用程序已成功处理消息,可以将其从队列中删除(accepted 结果)。
  • discard:应用程序无法处理消息,因为其无效,代理可以将其丢弃,或如果已配置,则将其发送到死信队列rejected 结果)。
  • requeue:应用程序未处理消息,代理可以将其重新入队并分发给相同或不同的消费者(released 结果)。

discardrequeue 有一个可选的消息注释参数,可以与消息头部分中现有的注释结合使用。此类消息注释可用于提供 discardrequeue 的原因详情。特定于应用程序的注释键必须以 x-opt- 前缀开头,而代理可识别的注释键仅以 x- 开头。discardrequeue 都使用带有消息注释参数的 modified 结果。

只有仲裁队列(quorum queues)支持使用 modified 结果修改消息注释

消费者优雅关闭

消费者通过接受、丢弃或重新入队来结算消息。

当消费者关闭时,未结算的消息会被重新入队。这可能导致消息的重复处理。

以下是一个示例:

  • 消费者对给定消息执行数据库操作。
  • 在消费者接受(结算)该消息之前,消费者被关闭。
  • 消息被重新入队。
  • 另一个消费者获取该消息并再次执行数据库操作。

完全避免重复消息是很困难的,这就是为什么处理应该是幂等的。消费者 API 提供了一种在消费者关闭时避免重复消息的方法。它包括暂停消息传递、获取未结算消息的数量以确保其最终达到 0,然后关闭消费者。这确保了消费者最终已静止(quiesced),并且所有收到的消息都已得到处理。

以下是消费者优雅关闭的示例:

优雅关闭消费者
// pause the delivery of messages
consumer.pause();
// ensure the number of unsettled messages reaches 0
long unsettledMessageCount = consumer.unsettledMessageCount();
// close the consumer
consumer.close();

应用程序仍然可以在不暂停的情况下关闭消费者,但这存在多次处理同一条消息的风险。

流支持

库在消费者配置中提供了对的开箱即用支持。

在从流中消费时,可以设置连接的位置:

附加到流的开头
Consumer consumer = connection.consumerBuilder()
.queue("some-stream")
.stream()
.offset(ConsumerBuilder.StreamOffsetSpecification.FIRST)
.builder()
.messageHandler((context, message) -> {
// message processing
})
.build();

还支持流过滤配置:

配置流过滤
Consumer consumer = connection.consumerBuilder()
.queue("some-stream")
.stream()
.filterValues("invoices", "orders")
.filterMatchUnfiltered(true)
.builder()
.messageHandler((ctx, msg) -> {
String filterValue = (String) msg.annotation("x-stream-filter-value");
// there must be some client-side filter logic
if ("invoices".equals(filterValue) || "orders".equals(filterValue)) {
// message processing
}
ctx.accept();
})
.build();

在处理流时,也请考虑使用首选编程语言的流客户端库及其原生流协议

基于 WebSocket 的 AMQP 1.0

VMware Tanzu RabbitMQ 支持基于 WebSocket 的 AMQP 1.0。文档可查看此处

注意

已在启用了 rabbitmq_web_amqp 插件的 VMware Tanzu RabbitMQ 4.1 上进行了测试。点击此处了解更多详情。

配置 ws 连接
// Not public yet

拓扑管理

应用程序可以管理 RabbitMQ 的 AMQ 0.9.1 模型:声明和删除交换机、队列和绑定。

为此,它们需要从连接中获取管理 API:

从环境中获取管理对象
Management management = connection.management();
// ...
// close the management instance when it is no longer needed
management.close();

管理 API 不再需要时应尽快关闭。应用程序通常在启动时创建其所需的拓扑,因此管理对象可以在此步骤后关闭。

交换器

以下是如何创建内置类型的交换机

创建内置类型的交换机
management.exchange()
.name("my-exchange")
.type(Management.ExchangeType.FANOUT) // enum for built-in type
.declare();

也可以将交换机类型指定为字符串(用于非内置类型的交换机):

创建非内置类型的交换机
management.exchange()
.name("my-exchange")
.type("x-delayed-message") // non-built-in type
.autoDelete(false)
.argument("x-delayed-type", "direct")
.declare();

以下是如何删除交换机:

删除交换机
management.exchangeDelete("my-exchange");

队列

以下是如何使用默认队列类型创建队列

创建经典队列
management.queue()
.name("my-queue")
.exclusive(true)
.autoDelete(false)
.declare();

管理 API 显式支持队列参数

创建带参数的队列
management.queue()
.name("my-queue")
.type(Management.QueueType.CLASSIC)
.messageTtl(Duration.ofMinutes(10))
.maxLengthBytes(ByteCapacity.MB(100))
.declare();

管理 API 还区分了所有队列类型共有的参数和仅对特定类型有效的参数。以下是创建仲裁队列的示例:

创建仲裁队列
management
.queue()
.name("my-quorum-queue")
.quorum() // set queue type to 'quorum'
.quorumInitialGroupSize(3) // specific to quorum queues
.deliveryLimit(3) // specific to quorum queues
.queue()
.declare();

可以查询有关队列的信息:

获取队列信息
Management.QueueInfo info = management.queueInfo("my-queue");
long messageCount = info.messageCount();
int consumerCount = info.consumerCount();
String leaderNode = info.leader();

此 API 也可用于检查队列是否存在。

以下是如何删除队列:

删除队列
management.queueDelete("my-queue");

绑定

管理 API 支持将队列绑定到交换机:

将队列绑定到交换机
management.binding()
.sourceExchange("my-exchange")
.destinationQueue("my-queue")
.key("foo")
.bind();

还支持交换机到交换机的绑定

将交换机绑定到另一个交换机
management.binding()
.sourceExchange("my-exchange")
.destinationExchange("my-other-exchange")
.key("foo")
.bind();

也可以解除实体绑定:

删除交换机和队列之间的绑定
management.unbind()
.sourceExchange("my-exchange")
.destinationQueue("my-queue")
.key("foo")
.unbind();

高级用法

生命周期监听器

应用程序可以通过添加监听器来对某些 API 组件的状态变化做出反应。应用程序可以向连接添加监听器,以便在连接恢复时停止发布消息。当连接恢复并再次打开时,应用程序可以恢复发布。

以下是如何在连接上设置监听器:

在连接上设置监听器
Connection connection = environment.connectionBuilder()
.listeners(context -> { // set one or several listeners
context.previousState(); // the previous state
context.currentState(); // the current (new) state
context.failureCause(); // the cause of the failure (in case of failure)
context.resource(); // the connection
}).build();

也可以在发布者实例上设置监听器:

在发布者上设置监听器
Publisher publisher = connection.publisherBuilder()
.listeners(context -> {
// ...
})
.exchange("foo").key("bar")
.build();

在消费者实例上也可以:

在消费者上设置监听器
Consumer consumer = connection.consumerBuilder()
.listeners(context -> {
// ...
})
.queue("my-queue")
.build();

自动连接恢复

自动连接恢复默认是开启的:客户端库将在意外关闭(如网络故障、节点重启等)后自动恢复连接。一旦启用连接恢复,自动拓扑恢复也会随之开启:客户端库将为恢复的连接重新创建 AMQP 实体以及发布者和消费者。开发人员不必过多担心网络稳定性和节点重启,因为客户端库会处理这些问题。

客户端每 5 秒尝试一次重新连接,直到成功。可以通过自定义退避延迟策略(back-off delay policy)来更改此行为:

为连接恢复设置退避策略
Connection connection = environment.connectionBuilder()
.recovery()
.backOffDelayPolicy(BackOffDelayPolicy.fixed(Duration.ofSeconds(2)))
.connectionBuilder().build();

如果拓扑恢复不适合特定的应用程序,也可以将其关闭。应用程序通常会注册一个连接生命周期监听器,以便在连接恢复时获知,并相应地恢复其自身状态。

关闭拓扑恢复
Connection connection = environment.connectionBuilder()
.recovery()
.topology(false) // deactivate topology recovery
.connectionBuilder()
.listeners(context -> {
// set listener that restores application state when connection is recovered
})
.build();

也可以完全关闭恢复功能:

关闭恢复
Connection connection = environment.connectionBuilder()
.recovery()
.activated(false) // deactivate recovery
.connectionBuilder().build();

队列亲和性

应用程序可以在打开连接时声明队列亲和性规则。客户端库将尽最大努力连接到适当的节点以强制执行这些规则。这可以减少集群内部流量,从而降低延迟并提高吞吐量。

在以下示例中,我们在给定连接上声明“发布”(publish)亲和性。库将连接到托管队列领导者(leader)的节点。如果环境已经创建了一个具有相同亲和性的连接,它也将重用该连接。

在连接上配置发布亲和性
Connection connection = environment.connectionBuilder()
.affinity()
.queue("my-queue")
.operation(PUBLISH)
.reuse(true)
.connection()
.build();

这是另一个关于“消费”(consume)操作的示例,库将优先选择托管队列副本(replica)的节点:

在连接上配置消费亲和性
Connection connection = environment.connectionBuilder()
.affinity()
.queue("my-queue")
.operation(CONSUME)
.connection()
.build();
© . This site is unofficial and not affiliated with VMware.