Springboot詳解RocketMQ實(shí)現(xiàn)廣播消息流程
RocketMQ消息模式主要有兩種:廣播模式、集群模式(負(fù)載均衡模式)
廣播模式是每個(gè)消費(fèi)者,都會(huì)消費(fèi)消息;
負(fù)載均衡模式是每一個(gè)消費(fèi)只會(huì)被某一個(gè)消費(fèi)者消費(fèi)一次;
我們業(yè)務(wù)上一般用的是負(fù)載均衡模式,當(dāng)然一些特殊場(chǎng)景需要用到廣播模式,比如發(fā)送一個(gè)信息到郵箱,手機(jī),站內(nèi)提示;
我們可以通過(guò)@RocketMQMessageListener的messageModel屬性值來(lái)設(shè)置,MessageModel.BROADCASTING是廣播模式,MessageModel.CLUSTERING是默認(rèn)集群負(fù)載均衡模式
下面來(lái)介紹下 springboot+rockermq 整合實(shí)現(xiàn) 廣播消息
- 創(chuàng)建Springboot項(xiàng)目,添加rockermq 依賴(lài)
<!--rocketMq依賴(lài)-->
<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
#生產(chǎn)者
producer:
#生產(chǎn)者組名,規(guī)定在一個(gè)應(yīng)用里面必須唯一
group: group1
#消息發(fā)送的超時(shí)時(shí)間 默認(rèn)3000ms
send-message-timeout: 3000
#消息達(dá)到4096字節(jié)的時(shí)候,消息就會(huì)被壓縮。默認(rèn) 4096
compress-message-body-threshold: 4096
#最大的消息限制,默認(rèn)為128K
max-message-size: 4194304
#同步消息發(fā)送失敗重試次數(shù)
retry-times-when-send-failed: 3
#在內(nèi)部發(fā)送失敗時(shí)是否重試其他代理,這個(gè)參數(shù)在有多個(gè)broker時(shí)才生效
retry-next-server: true
#異步消息發(fā)送失敗重試的次數(shù)
retry-times-when-send-async-failed: 3
- 生產(chǎn)端:新建一個(gè) controller 來(lái)做消息發(fā)送
生產(chǎn)端按正常發(fā)送邏輯發(fā)送消息即可
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;
/**
* 發(fā)送廣播消息
*/
@RequestMapping("/testBroadSend")
public void testSyncSend(){
//參數(shù)一:topic 如果想添加tag,可以使用"topic:tag"的寫(xiě)法
//參數(shù)二:消息內(nèi)容
for(int i=0;i<10;i++){
rocketMQTemplate.convertAndSend("test-topic-broad","test-message"+i);
}
}
}- 創(chuàng)建兩個(gè)消費(fèi)者來(lái)消費(fèi)消息
我們先集群負(fù)載均衡測(cè)試,加上messageModel=MessageModel.CLUSTERING
消費(fèi)者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監(jiān)聽(tīng)
* 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("集群模式 消費(fèi)者1,消費(fèi)消息:"+s);
}
}消費(fèi)者2: 與消費(fèi)者1在 同一個(gè)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監(jiān)聽(tīng)
* 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("集群模式 消費(fèi)者2,消費(fèi)消息:"+s);
}
}- 啟動(dòng)服務(wù),測(cè)試 集群模式消費(fèi)
集群模式測(cè)試: 兩個(gè)消費(fèi)者平攤 消息

- 把上面兩個(gè)消費(fèi)者的 messageModel 屬性值修改成 廣播模式
消費(fèi)者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監(jiān)聽(tīng)
* 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 廣播模式,消費(fèi)消息:"+s);
}
}消費(fèi)者2: 與消費(fèi)者1在 同一個(gè)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監(jiān)聽(tīng)
* 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 廣播模式,消費(fèi)消息:"+s);
}
}- 重啟服務(wù),測(cè)試 廣播模式消費(fèi)

廣播模式消費(fèi)下,兩個(gè)消費(fèi)者都消費(fèi)到Topic的所有消息。
測(cè)試成功!
到此這篇關(guān)于Springboot詳解RocketMQ實(shí)現(xiàn)廣播消息流程的文章就介紹到這了,更多相關(guān)Springboot廣播消息內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- Springboot詳解RocketMQ實(shí)現(xiàn)消息發(fā)送與接收流程
- Springboot詳細(xì)講解RocketMQ實(shí)現(xiàn)順序消息的發(fā)送與消費(fèi)流程
- SpringBoot整合RocketMQ實(shí)現(xiàn)消息發(fā)送和接收的詳細(xì)步驟
- 解決springboot集成rocketmq關(guān)于tag的坑
- 解決SpringBoot整合RocketMQ遇到的坑
- springboot整合rocketmq實(shí)現(xiàn)分布式事務(wù)
- SpringBoot定時(shí)監(jiān)聽(tīng)RocketMQ的NameServer問(wèn)題及解決方案
相關(guān)文章
java面向?qū)ο笾畬W(xué)生信息管理系統(tǒng)
這篇文章主要為大家詳細(xì)介紹了java面向?qū)ο笾畬W(xué)生信息管理系統(tǒng),文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2020-03-03
減小Maven項(xiàng)目生成的JAR包體積實(shí)現(xiàn)提升運(yùn)維效率
在Maven構(gòu)建Java項(xiàng)目過(guò)程中,減小JAR包體積可通過(guò)排除不必要的依賴(lài)和使依賴(lài)jar包獨(dú)立于應(yīng)用jar包來(lái)實(shí)現(xiàn),在pom.xml文件中使用<exclusions>標(biāo)簽排除不需要的依賴(lài),有助于顯著降低JAR包大小,此外,將依賴(lài)打包到應(yīng)用外,可減少應(yīng)用包的體積2024-10-10
關(guān)于SpringBoot在有Ajax時(shí)候不跳轉(zhuǎn)的問(wèn)題解決
最近在使用Ajax來(lái)發(fā)送一些數(shù)據(jù)給后臺(tái)一個(gè)Controller,但是遇到些問(wèn)題,所以下面這篇文章主要給大家介紹了關(guān)于SpringBoot在有Ajax時(shí)候不跳轉(zhuǎn)問(wèn)題的解決辦法,需要的朋友可以參考下2022-05-05
Spring?AOP利用切面實(shí)現(xiàn)日志保存的示例詳解
最近領(lǐng)導(dǎo)讓寫(xiě)個(gè)用切面實(shí)現(xiàn)日志保存,經(jīng)過(guò)調(diào)研和親測(cè),以完美解決。在這里分享給大家,給有需要的碼友直接使用,希望對(duì)大家有所幫助2022-11-11
mapstruct的用法之qualifiedByName示例詳解
qualifiedByName的意思就是使用這個(gè)Mapper接口中的指定的默認(rèn)方法去處理這個(gè)屬性的轉(zhuǎn)換,而不是簡(jiǎn)單的get?set,今天通過(guò)本文給大家介紹下mapstruct的用法之qualifiedByName示例詳解,感興趣的朋友一起看看吧2022-04-04

