最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

RabbitMQ消息確認(rèn)機(jī)制剖析

 更新時間:2022年08月18日 10:22:59   作者:劍圣無痕  
這篇文章主要為大家介紹了RabbitMQ消息確認(rèn)機(jī)制剖析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

前言

上一章講解了RabbitMq的三種Exchange消息發(fā)送的模式,但是在默認(rèn)情況下RabbitMQ并不能保證消息是否發(fā)送成功,以及是否能被成功消費,為了保證消息在傳遞過程中不丟失,需要對消息進(jìn)行確認(rèn)機(jī)制,來提高消息的可靠性。

消息確認(rèn)

基本流程

說明:

  • 生產(chǎn)者發(fā)送消息到RabbitMQ Server后,RabbitMQ Server需要對生產(chǎn)者進(jìn)行消息Confirm確認(rèn)。
  • 消費者消費消息后需要對 RabbitMQ Server進(jìn)行消息ACK確認(rèn)。

消息確認(rèn)模式

RabbitMq提供了兩種消息發(fā)送者確認(rèn)模式分別為: ConfirmCallback確認(rèn)模式和 ReturnCallback退回模式。

ConfirmCallback確認(rèn)模式

@Component
public class RabbitConfirmConfig implements ConfirmCallback
{
    private Logger logger = LoggerFactory.getLogger(RabbitConfirmConfig.class);
    public void confirm(CorrelationData correlationData, boolean ack,
            String cause)
    {
        logger.info("數(shù)據(jù)內(nèi)容:{}",correlationData);
        logger.info("是否確認(rèn)成功:{}",ack);
        logger.info("錯誤原因:{}",cause);
        if (!ack) 
        {
            logger.info("exchange produce confirm message send error" + cause);
        }
        else 
        {
            logger.info("exchange produce confirm message send success");
        }
    }
}

說明:ConfirmCallback模式確認(rèn),需要重寫confirm接方法,此方法的三個參數(shù)分別為:CorrelationData、ack、cause

  • CorrelationData:對象內(nèi)部只有一個id屬性,用來表示當(dāng)前消息的唯一性。
  • ack:消息投遞狀態(tài),true表示投遞成功
  • cause: 消息投遞失敗原因

雖然消息被broker接收到只能表示已經(jīng)到達(dá)MQ服務(wù)器,但是并不能保證消息一定會被投遞到目標(biāo) queue里。所以我們需要實現(xiàn)returnCallback來進(jìn)行相關(guān)處理。

ReturnCallback退回模式

@Component
public class RabbitReturnConfig implements ReturnCallback
{
    private Logger logger = LoggerFactory.getLogger(RabbitReturnConfig.class);
    public void returnedMessage(Message message, int replyCode,
            String replyText, String exchange, String routingKey)
    {
       logger.info("消息發(fā)送送到隊列信息:");
       logger.info("發(fā)生消息:{}",message);
       logger.info("回應(yīng)碼:{}",replyCode);
       logger.info("回應(yīng)信息:{}",replyText);
       logger.info("交換機(jī):{}",exchange);
       logger.info("路由鍵:{}",routingKey);
    }
}

說明:實現(xiàn)接口ReturnCallback重寫returnedMessage()方法,方法有五個參數(shù)message(消息體)、replyCode(響應(yīng)code)、replyText(響應(yīng)內(nèi)容)、exchange(交換機(jī))、routingKey(路由鍵)。

消息發(fā)送者確認(rèn)

@Component
public class MqConfirmProduce
{
    @Autowired
    private RabbitTemplate rabbitTemplate;
    @Autowired
    private RabbitConfirmConfig rabbitConfirmConfig;
    @Autowired
    private RabbitReturnConfig rabbitReturnConfig;
    /**
     * 
     * @param exchange 消息交互機(jī)名稱
     * @param routeKey 消息路由鍵的名稱
     * @param message  消息內(nèi)容
     */
    public void sendMessage(String exchange ,String routeKey,Object msg)
    {
        //確保消息發(fā)送失敗后可以重新返回到隊列中
        rabbitTemplate.setMandatory(true);
        // 消費者確認(rèn)收到消息后,手動ack回執(zhí)回調(diào)處理
        rabbitTemplate.setConfirmCallback(rabbitConfirmConfig);
        //消息投遞到隊列失敗回調(diào)處理
        rabbitTemplate.setReturnCallback(rabbitReturnConfig);
        //保證消息唯一性
        CorrelationData correlationData =new CorrelationData(UUID.randomUUID().toString());
        //發(fā)送消息
        rabbitTemplate.convertAndSend(exchange,routeKey,msg,
                message -> {
                    message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
                    return message;
                },
                correlationData);
    }
}

說明:注意需要開啟消息確認(rèn)的配置:

  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest
    virtual-host: /
    #開啟發(fā)送確認(rèn)
    publisher-confirms: true
    # 開啟發(fā)送失敗退回
    publisher-returns: true
    listener:
      simple:
       # 手動確認(rèn)
        acknowledge-mode: manual
        retry: 
          enabled: true

消息接收者確認(rèn)

@Component
@RabbitListener(queues = "testQueue")
public class MqConfirmConsumer
{
    private static final Logger logger = LoggerFactory.getLogger(MqConfirmConsumer.class);
    @RabbitHandler
    public void receive(String msg, Channel channel, Message message) throws IOException 
    {
        logger.info("receive message content:{}",message);
        try
        {
            logger.info("開始消息確認(rèn)");
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
            logger.info("消息確認(rèn)成功");
        }
        catch (Exception e)
        {
             logger.error("消息確認(rèn)失敗,即將再次返回隊列中");               channel.basicNack(message.getMessageProperties().getDeliveryTag(), true, true); 
        }
    }
}

說明:消息者確認(rèn)消息有三種模式,分別為basicAck、basicNack、basicReject。

basicAck模式

表示成功確認(rèn),使用此回執(zhí)方法后,消息會被rabbitmq broker刪除。

void basicAck(long deliveryTag, boolean multiple) 
  • deliveryTag:消息投遞序號,
  • multiple:是否批量確認(rèn),值為 true則會一次性ack所有小于當(dāng)前消息deliveryTag的消息。

basicNack模式

表示失敗確認(rèn),一般在消費消息異常時用到此方法,可以將消息重新投遞入隊列。

void basicNack(long deliveryTag, boolean multiple, boolean requeue)
  • deliveryTag:表示消息投遞序號。
  • requeue: 表示消息是否重新入隊列,true表示重新投入隊列中。
  • multiple:是否批量確認(rèn),true表示會一次性ack所有小于當(dāng)前消息deliveryTag的消息。

basicReject模式

basicReject:拒絕消息,與basicNack區(qū)別在于不能進(jìn)行批量操作,其他用法很相似。

void basicReject(long deliveryTag, boolean requeue)
  • deliveryTag:消息投遞序號。
  • requeue:值為true表示消息重新入隊列

測試

測試發(fā)送消息,消息發(fā)送者的確認(rèn)信息如下:

c.s.f.r.config.RabbitConfirmConfig - exchange produce confirm message send success
c.s.f.r.config.RabbitConfirmConfig - 數(shù)據(jù)內(nèi)容:CorrelationData [id=88ea47a5-726d-44c5-9839-1f2a6bf942ed]
c.s.f.r.config.RabbitConfirmConfig - 是否確認(rèn)成功:true
c.s.f.r.config.RabbitConfirmConfig - 錯誤原因:null
c.s.f.r.config.RabbitConfirmConfig - exchange produce confirm message send success

消費者的確認(rèn)信息如下:

receive message content:(Body:'this is test message' MessageProperties [headers={spring_listener_return_correlation=0fcefb6d-acea-4eb2-8484-e3a82f8c584f, spring_returned_message_correlation=88ea47a5-726d-44c5-9839-1f2a6bf942ed}, contentType=text/plain, contentEncoding=UTF-8, contentLength=0, receivedDeliveryMode=PERSISTENT, priority=0, redelivered=false, receivedExchange=testDirect, receivedRoutingKey=testDirectRouting, deliveryTag=2, consumerTag=amq.ctag-dOwkSPuI1e0HR_1Ufu3Erw, consumerQueue=testQueue])
 c.s.f.r.consumer.MqConfirmConsumer - 開始消息確認(rèn)
c.s.f.r.consumer.MqConfirmConsumer - 消息確認(rèn)成功

消費者確認(rèn)失敗

如果消息確認(rèn)在消費者確認(rèn)失敗,那么消息將會重寫投遞導(dǎo)導(dǎo)消息隊列的首部。模擬消費者確認(rèn)失敗場景:

@Component
@RabbitListener(queues = "testQueue")
public class MqConfirmConsumer
{
    private static final Logger logger = LoggerFactory.getLogger(MqConfirmConsumer.class);
    @RabbitHandler
    public void receive(String msg, Channel channel, Message message) throws IOException 
    {
        logger.info("receive message content:{}",message);
        try
        {
            logger.info("開始消息確認(rèn)");
            int c=1/0;
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
            logger.info("消息確認(rèn)成功");
        }
        catch (Exception e)
        {
          logger.error("消息確認(rèn)失敗,即將再次返回隊列中");                   channel.basicNack(message.getMessageProperties().getDeliveryTag(), true, true); 
        }
    }
}

查看執(zhí)行結(jié)果:

c.s.f.r.consumer.MqConfirmConsumer - receive message content:(Body:'this is test message' MessageProperties [headers={spring_listener_return_correlation=0fcefb6d-acea-4eb2-8484-e3a82f8c584f, spring_returned_message_correlation=39d4cdd1-cbeb-4090-91ea-9e5d0bed785c}, contentType=text/plain, contentEncoding=UTF-8, contentLength=0, receivedDeliveryMode=PERSISTENT, priority=0, redelivered=false, receivedExchange=testDirect, receivedRoutingKey=testDirectRouting, deliveryTag=1, consumerTag=amq.ctag-e5GtG455pkm7eWfY3xGleg, consumerQueue=testQueue])
c.s.f.r.consumer.MqConfirmConsumer - 開始消息確認(rèn)
c.s.f.r.consumer.MqConfirmConsumer - 消息確認(rèn)失敗,即將再次返回隊列中

消息已經(jīng)重新返回隊列中。我們查看隊列信息具體如下:

說明:我們可以看到消息為Unacked狀態(tài),消息又會重新會被消費,然后確認(rèn)失敗,又重新被消費,導(dǎo)致死循環(huán)。

解決辦法

針對這種情況,我們將如何處理呢?我們手動確認(rèn)失敗后,并將消息持久入到MySQL中通過定時任務(wù)做補(bǔ)償。然后刪除消息隊列。具體修改如下:

 @RabbitHandler
    public void receive(String msg, Channel channel, Message message) throws IOException 
    {
        logger.info("receive message content:{}",message);
        try
        {
            logger.info("開始消息確認(rèn)");
            int c=1/0;
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
            logger.info("消息確認(rèn)成功");
        }
        catch (Exception e)
        {
            if (message.getMessageProperties().getRedelivered()) 
            {
                logger.error("消息確認(rèn)失敗,拒絕處理");
              //執(zhí)行持久化處理                channel.basicReject(message.getMessageProperties().getDeliveryTag(), false); 
          }
            else 
            {
                logger.error("消息確認(rèn)失敗,即將再次返回隊列中");
                channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true); 
            }
        }
    }

修改后執(zhí)行結(jié)果如下:

總結(jié)

本文講解了RabbitMQ消息確認(rèn)機(jī)制,消息是否需要確認(rèn),我們需要根據(jù)業(yè)務(wù)的場景來分析,如有疑問,請隨時反饋,更多關(guān)于RabbitMQ消息確認(rèn)的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • idea out目錄與target目錄的區(qū)別詳解

    idea out目錄與target目錄的區(qū)別詳解

    這篇文章主要介紹了idea out目錄與target目錄的區(qū)別詳解,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2021-02-02
  • Java基礎(chǔ)知識之CharArrayWriter流的使用

    Java基礎(chǔ)知識之CharArrayWriter流的使用

    這篇文章主要介紹了Java基礎(chǔ)知識之CharArrayWriter流的使用,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-12-12
  • Java字符拼接成字符串的注意點詳解

    Java字符拼接成字符串的注意點詳解

    這篇文章主要介紹了Java字符拼接成字符串的注意點詳解,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2019-07-07
  • Spring整合junit的配置過程圖解

    Spring整合junit的配置過程圖解

    這篇文章主要介紹了Spring整合junit的配置過程圖解,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2020-02-02
  • Spring IOC原理詳解

    Spring IOC原理詳解

    這篇文章主要介紹了Spring IOC原理詳解,具有一定借鑒價值,需要的朋友可以參考下。
    2017-12-12
  • java實現(xiàn)分段讀取文件并通過HTTP上傳的方法

    java實現(xiàn)分段讀取文件并通過HTTP上傳的方法

    這篇文章主要介紹了java實現(xiàn)分段讀取文件并通過HTTP上傳的方法,實例分析了java分段讀取文件及使用http實現(xiàn)文件傳輸?shù)南嚓P(guān)技巧,具有一定參考借鑒價值,需要的朋友可以參考下
    2015-07-07
  • 基于Java的度分秒坐標(biāo)轉(zhuǎn)純經(jīng)緯度坐標(biāo)的漂亮國基地信息管理的方法

    基于Java的度分秒坐標(biāo)轉(zhuǎn)純經(jīng)緯度坐標(biāo)的漂亮國基地信息管理的方法

    本文以java語言為例,詳細(xì)介紹如何管理漂亮國的基地信息,為下一步全球的空間可視化打下堅實的基礎(chǔ),首先介紹如何對數(shù)據(jù)進(jìn)行去重處理,然后介紹在java當(dāng)中如何進(jìn)行度分秒位置的轉(zhuǎn)換,最后結(jié)合實現(xiàn)原型進(jìn)行詳細(xì)的說明,感興趣的朋友跟隨小編一起看看吧
    2024-06-06
  • SpringBoot中的響應(yīng)式web應(yīng)用詳解

    SpringBoot中的響應(yīng)式web應(yīng)用詳解

    這篇文章主要介紹了SpringBoot中的響應(yīng)式web應(yīng)用詳解,本文給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2020-11-11
  • springboot的Customizer源碼解析

    springboot的Customizer源碼解析

    這篇文章主要為大家介紹了springboot的Customizer源碼解析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-08-08
  • Spring Boot之內(nèi)嵌tomcat版本升級操作示例

    Spring Boot之內(nèi)嵌tomcat版本升級操作示例

    這篇文章主要為大家介紹了Spring Boot之內(nèi)嵌tomcat版本升級操作示例,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-06-06

最新評論

行唐县| 德令哈市| 阳朔县| 房产| 崇左市| 衡东县| 麟游县| 崇义县| 韶山市| 锡林浩特市| 安乡县| 平舆县| 罗源县| 浮梁县| 大新县| 星子县| 昌图县| 塘沽区| 称多县| 元谋县| 绥宁县| 长兴县| 双江| 芦山县| 原平市| 绥中县| 兴安县| 台州市| 五大连池市| 阳朔县| 天峻县| 四子王旗| 广汉市| 灌云县| 万年县| 峨边| 三门峡市| 扎鲁特旗| 澎湖县| 治多县| 景宁|