spring rocketmq集成方案
在 Spring 項(xiàng)目中集成 RocketMQ 是非常常見(jiàn)的消息隊(duì)列應(yīng)用場(chǎng)景,我會(huì)以 Spring Boot + RocketMQ 5.x(當(dāng)前主流版本)為例,提供完整、可直接運(yùn)行的集成方案,包括生產(chǎn)者、消費(fèi)者的核心代碼和配置說(shuō)明。
一、前置條件
- 已安裝并啟動(dòng) RocketMQ(NameServer + Broker),默認(rèn)端口:NameServer
9876 - Spring Boot 版本建議:
2.7.x或3.x(兼容 RocketMQ 官方 starter) - 開(kāi)發(fā)環(huán)境:JDK 8+
二、核心依賴引入
在 pom.xml 中添加 RocketMQ 與 Spring Boot 集成的官方 starter:
<!-- Spring Boot 基礎(chǔ)依賴 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
</dependency>
<!-- RocketMQ Spring Boot Starter(官方推薦) -->
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.2.3</version> <!-- 適配 RocketMQ 5.x,兼容 Spring Boot 2/3 -->
</dependency>
<!-- 可選:測(cè)試依賴 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>三、核心配置(application.yml)
在 resources 目錄下配置 RocketMQ 連接信息:
spring:
application:
name: rocketmq-demo # 應(yīng)用名稱
# RocketMQ 核心配置
rocketmq:
name-server: 127.0.0.1:9876 # NameServer 地址(集群用逗號(hào)分隔)
producer:
group: demo-producer-group # 生產(chǎn)者組(必填,標(biāo)識(shí)同一類生產(chǎn)者)
send-message-timeout: 3000 # 發(fā)送超時(shí)時(shí)間,默認(rèn)3000ms
compress-message-body-threshold: 4096 # 消息壓縮閾值,默認(rèn)4096字節(jié)
max-message-size: 4194304 # 最大消息大小,默認(rèn)4MB
retry-times-when-send-failed: 2 # 同步發(fā)送失敗重試次數(shù)
retry-times-when-send-async-failed: 2 # 異步發(fā)送失敗重試次數(shù)四、生產(chǎn)者實(shí)現(xiàn)(3 種發(fā)送方式)
1. 基礎(chǔ)同步發(fā)送(最常用)
適用于需要確認(rèn)發(fā)送結(jié)果的場(chǎng)景(如訂單創(chuàng)建、庫(kù)存扣減):
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
@Component
public class RocketMQProducer {
// 注入官方封裝的 RocketMQ 模板(類似 RabbitTemplate)
@Autowired
private RocketMQTemplate rocketMQTemplate;
/**
* 同步發(fā)送消息(阻塞等待結(jié)果)
* @param topic 消息主題(必填,需提前創(chuàng)建)
* @param message 消息內(nèi)容
* @return 發(fā)送結(jié)果
*/
public SendResult sendSyncMessage(String topic, String message) {
try {
// 發(fā)送格式:"topic:tag"(tag可選,用于消息過(guò)濾)
SendResult sendResult = rocketMQTemplate.syncSend(topic + ":demoTag", message);
System.out.println("同步發(fā)送成功,消息ID:" + sendResult.getMsgId());
return sendResult;
} catch (Exception e) {
System.err.println("同步發(fā)送失?。? + e.getMessage());
// 業(yè)務(wù)異常處理(如重試、記錄日志、告警)
throw new RuntimeException("消息發(fā)送失敗", e);
}
}
/**
* 異步發(fā)送消息(非阻塞,回調(diào)通知結(jié)果)
* @param topic 主題
* @param message 消息內(nèi)容
*/
public void sendAsyncMessage(String topic, String message) {
rocketMQTemplate.asyncSend(
topic + ":demoTag",
message,
// 發(fā)送成功回調(diào)
sendResult -> System.out.println("異步發(fā)送成功,消息ID:" + sendResult.getMsgId()),
// 發(fā)送失敗回調(diào)
throwable -> System.err.println("異步發(fā)送失?。? + throwable.getMessage())
);
}
/**
* 單向發(fā)送消息(無(wú)回調(diào),適用于日志、埋點(diǎn)等不關(guān)心結(jié)果的場(chǎng)景)
* @param topic 主題
* @param message 消息內(nèi)容
*/
public void sendOneWayMessage(String topic, String message) {
rocketMQTemplate.sendOneWay(topic + ":demoTag", message);
System.out.println("單向消息發(fā)送請(qǐng)求已提交");
}
}2. 發(fā)送自定義對(duì)象消息
如果需要發(fā)送 Java 對(duì)象(而非字符串),只需保證對(duì)象可序列化:
// 自定義消息實(shí)體(實(shí)現(xiàn) Serializable)
public class OrderMessage implements Serializable {
private Long orderId;
private String orderNo;
private BigDecimal amount;
// 省略 getter/setter/toString
}
// 生產(chǎn)者中新增方法
public SendResult sendObjectMessage(String topic, OrderMessage orderMessage) {
return rocketMQTemplate.syncSend(topic + ":orderTag", orderMessage);
}五、消費(fèi)者實(shí)現(xiàn)(2 種消費(fèi)模式)
1. 普通消費(fèi)(默認(rèn)集群模式)
適用于多實(shí)例負(fù)載均衡消費(fèi)(同一組消費(fèi)者分?jǐn)傁ⅲ?/p>
import org.apache.rocketmq.spring.annotation.ConsumeMode;
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;
/**
* RocketMQ 消費(fèi)者
* - topic:訂閱的主題(需與生產(chǎn)者一致)
* - consumerGroup:消費(fèi)者組(必填,同一組消費(fèi)同一主題)
* - messageModel:消費(fèi)模式(CLUSTERING 集群模式,BROADCASTING 廣播模式)
* - consumeMode:消費(fèi)方式(CONCURRENTLY 并發(fā)消費(fèi),ORDERLY 順序消費(fèi))
*/
@Component
@RocketMQMessageListener(
topic = "demo_topic", // 訂閱主題
consumerGroup = "demo-consumer-group", // 消費(fèi)者組
messageModel = MessageModel.CLUSTERING, // 集群模式(默認(rèn))
consumeMode = ConsumeMode.CONCURRENTLY // 并發(fā)消費(fèi)(默認(rèn))
)
public class RocketMQConsumer implements RocketMQListener<String> {
/**
* 消息消費(fèi)邏輯(接收到消息時(shí)觸發(fā))
* @param message 消息內(nèi)容(與生產(chǎn)者發(fā)送類型一致)
*/
@Override
public void onMessage(String message) {
try {
// 核心業(yè)務(wù)邏輯:如解析消息、處理訂單、更新庫(kù)存等
System.out.println("接收到消息:" + message);
// 消費(fèi)成功無(wú)需返回,拋出異常則會(huì)觸發(fā)重試
} catch (Exception e) {
System.err.println("消息消費(fèi)失?。? + e.getMessage());
// 異常拋出后,RocketMQ 會(huì)自動(dòng)重試(默認(rèn)最多16次)
throw new RuntimeException("消費(fèi)失敗", e);
}
}
}2. 消費(fèi)自定義對(duì)象消息
如果生產(chǎn)者發(fā)送的是自定義對(duì)象,消費(fèi)者需指定泛型為對(duì)應(yīng)類型:
@Component
@RocketMQMessageListener(
topic = "demo_topic",
consumerGroup = "order-consumer-group"
)
public class OrderMessageConsumer implements RocketMQListener<OrderMessage> {
@Override
public void onMessage(OrderMessage orderMessage) {
System.out.println("接收到訂單消息:" + orderMessage);
// 處理訂單業(yè)務(wù)邏輯
}
}六、測(cè)試驗(yàn)證
編寫(xiě)測(cè)試類,驗(yàn)證生產(chǎn)者發(fā)送、消費(fèi)者接收是否正常:
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
@SpringBootTest
public class RocketMQDemoTest {
@Autowired
private RocketMQProducer rocketMQProducer;
@Test
public void testSendSyncMessage() {
// 發(fā)送消息到 demo_topic 主題
rocketMQProducer.sendSyncMessage("demo_topic", "Hello RocketMQ + Spring Boot!");
// 暫停3秒,確保消費(fèi)者能接收到消息
try {
Thread.sleep(3000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}七、關(guān)鍵注意事項(xiàng)
- 主題 / 組命名規(guī)范:避免特殊字符,建議用
業(yè)務(wù)_模塊_主題格式(如order_pay_topic)。 - 重試機(jī)制:消費(fèi)失敗默認(rèn)重試 16 次,可通過(guò)
maxReconsumeTimes配置重試次數(shù)。 - 消息持久化:RocketMQ 默認(rèn)持久化消息,即使消費(fèi)者宕機(jī),重啟后仍能消費(fèi)未處理的消息。
- 順序消費(fèi):如需保證消息順序,需將
consumeMode設(shè)為ORDERLY,且生產(chǎn)者發(fā)送時(shí)指定同一消息隊(duì)列。 - 異常處理:生產(chǎn)環(huán)境建議對(duì)接告警(如釘釘、短信),避免消費(fèi)失敗無(wú)感知。
總結(jié)
- Spring Boot 集成 RocketMQ 的核心是引入官方
rocketmq-spring-boot-starter,配置 NameServer 地址和生產(chǎn) / 消費(fèi)組。 - 生產(chǎn)者通過(guò)
RocketMQTemplate實(shí)現(xiàn)同步 / 異步 / 單向發(fā)送,支持字符串和自定義對(duì)象消息。 - 消費(fèi)者通過(guò)
@RocketMQMessageListener注解聲明訂閱關(guān)系,實(shí)現(xiàn)RocketMQListener接口處理消息邏輯,默認(rèn)集群模式并發(fā)消費(fèi)。
核心關(guān)鍵點(diǎn):主題與消費(fèi)組必須配置正確,消費(fèi)失敗拋出異常會(huì)觸發(fā)自動(dòng)重試,生產(chǎn)環(huán)境需做好異常監(jiān)控和重試次數(shù)限制。
到此這篇關(guān)于spring rocketmq集成的文章就介紹到這了,更多相關(guān)spring rocketmq集成內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Spring Boot 菜單刪除實(shí)現(xiàn)代碼與事務(wù)管理
本文將詳細(xì)介紹Spring Boot環(huán)境下菜單刪除功能的實(shí)現(xiàn)邏輯,包括關(guān)聯(lián)數(shù)據(jù)處理、事務(wù)控制和異常處理等關(guān)鍵環(huán)節(jié),強(qiáng)調(diào)需處理多級(jí)嵌套、角色關(guān)聯(lián)及數(shù)據(jù)一致性,感興趣的朋友跟隨小編一起看看吧2025-08-08
關(guān)于Hibernate的一些學(xué)習(xí)心得總結(jié)
Hibernate是一個(gè)優(yōu)秀的Java 持久化層解決方案,是當(dāng)今主流的對(duì)象—關(guān)系映射(ORM)工具2013-07-07
詳解Java是如何通過(guò)接口來(lái)創(chuàng)建代理并進(jìn)行http請(qǐng)求
今天給大家?guī)?lái)的知識(shí)是關(guān)于Java的,文章圍繞Java是如何通過(guò)接口來(lái)創(chuàng)建代理并進(jìn)行http請(qǐng)求展開(kāi),文中有非常詳細(xì)的介紹及代碼示例,需要的朋友可以參考下2021-06-06
Spring Cloud實(shí)現(xiàn)5分鐘級(jí)區(qū)域切換的操作方法
Spring Cloud 2023.x通過(guò)智能路由預(yù)熱、多活數(shù)據(jù)同步和自動(dòng)化流量切換,實(shí)現(xiàn)5分鐘內(nèi)完成跨區(qū)域故障轉(zhuǎn)移,本文以某電商平臺(tái)從AWS亞太切換至阿里云華東的實(shí)戰(zhàn)為例,詳解關(guān)鍵技術(shù)路徑,需要的朋友可以參考下2025-04-04
springboot使用注解實(shí)現(xiàn)鑒權(quán)功能
這篇文章主要介紹了springboot使用注解實(shí)現(xiàn)鑒權(quán)功能,本文通過(guò)實(shí)例代碼給大家介紹的非常詳細(xì),感興趣的朋友跟隨小編一起看看吧2024-12-12
springboot+thymeleaf國(guó)際化之LocaleResolver接口的示例
本篇文章主要介紹了springboot+thymeleaf國(guó)際化之LocaleResolver的示例 ,小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧2017-11-11
Java實(shí)現(xiàn)map轉(zhuǎn)換成json的方法詳解
這篇文章主要為大家詳細(xì)介紹了Java語(yǔ)言實(shí)現(xiàn)map轉(zhuǎn)換成json的幾種方法,文中的示例代碼講解詳細(xì),對(duì)我們學(xué)習(xí)Java有一定幫助,需要的可以參考一下2022-05-05
JAVA學(xué)習(xí)筆記:注釋、變量的聲明和定義操作實(shí)例分析
這篇文章主要介紹了JAVA學(xué)習(xí)筆記:注釋、變量的聲明和定義操作,結(jié)合實(shí)例形式分析了Java注釋、變量的聲明和定義相關(guān)原理、實(shí)現(xiàn)方法及操作注意事項(xiàng),需要的朋友可以參考下2020-04-04

