【RabbitMQ】高级进阶,消费端限流、消息超时、死信队列、延时队列(重要)
目录
消费端限流
在 RabbitMQ 中,消息消费者的消息限流(message prefetching)是一种控制消费者接收消息数量的机制。这种机制可以避免消费者因处理能力不足而导致消息积压或丢失。
spring:
rabbitmq:
host: 10.0.70.102
port: 5672
username: guest
password: 123456
virtual-host: /
listener:
simple:
acknowledge-mode: manual
prefetch: 1 # 设置每次最多从消息队列服务器取回多少消息
通过设置预取值和手动应答,消费者端可以控制自身处理消息的速度,有效地实现消费端的限流。
消息超时
消息队列层面
消息层面
@RestController
@RequestMapping(value="/limiting")
public class LimitingController {
public static final String EXCHANGE_DIRECT = "limit.direct";
public static final String ROUTING_KEY = "limit";
@Resource
private RabbitTemplate rabbitTemplate;
@GetMapping( "/msgTimeout")
public String send(){
for (int i = 0; i < 10; i++) {
// 创建消息后置处理器对象
MessagePostProcessor postProcessor = message -> {
// 设置消息的过期时间,单位是毫秒
message.getMessageProperties().setExpiration("5000");
return message;
};
rabbitTemplate.convertAndSend(EXCHANGE_DIRECT, ROUTING_KEY, "测试设置消息过期时间" + i,postProcessor);
}
return "";
}
}
死信/死信队列
死信,当一个消息无法被消费掉,它就变成死信。
为了处理这些死信,RabbitMQ引入了死信队列的概念。当消息被标记为死信后,如果配置了死信队列,RabbitMQ会将该消息发送到死信交换机(Dead Letter Exchange)。死信交换机再根据配置的路由键(Routing Key)将消息投递到指定的死信队列中。
在死信队列中,可以对消息进行重新处理、记录或丢弃等操作。例如,可以将死信消息重新发送到另一个队列以供其他消费者再次尝试处理,或者将消息记录到日志中以供后续分析。
产生的原因大致有以下三种:
private static final String NORMAL_QUEUE = "normal.queue";
/**
* 监听确认队列当中的消息
*/
@RabbitListener(queues = NORMAL_QUEUE)
public void confirmMessage(String dataMsg,Message message, Channel channel) throws IOException{
// 1、获取当前消息的 deliveryTag 值备用
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
// 2、正常业务操作
log.info("消费端接收到消息内容:" + dataMsg);
System.out.println(10 / 0);
// 3、给 RabbitMQ 服务器返回 ACK 确认信息
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
// 4、获取信息,看当前消息是否曾经被投递过
Boolean redelivered = message.getMessageProperties().getRedelivered();
if (!redelivered) {
// 5、如果没有被投递过,那就重新放回队列,重新投递,再试一次
channel.basicNack(deliveryTag, false, true);
} else {
// 6、如果已经被投递过,且这一次仍然进入了 catch 块,那么返回拒绝且不再放回队列
channel.basicReject(deliveryTag, false);
}
}
}
查看输入结果,可以看到当消费者端,出现异常,进入异常逻辑后,再次投递的消息拒绝消费,那么本条消息会进入死信队列,有监听死信队列的消费者区处理。

/**
* 发布正常业务消息
*/
@GetMapping("/overLoad")
public String overLoadMessage(){
String message = "验证溢出的死信消息";
String date = new Date().toString();
log.info("生产者在:{},发布了消息:{}",date,message);
for (int i = 0; i < 12 ; i++) {
rabbitTemplate.convertAndSend(NORMAL_EXCHANGE,"businessKey",message+i,new CorrelationData(i+1+""));
}
return "生产者在:"+date+",发布了12条消息:"+message;
}
监听死信队列后,可以查看输出结果,原本先发送0和1两条消息溢出,被死信队列监听到:

@Configuration
public class DeadLetterConfig {
private static final String NORMAL_EXCHANGE = "normal.exchange";
private static final String NORMAL_QUEUE = "normal.queue";
private static final String DEAD_EXCHANGE = "dead.exchange";
private static final String DEAD_QUEUE = "dead.queue";
private static final String DEAD_ROUTING_KEY = "dead_routing";
/**
* 创建一个正常交换机
*/
@Bean(NORMAL_EXCHANGE)
public DirectExchange normalExchange() {
return ExchangeBuilder.directExchange(NORMAL_EXCHANGE).durable(true).build();
}
/**
* 创建死信交换机
*/
@Bean(DEAD_EXCHANGE)
public DirectExchange deadExchange() {
return ExchangeBuilder.directExchange(DEAD_EXCHANGE).durable(true).build();
}
/**
* 创建正常队列
* 设置消息过期参数 x-message-ttl:5000(5秒);
* 设置死信交换机 x-dead-letter-exchange:DEAD_EXCHANGE
*/
@Bean(NORMAL_QUEUE)
public Queue normalQueue() {
return QueueBuilder
.durable(NORMAL_QUEUE)
// 绑定死信交换机
.deadLetterExchange(DEAD_EXCHANGE)
// 死信队列路由关键字
.deadLetterRoutingKey(DEAD_ROUTING_KEY)
// 队列每条消息只能存活5s
.ttl(5000)
// 队列最大长度10
.maxLength(10)
.build();
}
/**
* 死信队列
*/
@Bean(DEAD_QUEUE)
public Queue backupQueue() {
return QueueBuilder.durable(DEAD_QUEUE).build();
}
/**
* 绑定正常交换机和正常队列
*
* @param queue 正常队列
* @param exchange 正常交换机
*/
@Bean
public Binding confirmQueueBindingConfirmExchange(
@Qualifier(NORMAL_QUEUE) Queue queue,
@Qualifier(NORMAL_EXCHANGE) Exchange exchange) {
return BindingBuilder.bind(queue).to(exchange).with("businessKey").noargs();
}
/**
* 绑定死信交换机和死信队列
*
* @param queue 死信队列
* @param exchange 死信交换机
*/
@Bean
public Binding backupQueueBindingBackupExchange(
@Qualifier(DEAD_QUEUE) Queue queue,
@Qualifier(DEAD_EXCHANGE) Exchange exchange) {
return BindingBuilder.bind(queue).to(exchange).with(DEAD_ROUTING_KEY).noargs();
}
}
发送一个正常消息:
private static final String NORMAL_EXCHANGE = "normal.exchange";
@Resource
private RabbitTemplate rabbitTemplate;
/**
* 发布正常业务消息
*/
@GetMapping("/dead")
public String sendDeadMessage(){
String message = "验证超时死信消息";
String date = new Date().toString();
log.info("生产者在:{},发布了消息:{}",date,message);
//发送一条路由key正确、id为1的消息
rabbitTemplate.convertAndSend(NORMAL_EXCHANGE,"businessKey",message,new CorrelationData("1"));
return "生产者在:"+date+",发布了一条消息:"+message;
}
添加死信队列监听,处理消费死信消息:
/**
* 监听死信队列当中的消息
*/
@RabbitListener(queues = DEAD_QUEUE)
public void deadMessage(Message message){
String date = new Date().toString();
log.info("死信消费者在:{},收到了消息:{}",date,new String(message.getBody()));
}
配置消息为手动应答,待5s超时时间过后查看输出,可以看到:

备份交换机和死信队列区别
备份交换机:备份交换器是为了实现没有路由到队列的消息,声明交换机的时候添加属性alternate-exchange,声明一个备用交换机,一般声明为fanout类型,这样交换机收到路由不到目标队列的消息时,就会发送到备用交换机绑定的队列中。---还没到队列
死信队列:当消息在一个队列中变成死信 (dead message) 之后,它能被重新发送到另一个交换机。死信 Exchange其实就是一种普通的Exchange,和创建其他Exchange没有两 样。只是在某一个设置为死信 Exchange的队列中有消息过期(死信)了,会自动触发消息的转发,发送到死信Exchange中去。---已经到队列
延时队列
延时队列,见名之意,就是放置到队列中的消息不是马上消费掉,而是在延缓一段时间后再去消费。RabbitMQ本身不直接支持延时消息,但可以通过两个特性结合使用来实现延时队列:死信队列(Dead Letter Queues)和消息的存活时间(Time-To-Live,TTL)。

基于死信队列和消息的超时时间实现延时队列功能大体步骤:
1、创建死信交换机和死信队列,绑定特定的路由键,如“dead_routing_key”;

2、创建正常消息交换机 和正常消息队列,给队列添加‘’x-dead-letter-exchange‘’参数绑定刚创建的死信交换机,添加‘x-dead-letter-routing-key’参数,绑定死信路由key。如果想让整个队列都有一个超时时间,可以再添加‘x-message-ttl’参数,设置队列消息的超时时间;

3、消费者监听死信队列

完成后,当生产者发送消息到NORMAL_QUEUE队列后,由于没有消费者监听处理,代超过延时时间5s后,就可以到信息队列中看到这条延时消息了;

RabbitMQ的延时队列还可以利用插件实现
插件地址:https://github.com/rabbitmq/rabbitmq-delayed-message-exchange
此插件最多支持延时两天的时间

具体操作可以参考官网,或待后续...
更多推荐



所有评论(0)