Ⅰ. 消息确认
一、消息确认机制
生产者发送消息之后,到达消费端之后,可能会有以下情况:
-
消息处理成功
-
消息处理异常

RabbitMQ 向消费者发送消息之后,就会把这条消息删掉,那么第两种情况,就会造成消息丢失。那么如何确保消费端已经成功接收了,并正确处理了呢?
为了保证消息从队列可靠地到达消费者,RabbitMQ 提供了消息确认机制!
消费者在订阅队列时,可以指定 autoAck 参数,根据这个参数设置,消息确认机制分为以下两种:
-
自动确认 : 当
autoAck等于true时,RabbitMQ 会自动把发送出去的消息置为确认,然后从内存(或者磁盘)中删除,而不管消费者是否真正地消费到了这些消息。适合对于消息可靠性要求不高 的场景。 -
手动确认 : 当
autoAck等于false时,RabbitMQ 会等待消费者显式地调用Basic.Ack命令,回复确认信号后才从内存(或者磁盘)中移去消息。适合对消息可靠性要求比较高 的场景。
String basicConsume(String queue, **boolean autoAck** , Consumer callback) throws IOException;代码示例:
DefaultConsumer consumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope,
AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("接收到消息: " + new String(body));
}
};
channel.basicConsume(Constants.TOPIC_QUEUE_NAME1, **true** , consumer);当 autoAck 参数置为 false,队列中的消息分成了两个部分:
-
Ready:等待投递给消费者的消息数 -
Unacked:已经投递给消费者,但是未收到消费者确认信号的消息数
如果 RabbitMQ 一直没有收到消费者的确认信号,并且消费此消息的消费者已经断开连接,则 RabbitMQ 会安排该消息重新进入队列,等待投递给下一个消费者,当然也有可能还是原来的那个消费者。

从管理平台上可以看到当前队列中 Ready 状态和 Unacked 状态的消息数:

二、手动确认方法
消费者在收到消息之后,可以选择确认,也可以选择直接拒绝或者跳过,RabbitMQ 也提供了不同的确认应答的方式,消费者客户端可以调用与其对应的 channel 的相关方法,共有以下三种:
① 肯定确认
RabbitMQ 已知道该消息并且成功的处理消息,可以将其丢弃了。
Channel.**basicAck** (long deliveryTag, boolean multiple);-
deliveryTag:消息的唯一标识,它是一个单调递增的 64 位的长整型值。deliveryTag是每个通道(Channel)独立维护的,所以在每个通道上都是唯一的 。当消费者确认(ack)一条消息时,必须使用对应的通道上进行确认。 -
multiple:是否批量确认。在某些情况下,为了减少网络流量,可以对一系列连续的deliveryTag进行批量确认。-
true:表示一次性 ack 所有小于等于指定deliveryTag的消息 -
false:表示确认当前指定deliveryTag的消息
-

deliveryTag是 RabbitMQ 中消息确认机制的一个重要组成部分,它确保了消息传递的可靠性和顺序性。
② 否定确认
RabbitMQ 在 2.0.0 版本开始引入了 Basic.Reject 这个命令,消费者客户端可以调用 channel.basicReject 方法来告诉 RabbitMQ 拒绝这个消息。
Channel.**basicReject** (long deliveryTag, boolean requeue);-
deliveryTag:参考channel.basicAck -
requeue:表示拒绝后,这条消息如何处理。-
true:会重新将这条消息存入队列,以便可以发送给下一个订阅的消费者。 -
false:会把消息从队列中移除,而不会把它发送给新的消费者。
-
③ 否定确认
Basic.Reject 命令一次只能拒绝一条消息,如果想要批量拒绝消息 ,则可以使用 Basic.Nack 这个命令,消费者客户端可以调用 channel.basicNack 方法来实现。
Channel.**basicNack** (long deliveryTag, boolean multiple, boolean requeue);三、代码示例
我们基于 SpringBoot 来演示消息的确认机制,使用方式和使用 RabbitMQ Java Client 库有一定差异。(主要体现在后者使用的是 channel.basicConsume 来接收消息以及做出回调处理,而前者用的是注解 @RabbitListener 来做回调处理)
Spring-AMQP 对消息确认机制提供了三种策略:
public enum AcknowledgeMode {
NONE,
MANUAL,
AUTO;
}-
NONE-
RabbitMQ 在投递消息后立即标记消息为已确认 。
-
不管消费者是否成功处理消息,Broker 都会立即从队列中移除消息。
-
如果消费者处理过程中宕机或抛出异常,消息无法再被重新投递 → 可能丢失 。
-
✅ 性能高,❌ 可靠性低
-
-
AUTO(默认)-
Spring 会在消息成功消费且未抛出异常时 自动发送
basicAck()。 -
如果消息处理过程中抛出异常,则不会确认,Spring 会根据重试机制 或死信策略 重新投递。
-
✅ 性能和可靠性平衡。
-
⚙️ 依赖 Spring 的异常处理机制判断是否 ack。
-
-
MANUAL-
开发者必须显式调用
channel.basicAck()或basicNack()。 -
若未确认,RabbitMQ 会认为消息 "尚未成功消费",并在消费者可用时重新投递。
-
✅ 最高可靠性,可实现精确控制(例如延迟 ack、批量 ack、失败重入队列)
-
❌ 代码复杂度稍高
-
下面以 AcknowledgeMode.NONE 为例,其它两种就是就是改一下配置文件中的 acknowledge-mode 即可!
-
配置确认机制:
spring: rabbitmq: addresses: amqp://liren:123123@127.0.0.1/lirendada listener: simple: **acknowledge-mode: ** ***none** ** # 设置确认机制为立刻确认* -
编写常量类:
public class Constant { public static final String ACK_EXCHANGE_NAME = "ack_exchange"; public static final String ACK_QUEUE = "ack_queue"; } -
配置与绑定队列和交换机:
@Configuration public class RabbitMQConfig { @Bean("ackQueue") public Queue ackQueue() { return QueueBuilder.*durable*(Constants.*ACK_QUEUE*).build(); } @Bean("ackExchange") public DirectExchange ackExchange() { return ExchangeBuilder.*directExchange*(Constants.*ACK_EXCHANGE_NAME*).durable(true).build(); } @Bean("ackBinding") public Binding ackBinding(@Qualifier("ackQueue") Queue queue, @Qualifier("ackExchange") DirectExchange exchange) { return BindingBuilder.*bind*(queue).to(exchange).with("ack"); } } -
发送消息:
@RequestMapping("/producer") @RestController public class producerController { @Autowired private RabbitTemplate rabbitTemplate; @RequestMapping("/ack") public String ack() { rabbitTemplate.convertAndSend(Constants.*ACK_EXCHANGE_NAME*, "ack", "consumer ack test..."); return "发送成功!"; } } -
消费消息:
@Component public class AckListener { @RabbitListener(queues = Constants.*ACK_QUEUE*) public void ListenerQueue(Message message, Channel channel) throws UnsupportedEncodingException { System.*out*.printf("接收到消息: %s, deliveryTag: %d%n", new String(message.getBody(),"UTF-8"), message.getMessageProperties().getDeliveryTag()); *// 模拟处理失败,会抛异常* * *int a = 3 / 0; System.*out*.println("处理完成"); } }
Ⅱ. 持久性
消费者处理消息时,消息如何不丢失呢?如何保证当 RabbitMQ 服务停掉以后,生产者发送的消息不丢失呢?
默认情况下,RabbitMQ 退出或者由于某种原因崩溃时,会忽视队列和消息,除非告知它不要这么做 。
RabbitMQ 的持久化分为三个部分:
-
交换器持久化
-
队列持久化
-
消息持久化
一、交换机持久化
交换器的持久化是在声明交换机时设置 durable 参数为 true,相当于将交换机的属性在服务器内部保存。
设置持久化之后,当 RabbitMQ 的服务器发生意外或关闭之后,进行重启时不需要重新去建立交换机,交换机会自动建立,相当于一直存在。
如果交换器不设置持久化,那么在 RabbitMQ 服务重启之后,相关的交换机元数据会丢失,对一个长期使用的交换器来说,建议将其置为持久化的。
ExchangeBuilder.topicExchange(Constant.ACK_EXCHANGE_NAME).**durable(true)** .build()二、队列持久化
队列的持久化是在声明队列时设置 durable 参数为 true。
如果队列不设置持久化,那么在 RabbitMQ 服务重启之后,该队列就会被删掉,此时数据也会丢失。(因为队列不存在了,那么消息也无处可存了)
但是设置了队列的持久化,也只能保证该队列本身的元数据不会因异常情况而丢失,并不能保证内部所存储的消息不会丢失 。要确保消息不会丢失,需要设置消息持久化 。
💥注意事项: 创建队列的时候默认 durable 为 true,即 RabbitMQ 默认开启队列持久化 。
三、消息持久化
消息实现持久化,需要在发送消息的时候,将消息的投递模式(MessageProperties 中的 deliveryMode)设置为 2,也就是 MessageDeliveryMode.PERSISTENT。
public enum MessageDeliveryMode {
NON_PERSISTENT, // 非持久化
PERSISTENT; // 持久化
}设置了队列以及消息的持久化,当 RabbitMQ 服务重启之后,消息依旧存在。
如果只设置队列持久化,重启之后消息会丢失;如果只设置消息持久化,重启之后队列消失,继而消息也丢失。所以单单设置消息持久化而不设置队列持久化显得毫无意义 。
// 非持久化信息
channel.basicPublish("", QUEUE_NAME, **null** , msg.getBytes());
// 持久化信息
channel.basicPublish("", QUEUE_NAME,
**MessageProperties.PERSISTENT_TEXT_PLAIN** , msg.getBytes());其中 MessageProperties.PERSISTENT_TEXT_PLAIN 实际就是封装了这个属性,其源码如下所示:
public static final BasicProperties PERSISTENT_TEXT_PLAIN =
new BasicProperties("text/plain",
null,
null,
2, // deliveryMode
0, null, null, null,
null, null, null, null,
null, null);如果是在 springboot 中使用 RabbitTemplate 发送持久化消息,操作如下所示:
// 要发送的消息内容
String message = "This is a persistent message";
// 创建一个Message对象,设置为持久化
Message messageObject = new Message(message.getBytes(), new MessageProperties());
messageObject.**getMessageProperties** ().**setDeliveryMode** (**MessageDeliveryMode.PERSISTENT** );
// 使用RabbitTemplate发送消息
rabbitTemplate.convertAndSend(Constant.ACK_EXCHANGE_NAME, "ack", **messageObject** );💥注意事项: RabbitMQ 默认情况下会将消息视为持久化的 ,除非队列被声明为非持久化,或者消息在发送时被标记为非持久化。
💡 将所有的消息都设置为持久化,会严重影响 RabbitMQ 的性能(随机)。写入磁盘的速度比写入内存的速度慢得不止一点点。对于可靠性不是那么高的消息可以不采用持久化处理以提高整体的吞吐量。在选择是否要将消息持久化时,需要在可靠性和吐吞量之间做一个权衡。
将交换器、队列、消息都设置了持久化之后就能百分之百保证数据不丢失了吗❓❓❓答案是否定的。
-
从消费者来说,如果在订阅消费队列时将
autoAck参数设置为true,那么当消费者接收到相关消息之后,还没来得及处理就宕机了,这样也算数据丢失。这种情况很好解决,将autoAck参数设置为false,进行手动确认。 -
在持久化消息被 RabbitMQ 接收后,并不会立即写入磁盘。RabbitMQ 并不会为每条消息都执行同步落盘(即调用操作系统的
fsync方法),而是先通过write()写入操作系统的页缓存中,等待系统或批量策略触发再真正写入磁盘。因此,如果在这段缓存尚未同步到磁盘的时间窗口内 RabbitMQ 节点发生宕机或重启,尚未落盘的消息仍可能丢失。
这个问题怎么解决呢❓❓❓
-
引入 RabbitMQ 的
仲裁队列,如果主节点(master)在此特殊时间内挂掉,可以自动切换到从节点(slave),这样有效地保证了高可用性,除非整个集群都挂掉。(此方法同样不能保证 100% 可靠,但是配置了仲裁队列要比没有配置仲裁队列的可靠性要高很多,实际生产环境中的关键业务队列一般都会设置仲裁队列) -
还可以在发送端引入
事务机制或者发布确认机制来保证消息已经正确地发送并存储至 RabbitMQ 中。
Ⅲ. 发布确认机制
在使用 RabbitMQ 的时候,可以通过消息持久化来解决因为服务器的异常崩溃而导致的消息丢失,但是还有一个问题,当消息的生产者将消息发送出去之后,消息到底有没有正确地到达服务器呢?如果在消息到达服务器之前已经丢失(比如 RabbitMQ 重启,那么 RabbitMQ 重启期间生产者消息投递失败),持久化操作也解决不了这个问题,因为消息根本没有到达服务器,何谈持久化?
RabbitMQ 为我们提供了两种解决方案:
-
通过** 事务机制** 实现**(比较消耗性能,在实际工作中使用也不多)**
-
通过 发布确认机制 ** 实现(这里主要介绍这种方案)**
RabbitMQ 为我们提供了两个方式来控制消息的可靠性投递
-
confirm确认模式 -
return退回模式
一、confirm确认模式
生产者在发送消息的时候,对发送端设置一个 ConfirmCallback 的监听,无论消息是否到达交换机,这个监听都会被执行。如果 Exchange 成功收到,ACK 为 true;如果没收到消息,ACK 就为 false。
RabbitTemplate.ConfirmCallback和ConfirmListener区别💥💥💥在 RabbitMQ 中,
ConfirmListener和ConfirmCallback都是用来处理消息确认的机制,但它们属于不同的客户端库,并且使用的场景和方式有所不同。
ConfirmListener是 RabbitMQ Java Client 库中的接口。这个库是 RabbitMQ 官方提供的一个直接与 RabbitMQ 服务器交互的客户端库。ConfirmListener接口提供了两个方法:handleAck和handleNack,用于处理消息确认和否定确认的事件。
ConfirmCallback是 Spring AMQP 框架中的一个接口,专门为 Spring 环境设计,用于简化与 RabbitMQ 交互的过程。它只包含一个confirm方法,用于处理消息确认的回调。在 SpringBoot 应用中,通常会使用
ConfirmCallback,因为它与 Spring 框架的其他部分更加整合,可以利用 Spring 的配置和依赖注入功能。而在使用 RabbitMQ Java Client 库时,则可能会直接实现ConfirmListener接口,更直接的与 RabbitMQ 的 Channel 交互。public interface ConfirmCallback { /** * 确认回调 * @param **correlationData** : 发送消息时的附加信息, 通常用于在确认回调中识别特定的消息 * @param **ack** : 交换机是否收到消息, 收到为true, 未收到为false * @param **cause** : 当消息确认失败时,这个字符串参数将提供失败的原因.这个原因可以用于调试和错误处理 * 成功时, cause为null */ void **confirm** (@Nullable CorrelationData correlationData, boolean ack, @Nullable String cause); }
- 配置RabbitMQ
spring: rabbitmq: addresses: amqp://liren:123123@127.0.0.1/lirendada listener: simple: acknowledge-mode: manual # 消息确认 **publisher-confirm-type: correlated** # 发布确认机制
其中发布确认机制这里有三种可选方式:(对应前面学习发布确认机制时候的 无确认、单独确认/批量确认、异步确认)
| 模式 | 是否启用确认 | 方式 | 阻塞性 | 性能 | 典型场景 |
|---|---|---|---|---|---|
| NONE | 否 | 无确认 | 非阻塞 | ⭐⭐⭐⭐ | 普通日志、监控消息 |
| SIMPLE | 是 | 同步等待 | 阻塞 | ⭐ | 关键事务型消息 |
| CORRELATED | 是 | 异步回调 | 非阻塞 | ⭐⭐⭐ | 高并发可靠投递 |
-
常量类
*// 发布确认机制* public static final String *CONFIRM_EXCHANGE_NAME *= "confirm_exchange"; public static final String *CONFIRM_QUEUE *= "confirm_queue"; -
设置确认回调逻辑并发送消息
- 无论消息确认成功还是失败,都会调用
ConfirmCallback的confirm方法 。-
如果消息发送成功,
ack=true -
如果消息发送失败,
ack=false,并且由参数 cause 提供失败的原因
-
public class RabbitTemplateConfig { @Bean("ackRabbitTemplate") public RabbitTemplate AckRabbitTemplate(ConnectionFactory connectionFactory) { RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory); return rabbitTemplate; } **@Bean("confirmRabbitTemplate")** public RabbitTemplate ConfirmRabbitTemplate(ConnectionFactory connectionFactory) { RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory); rabbitTemplate.**setConfirmCallback** (new RabbitTemplate.ConfirmCallback() { @Override public void confirm(CorrelationData correlationData, boolean ack, String cause) { if(ack) { System.*out*.printf("消息接收成功, id:%s", correlationData.getId()); } else { System.*out*.printf("消息接收失败, id:%s cause:%s", correlationData.getId(), cause); } } }); return rabbitTemplate; } } - 无论消息确认成功还是失败,都会调用
-
发送消息
@RequestMapping("/producer") @RestController public class producerController { @Resource(name = "ackRabbitTemplate") private RabbitTemplate ackRabbitTemplate; **@Resource(name = "confirmRabbitTemplate")** private RabbitTemplate confirmRabbitTemplate; @RequestMapping("/ack") public String ack() { ackRabbitTemplate.convertAndSend(Constants.*ACK_EXCHANGE_NAME*, "ack", "consumer ack test..."); return "发送成功!"; } **@RequestMapping("/confirm")** public String confirm() { CorrelationData correlationData = new CorrelationData("1"); confirmRabbitTemplate.convertAndSend(Constants.*CONFIRM_EXCHANGE_NAME*, "confirm", "consumer confirm test...", correlationData); return "发送成功!"; } }
💡 上面代码中有几处细节:
如果需要让不同的 RabbitTemplate 使用不同的配置的话,则需要对各自的 RabbitTemplate 进行配置,然后注册成不同的 Bean 对象交给 Spring 管理,在需要使用的时候利用
@Resource注入即可。为了防止多次调用 RabbitTemplate 的时候出现多次设置
ConfirmCallback导致报错的情况,通常将设置ConfirmCallback的操作放在配置类中完成 。
二、return退回模式
消息到达 Exchange 之后,会根据路由规则匹配,把消息放入 Queue 中。Exchange 到 Queue 的过程,如果一条消息无法被任何队列消费(即没有队列与消息的路由键匹配或队列不存在等),可以选择把消息退回给发送者。消息退回给发送者时,我们可以设置一个返回回调方法,对消息进行处理,这就是所谓的退回模式。
-
配置 RabbitMQ (同 confirm 模式)
spring: rabbitmq: addresses: amqp://liren:123123@127.0.0.1/lirendada listener: simple: acknowledge-mode: manual # 消息确认 publisher-confirm-type: correlated # 发布确认机制 -
设置返回回调逻辑: 当消息无法被路由到任何队列,它将返回给发送者,这时
setReturnCallback设置的回调将被触发@Bean("confirmRabbitTemplate") public RabbitTemplate ConfirmRabbitTemplate(ConnectionFactory connectionFactory) { RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory); *// 设置confirm回调(发送者 -> 交换机)* * *rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() { @Override public void confirm(CorrelationData correlationData, boolean ack, String cause) { if(ack) { System.*out*.printf("消息接收成功, id:%s\n", correlationData.getId()); } else { System.*out*.printf("消息接收失败, id:%s cause:%\n", correlationData.getId(), cause); } } }); ***// 设置return回调(交换机 -> 队列)** * * *rabbitTemplate.**setMandatory** (true); rabbitTemplate.**setReturnsCallback** (new RabbitTemplate.ReturnsCallback() { @Override public void returnedMessage(ReturnedMessage returned) { System.*out*.printf("消息被退回: %s\n", returned); } }); return rabbitTemplate; } -
发送消息:
@RequestMapping("/returns") public String returns() { CorrelationData correlationData = new CorrelationData("5"); confirmRabbitTemplate.convertAndSend(Constants.*CONFIRM_EXCHANGE_NAME*, "confirm", "consumer returns test...", correlationData); confirmRabbitTemplate.convertAndSend(Constants.*CONFIRM_EXCHANGE_NAME*, "confirm11", "consumer returns test...", correlationData); return "发送成功!"; }

使用 RabbitTemplate.setMandatory() 方法设置消息的 mandatory 属性为 true(默认为 false)。这个属性的作用是告诉 RabbitMQ,如果一条消息无法被任何队列消费,RabbitMQ 应该将消息返回给发送者,此时 ReturnCallback 会被触发。
其中该回调函数中有一个参数:ReturnedMessage,包含以下属性:
public class ReturnedMessage {
// 返回的消息对象,包含了消息体和消息属性
private final Message message;
// 由Broker提供的回复码, 表示消息无法路由的原因. 通常是一个数字代码,每个数字代表不同的含义.
private final int replyCode;
// 一个文本字符串, 提供了无法路由消息的额外信息或错误描述.
private final String replyText;
// 消息被发送到的交换机名称
private final String exchange;
// 消息的路由键,即发送消息时指定的键
private final String routingKey;
}三、常见面试题💥 -- 如何保证 RabbitMQ 消息的可靠传输?

从这个图中可以看出,消息可能丢失的场景以及解决方案:
-
生产者将消息发送到 RabbitMQ Server 失败
-
可能原因:网络问题等
-
解决办法:发布确认机制中的 confirm 模式
-
-
消息在交换机中无法路由到指定队列:
-
可能原因:代码或者配置层面错误,导致消息路由失败
-
解决办法:发布确认机制中的 return 模式
-
-
消息队列自身数据丢失
-
可能原因:消息到达 RabbitMQ 之后,RabbitMQ Server 宕机导致消息丢失
-
解决办法:持久性机制
- 开启 RabbitMQ 持久化,就是消息写入之后会持久化到磁盘,如果 RabbitMQ 挂了,恢复之后会自动读取之前存储的数据。(极端情况下,RabbitMQ 还未持久化就挂了,可能导致少量数据丢失,这个概率极低,也可以通过集群的方式提高可靠性)
-
-
消费者异常,导致消息丢失
-
可能原因:消息到达消费者,还没来得及消费、消费者宕机、消费者逻辑有问题
-
解决办法:消息确认机制
- RabbitMQ 提供了消费者应答机制来使 RabbitMQ 能够感知到消费者是否消费成功消息。默认情况下消费者应答机制是自动应答的,可以开启手动确认,当消费者确认消费成功后才会删除消息,从而避免消息丢失。除此之外,也可以配置重试机制(参考下一章节),当消息消费异常时,通过消息重试确保消息的可靠性。
-
Ⅳ. 重试机制
在消息传递过程中,可能会遇到各种问题,如网络故障、服务不可用、资源不足等,这些问题可能导致消息处理失败。为了解决这些问题,RabbitMQ 提供了重试机制,允许消息在处理失败后重新发送 。
但如果是程序逻辑引起的错误,那么多次重试也是没有用的,当然也可以设置重试次数。
💡 这里的 retry(重试)机制 通常指的是 消息从队列发送给消费者后,消费失败或者未 ack 时触发的重发 。而不是生产者发送给交换机时候失败触发的,这种触发是发布确认机制来解决的。
一、自动确认下的重试机制
- 重试配置
💥要启用 消费者消息重试机制 ,必须打开消息确认机制中的 AUTO 才行!
| 模式 | Spring自动重试机制 | 说明 |
|---|---|---|
| NONE | ❌ 无法重试(已确认) | RabbitMQ立即确认消息 |
| AUTO | ✅ 生效 | Spring 捕获异常自动重试 |
| MANUAL | ❌ 不生效 | 需自己实现重试逻辑 |
spring:
rabbitmq:
addresses: amqp://liren:123123@127.0.0.1/lirendada
listener:
simple:
acknowledge-mode: **auto** # 消息确认设置为auto,如果处理消息出现异常,会自动进行重试
**retry:**
enabled: true # 开启消费者失败重试
initial-interval: 5000ms # 初始失败等待时长为5秒
max-attempts: 5 # 最大重试次数(包括自身消费的一次)- 配置交换机 && 队列
首先是常量类:
*// 重试机制*
public static final String *RETRY_EXCHANGE_NAME *= "retry_exchange";
public static final String *RETRY_QUEUE *= "retry_queue";然后配置以及绑定交换机和队列:
*// 重试机制*
@Bean("retryQueue")
public Queue retryQueue() {
return QueueBuilder.*durable*(Constants.*RETRY_QUEUE*).build();
}
@Bean("retryExchange")
public DirectExchange retryExchange() {
return ExchangeBuilder.*directExchange*(Constants.*RETRY_EXCHANGE_NAME*).durable(true).build();
}
@Bean("retryBinding")
public Binding retryBinding(@Qualifier("retryQueue")Queue queue,
@Qualifier("retryExchange")Exchange exchange) {
return BindingBuilder.*bind*(queue).to(exchange).with("retry").noargs();
}- 发送消息
@RequestMapping("/retry")
public String retry() {
rabbitTemplate.convertAndSend(Constants.*RETRY_EXCHANGE_NAME*, "retry", "retry test...");
return "发送成功!";
}- 消费消息
import com.bite.rabbitmq.constant.Constant;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Component
public class RetryQueueListener {
//指定监听队列的名称
@RabbitListener(queues = Constant.RETRY_QUEUE)
public void ListenerQueue(Message message) throws Exception {
System.out.printf("接收到消息: %s, deliveryTag: %d%n", new String(message.getBody(),"UTF-8"),
message.getMessageProperties().getDeliveryTag());
//模拟处理失败
int num = 3/0;
System.out.println("处理完成");
}
}- 运行程序,观察结果
http://127.0.0.1:8080/product/retry

但是如果对异常进行捕获了,那么就不会进行重试!
二、手动确认下的重试机制
将消息确认机制改为手动 manual 模式,然后修改消费者代码:
@RabbitListener(queues = Constants.*RETRY_QUEUE*)
public void ListenerQueue(Message message, Channel channel) throws IOException, InterruptedException {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
System.*out*.printf("接收到消息:%s,deliveryTag:%d \n", new String(message.getBody()), deliveryTag);
try {
*// 模拟出现异常*
* *int a = 3 / 0;
System.*out*.println("处理完成!");
channel.basicAck(deliveryTag, true); *// 手动确认一下*
* *} catch (Exception e) {
System.*out*.println("出现异常!");
Thread.*sleep*(1000);
channel.basicNack(deliveryTag, true, true); *// 设置重新入队*
* *}
}
可以看到,手动确认模式时,重试次数的限制不会像在自动确认模式下那样直接生效,因为是否重试以及何时重试更多地取决于应用程序的逻辑和消费者的实现。
自动确认模式下,RabbitMQ 会在消息被投递给消费者后自动确认消息。如果消费者处理消息时抛出异常,RabbitMQ 根据配置的重试参数自动将消息重新入队,从而实现重试。重试次数和重试间隔等参数可以直接在 RabbitMQ 的配置中设定,并且 RabbitMQ 会负责执行这些重试策略。
而在手动确认模式下,消费者需要显式地对消息进行确认。如果消费者在处理消息时遇到异常,可以选择不确认消息使消息可以重新入队。重试的控制权在于应用程序本身,而不是 RabbitMQ 的内部机制。应用程序可以通过自己的逻辑和利用 RabbitMQ 的高级特性来实现有效的重试策略。
💡 使用重试机制时需要注意:
自动确认模式下:程序逻辑异常,多次重试还是失败,消息会自动确认,然后丢失 。
手动确认模式下:程序逻辑异常,多次重试消息依然处理失败,无法被确认,就一直是
unacked的状态,导致消息积压 。
Ⅴ. TTL
TTL(Time to Live)即过期时间。RabbitMQ 可以对消息和队列设置 TTL。
当消息达到存活时间之后,若还没有被消费,就会被自动清除。
咱们在网上购物,经常会遇到一个场景,当下单超过24小时还未付款,订单会被自动取消。还有类似的,申请退款之后,超过7天未被处理,则自动退款。
有两种方法设置消息的TTL:
-
队列级 TTL (队列中所有消息生效)
# 配置文件 arguments: x-message-ttl: 60000 # 毫秒 -
消息级 TTL (每条消息单独设置)
MessageProperties props = new MessageProperties(); props.setExpiration("60000"); // 毫秒
如果两种方法一起使用,则消息的 TTL 以两者之间较小的那个数值为准 。
一、单独设置消息的TTL
针对每条消息设置 TTL 的方法是在发送消息的方法中设置 expiration 属性参数,单位为毫秒 。
-
如果不设置 TTL,则表示此消息不会过期
-
如果将 TTL 设置为 0,则表示除非此时可以直接将消息投递到消费者,否则该消息会被立即丢弃
-
常量类:
*// TTL *public static final String TTL_TIME = "10000"; // 10s* *public static final String TTL_EXCHANGE_NAME = "ttl_exchange";* *public static final String TTL_QUEUE = "ttl_queue"; -
配置以及绑定交换机与队列:
*// TTL* @Bean("ttlQueue") public Queue ttlQueue() { return QueueBuilder.*durable*(Constants.*TTL_QUEUE*).build(); } @Bean("ttlExchange") public DirectExchange ttlExchange() { return ExchangeBuilder.*directExchange*(Constants.*TTL_EXCHANGE_NAME*).durable(true).build(); } @Bean("ttlBinding") public Binding ttlBinding(@Qualifier("ttlQueue")Queue queue, @Qualifier("ttlExchange")Exchange exchange) { return BindingBuilder.*bind*(queue).to(exchange).with("ttl").noargs(); } -
发送消息:
@RequestMapping("/ttl") public String ttl() { MessagePostProcessor messagePostProcessor = new MessagePostProcessor() { @Override public Message postProcessMessage(Message message) throws AmqpException { **message.getMessageProperties().setExpiration(Constants.** ***TTL_TIME** ***)** ; return message; } }; rabbitTemplate.convertAndSend(Constants.*TTL_EXCHANGE_NAME*, "ttl", "ttl test...", **messagePostProcessor** ); return "发送成功!"; } // 另一种写法:因为 MessagePostProcessor 是函数式接口,所以可以用lambda简化 @RequestMapping("/ttl") public String ttl() { rabbitTemplate.convertAndSend(Constants.*TTL_EXCHANGE_NAME*, "ttl", "ttl test...", **message -> {** **message.getMessageProperties().setExpiration(Constants.** ***TTL_TIME** ***);** **return message;** **}** ); return "发送成功!"; }
发送消息后可以到管理页面观察队列,可以发现 10s 后队列中的消息就消失了!
二、设置队列的TTL
直接在创建队列的时候使用封装好的 ttl() 方法即可设置队列中的消息过期时间,单位是毫秒 。
@Bean("ttlQueue2")
public Queue ttlQueue2() {
return QueueBuilder.*durable*(Constants.*TTL_QUEUE*).**ttl** (10000).build();
}实际上设置队列 TTL 的原理,是在创建队列时加入 x-message-ttl 参数实现的,下面是源码:

三、两者区别
-
设置队列 TTL 属性的方法,一旦消息过期,就会从队列中删除
-
单独设置消息 TTL 的方法,即使消息过期,也不会马上从队列中删除,而是当投递到消费者之前进行判定为过期了才删除的
为什么这两种方法处理的方式不一样???
因为 RabbitMQ 的队列不是扫描式的,而是顺序读取式队列 。
消息存放在队列中时,只有在它到达队列头部(准备被投递给消费者)时 ,RabbitMQ 才会检查它是否已过期。
换句话说:RabbitMQ 不会遍历整个队列去找 "已经过期" 的消息。只有 "轮到要被投递的消息" 时,才判断是否过期。
因此两者处理结果不同:
-
如果消息在队列中排得很靠后,它可能在过期后很久 才被清理掉。
-
清理时是 "惰性删除":当检测到已过期 → 丢弃并(可选)发送到死信交换机(DLX)。
Ⅵ. 死信队列
一、死信的概念
死信就是因为种种原因,而导致的无法被消费的信息 。
有死信,自然就有死信队列。当消息在一个队列中变成死信之后,它能被重新被发送到另一个交换器中,这个交换器就是 DLX(Dead Letter Exchange),绑定 DLX 的队列,就称为死信队列 DLQ(Dead Letter Queue)。

消息变成死信通常有以下几种可能:
-
消息被拒绝(
Basic.Reject/Basic.Nack),并且设置requeue参数为false -
消息过期
-
队列达到最大长度,消息溢出
二、代码示例
1. 声明配置队列和交换机
包含两部分:
-
声明正常的队列和正常的交换机
-
声明死信队列和死信交换机
死信交换机/队列 和 普通的交换机/队列 没有区别,只是处理的事情不同罢了!
首先是常量类:
*// 死信*
public static final String *NORMAL_EXCHANGE *= "normal_exchange";
public static final String *NORMAL_QUEUE *= "normal_queue";
public static final String *DL_EXCHANGE *= "dl_exchange";
public static final String *DL_QUEUE *= "dl_queue";然后声明以及配置交换机和队列:
**// 正常队列**
@Bean("normalQueue")
public Queue normalQueue() {
return QueueBuilder
.*durable*(Constants.*NORMAL_QUEUE*)
.**deadLetterExchange** (Constants.*DL_EXCHANGE*) *// 绑定死信交换机*
* *.**deadLetterRoutingKey** ("dlk") *// 绑定死信路由键*
* *.**ttl** (10000) *// 过期时间设置10s,方便测试*
* *.**maxLength** (10L) *// 队列最大长度设为10,方便测试*
* *.build();
}
@Bean("normalExchange")
public DirectExchange normalExchange() {
return ExchangeBuilder.*directExchange*(Constants.*NORMAL_EXCHANGE*).durable(true).build();
}
@Bean("normalBinding")
public Binding normalBinding(@Qualifier("normalQueue")Queue queue,
@Qualifier("normalExchange")Exchange exchange) {
return BindingBuilder.*bind*(queue).to(exchange).with("normal").noargs();
}
**// 死信队列**
@Bean("dlQueue")
public Queue dlQueue() {
return QueueBuilder.*durable*(Constants.*DL_QUEUE*).build();
}
@Bean("dlExchange")
public DirectExchange dlExchange() {
return ExchangeBuilder.*directExchange*(Constants.*DL_EXCHANGE*).durable(true).build();
}
@Bean("dlBinding")
public Binding dlBinding(@Qualifier("dlQueue")Queue queue,
@Qualifier("dlExchange")Exchange exchange) {
return BindingBuilder.*bind*(queue).to(exchange).with("dlk").noargs();
}2. 发送消息
@RequestMapping("/dlx")
public String dlx() {
*// 1. 测试过期时间, 当时间达到TTL, 消息自动进入到死信队列*
* *rabbitTemplate.convertAndSend(Constants.*NORMAL_EXCHANGE*, "normal", "dlx test...");
*// 2. 测试队列长度溢出,消息自动进入到死信队列*
* *for(int i = 0; i < 20; ++i) {
rabbitTemplate.convertAndSend(Constants.*NORMAL_EXCHANGE*, "normal", "dlx test...");
}
return "发送成功!";
}3. 测试死信
① 程序启动之后,观察队列

-
D:队列设置了持久化机制 -
TTL:队列设置了消息过期时间 -
Lim:队列设置了长度(x-max-length) -
DLX:队列设置了死信交换机(x-dead-letter-exchange) -
DLK:队列设置了死信路由键(x-dead-letter-routing-key)
② 测试过期时间,到达过期时间之后,进入死信队列
发送消息:http://127.0.0.1:8080/product/dlx
发送之后:

10秒后,消息进入到死信队列:

生产者首先发送一条消息,然后经过交换器(normal_exchange)顺利地存储到队列(normal_queue)中。由于队列 normal_queue 设置了过期时间为 10s,在这 10s 内没有消费者消费这条消息,那么判定这条消息过期。由于设置了 DLX,过期之时,消息会被丢给交换器(dl_exchange)中,这时根据 RoutingKey 匹配,找到匹配的队列(dl_queue),最后消息被存储在 queue.dlx 这个死信队列中。
③ 测试达到队列长度,消息进入死信队列
队列长度设置为 10,我们发送 20 条数据,会有 10 条数据直接进入到死信队列
发送前,死信队列只有一条数据:

运行后,可以看到死信队列变成了 11 条:

过期之后,正常队列的 10 条也会进入到死信队列:

④ 测试消息拒收
写消费者代码,并强制异常,测试拒绝签收:
@Component
public class DLListener {
*// 监听正常队列*
* ***@RabbitListener(queues = Constants.** ***NORMAL_QUEUE** ***)**
public void normalQueue(Message message, Channel channel) throws InterruptedException, IOException {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
System.*out*.printf("接收到消息: %s, deliveryTag: %d\n", new String(message.getBody()), deliveryTag);
*// 模拟处理失败*
* *int num = 3/0;
System.*out*.println("处理完成");
*// 手动确认*
* *channel.basicAck(deliveryTag, true);
}catch (Exception e){
***// 第三个参数requeue决定是否重新入队,如果为true,则会重新发送;若为false,则直接丢弃,若此时设置了死信,会进入到死信队列** *
* *channel.basicNack(deliveryTag, true,false);
}
}
*// 监听死信队列*
* ***@RabbitListener(queues = Constants.** ***DL_QUEUE** ***)**
public void dlQueue(Message message) {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
System.*out*.printf("死信队列接收到消息: %s, deliveryTag: %d\n", new String(message.getBody()), deliveryTag);
}
}
三、常见面试题💥
1. 死信队列的概念
死信就是因为种种原因,而导致的无法被消费的信息 。
2. 死信的来源
-
消息被拒绝(
Basic.Reject/Basic.Nack),并且设置requeue参数为false -
消息过期
-
队列达到最大长度,消息溢出
3. 死信队列的应用场景
对于 RabbitMQ 来说,死信队列是一个非常有用的特性。它可以处理异常情况下,消息不能够被消费者正确消费而被置入死信队列中的情况,应用程序可以通过消费这个死信队列中的内容来分析当时所遇到的异常情况,进而可以改善和优化系统。
比如:用户支付订单之后,支付系统会给订单系统返回当前订单的支付状态。
为了保证支付信息不丢失,需要使用到死信队列机制。当消息消费异常时,将消息投入到死信队列中,由订单系统的其他消费者来监听这个队列,并对数据进行处理 (比如发送工单等,进行人工确认)。
场景的应用场景还有:
-
消息重试:将死信消息重新发送到原队列或另一个队列进行重试处理。
-
消息丢弃:直接丢弃这些无法处理的消息,以避免它们占用系统资源。
-
日志收集:将死信消息作为日志收集起来,用于后续分析和问题定位。
Ⅶ. 延迟队列
一、概念 && 应用场景
延迟队列(Delayed Queue)即消息被发送以后,并不想让消费者立刻拿到消息,而是等待特定时间后,消费者才能拿到这个消息进行消费。
延迟队列的使用场景有很多,比如:
-
智能家居:用户希望通过手机远程遥控家里的智能设备在指定的时间进行工作。这时候就可以将用户指令发送到延迟队列,当指令设定的时间到了再将指令推送到智能设备。
-
日常管理:预定会议后,需要在会议开始前十五分钟提醒参会人参加会议。
-
用户注册成功后,7天后发送短信,提高用户活跃度等。
-
......
RabbitMQ 本身没有直接支持延迟队列的功能 ,但是可以通过 TTL+死信队列 的组合模拟出延迟队列的功能,所以死信队列章节展示的也是延迟队列的使用。
假设一个应用中需要将每条消息都设置为 10 秒的延迟,生产者通过 normal_exchange 这个交换器将发送的消息存储在 normal_queue 这个队列中。消费者订阅的并非是 normal_queue 这个队列,而是 dl_queue 死信队列。当消息从 normal_queue 这个队列中过期之后被存入 dl_queue 这个队列中,消费者就恰巧消费到了延迟 10 秒的这条消息。

二、TTL+死信队列实现
延迟队列,就是希望等待特定的时间之后,消费者才能拿到这个消息。TTL 刚好可以让消息延迟一段时间成为死信,成为死信的消息会被投递到死信队列里,这样消费者一直消费死信队列里的消息就可以了。
-
声明以及配置队列: (沿用前面死信队列的配置,只不过略做修改)
@Bean("normalQueue") public Queue normalQueue() { return QueueBuilder .*durable*(Constants.*NORMAL_QUEUE*) .deadLetterExchange(Constants.*DL_EXCHANGE*) *// 绑定死信交换机* * *.deadLetterRoutingKey("dlk") *// 绑定死信路由键* * *.build(); } -
发送消息: 发送两条消息,一条消息 10s 后过期,第二条 20s 后过期
@RequestMapping("/delay") public String delay() { *// 发送两条单独带TTL的消息* * *rabbitTemplate.convertAndSend(Constants.*NORMAL_EXCHANGE*, "normal", "delay test 10s..." + new Date(), message -> { message.getMessageProperties().setExpiration("10000"); *// 延迟10s到达死信队列* * *return message; }); rabbitTemplate.convertAndSend(Constants.*NORMAL_EXCHANGE*, "normal", "delay test 20s..." + new Date(), message -> { message.getMessageProperties().setExpiration("20000"); *// 延迟20s到达死信队列* * *return message; }); return "发送成功!"; } -
消费者: 监听死信队列,打印信息,观察现象
*// 监听死信队列* @RabbitListener(queues = Constants.*DL_QUEUE*) public void dlQueue(Message message) { System.*out*.printf("%tc 死信队列接收到消息: %s\n", new Date(), new String(message.getBody())); }

该实现方式存在的问题🐔
把生产消息的顺序修改一下:先发送 20s 过期数据,再发送 10s 过期数据:
@RequestMapping("/delay")
public String delay() {
*// 发送两条单独带TTL的消息*
* *rabbitTemplate.convertAndSend(Constants.*NORMAL_EXCHANGE*, "normal", "delay test 20s..." + new Date(), message -> {
message.getMessageProperties().setExpiration("**20000** "); *// 延迟20s到达死信队列*
* *return message;
});
rabbitTemplate.convertAndSend(Constants.*NORMAL_EXCHANGE*, "normal", "delay test 10s..." + new Date(), message -> {
message.getMessageProperties().setExpiration("**10000** "); *// 延迟10s到达死信队列*
* *return message;
});
return "发送成功!";
}
这时会发现:10s 过期的消息在 20s 后才进入到死信队列??
这是因为消息过期之后,不一定会被马上丢弃。因为 RabbitMQ 只会检查队首消息是否过期 ,如果过期则丢到死信队列,此时就会造成一个问题,如果第一个消息的延时时间很长,第二个消息的延时时间很短,那第二个消息并不会优先得到执行。
所以在考虑使用 TTL+死信队列 实现延迟任务队列的时候,需要确认业务上每个任务的延迟时间是一致的,如果遇到不同的任务类型需要不同的延迟的话,需要为每一种不同延迟时间的消息建立单独的消息队列。
三、延迟队列插件
RabbitMQ 官方提供了一个延迟的插件来实现延迟的功能
参考:https://www.rabbitmq.com/blog/2015/04/16/scheduling-messages-with-rabbitmq
① 安装延迟队列插件

根据自己的 RabbitMQ 版本选择相应版本的延迟插件,下载后上传到服务器或者放到本地的 RabbitMQ 的 plugins 目录中,可以参考下图解释:

-
启动插件 (下面是 linux 系统指令,其它系统指令直接问 gpt 即可)
# 查看插件列表 rabbitmq-plugins list # 启动插件 rabbitmq-plugins enable rabbitmq_delayed_message_exchange # 重启服务 service rabbitmq-server restart -
验证插件
在 RabbitMQ 管理平台查看,新建交换机时是否有延迟消息选项,如果有就说明延迟消息插件已经正常运行了。

② 基于插件延迟队列实现
-
声明与绑定交换机、队列
*// 常量* public static final String *DELAY_EXCHANGE *= "delay_exchange"; public static final String *DELAY_QUEUE *= "delay_queue"; *// 延迟队列* @Bean("delayQueue") public Queue delayQueue() { return QueueBuilder.*durable*(Constants.*DELAY_QUEUE*).build(); // 队列正常设置 } @Bean("delayExchange") public DirectExchange delayExchange() { return ExchangeBuilder.*directExchange*(Constants.*DELAY_EXCHANGE*).**delayed()** .build(); } @Bean("delayBinding") public Binding delayBinding(@Qualifier("delayQueue")Queue queue, @Qualifier("delayExchange")Exchange exchange) { return BindingBuilder.*bind*(queue).to(exchange).with("delay").noargs(); } -
生产者发送两条消息,并设置延迟时间
@RequestMapping("/delay") public String delay() { *// 发送两条单独带TTL的消息* * *rabbitTemplate.convertAndSend(Constants.*DELAY_EXCHANGE*, "delay", "delay test 10s..." + new Date(), message -> { message.getMessageProperties().**setDelayLong(20000L)** ; *// 延迟20s到达死信队列* * *return message; }); rabbitTemplate.convertAndSend(Constants.*DELAY_EXCHANGE*, "delay", "delay test 20s..." + new Date(), message -> { message.getMessageProperties().**setDelayLong(10000L)** ; *// 延迟10s到达死信队列* * *return message; }); return "发送成功!"; } -
消费者监听延迟队列,打印并观察消息
@Component public class DelayListener { @RabbitListener(queues = Constants.*DELAY_QUEUE*) public void delayQueue(Message message) { System.*out*.printf("%tc 延迟队列接收到消息: %s\n", new Date(), new String(message.getBody())); } }

从结果可以看出,使用延迟队列,可以保证消息按照延迟时间到达消费者。
四、两种实现方式的区别
| 实现方式 | 优点 | 缺点 |
|---|---|---|
| TTL+死信 | ① 灵活,不依赖额外插件 ② 适用于任何标准 RabbitMQ 环境 |
① 存在消息顺序问题(先到期的消息可能被后到期的阻塞) ② 需要额外逻辑处理死信消息,系统复杂度提高 |
| 插件 | ① 插件可直接创建延迟队列,实现简单 ② 避免 DLX 的时序问题,顺序更可靠 |
① 依赖特定插件(需安装维护) ② 只支持部分 RabbitMQ 版本,兼容性有限 |
Ⅷ. 事务
RabbitMQ 是基于 AMQP 协议实现的,该协议实现了事务机制,因此 RabbitMQ 也支持事务机制。Spring AMQP 也提供了对事务相关的操作。
RabbitMQ 事务允许开发者确保消息的发送和接收是原子性的,要么全部成功,要么全部失败 。
要使用 RabbitMQ 事务,需要同时完成下面三步操作 !
一、配置事务管理器
因为需要配置事务管理器,所以通常单独配置 RabbitTemplate,然后配置时候调用 rabbitTemplate.``setChannelTransacted``(true) 打开事务管理器,并且需要配置一下事务管理器 RabbitTransactionManager,如下所示:
@Configuration
public class TransactionConfig {
@Bean
public **RabbitTransactionManager** transactionManager(ConnectionFactory connectionFactory) {
return new RabbitTransactionManager(connectionFactory);
}
@Bean("transRabbitTemplate")
public RabbitTemplate transRabbitTemplate(ConnectionFactory connectionFactory) {
RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
rabbitTemplate.**setChannelTransacted** (true); *// 开启事务*
* *return rabbitTemplate;
}
}二、声明队列
声明队列就和普通队列一样,不需要什么特殊设置:
*// 事务*
@Bean("transQueue")
public Queue transQueue() {
return QueueBuilder.*durable*("transQueue").build();
}三、发送消息时打开事务
@Resource(name = "transRabbitTemplate")
private RabbitTemplate transRabbitTemplate; // 注入transRabbitTemplate
**@Transactional**
@RequestMapping("/trans")
public String trans() {
transRabbitTemplate.convertAndSend("", "transQueue", "test trans 1...");
int a = 5 / 0; *// 模拟出现异常*
* *transRabbitTemplate.convertAndSend("", "transQueue", "test trans 2...");
return "发送成功!";
}
如果三个步骤中没做其中的任何一个,都没办法保证事务机制的启动! (自行测试)
Ⅸ. 消息分发
一、概念
当 RabbitMQ 队列拥有多个消费者时,队列会把收到的消息分派给不同的消费者。每条消息只会发送给订阅列表里的一个消费者**(普通队列的点对点消费)** 。这种方式非常适合扩展,如果现在负载加重,那么只需要创建更多的消费者来消费处理消息即可。
默认情况下,RabbitMQ 是以 轮询的方法进行分发的 ,而不管消费者是否已经消费并已经确认了消息。这种方式是不太合理的,试想一下,如果某些消费者消费速度慢,而某些消费者消费速度快,就可能会导致某些消费者消息积压,某些消费者空闲,进而应用整体的吞吐量下降。
如何解决❓❓❓
可以使用 channel.``basicQos``(int prefetchCount),限制当前信道上的消费者所能保持的最大未确认消息的数量 。
其中参数 prefetchCount 设置为 0 时表示没有上限。
比如:消费端调用了
channel.basicQos(5),RabbitMQ 会为该消费者计数,发送一条消息计数+1,消费一条消息计数-1,当达到了设定的上限,RabbitMQ 就不会再向它发送消息了,直到消费者确认了某条消息。类似 TCP/IP 中的 "滑动窗口"。
💥注意事项: basicQos()对拉模式的消费无效 。
二、应用场景
消息分发的常见应用场景有如下:
-
限流
-
非公平分发
① 限流
如下场景:
订单系统每秒最多处理 5000 个请求,正常情况下,订单系统可以正常满足需求。
但是在秒杀时间点,请求瞬间增多,每秒 1w 个请求,如果这些请求全部通过 MQ 发送到订单系统,无疑会把订单系统压垮。

所以 RabbitMQ 提供了限流机制,可以控制消费端一次只拉取 N 个请求,保证消费端的正常运行。
操作:设置 prefetchCount参数 ,同时设置消息确认机制为手动应答 manual。
-
配置
prefetch参数,设置应答方式为手动应答spring: rabbitmq: addresses: amqp://liren:123123@127.0.0.1/lirendada listener: simple: * *acknowledge-mode: ***manual** ** # 手动确认* * ***prefetch** : 5 -
配置交换机,队列
*// 常量* public static final String *QOS_EXCHANGE *= "qos_exchange"; public static final String *QOS_QUEUE *= "qos_queue"; *// 消息分发* @Bean("qosQueue") public Queue qosQueue() { return QueueBuilder.*durable*(Constants.*QOS_QUEUE*).build(); } @Bean("qosExchange") public DirectExchange qosExchange() { return ExchangeBuilder.*directExchange*(Constants.*QOS_EXCHANGE*).build(); } @Bean("qosBinding") public Binding qosBinding(@Qualifier("qosQueue")Queue queue, @Qualifier("qosExchange")Exchange exchange) { return BindingBuilder.*bind*(queue).to(exchange).with("qos").noargs(); } -
发送消息,一次发送20条消息
@RequestMapping("/qos") public String qos() { for(int i = 0; i < 20; ++i) { rabbitTemplate.convertAndSend(Constants.*QOS_EXCHANGE*, "qos", "test qos..." + i); } return "发送成功!"; } -
消费者监听,进行手动确认
@Configuration public class QosListener { @RabbitListener(queues = Constants.*QOS_QUEUE*) public void qosQueue(Message message, Channel channel) throws IOException { long deliveryTag = message.getMessageProperties().getDeliveryTag(); System.*out*.printf("接收到消息:%s,deliveryTag:%d%n", new String(message.getBody()), deliveryTag); *// channel.basicAck(deliveryTag, true); // 注释掉,不进行确认,观察现象* * *} }
发送消息时,需要先把手动确认注释掉,不然会直接消费掉

将 prefetch 注释掉,然后重新启动程序观察现象:

可以看到消息一次性都被消费者拿到了,就没有限流效果了!
② 负载均衡
如下图所示,在有两个消费者的情况下,一个消费者处理任务非常快,另一个非常慢,就会造成一个消费者会一直很忙,而另一个消费者很闲。这是因为 RabbitMQ 只是在消息进入队列时分派消息,它不考虑消费者未确认消息的数量。

我们可以使用设置 prefetch=1 的方式,告诉 RabbitMQ 一次只给一个消费者一条消息,也就是说,在处理并确认前一条消息之前,不要向该消费者发送新消息 。此时,它会将它分派给下一个不忙的消费者。
-
配置
prefetch参数,设置应答方式为手动应答spring: rabbitmq: addresses: amqp://liren:123123@127.0.0.1/lirendada listener: simple: * *acknowledge-mode: ***manual** ** # 手动确认消息* * ***prefetch** : 1 -
启动两个消费者(用休眠模拟业务处理耗时的不同)
@Configuration public class QosListener { @RabbitListener(queues = Constants.*QOS_QUEUE*) public void qosQueue1(Message message, Channel channel) throws IOException, InterruptedException { long deliveryTag = message.getMessageProperties().getDeliveryTag(); System.*out*.printf("qosQueue1 接收到消息:%s,deliveryTag:%d%n", new String(message.getBody()), deliveryTag); Thread.*sleep*(1000); *// 模拟快业务处理,1s* * *channel.basicAck(deliveryTag, true); } @RabbitListener(queues = Constants.*QOS_QUEUE*) public void qosQueue2(Message message, Channel channel) throws IOException, InterruptedException { long deliveryTag = message.getMessageProperties().getDeliveryTag(); System.*out*.printf("qosQueue2 接收到消息:%s,deliveryTag:%d%n", new String(message.getBody()), deliveryTag); Thread.*sleep*(2000); *// 模拟慢业务处理,2s* * *channel.basicAck(deliveryTag, true); } }

💥注意: deliveryTag 有重复是因为两个消费者使用的是不同的 Channel,每个 Channel 上的 deliveryTag 是独立计数的。
