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

Java消息隊(duì)列RabbitMQ之消息回調(diào)詳解

 更新時(shí)間:2023年07月31日 09:43:39   作者:迷鹿小女子  
這篇文章主要介紹了Java消息隊(duì)列RabbitMQ之消息回調(diào)詳解,消息回調(diào),其實(shí)就是消息確認(rèn)(生產(chǎn)者推送消息成功,消費(fèi)者接收消息成功)  , 對(duì)于程序來(lái)說(shuō),發(fā)送者沒(méi)法確認(rèn)是否發(fā)送成功,需要的朋友可以參考下

消息100%的投遞

消息如何保障100%的投遞成功? 什么是生產(chǎn)端的可靠性投遞?

  • 保障消息的成功發(fā)出
  • 保障MQ節(jié)點(diǎn)的成功接收
  • 發(fā)送端收到MQ節(jié)點(diǎn)Broker確認(rèn)應(yīng)答
  • 完善的消息進(jìn)行補(bǔ)償機(jī)制

BAT/TMD互聯(lián)網(wǎng)大廠的解決方案:

  • 消息落庫(kù),對(duì)消息狀態(tài)進(jìn)行打標(biāo)
  • 消息的延遲投遞,做二次確認(rèn),回調(diào)檢查

在這里插入圖片描述

在這里插入圖片描述

冪等性概念

冪等性是什么?

  • 我們可以借鑒數(shù)據(jù)庫(kù)的樂(lè)觀鎖機(jī)制
  • 比如我們執(zhí)行一條更新庫(kù)存的SQL語(yǔ)句
  • pdate t_repository set cont = cont -1,version = version + 1 where version = 1
  • Elasticsearch也是嚴(yán)格遵循冪等性概念,每次數(shù)據(jù)更新,version+1

消費(fèi)端-冪等性保障

在海量訂單產(chǎn)生的業(yè)務(wù)高峰期,如何避免消息的重復(fù)消費(fèi)問(wèn)題?

消費(fèi)實(shí)現(xiàn)冪等性,就意味著,我們的消息永遠(yuǎn)不會(huì)消費(fèi)多次,即使我們收到了多條一樣的消息

業(yè)界主流的冪等性操作

唯一ID+指紋碼機(jī)制,利用數(shù)據(jù)庫(kù)主鍵去重 利用Redis的原子性去實(shí)現(xiàn)

  • 唯一ID+指紋碼機(jī)制,利用數(shù)據(jù)庫(kù)主鍵去重
  • Select count(1) from T_order where ID=唯一ID+指紋碼

好處:實(shí)現(xiàn)簡(jiǎn)單

壞處:高并發(fā)下有數(shù)據(jù)庫(kù)寫(xiě)入的性能瓶頸

解決方案:根據(jù)ID進(jìn)行分庫(kù)分表進(jìn)行算法路由

利用Redis的原子性去實(shí)現(xiàn)

使用Redis進(jìn)行冪等,需要考慮的問(wèn)題

第一:我們是否要進(jìn)行數(shù)據(jù)落庫(kù),如果落庫(kù)的話(huà),關(guān)鍵解決的問(wèn)題是數(shù)據(jù)庫(kù)和緩存如何做到原子性?

第二:如果不進(jìn)行落庫(kù),那么都存儲(chǔ)到緩存中,如何設(shè)置定時(shí)同步策略?

Confirm確認(rèn)消息

理解Confirm消息確認(rèn)機(jī)制 消息的確認(rèn),是指生產(chǎn)者投遞消息后,如果Broker收到消息,則會(huì)給我們生產(chǎn)者 一個(gè)應(yīng)答。

生產(chǎn)者進(jìn)行接收應(yīng)答,用來(lái)確定這條消息是否正常的發(fā)送到Broker,這種方式也是消息的可靠性投遞的核心保障

在這里插入圖片描述

如何實(shí)現(xiàn)Confirm確認(rèn)消息?

第一步:在Channel上開(kāi)啟確認(rèn)模式:channel.confirmSelect()

第二步:在channel上添加監(jiān)聽(tīng):addConfirmListener,監(jiān)聽(tīng)成功和失敗的返回結(jié)果,根據(jù)具體的結(jié)果對(duì)消息進(jìn)行重新發(fā)送、或記錄日志等后續(xù)處理!

消費(fèi)端代碼

package com.xieminglu.rabbitmqapi.confirm;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.QueueingConsumer;
public class Consumer {
    public static void main(String[] args) throws Exception {
        //1 創(chuàng)建ConnectionFactory
        ConnectionFactory connectionFactory = new ConnectionFactory();
        connectionFactory.setHost("192.168.248.132");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        //2 獲取C    onnection
        Connection connection = connectionFactory.newConnection();
        //3 通過(guò)Connection創(chuàng)建一個(gè)新的Channel
        Channel channel = connection.createChannel();
        String exchangeName = "test_confirm_exchange";
        String routingKey = "confirm.#";
        String queueName = "test_confirm_queue";
        //4 聲明交換機(jī)和隊(duì)列 然后進(jìn)行綁定設(shè)置, 最后制定路由Key
        channel.exchangeDeclare(exchangeName, "topic", true);
        channel.queueDeclare(queueName, true, false, false, null);
        channel.queueBind(queueName, exchangeName, routingKey);
        //5 創(chuàng)建消費(fèi)者
        QueueingConsumer queueingConsumer = new QueueingConsumer(channel);
        channel.basicConsume(queueName, true, queueingConsumer);
        while(true){
            QueueingConsumer.Delivery delivery = queueingConsumer.nextDelivery();
            String msg = new String(delivery.getBody());
            System.err.println("消費(fèi)端: " + msg);
        }
    }
}

服務(wù)提供方代碼

package com.xieminglu.rabbitmqapi.confirm;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.ConfirmListener;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.io.IOException;
public class Producer {
    public static void main(String[] args) throws Exception {
        //1 創(chuàng)建ConnectionFactory
        ConnectionFactory connectionFactory = new ConnectionFactory();
        connectionFactory.setHost("192.168.248.132");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        //2 獲取C    onnection
        Connection connection = connectionFactory.newConnection();
        //3 通過(guò)Connection創(chuàng)建一個(gè)新的Channel
        Channel channel = connection.createChannel();
        //4 指定我們的消息投遞模式: 消息的確認(rèn)模式
        channel.confirmSelect();
        String exchangeName = "test_confirm_exchange";
        String routingKey = "confirm.save";
        //5 發(fā)送一條消息
        String msg = "Hello RabbitMQ Send confirm message!";
        channel.basicPublish(exchangeName, routingKey, null, msg.getBytes());
        //6 添加一個(gè)確認(rèn)監(jiān)聽(tīng)
        channel.addConfirmListener(new ConfirmListener() {
            @Override
            public void handleNack(long deliveryTag, boolean multiple) throws IOException {
                System.err.println("-------no ack!-----------");
            }
            @Override
            public void handleAck(long deliveryTag, boolean multiple) throws IOException {
                System.err.println("-------ack!-----------");
            }
        });
    }
}

Retrn返回消息

Retrn Listener用于處理一些不可路由的消息!

  • 正常情況:我們的消息生產(chǎn)者,通過(guò)指定一個(gè)Exchange和RotingKey,把消息送達(dá)到某一個(gè)隊(duì)列中去,然后我們的消費(fèi)者監(jiān)聽(tīng)隊(duì)列,進(jìn)行消費(fèi)處理操作!
  • 異常情況:在某些情況下,如果我們?cè)诎l(fā)送消息的時(shí)候,當(dāng)前的Exchange不存在或者指定的路由key路由不到,這個(gè)時(shí)候如果我們需要監(jiān)聽(tīng)這種不可達(dá)的消息,就需要使用Retrn Listener!

在基礎(chǔ)API中有一個(gè)關(guān)鍵的配置項(xiàng) Mandatory:如果為tre,則監(jiān)聽(tīng)器會(huì)接收到路由不可達(dá)的消息,然后進(jìn)行后續(xù)處理,如果為false,那么Broker端自動(dòng)刪除該消息!

在這里插入圖片描述

消費(fèi)端代碼

package com.xieminglu.rabbitmqapi.returnlistener;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.QueueingConsumer;
public class Consumer {
    public static void main(String[] args) throws Exception {
        ConnectionFactory connectionFactory = new ConnectionFactory();
        connectionFactory.setHost("192.168.248.132");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        Connection connection = connectionFactory.newConnection();
        Channel channel = connection.createChannel();
        String exchangeName = "test_return_exchange";
        String routingKey = "return.#";
        String queueName = "test_return_queue";
        channel.exchangeDeclare(exchangeName, "topic", true, false, null);
        channel.queueDeclare(queueName, true, false, false, null);
        channel.queueBind(queueName, exchangeName, routingKey);
        QueueingConsumer queueingConsumer = new QueueingConsumer(channel);
        channel.basicConsume(queueName, true, queueingConsumer);
        while(true){
            QueueingConsumer.Delivery delivery = queueingConsumer.nextDelivery();
            String msg = new String(delivery.getBody());
            System.err.println("消費(fèi)者: " + msg);
        }
    }
}

生產(chǎn)端代碼

package com.xieminglu.rabbitmqapi.returnlistener;
import com.rabbitmq.client.*;
import java.io.IOException;
public class Producer {
    public static void main(String[] args) throws Exception {
        ConnectionFactory connectionFactory = new ConnectionFactory();
        connectionFactory.setHost("192.168.248.132");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        Connection connection = connectionFactory.newConnection();
        Channel channel = connection.createChannel();
        String exchange = "test_return_exchange";
        String routingKey = "return.save";
        String routingKeyError = "abc.save";
        String msg = "Hello RabbitMQ Return Message";
        channel.addReturnListener(new ReturnListener() {
            @Override
            public void handleReturn(int replyCode, String replyText, String exchange,
                                     String routingKey, AMQP.BasicProperties properties, byte[] body) throws IOException {
                System.err.println("---------handle  return----------");
                System.err.println("replyCode: " + replyCode);
                System.err.println("replyText: " + replyText);
                System.err.println("exchange: " + exchange);
                System.err.println("routingKey: " + routingKey);
                System.err.println("properties: " + properties);
                System.err.println("body: " + new String(body));
            }
        });
        //消息投遞成功,會(huì)被消費(fèi)者所消費(fèi)
        channel.basicPublish(exchange, routingKey, true, null, msg.getBytes());
        //消息不可達(dá),將觸發(fā)ReturnListener
//         channel.basicPublish(exchange, routingKeyError, true, null, msg.getBytes());
    }
}

自定義消費(fèi)者

我們一般就是在代碼中編寫(xiě)while循環(huán),進(jìn)行consmer.nextDelivery方法進(jìn)行獲取下一條消息,然后進(jìn)行消費(fèi)處理!

但是我們使用自定義的Consmer更加的方便,解耦性更加的強(qiáng),也是實(shí)際工作中最常用的使用方式!

自定義消費(fèi)端代碼

package com.xieminglu.rabbitmqapi.consumer;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;
import java.io.IOException;
public class MyConsumer extends DefaultConsumer {
    public MyConsumer(Channel channel) {
        super(channel);
    }
    @Override
    public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
        System.err.println("-----------consume message----------");
        System.err.println("consumerTag: " + consumerTag);
        System.err.println("envelope: " + envelope);
        System.err.println("properties: " + properties);
        System.err.println("body: " + new String(body));
    }
}

消費(fèi)端調(diào)用

package com.xieminglu.rabbitmqapi.consumer;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
public class Consumer {
    public static void main(String[] args) throws Exception {
        ConnectionFactory connectionFactory = new ConnectionFactory();
        connectionFactory.setHost("192.168.248.132");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        Connection connection = connectionFactory.newConnection();
        Channel channel = connection.createChannel();
        String exchangeName = "test_consumer_exchange";
        String routingKey = "consumer.#";
        String queueName = "test_consumer_queue";
        channel.exchangeDeclare(exchangeName, "topic", true, false, null);
        channel.queueDeclare(queueName, true, false, false, null);
        channel.queueBind(queueName, exchangeName, routingKey);
        channel.basicConsume(queueName, true, new MyConsumer(channel));
        }
        }

生產(chǎn)端調(diào)用

package com.xieminglu.rabbitmqapi.consumer;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
public class Producer {
    public static void main(String[] args) throws Exception {
        ConnectionFactory connectionFactory = new ConnectionFactory();
        connectionFactory.setHost("192.168.248.132");
        connectionFactory.setPort(5672);
        connectionFactory.setVirtualHost("/");
        Connection connection = connectionFactory.newConnection();
        Channel channel = connection.createChannel();
        String exchange = "test_consumer_exchange";
        String routingKey = "consumer.save";
        String msg = "Hello RabbitMQ Consumer Message";
        for(int i =0; i<5; i ++){
            channel.basicPublish(exchange, routingKey, true, null, msg.getBytes());
        }
    }
}

到此這篇關(guān)于Java消息隊(duì)列RabbitMQ之消息回調(diào)詳解的文章就介紹到這了,更多相關(guān)RabbitMQ消息回調(diào)內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • JVM回收跨代垃圾的方式詳解

    JVM回收跨代垃圾的方式詳解

    在Java堆內(nèi)存中,年輕代和老年代之間存在的對(duì)象相互引用,假設(shè)現(xiàn)在要進(jìn)行一次新生代的YGC,但新生代中的對(duì)象可能被老年代所引用的,為了找到新生代中的存活對(duì)象,不得不遍歷整個(gè)老年代,這樣明顯效率很低下,那么如何快速識(shí)別并回收這種引用對(duì)象呢
    2024-02-02
  • Java實(shí)現(xiàn)生成n個(gè)不重復(fù)的隨機(jī)數(shù)

    Java實(shí)現(xiàn)生成n個(gè)不重復(fù)的隨機(jī)數(shù)

    這篇文章主要為大家詳細(xì)介紹了Java實(shí)現(xiàn)生成n個(gè)不重復(fù)的隨機(jī)數(shù),文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2020-05-05
  • java如何實(shí)現(xiàn)基于opencv全景圖合成實(shí)例代碼

    java如何實(shí)現(xiàn)基于opencv全景圖合成實(shí)例代碼

    全景圖相信大家應(yīng)該都不陌生,下面這篇文章主要給大家介紹了關(guān)于java如何實(shí)現(xiàn)基于opencv全景圖合成的相關(guān)資料,文中通過(guò)示例代碼介紹的非常詳細(xì),需要的朋友可以參考借鑒,下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2018-07-07
  • Java 你知道什么是耦合、如何解(降低)耦合

    Java 你知道什么是耦合、如何解(降低)耦合

    這篇文章主要介紹了Java 你知道什么是耦合、如何解(降低)耦合的操作,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-09-09
  • Java的封裝類(lèi)和裝箱拆箱詳解

    Java的封裝類(lèi)和裝箱拆箱詳解

    Java中存在基礎(chǔ)數(shù)據(jù)類(lèi)型,但是在某些情況下,我們要對(duì)基礎(chǔ)數(shù)據(jù)類(lèi)型進(jìn)行對(duì)象的操作,例如,集合中只能存對(duì)象,而不能存在基礎(chǔ)數(shù)據(jù)類(lèi)型,于是便出現(xiàn)了封裝類(lèi),本文將詳細(xì)給大家介紹Java封裝類(lèi)和裝箱拆箱,需要的朋友可以參考下
    2023-05-05
  • 關(guān)于spring中單例Bean引用原型Bean產(chǎn)生的問(wèn)題及解決

    關(guān)于spring中單例Bean引用原型Bean產(chǎn)生的問(wèn)題及解決

    這篇文章主要介紹了關(guān)于spring中單例Bean引用原型Bean產(chǎn)生的問(wèn)題及解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-06-06
  • springboot Quartz動(dòng)態(tài)修改cron表達(dá)式的方法

    springboot Quartz動(dòng)態(tài)修改cron表達(dá)式的方法

    這篇文章主要介紹了springboot Quartz動(dòng)態(tài)修改cron表達(dá)式的方法,小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2018-09-09
  • Java實(shí)現(xiàn)將PDF轉(zhuǎn)換為Word的示例詳解

    Java實(shí)現(xiàn)將PDF轉(zhuǎn)換為Word的示例詳解

    在日常的數(shù)據(jù)處理、文檔編輯和系統(tǒng)集成工作中,將不可編輯的PDF文檔轉(zhuǎn)換為可編輯的Word文檔是一項(xiàng)常見(jiàn)且重要的需求,下面我們就來(lái)看看如何使用java實(shí)現(xiàn)這一操作吧
    2025-08-08
  • Spring Task定時(shí)任務(wù)的實(shí)現(xiàn)詳解

    Spring Task定時(shí)任務(wù)的實(shí)現(xiàn)詳解

    這篇文章主要介紹了SpringBoot定時(shí)任務(wù)功能詳細(xì)解析,這次的功能開(kāi)發(fā)過(guò)程中也算是對(duì)其內(nèi)涵的進(jìn)一步了解,以后遇到定時(shí)任務(wù)的處理也更清晰,更有效率了,對(duì)SpringBoot定時(shí)任務(wù)相關(guān)知識(shí)感興趣的朋友一起看看吧
    2022-08-08
  • Spring中自定義Schema如何解析生效詳解

    Spring中自定義Schema如何解析生效詳解

    Spring2.5在2.0的基于Schema的Bean配置的基礎(chǔ)之上,再增加了擴(kuò)展XML配置的機(jī)制。通過(guò)該機(jī)制,我們可以編寫(xiě)自己的Schema,并根據(jù)自定義的Schema用自定的標(biāo)簽配置Bean,下面這篇文章主要介紹了關(guān)于Spring中自定義Schema如何解析生效的相關(guān)資料,需要的朋友可以參考下
    2018-07-07

最新評(píng)論

卢龙县| 桦南县| 镇康县| 河东区| 堆龙德庆县| 石景山区| 招远市| 团风县| 嘉义县| 临泽县| 稷山县| 霍城县| 博野县| 瓦房店市| 宁晋县| 安丘市| 高碑店市| 玉山县| 曲阜市| 兴化市| 尼勒克县| 仙居县| 绩溪县| 高阳县| 全州县| 萨迦县| 志丹县| 正安县| 文昌市| 老河口市| 常熟市| 城口县| 远安县| 海阳市| 商洛市| 扎囊县| 宁海县| 大兴区| 大丰市| 沅陵县| 芜湖县|