自拍偷在线精品自拍偷,亚洲欧美中文日韩v在线观看不卡

RabbitMQ如何保證消息的可靠投遞?

開發(fā) 前端
總而言之,在生產(chǎn)環(huán)境中,我們一般都是單條手動ack,消費(fèi)失敗后不會重新入隊(duì)(因?yàn)楹艽蟾怕蔬€會再次失敗),而是將消息重新投遞到死信隊(duì)列,方便以后排查問題。

 

Spring Boot整合RabbitMQ

github地址:

https://github.com/erlieStar/rabbitmq-examples

Spring有三種配置方式

  1. 基于XML
  2. 基于JavaConfig
  3. 基于注解

當(dāng)然現(xiàn)在已經(jīng)很少使用XML來做配置了,只介紹一下用JavaConfig和注解的配置方式

RabbitMQ整合Spring Boot,我們只需要增加對應(yīng)的starter即可

  1. <dependency> 
  2.   <groupId>org.springframework.boot</groupId> 
  3.   <artifactId>spring-boot-starter-amqp</artifactId> 
  4. </dependency> 

基于注解

在application.yaml的配置如下

  1. spring: 
  2.   rabbitmq: 
  3.     host: myhost 
  4.     port: 5672 
  5.     username: guest 
  6.     password: guest 
  7.     virtual-host: / 
  8.  
  9. log: 
  10.   exchange: log.exchange 
  11.   info: 
  12.     queue: info.log.queue 
  13.     binding-key: info.log.key 
  14.   error: 
  15.     queue: error.log.queue 
  16.     binding-key: error.log.key 
  17.   all
  18.     queue: all.log.queue 
  19.     binding-key'*.log.key' 

消費(fèi)者代碼如下

  1. @Slf4j 
  2. @Component 
  3. public class LogReceiverListener { 
  4.  
  5.     /** 
  6.      * 接收info級別的日志 
  7.      */ 
  8.     @RabbitListener( 
  9.             bindings = @QueueBinding( 
  10.                     value = @Queue(value = "${log.info.queue}", durable = "true"), 
  11.                     exchange = @Exchange(value = "${log.exchange}", type = ExchangeTypes.TOPIC), 
  12.                     key = "${log.info.binding-key}" 
  13.             ) 
  14.     ) 
  15.     public void infoLog(Message message) { 
  16.         String msg = new String(message.getBody()); 
  17.         log.info("infoLogQueue 收到的消息為: {}", msg); 
  18.     } 
  19.  
  20.     /** 
  21.      * 接收所有的日志 
  22.      */ 
  23.     @RabbitListener( 
  24.             bindings = @QueueBinding( 
  25.                     value = @Queue(value = "${log.all.queue}", durable = "true"), 
  26.                     exchange = @Exchange(value = "${log.exchange}", type = ExchangeTypes.TOPIC), 
  27.                     key = "${log.all.binding-key}" 
  28.             ) 
  29.     ) 
  30.     public void allLog(Message message) { 
  31.         String msg = new String(message.getBody()); 
  32.         log.info("allLogQueue 收到的消息為: {}", msg); 
  33.     } 

生產(chǎn)者如下

  1. @RunWith(SpringRunner.class) 
  2. @SpringBootTest 
  3. public class MsgProducerTest { 
  4.  
  5.     @Autowired 
  6.     private AmqpTemplate amqpTemplate; 
  7.     @Value("${log.exchange}"
  8.     private String exchange; 
  9.     @Value("${log.info.binding-key}"
  10.     private String routingKey; 
  11.  
  12.     @SneakyThrows 
  13.     @Test 
  14.     public void sendMsg() { 
  15.         for (int i = 0; i < 5; i++) { 
  16.             String message = "this is info message " + i; 
  17.             amqpTemplate.convertAndSend(exchange, routingKey, message); 
  18.         } 
  19.  
  20.         System.in.read(); 
  21.     } 

Spring Boot針對消息ack的方式和原生api針對消息ack的方式有點(diǎn)不同

原生api消息ack的方式

消息的確認(rèn)方式有2種

自動確認(rèn)(autoAck=true)

手動確認(rèn)(autoAck=false)

消費(fèi)者在消費(fèi)消息的時(shí)候,可以指定autoAck參數(shù)

String basicConsume(String queue, boolean autoAck, Consumer callback)

autoAck=false: RabbitMQ會等待消費(fèi)者顯示回復(fù)確認(rèn)消息后才從內(nèi)存(或者磁盤)中移出消息

autoAck=true: RabbitMQ會自動把發(fā)送出去的消息置為確認(rèn),然后從內(nèi)存(或者磁盤)中刪除,而不管消費(fèi)者是否真正的消費(fèi)了這些消息

手動確認(rèn)的方法如下,有2個(gè)參數(shù)

basicAck(long deliveryTag, boolean multiple)

deliveryTag: 用來標(biāo)識信道中投遞的消息。RabbitMQ 推送消息給Consumer時(shí),會附帶一個(gè)deliveryTag,以便Consumer可以在消息確認(rèn)時(shí)告訴RabbitMQ到底是哪條消息被確認(rèn)了。

RabbitMQ保證在每個(gè)信道中,每條消息的deliveryTag從1開始遞增

multiple=true: 消息id<=deliveryTag的消息,都會被確認(rèn)

myltiple=false: 消息id=deliveryTag的消息,都會被確認(rèn)

消息一直不確認(rèn)會發(fā)生啥?

如果隊(duì)列中的消息發(fā)送到消費(fèi)者后,消費(fèi)者不對消息進(jìn)行確認(rèn),那么消息會一直留在隊(duì)列中,直到確認(rèn)才會刪除。

如果發(fā)送到A消費(fèi)者的消息一直不確認(rèn),只有等到A消費(fèi)者與rabbitmq的連接中斷,rabbitmq才會考慮將A消費(fèi)者未確認(rèn)的消息重新投遞給另一個(gè)消費(fèi)者

Spring Boot中針對消息ack的方式

有三種方式,定義在AcknowledgeMode枚舉類中

方式 解釋
NONE 沒有ack,等價(jià)于原生api中的autoAck=true
MANUAL 用戶需要手動發(fā)送ack或者nack
AUTO 方法正常結(jié)束,spring boot 框架返回ack,發(fā)生異常spring boot框架返回nack

spring boot針對消息默認(rèn)的ack的方式為AUTO。

在實(shí)際場景中,我們一般都是手動ack。

application.yaml的配置改為如下

  1. spring: 
  2.   rabbitmq: 
  3.     host: myhost 
  4.     port: 5672 
  5.     username: guest 
  6.     password: guest 
  7.     virtual-host: / 
  8.     listener: 
  9.       simple: 
  10.         acknowledge-mode: manual # 手動ack,默認(rèn)為auto 

相應(yīng)的消費(fèi)者代碼改為

  1. @Slf4j 
  2. @Component 
  3. public class LogListenerManual { 
  4.  
  5.     /** 
  6.      * 接收info級別的日志 
  7.      */ 
  8.     @RabbitListener( 
  9.             bindings = @QueueBinding( 
  10.                     value = @Queue(value = "${log.info.queue}", durable = "true"), 
  11.                     exchange = @Exchange(value = "${log.exchange}", type = ExchangeTypes.TOPIC), 
  12.                     key = "${log.info.binding-key}" 
  13.             ) 
  14.     ) 
  15.     public void infoLog(Message message, Channel channel) throws Exception { 
  16.         String msg = new String(message.getBody()); 
  17.         log.info("infoLogQueue 收到的消息為: {}", msg); 
  18.         try { 
  19.             // 這里寫各種業(yè)務(wù)邏輯 
  20.             channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); 
  21.         } catch (Exception e) { 
  22.             channel.basicNack(message.getMessageProperties().getDeliveryTag(), falsefalse); 
  23.         } 
  24.     } 

我們上面用到的注解,作用如下

注解 作用
RabbitListener 消費(fèi)消息,可以定義在類上,方法上,當(dāng)定義在類上時(shí)需要和RabbitHandler配合使用
QueueBinding 定義綁定關(guān)系
Queue 定義隊(duì)列
Exchange 定義交換機(jī)
RabbitHandler RabbitListener定義在類上時(shí),需要用RabbitHandler指定處理的方法

基于JavaConfig

既然用注解這么方便,為啥還需要JavaConfig的方式呢?

JavaConfig方便自定義各種屬性,比如同時(shí)配置多個(gè)virtual host等

具體代碼看GitHub把

RabbitMQ如何保證消息的可靠投遞

一個(gè)消息往往會經(jīng)歷如下幾個(gè)階段

在這里插入圖片描述

所以要保證消息的可靠投遞,只需要保證這3個(gè)階段的可靠投遞即可

生產(chǎn)階段

這個(gè)階段的可靠投遞主要靠ConfirmListener(發(fā)布者確認(rèn))和ReturnListener(失敗通知)

前面已經(jīng)介紹過了,一條消息在RabbitMQ中的流轉(zhuǎn)過程為

producer -> rabbitmq broker cluster -> exchange -> queue -> consumer

ConfirmListener可以獲取消息是否從producer發(fā)送到broker

ReturnListener可以獲取從exchange路由不到queue的消息

我用Spring Boot Starter 的api來演示一下效果

application.yaml

  1. spring: 
  2.   rabbitmq: 
  3.     host: myhost 
  4.     port: 5672 
  5.     username: guest 
  6.     password: guest 
  7.     virtual-host: / 
  8.     listener: 
  9.       simple: 
  10.         acknowledge-mode: manual # 手動ack,默認(rèn)為auto 
  11.  
  12. log: 
  13.   exchange: log.exchange 
  14.   info: 
  15.     queue: info.log.queue 
  16.     binding-key: info.log.key 

發(fā)布者確認(rèn)回調(diào)

  1. @Component 
  2. public class ConfirmCallback implements RabbitTemplate.ConfirmCallback { 
  3.  
  4.     @Autowired 
  5.     private MessageSender messageSender; 
  6.  
  7.     @Override 
  8.     public void confirm(CorrelationData correlationData, boolean ack, String cause) { 
  9.         String msgId = correlationData.getId(); 
  10.         String msg = messageSender.dequeueUnAckMsg(msgId); 
  11.         if (ack) { 
  12.             System.out.println(String.format("消息 {%s} 成功發(fā)送給mq", msg)); 
  13.         } else { 
  14.             // 可以加一些重試的邏輯 
  15.             System.out.println(String.format("消息 {%s} 發(fā)送mq失敗", msg)); 
  16.         } 
  17.     } 

失敗通知回調(diào)

  1. @Component 
  2. public class ReturnCallback implements RabbitTemplate.ReturnCallback { 
  3.  
  4.     @Override 
  5.     public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) { 
  6.         String msg = new String(message.getBody()); 
  7.         System.out.println(String.format("消息 {%s} 不能被正確路由,routingKey為 {%s}", msg, routingKey)); 
  8.     } 
  1. @Configuration 
  2. public class RabbitMqConfig { 
  3.  
  4.     @Bean 
  5.     public ConnectionFactory connectionFactory( 
  6.             @Value("${spring.rabbitmq.host}") String host, 
  7.             @Value("${spring.rabbitmq.port}"int port, 
  8.             @Value("${spring.rabbitmq.username}") String username, 
  9.             @Value("${spring.rabbitmq.password}") String password
  10.             @Value("${spring.rabbitmq.virtual-host}") String vhost) { 
  11.         CachingConnectionFactory connectionFactory = new CachingConnectionFactory(host); 
  12.         connectionFactory.setPort(port); 
  13.         connectionFactory.setUsername(username); 
  14.         connectionFactory.setPassword(password); 
  15.         connectionFactory.setVirtualHost(vhost); 
  16.         connectionFactory.setPublisherConfirms(true); 
  17.         connectionFactory.setPublisherReturns(true); 
  18.         return connectionFactory; 
  19.     } 
  20.  
  21.     @Bean 
  22.     public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory, 
  23.                                          ReturnCallback returnCallback, ConfirmCallback confirmCallback) { 
  24.         RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory); 
  25.         rabbitTemplate.setReturnCallback(returnCallback); 
  26.         rabbitTemplate.setConfirmCallback(confirmCallback); 
  27.         // 要想使 returnCallback 生效,必須設(shè)置為true 
  28.         rabbitTemplate.setMandatory(true); 
  29.         return rabbitTemplate; 
  30.     } 

這里我對RabbitTemplate做了一下包裝,主要就是發(fā)送的時(shí)候增加消息id,并且保存消息id和消息的對應(yīng)關(guān)系,因?yàn)镽abbitTemplate.ConfirmCallback只能拿到消息id,并不能拿到消息內(nèi)容,所以需要我們自己保存這種映射關(guān)系。在一些可靠性要求比較高的系統(tǒng)中,你可以將這種映射關(guān)系存到數(shù)據(jù)庫中,成功發(fā)送刪除映射關(guān)系,失敗則一直發(fā)送

  1. @Component 
  2. public class MessageSender { 
  3.  
  4.     @Autowired 
  5.     private RabbitTemplate rabbitTemplate; 
  6.  
  7.     public final Map<String, String> unAckMsgQueue = new ConcurrentHashMap<>(); 
  8.  
  9.     public void convertAndSend(String exchange, String routingKey, String message) { 
  10.         String msgId = UUID.randomUUID().toString(); 
  11.         CorrelationData correlationData = new CorrelationData(); 
  12.         correlationData.setId(msgId); 
  13.         rabbitTemplate.convertAndSend(exchange, routingKey, message, correlationData); 
  14.         unAckMsgQueue.put(msgId, message); 
  15.     } 
  16.  
  17.     public String dequeueUnAckMsg(String msgId) { 
  18.         return unAckMsgQueue.remove(msgId); 
  19.     } 
  20.  

測試代碼為

  1. @RunWith(SpringRunner.class) 
  2. @SpringBootTest 
  3. public class MsgProducerTest { 
  4.  
  5.     @Autowired 
  6.     private MessageSender messageSender; 
  7.     @Value("${log.exchange}"
  8.     private String exchange; 
  9.     @Value("${log.info.binding-key}"
  10.     private String routingKey; 
  11.  
  12.     /** 
  13.      * 測試失敗通知 
  14.      */ 
  15.     @SneakyThrows 
  16.     @Test 
  17.     public void sendErrorMsg() { 
  18.         for (int i = 0; i < 3; i++) { 
  19.             String message = "this is error message " + i; 
  20.             messageSender.convertAndSend(exchange, "test", message); 
  21.         } 
  22.         System.in.read(); 
  23.     } 
  24.  
  25.     /** 
  26.      * 測試發(fā)布者確認(rèn) 
  27.      */ 
  28.     @SneakyThrows 
  29.     @Test 
  30.     public void sendInfoMsg() { 
  31.         for (int i = 0; i < 3; i++) { 
  32.             String message = "this is info message " + i; 
  33.             messageSender.convertAndSend(exchange, routingKey, message); 
  34.         } 
  35.         System.in.read(); 
  36.     } 

先來測試失敗者通知

輸出為

  1. 消息 {this is error message 0} 不能被正確路由,routingKey為 {test} 
  2. 消息 {this is error message 0} 成功發(fā)送給mq 
  3. 消息 {this is error message 2} 不能被正確路由,routingKey為 {test} 
  4. 消息 {this is error message 2} 成功發(fā)送給mq 
  5. 消息 {this is error message 1} 不能被正確路由,routingKey為 {test} 
  6. 消息 {this is error message 1} 成功發(fā)送給mq 

消息都成功發(fā)送到broker,但是并沒有被路由到queue中

再來測試發(fā)布者確認(rèn)

輸出為

  1. 消息 {this is info message 0} 成功發(fā)送給mq 
  2. infoLogQueue 收到的消息為: {this is info message 0} 
  3. infoLogQueue 收到的消息為: {this is info message 1} 
  4. 消息 {this is info message 1} 成功發(fā)送給mq 
  5. infoLogQueue 收到的消息為: {this is info message 2} 
  6. 消息 {this is info message 2} 成功發(fā)送給mq 

消息都成功發(fā)送到broker,也成功被路由到queue中

存儲階段

這個(gè)階段的高可用還真沒研究過,畢竟集群都是運(yùn)維搭建的,后續(xù)有時(shí)間的話會把這快的內(nèi)容補(bǔ)充一下

消費(fèi)階段

消費(fèi)階段的可靠投遞主要靠ack來保證。

總而言之,在生產(chǎn)環(huán)境中,我們一般都是單條手動ack,消費(fèi)失敗后不會重新入隊(duì)(因?yàn)楹艽蟾怕蔬€會再次失敗),而是將消息重新投遞到死信隊(duì)列,方便以后排查問題

總結(jié)一下各種情況

  1. ack后消息從broker中刪除
  2. nack或者reject后,分為如下2種情況

(1) reque=true,則消息會被重新放入隊(duì)列

(2) reque=fasle,消息會被直接丟棄,如果指定了死信隊(duì)列的話,會被投遞到死信隊(duì)列

本文轉(zhuǎn)載自微信公眾號「Java識堂」,可以通過以下二維碼關(guān)注。轉(zhuǎn)載本文請聯(lián)系Java識堂公眾號。

 

責(zé)任編輯:武曉燕 來源: Java識堂
相關(guān)推薦

2023-03-06 08:16:04

SpringRabbitMQ

2021-04-27 07:52:18

RocketMQ消息投遞

2024-05-09 08:04:23

RabbitMQ消息可靠性

2021-02-02 11:01:31

RocketMQ消息分布式

2020-09-27 07:44:08

RabbitMQ投遞消息

2024-12-18 07:43:49

2023-12-04 09:23:49

分布式消息

2023-11-30 18:03:02

TCP傳輸

2024-08-12 12:17:03

2024-02-20 11:30:23

光纖

2022-07-26 20:00:35

場景RabbitMQMQ

2022-08-02 11:27:25

RabbitMQ消息路由

2023-10-17 16:30:00

TCP

2024-06-27 08:00:17

2024-08-06 09:55:25

2023-11-27 17:29:43

Kafka全局順序性

2024-05-23 12:11:39

2024-12-04 15:38:43

2021-03-08 10:19:59

MQ消息磁盤

2021-08-10 09:59:15

RabbitMQ消息微服務(wù)
點(diǎn)贊
收藏

51CTO技術(shù)棧公眾號