首頁 > 軟體

Springboot詳解RocketMQ實現廣播訊息流程

2022-06-22 14:02:01

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);
    }
}
  • 啟動服務,測試 叢集模式消費

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

  • 把上面兩個消費者的 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);
    }
}
  • 重啟服務,測試 廣播模式消費

廣播模式消費下,兩個消費者都消費到Topic的所有訊息。

測試成功!

到此這篇關於Springboot詳解RocketMQ實現廣播訊息流程的文章就介紹到這了,更多相關Springboot廣播訊息內容請搜尋it145.com以前的文章或繼續瀏覽下面的相關文章希望大家以後多多支援it145.com!


IT145.com E-mail:sddin#qq.com