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

Spring?Boot?整合RocketMq實(shí)現(xiàn)消息過(guò)濾功能

 更新時(shí)間:2022年06月07日 14:29:43   作者:劍圣無(wú)痕  
這篇文章主要介紹了Spring?Boot?整合RocketMq實(shí)現(xiàn)消息過(guò)濾,本文講解了RocketMQ實(shí)現(xiàn)消息過(guò)濾,針對(duì)不同的業(yè)務(wù)場(chǎng)景選擇合適的方案即可,需要的朋友可以參考下

簡(jiǎn)介

消息過(guò)濾是指消費(fèi)者一端在消費(fèi)消息時(shí),對(duì)消息進(jìn)行選擇性過(guò)濾,只消費(fèi)符合過(guò)濾條件的消息。 RocketMQ的消息過(guò)濾機(jī)制大致分為兩種:標(biāo)簽過(guò)濾和類過(guò)濾。其中標(biāo)簽過(guò)濾又分為T(mén)ag過(guò)濾和SQL92過(guò)濾。

根據(jù)TAG過(guò)濾消息

消息發(fā)送端只能設(shè)置一個(gè)tag,消息接收端可以設(shè)置多個(gè)tag。

生產(chǎn)者

 public void sendTagMessage()
   {
       String[] tags = new String[]{"TagA", "TagB", "TagC", "TagD", "TagE"};
       for(int i=0;i<10;i++)
       {
           String tag = tags[i % tags.length];
           logger.info("sendTagMessage tag is :{}",tag);
           String msg = "hello, 這是第" + (i + 1) + "條消息";
           org.springframework.messaging.Message<String> msg1 = MessageBuilder.withPayload(msg).build(); 
           rocketMQTemplate.convertAndSend("test-tag-rocketmq" + ":" + tag, msg1);
       }
   }

說(shuō)明:示例中循環(huán)發(fā)送了10條消息,每條消息設(shè)置了一個(gè)tag發(fā)送過(guò)濾消息的格式為:topic:tag的形式,注意發(fā)送端只能設(shè)定一個(gè)tag。

消費(fèi)者

@Component
@RocketMQMessageListener(consumerGroup="test-tagrocketmq-group",topic="test-tag-rocketmq",selectorExpression="TagA || TagC || TagD",selectorType=SelectorType.TAG, messageModel = MessageModel.CLUSTERING)
public class TagConsumer implements RocketMQListener<Object>
{
    private Logger logger =LoggerFactory.getLogger(getClass());
    @Override
    public void onMessage(Object o)
    {
        String msg=JSON.toJSONString(o);
        logger.info("send TagA || TagC || TagD  succss content is:{}", msg);
    }
}

說(shuō)明:

  • selectorType:指定消息通過(guò)的tag的方式,默認(rèn)為SelectorType.TAG
  • messageModel:指定消息的消費(fèi)模式,默認(rèn)為MessageModel.CLUSTERING模式每條消息只能由一個(gè)消費(fèi)者消費(fèi),而MessageModel.BROADCASTING模式為廣播模式,所有訂閱者都能消費(fèi)。
  • selectorExpression :指定那些Tag消息能夠被消費(fèi),多個(gè)采用||分割。

測(cè)試結(jié)果

從結(jié)果我可以看出第2條為T(mén)AGC、第7條為T(mén)AGC、第8條為T(mén)AGD,第3條為T(mén)AGD,第5條為T(mén)AGA,第0條為T(mén)AGA,而消費(fèi)端監(jiān)聽(tīng)的TAG為T(mén)AGA、TAGC、TAGD所以對(duì)于不符合條件的消息進(jìn)行了過(guò)濾。

根據(jù)SQL表達(dá)式過(guò)濾消息

SQL表達(dá)式方式可以根據(jù)發(fā)送消息時(shí)輸入的屬性進(jìn)行一些計(jì)算。

RocketMQ的SQL表達(dá)式語(yǔ)法 只定義了一些基本的語(yǔ)法功能。

  • 數(shù)字比較,如>,>=,<,<=,BETWEEN,=;
  • 字符比較,如:=,<>,IN;IS NULL or IS NOT NULL;
  • 邏輯運(yùn)算符:AND, OR, NOT;
  • 常量類型:
  • 數(shù)值,如:123, 3.1415;
  • 字符, 如:‘abc’, 必須使用單引號(hào);
  • NULL,特殊常量
  • Boolean, TRUE or FALSE;

生產(chǎn)者

   public void sendSQLMessage()
   {
       String msg = "hello, 這是第1條消息";
       org.springframework.messaging.Message<String> message = MessageBuilder.withPayload(msg).build() ;
       Map<String, Object> headers = new HashMap<>() ;
       headers.put("i", 5) ;
       rocketMQTemplate.convertAndSend("test-sql-rocketmq", message, headers);
   }

說(shuō)明:傳遞了參數(shù)為5進(jìn)行條件判斷。

消費(fèi)者

@Component
@RocketMQMessageListener(consumerGroup="test-sqlrocketmq-group",topic="test-sql-rocketmq",selectorExpression = "i=5",selectorType=SelectorType.SQL92, messageModel = MessageModel.CLUSTERING)
public class SQLConsumer implements RocketMQListener<MessageExt>
{
    private Logger logger =LoggerFactory.getLogger(getClass());
    @Override
    public void onMessage(MessageExt message)
    {
        String msg=new String(message.getBody());
        String paramStr=JSON.toJSONString(message.getProperties());
        //消息內(nèi)容
        logger.info("send succss content is:{}", msg);
        //消息參數(shù)
        logger.info("send mssage parma is:{}", paramStr);
    }
}

說(shuō)明:

  • selectorType:指定消息通過(guò)的tag的方式,默認(rèn)為SelectorType.CLUSTERING
  • messageModel:指定消息的消費(fèi)模式,默認(rèn)為MessageModel.CLUSTERING模式每條消息只能由一個(gè)消費(fèi)者消費(fèi),而MessageModel.BROADCASTING模式為廣播模式,所有訂閱者都能消費(fèi)。
  • selectorExpression : 采用rocketMQ支持的表達(dá)式。例如i=5

啟動(dòng)程序報(bào)錯(cuò)The broker does not support consumer to filter message by SQL92

原因:默認(rèn)情況下broke沒(méi)有開(kāi)啟對(duì)SQL語(yǔ)法的支持,需要修改配置

1.打開(kāi)rocketmq服務(wù)下的broke.conf文件,添加如下配置即可。

2.重啟broke服務(wù)即可.

測(cè)試結(jié)果

說(shuō)明:只有滿足SQL條件能進(jìn)行消費(fèi)。

總結(jié)

本文講解了RocketMQ實(shí)現(xiàn)消息過(guò)濾,針對(duì)不同的業(yè)務(wù)場(chǎng)景選擇合適的方案即可,如果疑問(wèn),請(qǐng)隨時(shí)反饋,

到此這篇關(guān)于Spring Boot 整合RocketMq實(shí)現(xiàn)消息過(guò)濾的文章就介紹到這了,更多相關(guān)Spring Boot消息過(guò)濾內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 新手初學(xué)Java List 接口

    新手初學(xué)Java List 接口

    這篇文章主要介紹了Java集合操作之List接口及其實(shí)現(xiàn)方法,詳細(xì)分析了Java集合操作中List接口原理、功能、用法及操作注意事項(xiàng),需要的朋友可以參考下
    2021-07-07
  • 強(qiáng)烈推薦這些提升代碼效率的IDEA使用技巧

    強(qiáng)烈推薦這些提升代碼效率的IDEA使用技巧

    在平常的開(kāi)發(fā)中,發(fā)現(xiàn)一些同事對(duì)Idea 使用的不是很熟練,僅僅用來(lái)編輯,編譯,不能很好的發(fā)揮Idea 的神奇.整理了下我平常用的一些技巧,希望你能從中學(xué)習(xí)到一些.需要的朋友可以參考下
    2021-05-05
  • 基于Security實(shí)現(xiàn)OIDC單點(diǎn)登錄的詳細(xì)流程

    基于Security實(shí)現(xiàn)OIDC單點(diǎn)登錄的詳細(xì)流程

    本文主要是給大家介紹 OIDC 的核心概念以及如何通過(guò)對(duì) Spring Security 的授權(quán)碼模式進(jìn)行擴(kuò)展來(lái)實(shí)現(xiàn) OIDC 的單點(diǎn)登錄。對(duì)Security實(shí)現(xiàn)OIDC單點(diǎn)登錄的詳細(xì)過(guò)程感興趣的朋友,一起看看吧
    2021-09-09
  • java中各種對(duì)象的比較方法

    java中各種對(duì)象的比較方法

    Java對(duì)象的比較是初學(xué)者不易掌握的,下面這篇文章主要給大家介紹了關(guān)于java中各種對(duì)象的比較方法,文中通過(guò)實(shí)例代碼以及圖文介紹的非常詳細(xì),需要的朋友可以參考下
    2023-04-04
  • 高吞吐、線程安全的LRU緩存詳解

    高吞吐、線程安全的LRU緩存詳解

    這篇文章主要介紹了高吞吐、線程安全的LRU緩存詳解,分享了相關(guān)代碼示例,小編覺(jué)得還是挺不錯(cuò)的,具有一定借鑒價(jià)值,需要的朋友可以參考下
    2018-02-02
  • Java通過(guò)JsApi方式實(shí)現(xiàn)微信支付

    Java通過(guò)JsApi方式實(shí)現(xiàn)微信支付

    本文講解了Java如何實(shí)現(xiàn)JsApi方式的微信支付,代碼內(nèi)容詳細(xì),文章思路清晰,需要的朋友可以參考下
    2015-07-07
  • java實(shí)現(xiàn)學(xué)籍管理系統(tǒng)

    java實(shí)現(xiàn)學(xué)籍管理系統(tǒng)

    這篇文章主要為大家詳細(xì)介紹了java實(shí)現(xiàn)學(xué)籍管理系統(tǒng),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2016-12-12
  • springboot配置多數(shù)據(jù)源(靜態(tài)和動(dòng)態(tài)數(shù)據(jù)源)

    springboot配置多數(shù)據(jù)源(靜態(tài)和動(dòng)態(tài)數(shù)據(jù)源)

    在開(kāi)發(fā)過(guò)程中,很多時(shí)候都會(huì)有垮數(shù)據(jù)庫(kù)操作數(shù)據(jù)的情況,需要同時(shí)配置多套數(shù)據(jù)源,本文主要介紹了springboot配置多數(shù)據(jù)源(靜態(tài)和動(dòng)態(tài)數(shù)據(jù)源),感興趣的可以了解一下
    2023-09-09
  • java導(dǎo)出Excel(非模板)可導(dǎo)出多個(gè)sheet方式

    java導(dǎo)出Excel(非模板)可導(dǎo)出多個(gè)sheet方式

    Java開(kāi)發(fā)中,導(dǎo)出Excel是常見(jiàn)需求,有時(shí)需要支持多個(gè)Sheet導(dǎo)出,此技巧介紹非模板方式實(shí)現(xiàn)單標(biāo)題單Sheet以及多Sheet導(dǎo)出,標(biāo)題一致或不一致均可,可換成Map使用,適合個(gè)人開(kāi)發(fā)者和需要Excel導(dǎo)出功能的場(chǎng)景
    2024-09-09
  • Maven構(gòu)建Hadoop項(xiàng)目的實(shí)踐步驟

    Maven構(gòu)建Hadoop項(xiàng)目的實(shí)踐步驟

    本文主要介紹了Maven構(gòu)建Hadoop項(xiàng)目的實(shí)踐步驟,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2023-06-06

最新評(píng)論

新泰市| 陇川县| 原平市| 涿鹿县| 徐州市| 原平市| 大悟县| 天气| 眉山市| 福州市| 辛集市| 周口市| 阿尔山市| 涟水县| 富平县| 沽源县| 桐庐县| 万年县| 徐闻县| 香港| 涞源县| 定陶县| 楚雄市| 长顺县| 阳东县| 常山县| 阿拉尔市| 宁武县| 延川县| 修文县| 兴宁市| 渝中区| 莱芜市| 庄河市| 什邡市| 获嘉县| 香格里拉县| 通化县| 钦州市| 休宁县| 崇信县|