SpringBoot集成RocketMQ的使用示例
一、RocketMQ基本概念
消息模型(Message Model)
RocketMQ主要由Producer、Broker、Consumer三部分組成,其中Producer負(fù)責(zé)生產(chǎn)消息,Consumer負(fù)責(zé)消費(fèi)消息,Broker負(fù)責(zé)存儲(chǔ)消息。Broker在實(shí)際部署過程中對(duì)應(yīng)一臺(tái)服務(wù)器,每個(gè)Broker可以存儲(chǔ)多個(gè)Topic的消息,每個(gè)Topic的消息也可以分片存儲(chǔ)于不同的Broker。MessageQueue用于存儲(chǔ)消息的物理地址,每個(gè)Topic中的消息地址存儲(chǔ)于多個(gè)MessageQueue中。ConsumerGroup由多個(gè)Consumer實(shí)例構(gòu)成。
1、在springBoot項(xiàng)目中添加Maven依賴
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.0.4</version>
</dependency>
2、添加配置:
application.yml 文件中添加如下配置:
rocketmq:
name-server: 192.168.152.165:9876
producer:
group: my-groupSpringBoot 集成 RocketMQ代碼:
生產(chǎn)者: 消息發(fā)送的三種方式
package com.rocketmq.springbootrocketmq;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.test.context.junit4.SpringRunner;
import java.util.concurrent.TimeUnit;
@RunWith(SpringRunner.class)
@SpringBootTest
public class T {
@Autowired
private RocketMQTemplate rocketMQTemplate;
//同步消息
@Test
public void testRocketMQ() {
Message msg = MessageBuilder.withPayload("boot發(fā)送同步消息").build();
rocketMQTemplate.send("helloTopicBoot", msg);
System.out.println("success send");
}
//異步消息
@Test
public void sendASYCMsg() throws InterruptedException {
Message message = MessageBuilder.withPayload("boot發(fā)送異步消息").build();
rocketMQTemplate.asyncSend("helloTopicBoot", message, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
System.out.println("發(fā)送狀態(tài):"+sendResult.getSendStatus());
}
@Override
public void onException(Throwable throwable) {
System.out.println("消息發(fā)送失敗");
}
});
TimeUnit.SECONDS.sleep(5);
}
//一次性消息
@Test
public void sendOneWayRocketMQ() {
Message msg = MessageBuilder.withPayload("boot發(fā)送一次性消息").build();
rocketMQTemplate.sendOneWay("helloTopicBoot", msg);
}
}
消費(fèi)者:
package com.example.springbooTRocketMQConsumer.listener;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
import java.nio.charset.Charset;
@Component
@RocketMQMessageListener(consumerGroup = "htbConsumerGroup",topic = "helloTopicBoot")
public class HelloTopicListener implements RocketMQListener<MessageExt> {
@Override
public void onMessage(MessageExt messageExt) {
System.out.println("success get:"+new String(messageExt.getBody(), Charset.defaultCharset()));
}
}
消息消費(fèi)的兩種模式
集群模式:默認(rèn)模式
廣播模式:
消費(fèi)者:messageModel = MessageModel.BROADCASTING
package com.example.springbooTRocketMQConsumer.listener;
import org.apache.rocketmq.common.message.MessageExt;
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.Component;
import java.nio.charset.Charset;
@Component
@RocketMQMessageListener(consumerGroup = "htbConsumerGroup",topic = "helloTopicBoot",messageModel = MessageModel.BROADCASTING)
public class HelloTopicListener implements RocketMQListener<MessageExt> {
@Override
public void onMessage(MessageExt messageExt) {
System.out.println("success get:"+new String(messageExt.getBody(), Charset.defaultCharset()));
}
}
順序消息
生產(chǎn)者:
//順序消息
@Test
public void sendOrderlyMsg(){
//設(shè)置隊(duì)列選擇器
rocketMQTemplate.setMessageQueueSelector(new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> list, org.apache.rocketmq.common.message.Message message, Object o) {
String orderIdStr = (String) o;
long orderId = Long.parseLong(orderIdStr);
int index = (int)orderId % list.size();
return list.get(index);
}
});
List<OrderStep> orderSteps = OrderUtil.buildOrders();
for (OrderStep orderStep : orderSteps) {
Message msg = MessageBuilder.withPayload(orderStep.toString()).build();
rocketMQTemplate.sendOneWayOrderly("orderlyTopicBoot",msg,String.valueOf(orderStep.getOrderId()));
}
}消費(fèi)者:
package com.example.springbooTRocketMQConsumer.listener;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.spring.annotation.ConsumeMode;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
import java.nio.charset.Charset;
@Component
@RocketMQMessageListener(consumerGroup = "orderlyConsumerBoot",topic = "orderlyTopicBoot",consumeMode = ConsumeMode.ORDERLY)
public class OrderlyTopicListener implements RocketMQListener<MessageExt> {
@Override
public void onMessage(MessageExt messageExt) {
System.out.println("當(dāng)前線程:" + Thread.currentThread() + "隊(duì)列ID"+messageExt.getQueueId() + ",消息內(nèi)容:" + new String(messageExt.getBody(),Charset.defaultCharset()));
}
}
延遲消息
生產(chǎn)者:
//延遲消息
@Test
public void sendDelayRocketMQ() {
Message msg = MessageBuilder.withPayload("boot發(fā)送延時(shí)消息,發(fā)送時(shí)間:"+new Date()).build();
rocketMQTemplate.syncSend("helloTopicBoot", msg,3000,3);
}
消費(fèi)者:
package com.example.springbooTRocketMQConsumer.listener;
import org.apache.rocketmq.common.message.MessageExt;
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.Component;
import java.nio.charset.Charset;
import java.util.Date;
@Component
@RocketMQMessageListener(consumerGroup = "htbConsumerGroup",topic = "helloTopicBoot")
public class DelayTopicListener implements RocketMQListener<MessageExt> {
@Override
public void onMessage(MessageExt messageExt) {
System.out.println("success get:發(fā)送時(shí)間"+new Date()+new String(messageExt.getBody(), Charset.defaultCharset()));
}
}
消息Tag條件過濾
生成者
//Tag消息
@Test
public void sendTagFilterRocketMQ() {
Message msg1 = MessageBuilder.withPayload("消息A").build();
rocketMQTemplate.sendOneWay("tagFilterBoot:TagA", msg1);
Message msg2 = MessageBuilder.withPayload("消息B").build();
rocketMQTemplate.sendOneWay("tagFilterBoot:TagB", msg2);
Message msg3 = MessageBuilder.withPayload("消息C").build();
rocketMQTemplate.sendOneWay("tagFilterBoot:TagC", msg3);
}消費(fèi)者:
package com.example.springbooTRocketMQConsumer.listener;
import org.apache.rocketmq.common.message.MessageExt;
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.Component;
import java.nio.charset.Charset;
import java.util.Date;
@Component
@RocketMQMessageListener(consumerGroup = "tagFilterGroupBoot",topic = "tagFilterBoot",selectorExpression = "TagA || TagC")
public class TagFilterTopicListener implements RocketMQListener<MessageExt> {
@Override
public void onMessage(MessageExt messageExt) {
System.out.println("success get:發(fā)送時(shí)間"+new Date()+new String(messageExt.getBody(), Charset.defaultCharset()));
}
}
SQL92消息過濾
生產(chǎn)者:
//SQL92消息
@Test
public void sendSQL92FilterRocketMQ() {
Message msg1 = MessageBuilder.withPayload("小紅,年齡22,體重45").setHeader("age","22").setHeader("weight",45).build();
rocketMQTemplate.sendOneWay("SQL92FilterBoot", msg1);
Message msg2 = MessageBuilder.withPayload("小明,年齡25,體重60").setHeader("age","25").setHeader("weight",60).build();
rocketMQTemplate.sendOneWay("SQL92FilterBoot", msg2);
Message msg3 = MessageBuilder.withPayload("小藍(lán),年齡40,體重70").setHeader("age","40").setHeader("weight",70).build();
rocketMQTemplate.sendOneWay("SQL92FilterBoot", msg3);
}
消費(fèi)者:
package com.example.springbooTRocketMQConsumer.listener;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.annotation.SelectorType;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
import java.nio.charset.Charset;
import java.util.Date;
@Component
@RocketMQMessageListener(consumerGroup = "SQL92FilterGroupBoot",topic = "SQL92FilterBoot",selectorType = SelectorType.SQL92,selectorExpression = "age > 23 and weight > 60")
public class SQL92FilterTopicListener implements RocketMQListener<MessageExt> {
@Override
public void onMessage(MessageExt messageExt) {
System.out.println("success get:發(fā)送時(shí)間"+new Date()+new String(messageExt.getBody(), Charset.defaultCharset()));
}
}到此這篇關(guān)于SpringBoot集成RocketMQ的使用示例的文章就介紹到這了,更多相關(guān)SpringBoot集成RocketMQ內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- 淺談Springboot整合RocketMQ使用心得
- springBoot整合RocketMQ及坑的示例代碼
- Springboot RocketMq實(shí)現(xiàn)過程詳解
- Springboot詳解RocketMQ實(shí)現(xiàn)廣播消息流程
- 解決SpringBoot整合RocketMQ遇到的坑
- SpringBoot整合RocketMQ實(shí)現(xiàn)發(fā)送同步消息
- 解決springboot集成rocketmq關(guān)于tag的坑
- 在SpringBoot中利用RocketMQ實(shí)現(xiàn)批量消息消費(fèi)功能
- SpringBoot項(xiàng)目嵌入RocketMQ的實(shí)現(xiàn)示例
- RocketMQ在Spring Boot上的基礎(chǔ)使用
相關(guān)文章
Spring Boot 集成 RocketMQ 全流程指南(從依賴引入到消息收發(fā)
本文將通過 手動(dòng)連接 和 配置連接 兩種方式,詳細(xì)講解如何在 Spring Boot 中集成 RocketMQ,實(shí)現(xiàn)消息的同步與異步發(fā)送,并提供完整示例代碼,感興趣的朋友一起看看吧2025-04-04
MyBatisPlus報(bào)錯(cuò):Failed to process,please exclud
這篇文章主要介紹了MyBatisPlus報(bào)錯(cuò):Failed to process,please exclude the tableName or statementId問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2024-08-08
Mybatis一對(duì)多關(guān)聯(lián)關(guān)系映射實(shí)現(xiàn)過程解析
這篇文章主要介紹了Mybatis一對(duì)多關(guān)聯(lián)關(guān)系映射實(shí)現(xiàn)過程解析,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2020-02-02
SpringCloud GateWay動(dòng)態(tài)路由用法
網(wǎng)關(guān)作為所有項(xiàng)目的入口,不希望重啟,因此動(dòng)態(tài)路由是必須的,動(dòng)態(tài)路由主要通過RouteDefinitionRepository接口實(shí)現(xiàn),其默認(rèn)的實(shí)現(xiàn)是InMemoryRouteDefinitionRepository,即在內(nèi)存中存儲(chǔ)路由配置,可基于這個(gè)map對(duì)象操作,動(dòng)態(tài)路由的實(shí)現(xiàn)方案有兩種2024-10-10
如何查找YUM安裝的JAVA_HOME環(huán)境變量詳解
這篇文章主要給大家介紹了關(guān)于如何查找YUM安裝的JAVA_HOME環(huán)境變量的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧。2017-10-10
SpringBoot實(shí)現(xiàn)熱部署Community的示例代碼
本文主要介紹了SpringBoot實(shí)現(xiàn)熱部署Community的示例代碼,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2023-06-06
SpringBoot實(shí)現(xiàn)異步事件Event詳解
這篇文章主要介紹了SpringBoot實(shí)現(xiàn)異步事件Event詳解,異步事件的模式,通常將一些非主要的業(yè)務(wù)放在監(jiān)聽器中執(zhí)行,因?yàn)楸O(jiān)聽器中存在失敗的風(fēng)險(xiǎn),所以使用的時(shí)候需要注意,需要的朋友可以參考下2023-11-11

