中文字幕av专区_日韩电影在线播放_精品国产精品久久一区免费式_av在线免费观看网站

溫馨提示×

溫馨提示×

您好,登錄后才能下訂單哦!

密碼登錄×
登錄注冊×
其他方式登錄
點擊 登錄注冊 即表示同意《億速云用戶服務條款》

Springboot中RocketMQ怎么實現廣播消息

發布時間:2022-06-22 17:32:11 來源:億速云 閱讀:423 作者:iii 欄目:開發技術

這篇文章主要介紹“Springboot中RocketMQ怎么實現廣播消息”,在日常操作中,相信很多人在Springboot中RocketMQ怎么實現廣播消息問題上存在疑惑,小編查閱了各式資料,整理出簡單好用的操作方法,希望對大家解答”Springboot中RocketMQ怎么實現廣播消息”的疑惑有所幫助!接下來,請跟著小編一起來學習吧!

RocketMQ消息模式主要有兩種:廣播模式、集群模式(負載均衡模式)

廣播模式是每個消費者,都會消費消息;

負載均衡模式是每一個消費只會被某一個消費者消費一次;

我們業務上一般用的是負載均衡模式,當然一些特殊場景需要用到廣播模式,比如發送一個信息到郵箱,手機,站內提示;

我們可以通過@RocketMQMessageListenermessageModel屬性值來設置,MessageModel.BROADCASTING是廣播模式,MessageModel.CLUSTERING是默認集群負載均衡模式

下面來介紹下 springboot+rockermq 整合實現 廣播消息

  • 創建Springboot項目,添加rockermq 依賴

<!--rocketMq依賴-->
<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <version>2.2.1</version>
</dependency>
  • 配置rocketmq

# 端口
server:
  port: 8083

# 配置 rocketmq
rocketmq:
  name-server: 127.0.0.1:9876
  #生產者
  producer:
    #生產者組名,規定在一個應用里面必須唯一
    group: group1
    #消息發送的超時時間 默認3000ms
    send-message-timeout: 3000
    #消息達到4096字節的時候,消息就會被壓縮。默認 4096
    compress-message-body-threshold: 4096
    #最大的消息限制,默認為128K
    max-message-size: 4194304
    #同步消息發送失敗重試次數
    retry-times-when-send-failed: 3
    #在內部發送失敗時是否重試其他代理,這個參數在有多個broker時才生效
    retry-next-server: true
    #異步消息發送失敗重試的次數
    retry-times-when-send-async-failed: 3

  • 生產端:新建一個 controller 來做消息發送

生產端按正常發送邏輯發送消息即可

package com.example.springbootrocketdemo.controller;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
/**
 * 廣播消息
 * @author qzz
 */
@RestController
public class RocketMQBroadCOntroller {
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    /**
     * 發送廣播消息
     */
    @RequestMapping("/testBroadSend")
    public void testSyncSend(){
        //參數一:topic   如果想添加tag,可以使用"topic:tag"的寫法
        //參數二:消息內容
        for(int i=0;i<10;i++){
            rocketMQTemplate.convertAndSend("test-topic-broad","test-message"+i);
        }
    }
}
  • 創建兩個消費者來消費消息

我們先集群負載均衡測試,加上messageModel=MessageModel.CLUSTERING

消費者1:

package com.example.springbootrocketdemo.config;
import org.apache.rocketmq.spring.annotation.MessageModel;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Service;
/**
 * 廣播消息
 * 配置RocketMQ監聽
 * MessageModel.CLUSTERING:集群模式
 * MessageModel.BROADCASTING:廣播模式
 * @author qzz
 */
@Service
@RocketMQMessageListener(consumerGroup = "test-broad",topic = "test-topic-broad",messageModel = MessageModel.CLUSTERING)
public class RocketMQBroadConsumerListener implements RocketMQListener<String> {
    @Override
    public void onMessage(String s) {
        System.out.println("集群模式 消費者1,消費消息:"+s);
    }
}

消費者2: 與消費者1在 同一個consumerGroup 和 topic

package com.example.springbootrocketdemo.config;
import org.apache.rocketmq.spring.annotation.MessageModel;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Service;
/**
 * 廣播消息
 * 配置RocketMQ監聽
 * MessageModel.CLUSTERING:集群模式
 * MessageModel.BROADCASTING:廣播模式
 * @author qzz
 */
@Service
@RocketMQMessageListener(consumerGroup = "test-broad",topic = "test-topic-broad",messageModel = MessageModel.CLUSTERING)
public class RocketMQBroadConsumerListener2 implements RocketMQListener<String> {
    @Override
    public void onMessage(String s) {
        System.out.println("集群模式 消費者2,消費消息:"+s);
    }
}
  • 啟動服務,測試 集群模式消費

集群模式測試: 兩個消費者平攤 消息

Springboot中RocketMQ怎么實現廣播消息

  • 把上面兩個消費者的 messageModel 屬性值修改成 廣播模式

消費者1:

package com.example.springbootrocketdemo.config;
import org.apache.rocketmq.spring.annotation.MessageModel;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Service;
/**
 * 廣播消息
 * 配置RocketMQ監聽
 * MessageModel.CLUSTERING:集群模式
 * MessageModel.BROADCASTING:廣播模式
 * @author qzz
 */
@Service
@RocketMQMessageListener(consumerGroup = "test-broad",topic = "test-topic-broad",messageModel = MessageModel.BROADCASTING)
public class RocketMQBroadConsumerListener implements RocketMQListener<String> {
    @Override
    public void onMessage(String s) {
        System.out.println("廣播消息1 廣播模式,消費消息:"+s);
    }
}

消費者2: 與消費者1在 同一個consumerGroup 和 topic

package com.example.springbootrocketdemo.config;
import org.apache.rocketmq.spring.annotation.MessageModel;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Service;
/**
 * 廣播消息
 * 配置RocketMQ監聽
 * MessageModel.CLUSTERING:集群模式
 * MessageModel.BROADCASTING:廣播模式
 * @author qzz
 */
@Service
@RocketMQMessageListener(consumerGroup = "test-broad",topic = "test-topic-broad",messageModel = MessageModel.BROADCASTING)
public class RocketMQBroadConsumerListener2 implements RocketMQListener<String> {
    @Override
    public void onMessage(String s) {
        System.out.println("廣播消息2 廣播模式,消費消息:"+s);
    }
}
  • 重啟服務,測試 廣播模式消費

Springboot中RocketMQ怎么實現廣播消息

到此,關于“Springboot中RocketMQ怎么實現廣播消息”的學習就結束了,希望能夠解決大家的疑惑。理論與實踐的搭配能更好的幫助大家學習,快去試試吧!若想繼續學習更多相關知識,請繼續關注億速云網站,小編會繼續努力為大家帶來更多實用的文章!

向AI問一下細節

免責聲明:本站發布的內容(圖片、視頻和文字)以原創、轉載和分享為主,文章觀點不代表本網站立場,如果涉及侵權請聯系站長郵箱:is@yisu.com進行舉報,并提供相關證據,一經查實,將立刻刪除涉嫌侵權內容。

AI

潮州市| 大田县| 昌图县| 弥勒县| 攀枝花市| 武夷山市| 泊头市| 綦江县| 承德县| 阳春市| 满城县| 延津县| 正安县| 玛纳斯县| 望江县| 大方县| 南宫市| 沂源县| 高要市| 岳阳县| 施甸县| 织金县| 普格县| 巩义市| 长海县| 伊通| 中方县| 涟水县| 桦甸市| 寿阳县| 仲巴县| 永春县| 苍溪县| 通河县| 中西区| 平阴县| 丹凤县| 青田县| 日喀则市| 山阴县| 鄢陵县|