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

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

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

簡介

消息過濾是指消費者一端在消費消息時,對消息進行選擇性過濾,只消費符合過濾條件的消息。 RocketMQ的消息過濾機制大致分為兩種:標簽過濾和類過濾。其中標簽過濾又分為Tag過濾和SQL92過濾。

根據(jù)TAG過濾消息

消息發(fā)送端只能設置一個tag,消息接收端可以設置多個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);
       }
   }

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

消費者

@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);
    }
}

說明:

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

測試結果

從結果我可以看出第2條為TAGC、第7條為TAGC、第8條為TAGD,第3條為TAGD,第5條為TAGA,第0條為TAGA,而消費端監(jiān)聽的TAG為TAGA、TAGC、TAGD所以對于不符合條件的消息進行了過濾。

根據(jù)SQL表達式過濾消息

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

RocketMQ的SQL表達式語法 只定義了一些基本的語法功能。

  • 數(shù)字比較,如>,>=,<,<=,BETWEEN,=;
  • 字符比較,如:=,<>,IN;IS NULL or IS NOT NULL;
  • 邏輯運算符:AND, OR, NOT;
  • 常量類型:
  • 數(shù)值,如:123, 3.1415;
  • 字符, 如:‘abc’, 必須使用單引號;
  • 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);
   }

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

消費者

@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());
        //消息內容
        logger.info("send succss content is:{}", msg);
        //消息參數(shù)
        logger.info("send mssage parma is:{}", paramStr);
    }
}

說明:

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

啟動程序報錯The broker does not support consumer to filter message by SQL92

原因:默認情況下broke沒有開啟對SQL語法的支持,需要修改配置

1.打開rocketmq服務下的broke.conf文件,添加如下配置即可。

2.重啟broke服務即可.

測試結果

說明:只有滿足SQL條件能進行消費。

總結

本文講解了RocketMQ實現(xiàn)消息過濾,針對不同的業(yè)務場景選擇合適的方案即可,如果疑問,請隨時反饋,

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

相關文章

  • 新手初學Java List 接口

    新手初學Java List 接口

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

    強烈推薦這些提升代碼效率的IDEA使用技巧

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

    基于Security實現(xiàn)OIDC單點登錄的詳細流程

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

    java中各種對象的比較方法

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

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

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

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

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

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

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

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

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

    java導出Excel(非模板)可導出多個sheet方式

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

    Maven構建Hadoop項目的實踐步驟

    本文主要介紹了Maven構建Hadoop項目的實踐步驟,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2023-06-06

最新評論

巢湖市| 邹平县| 合阳县| 特克斯县| 天台县| 苍山县| 房山区| 霍州市| 龙胜| 鹤山市| 桐城市| 邯郸市| 铅山县| 昆明市| 灌云县| 六盘水市| 山阳县| 马龙县| 南乐县| 遵义县| 红原县| 宜阳县| 雷州市| 西乌珠穆沁旗| 易门县| 光山县| 岳阳县| 岳普湖县| 鲁山县| 加查县| 乐至县| 临武县| 大石桥市| 台东市| 仁怀市| 长武县| 栖霞市| 砚山县| 云南省| 阿坝| 新安县|