SpringBoot整合RocketMQ極速實(shí)戰(zhàn)教程
一、前期準(zhǔn)備
1.1 環(huán)境依賴
- SpringBoot 2.x / 3.x(本文代碼全兼容)
- RocketMQ 服務(wù)端(4.x/5.x/7.x 均可)
- maven/gradle 項(xiàng)目
1.2 核心依賴引入(Maven)
Apache 官方適配 SpringBoot Starter,無(wú)需手動(dòng)適配版本,自動(dòng)兼容。
<!-- SpringBoot 整合 RocketMQ 官方依賴 -->
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.2.3</version>
</dependency>二、全局配置文件(application.yml)
配置服務(wù)地址、生產(chǎn)組、消費(fèi)組,后續(xù)所有功能統(tǒng)一復(fù)用該配置。
# RocketMQ 配置
rocketmq:
# NameServer 集群地址
name-server: 127.0.0.1:9876
# 生產(chǎn)者配置
producer:
# 生產(chǎn)者組名
group: demo-producer-group
# 消息發(fā)送超時(shí)時(shí)間
send-message-timeout: 3000
# 最大重試次數(shù)
retry-times-when-send-failed: 2
# 消費(fèi)者默認(rèn)配置
consumer:
group: demo-consumer-group三、SpringBoot 集成 RocketMQ 核心原理
在進(jìn)行各類消息實(shí)戰(zhàn)之前,先深度講解 SpringBoot 與 RocketMQ 的底層集成原理,搞懂自動(dòng)裝配、核心組件、通信機(jī)制,知其然更知其所以然,解決開發(fā)中組件失效、連接失敗、消息發(fā)送異常等底層問題。
3.1 自動(dòng)裝配原理(核心)
RocketMQ 官方提供的 rocketmq-spring-boot-starter 遵循 SpringBoot 自動(dòng)裝配機(jī)制,無(wú)需手動(dòng)創(chuàng)建生產(chǎn)者、消費(fèi)者實(shí)例,框架自動(dòng)初始化托管Bean。
- 配置加載:項(xiàng)目啟動(dòng)時(shí)自動(dòng)讀取
application.yml中 rocketmq 全局配置,包含NameServer地址、生產(chǎn)/消費(fèi)組、超時(shí)時(shí)間、AK/SK權(quán)限信息等。 - Bean自動(dòng)注冊(cè):starter 內(nèi)置自動(dòng)配置類,自動(dòng)初始化 RocketMQTemplate 模板類,交由Spring容器管理,開發(fā)者可直接@Autowired注入使用。
- 消費(fèi)者動(dòng)態(tài)注冊(cè):被
@RocketMQMessageListener注解的消費(fèi)者類,項(xiàng)目啟動(dòng)時(shí)會(huì)被Spring掃描,自動(dòng)根據(jù)注解參數(shù)創(chuàng)建消費(fèi)者實(shí)例、訂閱對(duì)應(yīng)Topic、綁定消費(fèi)組與消費(fèi)模式。 - 資源自動(dòng)銷毀:項(xiàng)目關(guān)閉時(shí),Spring容器自動(dòng)銷毀生產(chǎn)、消費(fèi)實(shí)例,優(yōu)雅關(guān)閉連接,避免消息堆積、連接殘留問題。
3.2 核心集成組件說明
- RocketMQTemplate:SpringBoot整合的核心操作模板,封裝了原生API的所有發(fā)送能力,包含同步、異步、單向、有序、批量、事務(wù)消息發(fā)送,是生產(chǎn)者唯一操作入口,屏蔽原生底層復(fù)雜API。
- @RocketMQMessageListener:消費(fèi)者核心注解,承載所有消費(fèi)配置,可配置Topic、消費(fèi)組、集群/廣播模式、消息過濾表達(dá)式、并發(fā)消費(fèi)數(shù)等參數(shù),實(shí)現(xiàn)零配置快速訂閱。
- RocketMQLocalTransactionListener:事務(wù)消息專屬監(jiān)聽接口,框架自動(dòng)注冊(cè)事務(wù)監(jiān)聽器,實(shí)現(xiàn)本地事務(wù)執(zhí)行、事務(wù)狀態(tài)回查兩大核心能力。
3.3 完整通信執(zhí)行流程
- 啟動(dòng)初始化:SpringBoot項(xiàng)目啟動(dòng),自動(dòng)裝配機(jī)制加載RocketMQ配置,初始化RocketMQTemplate、消費(fèi)者實(shí)例。
- 路由拉取:生產(chǎn)者、消費(fèi)者自動(dòng)連接NameServer,定時(shí)拉取Topic對(duì)應(yīng)的Broker路由信息,緩存至本地。
- 消息投遞:業(yè)務(wù)調(diào)用RocketMQTemplate方法發(fā)送消息,框架基于負(fù)載均衡選擇最優(yōu)Broker節(jié)點(diǎn)投遞。
- 消息存儲(chǔ):Broker接收消息、持久化落地磁盤,返回投遞結(jié)果給生產(chǎn)者。
- 消息消費(fèi):消費(fèi)者主動(dòng)拉取對(duì)應(yīng)Topic消息,執(zhí)行業(yè)務(wù)邏輯,框架自動(dòng)維護(hù)消費(fèi)位點(diǎn)、失敗重試機(jī)制。
3.4 SpringBoot集成優(yōu)勢(shì)
- 極簡(jiǎn)開發(fā):摒棄原生繁瑣的代碼創(chuàng)建實(shí)例、配置綁定,注解+配置文件即可快速開發(fā)。
- 容器托管:所有MQ組件交由Spring容器統(tǒng)一管理,生命周期可控、優(yōu)雅啟停。
- 能力全覆蓋:封裝所有高級(jí)消息特性,適配普通、延遲、順序、批量、事務(wù)等全場(chǎng)景。
- 高容錯(cuò)性:框架內(nèi)置發(fā)送重試、異常捕獲、位點(diǎn)維護(hù)、自動(dòng)重連機(jī)制,大幅降低開發(fā)容錯(cuò)成本。
四、基礎(chǔ)消息實(shí)戰(zhàn):普通生產(chǎn)與消費(fèi)
三、基礎(chǔ)消息實(shí)戰(zhàn):普通生產(chǎn)與消費(fèi)
3.1 普通消息實(shí)現(xiàn)原理
原理說明:普通消息是RocketMQ最基礎(chǔ)的消息模型,采用生產(chǎn)者主動(dòng)推送、消費(fèi)者主動(dòng)拉取模式。生產(chǎn)者通過NameServer獲取Broker路由信息,基于負(fù)載均衡策略選擇Broker節(jié)點(diǎn)投遞消息,消息落地Broker磁盤持久化;消費(fèi)者定時(shí)拉取消息,業(yè)務(wù)執(zhí)行成功自動(dòng)提交消費(fèi)位點(diǎn),異常觸發(fā)重試,保障At-Least-Once至少一次投遞語(yǔ)義。支持同步、異步、單向三種發(fā)送模式,適配不同吞吐與可靠性需求。
3.2 普通消息生產(chǎn)者
支持三種發(fā)送模式:同步、異步、單向,覆蓋絕大多數(shù)業(yè)務(wù)場(chǎng)景。
import org.apache.rocketmq.client.producer.SendCallback;
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 MqProducer {
@Autowired
private RocketMQTemplate rocketMqTemplate;
// 定義全局Topic
private static final String TOPIC = "demo-normal-topic";
/**
* 同步發(fā)送消息(可靠,適合核心業(yè)務(wù))
*/
public void sendSyncMsg(String msg) {
SendResult sendResult = rocketMqTemplate.syncSend(TOPIC, msg);
System.out.println("同步發(fā)送結(jié)果:" + sendResult.getSendStatus());
}
/**
* 異步發(fā)送消息(高吞吐,適合非核心業(yè)務(wù))
*/
public void sendAsyncMsg(String msg) {
rocketMqTemplate.asyncSend(TOPIC, msg, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
System.out.println("異步發(fā)送成功");
}
@Override
public void onException(Throwable e) {
System.err.println("異步發(fā)送失?。? + e.getMessage());
}
});
}
/**
* 單向發(fā)送(極致高性能,無(wú)需響應(yīng))
*/
public void sendOneWayMsg(String msg) {
rocketMqTemplate.sendOneWay(TOPIC, msg);
}
}3.3 普通消息消費(fèi)者
默認(rèn)集群消費(fèi)模式,自動(dòng)負(fù)載均衡,異常自動(dòng)重試。
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
@Component
@RocketMQMessageListener(
topic = "demo-normal-topic",
consumerGroup = "demo-consumer-group"
)
public class NormalMsgConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
// 執(zhí)行業(yè)務(wù)邏輯
System.out.println("收到普通消息:" + message);
// 異常自動(dòng)進(jìn)入重試隊(duì)列,無(wú)需手動(dòng)處理
// int a = 1 / 0;
}
}五、廣播消息實(shí)戰(zhàn)
實(shí)現(xiàn)原理:廣播消息核心是消費(fèi)組級(jí)別的全量投遞。集群消費(fèi)是組內(nèi)分?jǐn)傁?,而廣播消費(fèi)會(huì)讓Broker將同一條消息推送給當(dāng)前消費(fèi)組內(nèi)所有在線消費(fèi)者實(shí)例,每個(gè)實(shí)例獨(dú)立消費(fèi)、獨(dú)立維護(hù)位點(diǎn)。為避免集群異常刷屏,廣播消費(fèi)關(guān)閉重試機(jī)制,消息消費(fèi)失敗不會(huì)進(jìn)入重試隊(duì)列。適用于全服務(wù)緩存刷新、全局配置更新等全員同步場(chǎng)景。
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;
@Component
@RocketMQMessageListener(
topic = "demo-broadcast-topic",
consumerGroup = "demo-broadcast-group",
consumeMode = ConsumeMode.BROADCASTING // 開啟廣播消費(fèi)
)
public class BroadcastConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
System.out.println("廣播消費(fèi)消息:" + message);
}
}注意:廣播消費(fèi)不支持重試,適合全局緩存刷新、配置更新場(chǎng)景。
六、順序消息實(shí)戰(zhàn)(分區(qū)有序)
實(shí)現(xiàn)原理:RocketMQ順序消息僅支持分區(qū)有序(局部有序),不支持全局有序。核心實(shí)現(xiàn)邏輯:Topic會(huì)被拆分為多個(gè)消息隊(duì)列,生產(chǎn)者通過自定義業(yè)務(wù)Key哈希取模,將同一業(yè)務(wù)維度(同一訂單/同一用戶)的消息固定投遞到同一個(gè)MessageQueue;消費(fèi)者單線程消費(fèi)單個(gè)隊(duì)列,嚴(yán)格保證隊(duì)列內(nèi)消息FIFO先進(jìn)先出,從而實(shí)現(xiàn)業(yè)務(wù)時(shí)序一致。不同隊(duì)列消息并行消費(fèi),兼顧有序性與吞吐量。
5.1 有序消息生產(chǎn)者原理與代碼
public void sendOrderMsg(String bizKey, String msg) {
// bizKey相同 → 固定發(fā)送同一隊(duì)列 → 保證順序
rocketMqTemplate.syncSendOrderly("demo-order-topic", msg, bizKey);
}5.2 有序消息消費(fèi)者原理與代碼
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;
@Component
@RocketMQMessageListener(
topic = "demo-order-topic",
consumerGroup = "demo-order-group",
messageModel = MessageModel.CLUSTERING
)
public class OrderMsgConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
System.out.println("順序消費(fèi):" + message);
}
}七、延遲消息實(shí)戰(zhàn)
實(shí)現(xiàn)原理:延遲消息核心是系統(tǒng)隊(duì)列中轉(zhuǎn)延時(shí)機(jī)制。生產(chǎn)者發(fā)送消息時(shí)指定延遲等級(jí),Broker接收消息后不會(huì)存入目標(biāo)Topic隊(duì)列,而是轉(zhuǎn)入內(nèi)置的SCHEDULE_TOPIC_XXXX延遲隊(duì)列。Broker后臺(tái)定時(shí)線程掃描延遲隊(duì)列,倒計(jì)時(shí)結(jié)束后,將消息重新路由至用戶指定的普通Topic隊(duì)列,消費(fèi)者方可正常拉取消費(fèi)。RocketMQ不支持自定義任意時(shí)間,僅支持官方預(yù)設(shè)18個(gè)固定延遲等級(jí),保證服務(wù)性能穩(wěn)定。
/**
* 發(fā)送延遲消息
* 延遲等級(jí):1~18級(jí) → 1s、5s、10s、30s、1m、2m、3m、4m、5m、6m、7m、8m、9m、10m、20m、30m、1h、2h
*/
public void sendDelayMsg(String msg) {
// 3級(jí)延遲 = 10秒后消費(fèi)
rocketMqTemplate.syncSend("demo-delay-topic", msg, 3);
}八、批量消息實(shí)戰(zhàn)
實(shí)現(xiàn)原理:批量消息核心是合并網(wǎng)絡(luò)IO、減少請(qǐng)求次數(shù)。多條同Topic、同配置的消息封裝為一個(gè)消息列表,通過一次網(wǎng)絡(luò)請(qǐng)求提交至Broker,Broker批量持久化存儲(chǔ)。大幅降低頻繁網(wǎng)絡(luò)連接、請(qǐng)求握手帶來的性能開銷,極大提升高吞吐場(chǎng)景的并發(fā)能力。批量消息為整體事務(wù)機(jī)制,整批消息要么全部成功存儲(chǔ),要么全部失敗,不支持部分成功。
import org.apache.rocketmq.common.message.Message;
import java.util.ArrayList;
import java.util.List;
public void sendBatchMsg() {
List<Message> messageList = new ArrayList<>();
for (int i = 0; i < 10; i++) {
Message message = new Message("demo-batch-topic", ("批量消息" + i).getBytes());
messageList.add(message);
}
// 批量發(fā)送
rocketMqTemplate.syncSend(messageList);
}注意:?jiǎn)闻慰偞笮〔怀^4MB,整體成功/整體失敗。
九、消息過濾實(shí)戰(zhàn)(Tag過濾)
實(shí)現(xiàn)原理:消息過濾采用服務(wù)端預(yù)過濾機(jī)制,避免無(wú)效消息拉取浪費(fèi)網(wǎng)絡(luò)資源。Tag過濾是輕量級(jí)過濾方案,生產(chǎn)者為消息綁定業(yè)務(wù)標(biāo)簽,Broker存儲(chǔ)時(shí)關(guān)聯(lián)Tag標(biāo)識(shí);消費(fèi)者訂閱指定Tag表達(dá)式,Broker在服務(wù)端直接匹配過濾,僅推送符合條件的消息至消費(fèi)者,無(wú)性能損耗。復(fù)雜場(chǎng)景可使用SQL過濾,基于消息自定義屬性做多條件篩選。
8.1 Tag過濾實(shí)現(xiàn)原理與代碼
// 發(fā)送訂單消息,tag = ORDER
public void sendTagMsg() {
rocketMqTemplate.syncSend("demo-tag-topic:ORDER", "訂單創(chuàng)建消息");
}8.2 消費(fèi)者訂閱指定Tag
@RocketMQMessageListener(
topic = "demo-tag-topic",
consumerGroup = "demo-tag-group",
selectorExpression = "ORDER||PAY" // 只消費(fèi)ORDER、PAY標(biāo)簽消息
)
@Component
public class TagFilterConsumer implements RocketMQListener<String>{
@Override
public void onMessage(String s) {
System.out.println("過濾消費(fèi)消息:" + s);
}
}十、事務(wù)消息實(shí)戰(zhàn)(核心重點(diǎn))
實(shí)現(xiàn)原理:事務(wù)消息基于兩階段提交+超時(shí)回查機(jī)制,解決本地?cái)?shù)據(jù)庫(kù)事務(wù)與消息發(fā)送的分布式一致性問題。第一階段發(fā)送半消息預(yù)提交,消息持久化但對(duì)消費(fèi)者不可見;第二階段執(zhí)行本地事務(wù),根據(jù)事務(wù)結(jié)果提交或回滾消息。針對(duì)生產(chǎn)者宕機(jī)、網(wǎng)絡(luò)超時(shí)等異常,Broker提供定時(shí)回查兜底,主動(dòng)校驗(yàn)本地事務(wù)狀態(tài),徹底杜絕消息懸掛、數(shù)據(jù)不一致問題。
實(shí)現(xiàn)本地事務(wù)與消息發(fā)送最終一致性。
9.1 事務(wù)消息完整代碼實(shí)現(xiàn)
import org.apache.rocketmq.spring.annotation.RocketMQTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionState;
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.Component;
@Component
@RocketMQTransactionListener
public class TransactionMsgListener implements RocketMQLocalTransactionListener {
@Autowired
private RocketMQTemplate rocketMqTemplate;
private static final String TRANS_TOPIC = "demo-trans-topic";
// 發(fā)送事務(wù)消息入口
public void sendTransMsg(String content) {
Message<String> message = MessageBuilder.withPayload(content).build();
rocketMqTemplate.sendMessageInTransaction(TRANS_TOPIC, message, content);
}
// 執(zhí)行本地事務(wù)
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message message, Object arg) {
try {
// 模擬本地?cái)?shù)據(jù)庫(kù)事務(wù)
System.out.println("執(zhí)行本地事務(wù):" + arg);
// 成功則提交消息
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
// 異?;貪L消息
return RocketMQLocalTransactionState.ROLLBACK;
}
}
// 事務(wù)回查
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message message) {
// 查詢數(shù)據(jù)庫(kù)事務(wù)狀態(tài),此處模擬成功
return RocketMQLocalTransactionState.COMMIT;
}
}到此這篇關(guān)于SpringBoot整合RocketMQ極速實(shí)戰(zhàn)教程的文章就介紹到這了,更多相關(guān)SpringBoot整合RocketMQ內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- springboot 3.x 整合 RocketMQ 5.x的詳細(xì)過程
- 在SpringBoot中利用RocketMQ實(shí)現(xiàn)批量消息消費(fèi)功能
- 解決SpringBoot2.1.0+RocketMQ版本沖突問題
- SpringBoot整合RocketMQ批量發(fā)送消息的實(shí)現(xiàn)代碼
- SpringBoot配置RocketMQ的詳細(xì)過程
- SpringBoot定時(shí)監(jiān)聽RocketMQ的NameServer問題及解決方案
- SpringBoot集成RocketMQ的使用示例
- SpringBoot集成RocketMQ實(shí)現(xiàn)消息發(fā)送的三種方式
- springboot集成RocketMQ過程及使用示例詳解
相關(guān)文章
MyBatis-Plus高效開發(fā)實(shí)戰(zhàn)指南
MyBatis-Plus(簡(jiǎn)稱?MP)是?個(gè)MyBatis的增強(qiáng)工具,在MyBatis的基礎(chǔ)上只做增強(qiáng)不做改變,為簡(jiǎn)化開發(fā),本文介紹MyBatis-Plus高效開發(fā)全攻略,感興趣的朋友跟隨小編一起看看吧2026-03-03
MyBatis中RowBounds實(shí)現(xiàn)內(nèi)存分頁(yè)
RowBounds是MyBatis提供的一種內(nèi)存分頁(yè)方式,適用于小數(shù)據(jù)量的分頁(yè)場(chǎng)景,本文就來詳細(xì)的介紹一下,具有一定的參考價(jià)值,感興趣的可以了解一下2024-12-12
Java基礎(chǔ)教程之final關(guān)鍵字淺析
這篇文章主要給大家介紹了關(guān)于Java基礎(chǔ)教程之final關(guān)鍵字的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面來一起學(xué)習(xí)學(xué)習(xí)吧2019-06-06
sqlite數(shù)據(jù)庫(kù)的介紹與java操作sqlite的實(shí)例講解
今天小編就為大家分享一篇關(guān)于sqlite數(shù)據(jù)庫(kù)的介紹與java操作sqlite的實(shí)例講解,小編覺得內(nèi)容挺不錯(cuò)的,現(xiàn)在分享給大家,具有很好的參考價(jià)值,需要的朋友一起跟隨小編來看看吧2019-02-02
mybatis定義sql語(yǔ)句標(biāo)簽之delete標(biāo)簽解析
這篇文章主要介紹了mybatis定義sql語(yǔ)句標(biāo)簽之delete標(biāo)簽解析,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2022-03-03

