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

