在SpringBoot中利用RocketMQ實(shí)現(xiàn)批量消息消費(fèi)功能
準(zhǔn)備工作
在開始之前,請(qǐng)確保你已經(jīng)安裝和配置好 RocketMQ。如果還沒安裝,請(qǐng)參考 RocketMQ 官網(wǎng) 獲取安裝指南。
項(xiàng)目依賴
首先,我們需要在 Spring Boot 項(xiàng)目中添加 RocketMQ 的依賴。打開 pom.xml 文件,添加以下內(nèi)容:
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.1.0</version>
</dependency>
這個(gè)依賴包包含了與 RocketMQ 集成所需的所有內(nèi)容。
配置 RocketMQ
在 application.yml 文件中添加 RocketMQ 的相關(guān)配置:
rocketmq:
name-server: 127.0.0.1:9876
consumer:
group: batchConsumerGroup
producer:
group: batchProducerGroup
name-server:RocketMQ 服務(wù)的地址consumer.group:消息消費(fèi)的分組producer.group:消息生產(chǎn)的分組
確保 name-server 地址是正確的,指向你的 RocketMQ 服務(wù)。
生產(chǎn)批量消息
創(chuàng)建一個(gè)消息生產(chǎn)者,用于發(fā)送批量消息。以下是 BatchProducer.java 的示例代碼:
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.List;
@Service
public class BatchProducer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
public void sendBatchMessages() {
List<Message<String>> messages = new ArrayList<>();
for (int i = 0; i < 10; i++) {
Message<String> message = MessageBuilder.withPayload("Hello RocketMQ " + i).build();
messages.add(message);
}
rocketMQTemplate.syncSend("BatchTopic", messages, 10000);
System.out.println("批量消息發(fā)送成功!");
}
}
- 這里,我們創(chuàng)建了 10 條消息并將它們添加到列表
messages中。 - 調(diào)用
rocketMQTemplate.syncSend方法將消息批量發(fā)送到主題BatchTopic。
消費(fèi)批量消息
接下來,我們創(chuàng)建一個(gè)消息消費(fèi)者,用于批量消費(fèi)消息。以下是 BatchConsumer.java 的示例代碼:
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Service;
import java.util.List;
@Service
@RocketMQMessageListener(topic = "BatchTopic", consumerGroup = "batchConsumerGroup", selectorExpression = "*", consumeMessageBatchMaxSize = 10)
public class BatchConsumer implements RocketMQListener<List<String>> {
@Override
public void onMessage(List<String> messages) {
System.out.println("批量接收到消息:");
messages.forEach(message -> System.out.println("消息內(nèi)容:" + message));
}
}
在這段代碼中:
@RocketMQMessageListener注解用于標(biāo)識(shí)這是一個(gè) RocketMQ 的消息監(jiān)聽器,指定了監(jiān)聽的主題BatchTopic和消費(fèi)分組batchConsumerGroup。consumeMessageBatchMaxSize = 10表示每次批量消費(fèi)最多 10 條消息。onMessage方法會(huì)處理接收到的消息列表,并逐條打印出消息內(nèi)容。
測(cè)試批量消息發(fā)送和消費(fèi)
創(chuàng)建一個(gè)簡(jiǎn)單的 Spring Boot 控制器,用于觸發(fā)批量消息發(fā)送。以下是 MessageController.java 的代碼:
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class MessageController {
@Autowired
private BatchProducer batchProducer;
@GetMapping("/sendBatchMessages")
public String sendBatchMessages() {
batchProducer.sendBatchMessages();
return "批量消息已發(fā)送";
}
}
通過訪問 http://localhost:8080/sendBatchMessages 觸發(fā)消息發(fā)送。
- 調(diào)用這個(gè)接口會(huì)將批量消息發(fā)送到 RocketMQ 主題
BatchTopic。 BatchConsumer會(huì)自動(dòng)接收并批量處理這些消息。
總結(jié)
我們成功在 Spring Boot 中實(shí)現(xiàn)了 RocketMQ 的批量消息發(fā)送與消費(fèi):
- 使用 BatchProducer 類批量發(fā)送消息。
- 使用 BatchConsumer 類批量消費(fèi)消息,并設(shè)置最大批量大小。
- 通過簡(jiǎn)單的 REST API 控制消息發(fā)送,確保一切順利。
批量消息處理可以提高消息傳遞的效率,適合高并發(fā)場(chǎng)景。這種方式可以減少網(wǎng)絡(luò)開銷,并有效利用系統(tǒng)資源。
到此這篇關(guān)于在SpringBoot中利用RocketMQ實(shí)現(xiàn)批量消息消費(fèi)功能的文章就介紹到這了,更多相關(guān)SpringBoot RocketMQ消息消費(fèi)內(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的使用示例
- SpringBoot項(xiàng)目嵌入RocketMQ的實(shí)現(xiàn)示例
- RocketMQ在Spring Boot上的基礎(chǔ)使用
相關(guān)文章
Java中CountDownLatch工具類詳細(xì)解析
這篇文章主要介紹了Java中CountDownLatch工具類詳細(xì)解析,創(chuàng)建CountDownLatch對(duì)象時(shí),會(huì)傳入一個(gè)count數(shù)值,該對(duì)象每次調(diào)用countDown()方法會(huì)使count?--?,就是count每次減1,需要的朋友可以參考下2023-11-11
PowerJob分布式任務(wù)調(diào)度源碼流程解讀
這篇文章主要為大家介紹了PowerJob分布式任務(wù)調(diào)度源碼流程解讀,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2024-02-02
java中為什么要謹(jǐn)慎使用Arrays.asList、ArrayList的subList
這篇文章主要介紹了java中為什么要謹(jǐn)慎使用Arrays.asList、ArrayList的subList,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2021-02-02
Maven中的dependencyManagement 實(shí)例詳解
dependencyManagement的中文意思就是依賴關(guān)系管理,它就是為了能通更好統(tǒng)一管理項(xiàng)目的版本號(hào)和各種jar版本號(hào),可以更加方便升級(jí),解決包沖突問題,這篇文章主要介紹了Maven中的dependencyManagement 實(shí)例詳解,需要的朋友可以參考下2024-02-02
java調(diào)用未知類的指定方法簡(jiǎn)單實(shí)例
這篇文章介紹了java調(diào)用未知類的指定方法簡(jiǎn)單實(shí)例,有需要的朋友可以參考一下2013-09-09
Java嵌套for循環(huán)的幾種常見優(yōu)化方案
這篇文章主要給大家介紹了關(guān)于Java嵌套for循環(huán)的幾種常見優(yōu)化,在Java中優(yōu)化嵌套for循環(huán)可以通過以下幾種方式來提高性能和效率,文中通過代碼介紹的非常詳細(xì),需要的朋友可以參考下2024-07-07
Java?SpringBoot?@Async實(shí)現(xiàn)異步任務(wù)的流程分析
這篇文章主要介紹了Java?SpringBoot?@Async實(shí)現(xiàn)異步任務(wù),主要包括@Async?異步任務(wù)-無返回值,@Async?異步任務(wù)-有返回值,@Async?+?自定義線程池的操作代碼,需要的朋友可以參考下2022-12-12
SpringBoot 如何使用RestTemplate發(fā)送Post請(qǐng)求
這篇文章主要介紹了SpringBoot 如何使用RestTemplate發(fā)送Post請(qǐng)求的操作,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2021-08-08

