RocketMQ在Spring Boot上的基礎(chǔ)使用
一、整合前置準(zhǔn)備
1. 環(huán)境依賴
- Spring Boot 版本:
2.x或3.x - RocketMQ 版本:
4.9.x或5.x - JDK 版本:1.8 及以上
2. 引入 Maven 依賴
在 pom.xml 中添加 RocketMQ 官方 Spring Boot Starter 依賴:
<!-- RocketMQ Spring Boot Starter 核心依賴 -->
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.2.3</version> <!-- 與 RocketMQ 服務(wù)端版本匹配 -->
</dependency>
<!-- 可選:測試依賴 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
版本匹配建議:
RocketMQ 4.9.x → starter 2.2.3
RocketMQ 5.x → starter 2.3.0+
二、核心配置
1. 基礎(chǔ)配置(application.yml)
spring:
application:
name: rocketmq-spring-boot-demo
# RocketMQ 核心配置
rocketmq:
# NameServer 地址(集群用分號分隔)
name-server: 127.0.0.1:9876
# 生產(chǎn)者配置
producer:
# 生產(chǎn)者組名(必填,建議按業(yè)務(wù)命名)
group: demo-producer-group
# 發(fā)送超時時間(默認 3000ms)
send-message-timeout: 3000
# 同步發(fā)送失敗重試次數(shù)(默認 2)
retry-times-when-send-failed: 2
# 異步發(fā)送失敗重試次數(shù)(默認 2)
retry-times-when-send-async-failed: 2
# 消息最大長度(默認 4194304 字節(jié) = 4MB)
max-message-size: 4194304
# 壓縮閾值(默認 4096 字節(jié),超過自動壓縮)
compress-message-body-threshold: 4096
# 消費者配置(全局默認,可在消費端注解覆蓋)
consumer:
group: demo-consumer-group
# 消費線程數(shù)(默認 20)
consume-thread-min: 10
consume-thread-max: 20
# 批量消費最大條數(shù)(默認 1)
consume-message-batch-max-size: 1
# 最大重試次數(shù)(默認 -1 表示 16 次)
max-reconsume-times: 3
2. 核心配置項說明
| 配置項 | 作用 | 生產(chǎn)建議 |
|---|---|---|
| name-server | NameServer 地址 | 集群環(huán)境配置多個(用分號分隔),避免單點 |
| producer.group | 生產(chǎn)者組名 | 按業(yè)務(wù)劃分(如 order-producer-group) |
| producer.send-message-timeout | 發(fā)送超時 | 核心業(yè)務(wù)設(shè)為 5000ms,避免超時過短 |
| consumer.group | 消費者組名 | 一個業(yè)務(wù)邏輯對應(yīng)一個組,不可復(fù)用 |
| consumer.max-reconsume-times | 最大重試次數(shù) | 非核心業(yè)務(wù)設(shè)為 3 次,核心業(yè)務(wù)設(shè)為 5 次 |
三、各消息類型代碼示例
1. 普通消息(最基礎(chǔ))
1.1 生產(chǎn)者(同步發(fā)送)
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.apache.rocketmq.spring.support.RocketMQHeaders;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.nio.charset.StandardCharsets;
@Component
public class NormalMessageProducer {
@Resource
private RocketMQTemplate rocketMQTemplate;
/**
* 同步發(fā)送普通消息
* @param topic 消息主題
* @param msgContent 消息內(nèi)容
* @param msgKey 消息唯一鍵(用于追蹤/冪等)
*/
public void sendNormalMessage(String topic, String msgContent, String msgKey) {
// 構(gòu)建消息(支持自定義 Header)
Message<String> message = MessageBuilder
.withPayload(msgContent)
// 設(shè)置消息 Key(必填,用于冪等/追蹤)
.setHeader(RocketMQHeaders.KEYS, msgKey)
// 可選:設(shè)置 Tag
.setHeader(RocketMQHeaders.TAGS, "normal_tag")
.build();
// 同步發(fā)送(topic:tag 格式指定 Tag)
rocketMQTemplate.syncSend(topic + ":normal_tag", message);
System.out.println("普通消息發(fā)送成功:" + msgKey);
}
}
1.2 消費者
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;
/**
* 普通消息消費者
* - consumerGroup:消費者組(必填)
* - topic:訂閱主題(必填)
* - selectorExpression:Tag 過濾(* 表示所有)
* - messageModel:消費模式(CLUSTERING 集群/ BROADCASTING 廣播)
* - consumeMode:消費模式(CONCURRENTLY 并發(fā)/ ORDERLY 順序)
*/
@Component
@RocketMQMessageListener(
consumerGroup = "demo-consumer-group",
topic = "normal_topic",
selectorExpression = "normal_tag",
messageModel = MessageModel.CLUSTERING,
consumeMode = ConsumeMode.CONCURRENTLY
)
public class NormalMessageConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String msg) {
// 消費邏輯(需保證冪等)
System.out.println("收到普通消息:" + msg);
// 業(yè)務(wù)處理...
}
}
2. 順序消息
2.1 生產(chǎn)者
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
@Component
public class OrderedMessageProducer {
@Resource
private RocketMQTemplate rocketMQTemplate;
/**
* 發(fā)送順序消息(按業(yè)務(wù)鍵哈希選擇隊列)
* @param topic 主題
* @param msgContent 內(nèi)容
* @param orderId 業(yè)務(wù)唯一鍵(如訂單ID,保證同ID入同隊列)
*/
public void sendOrderedMessage(String topic, String msgContent, String orderId) {
// 構(gòu)建消息
String message = "訂單" + orderId + ":" + msgContent;
// 發(fā)送順序消息(指定 hashKey 為 orderId)
rocketMQTemplate.syncSendOrderly(
topic + ":order_tag",
MessageBuilder.withPayload(message).build(),
orderId // 關(guān)鍵:hashKey,保證同值入同隊列
);
System.out.println("順序消息發(fā)送成功:" + orderId);
}
}
2.2 消費者(必須設(shè)為順序消費)
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(
consumerGroup = "order-consumer-group",
topic = "order_topic",
selectorExpression = "order_tag",
// 核心:順序消費必須設(shè)為 ORDERLY
consumeMode = ConsumeMode.ORDERLY
)
public class OrderedMessageConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String msg) {
// 單線程消費,保證順序
System.out.println("收到順序消息:" + msg);
// 業(yè)務(wù)處理(如訂單創(chuàng)建→支付→發(fā)貨)
}
}
3. 延遲消息
3.1 生產(chǎn)者
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
@Component
public class DelayMessageProducer {
@Resource
private RocketMQTemplate rocketMQTemplate;
/**
* 發(fā)送延遲消息(固定等級)
* @param topic 主題
* @param msgContent 內(nèi)容
* @param delayLevel 延遲等級(1=1s,5=1m,18=2h)
*/
public void sendDelayMessage(String topic, String msgContent, int delayLevel) {
// 構(gòu)建消息
org.apache.rocketmq.common.message.Message rocketMsg = new org.apache.rocketmq.common.message.Message(
topic,
"delay_tag",
msgContent.getBytes()
);
// 設(shè)置延遲等級
rocketMsg.setDelayTimeLevel(delayLevel);
// 發(fā)送延遲消息
rocketMQTemplate.getProducer().send(rocketMsg);
System.out.println("延遲消息發(fā)送成功(等級" + delayLevel + "):" + msgContent);
}
}
3.2 消費者(與普通消息一致)
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
@Component
@RocketMQMessageListener(
consumerGroup = "delay-consumer-group",
topic = "delay_topic",
selectorExpression = "delay_tag"
)
public class DelayMessageConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String msg) {
System.out.println("收到延遲消息:" + msg);
// 業(yè)務(wù)處理(如訂單超時取消)
}
}
4. 事務(wù)消息
4.1 事務(wù)監(jiān)聽器(核心)
import org.apache.rocketmq.spring.annotation.RocketMQTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionState;
import org.springframework.messaging.Message;
import org.springframework.stereotype.Component;
/**
* 事務(wù)消息監(jiān)聽器
* - txProducerGroup:對應(yīng)生產(chǎn)者組名
*/
@RocketMQTransactionListener(txProducerGroup = "tx-producer-group")
@Component
public class TransactionListener implements RocketMQLocalTransactionListener {
/**
* 執(zhí)行本地事務(wù)
*/
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 獲取消息內(nèi)容
String msgContent = new String((byte[]) msg.getPayload());
String orderId = msg.getHeaders().get("KEYS").toString();
try {
// 執(zhí)行本地事務(wù)(如扣減庫存、創(chuàng)建訂單)
boolean success = executeLocalDBTransaction(orderId);
if (success) {
// 提交消息
return RocketMQLocalTransactionState.COMMIT;
} else {
// 回滾消息
return RocketMQLocalTransactionState.ROLLBACK;
}
} catch (Exception e) {
// 未知狀態(tài),等待回查
return RocketMQLocalTransactionState.UNKNOWN;
}
}
/**
* 事務(wù)回查(Broker 主動調(diào)用)
*/
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
String orderId = msg.getHeaders().get("KEYS").toString();
// 查詢本地事務(wù)狀態(tài)
boolean isSuccess = queryLocalTransactionStatus(orderId);
return isSuccess ? RocketMQLocalTransactionState.COMMIT : RocketMQLocalTransactionState.ROLLBACK;
}
// 模擬本地事務(wù)執(zhí)行
private boolean executeLocalDBTransaction(String orderId) {
// 實際業(yè)務(wù)邏輯:操作數(shù)據(jù)庫/緩存等
return true;
}
// 模擬查詢本地事務(wù)狀態(tài)
private boolean queryLocalTransactionStatus(String orderId) {
// 實際查詢數(shù)據(jù)庫狀態(tài)
return true;
}
}
4.2 事務(wù)消息生產(chǎn)者
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
@Component
public class TransactionMessageProducer {
@Resource
private RocketMQTemplate rocketMQTemplate;
/**
* 發(fā)送事務(wù)消息
* @param topic 主題
* @param msgContent 內(nèi)容
* @param orderId 訂單ID(消息Key)
*/
public void sendTransactionMessage(String topic, String msgContent, String orderId) {
// 構(gòu)建消息
Message<String> message = MessageBuilder
.withPayload(msgContent)
.setHeader("KEYS", orderId)
.setHeader("TAGS", "tx_tag")
.build();
// 發(fā)送事務(wù)消息(最后一個參數(shù)為 arg,會傳給監(jiān)聽器)
rocketMQTemplate.sendMessageInTransaction(
topic + ":tx_tag",
message,
null // 自定義參數(shù),可傳業(yè)務(wù)對象
);
System.out.println("事務(wù)消息發(fā)送成功(半消息):" + orderId);
}
}
4.3 事務(wù)消息消費者(與普通消息一致)
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
@Component
@RocketMQMessageListener(
consumerGroup = "tx-consumer-group",
topic = "tx_topic",
selectorExpression = "tx_tag"
)
public class TransactionMessageConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String msg) {
System.out.println("收到事務(wù)消息:" + msg);
// 業(yè)務(wù)處理(如通知物流、更新積分)
}
}
5. 批量消息
5.1 生產(chǎn)者
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.util.ArrayList;
import java.util.List;
@Component
public class BatchMessageProducer {
@Resource
private RocketMQTemplate rocketMQTemplate;
/**
* 發(fā)送批量消息
* @param topic 主題
*/
public void sendBatchMessage(String topic) {
// 構(gòu)建批量消息列表(必須同Topic、同Tag、無延遲/事務(wù))
List<Message> msgList = new ArrayList<>();
for (int i = 0; i < 10; i++) {
Message msg = new Message(
topic,
"batch_tag",
("批量消息" + i).getBytes()
);
msgList.add(msg);
}
// 發(fā)送批量消息
SendResult result = rocketMQTemplate.getProducer().send(msgList);
System.out.println("批量消息發(fā)送成功:" + result.getMsgId());
}
}
5.2 批量消息消費者(支持批量消費)
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQBatchListener;
import org.springframework.stereotype.Component;
import java.util.List;
@Component
@RocketMQMessageListener(
consumerGroup = "batch-consumer-group",
topic = "batch_topic",
selectorExpression = "batch_tag",
// 批量消費最大條數(shù)(需與配置項一致)
consumeMessageBatchMaxSize = 10
)
public class BatchMessageConsumer implements RocketMQBatchListener<String> {
@Override
public void onMessage(List<String> msgs) {
// 批量處理消息
System.out.println("收到批量消息,共" + msgs.size() + "條:");
msgs.forEach(msg -> System.out.println("- " + msg));
}
}
四、測試用例(完整示例)
import org.junit.jupiter.api.Test;
import org.springframework.boot.test.context.SpringBootTest;
import javax.annotation.Resource;
@SpringBootTest
public class RocketMQTest {
@Resource
private NormalMessageProducer normalMessageProducer;
@Resource
private OrderedMessageProducer orderedMessageProducer;
@Resource
private DelayMessageProducer delayMessageProducer;
@Resource
private TransactionMessageProducer transactionMessageProducer;
@Resource
private BatchMessageProducer batchMessageProducer;
// 測試普通消息
@Test
public void testNormalMessage() {
normalMessageProducer.sendNormalMessage("normal_topic", "Hello RocketMQ", "MSG_001");
}
// 測試順序消息
@Test
public void testOrderedMessage() {
String orderId = "ORDER_1001";
orderedMessageProducer.sendOrderedMessage("order_topic", "創(chuàng)建", orderId);
orderedMessageProducer.sendOrderedMessage("order_topic", "支付", orderId);
orderedMessageProducer.sendOrderedMessage("order_topic", "發(fā)貨", orderId);
}
// 測試延遲消息(等級5=1分鐘)
@Test
public void testDelayMessage() {
delayMessageProducer.sendDelayMessage("delay_topic", "訂單超時取消", 5);
}
// 測試事務(wù)消息
@Test
public void testTransactionMessage() {
transactionMessageProducer.sendTransactionMessage("tx_topic", "訂單支付成功", "ORDER_10086");
}
// 測試批量消息
@Test
public void testBatchMessage() {
batchMessageProducer.sendBatchMessage("batch_topic");
}
}
五、生產(chǎn)級優(yōu)化建議
1. 冪等性保障
- 消費端必須基于
msgKey做冪等(如 Redis 分布式鎖、數(shù)據(jù)庫唯一鍵); - 避免重復(fù)消費導(dǎo)致業(yè)務(wù)異常。
2. 異常處理
- 生產(chǎn)者:捕獲
MQClientException,實現(xiàn)失敗重試/降級; - 消費者:消費失敗時拋出異常,觸發(fā)重試(或手動返回失敗)。
3. 監(jiān)控告警
- 監(jiān)控消息堆積量(
brokerOffset - consumerOffset); - 監(jiān)控發(fā)送/消費失敗率、延遲時間。
4. 死信隊列處理
- 配置死信隊列(默認
%DLQ%{consumerGroup}); - 定期處理死信消息,避免消息丟失。
總結(jié)
核心關(guān)鍵點
- 配置核心:
name-server和producer/consumer.group是必填項,需按業(yè)務(wù)規(guī)范命名; - 消息類型:
- 普通消息:直接發(fā)送,適配大部分場景;
- 順序消息:需指定
hashKey,消費端設(shè)為ORDERLY; - 延遲消息:僅支持固定等級,需設(shè)置
delayTimeLevel; - 事務(wù)消息:需實現(xiàn)
RocketMQLocalTransactionListener,處理本地事務(wù)和回查; - 批量消息:需保證同Topic/Tag,消費端實現(xiàn)
RocketMQBatchListener;
- 生產(chǎn)規(guī)范:消息 Key 必設(shè)、消費冪等必做、重試次數(shù)合理配置。
到此這篇關(guān)于RocketMQ在Spring Boot上的基礎(chǔ)使用的文章就介紹到這了,更多相關(guān)SpringBoot RocketMQ使用內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- 淺談Springboot整合RocketMQ使用心得
- springBoot整合RocketMQ及坑的示例代碼
- Springboot RocketMq實現(xiàn)過程詳解
- Springboot詳解RocketMQ實現(xiàn)廣播消息流程
- 解決SpringBoot整合RocketMQ遇到的坑
- SpringBoot整合RocketMQ實現(xiàn)發(fā)送同步消息
- 解決springboot集成rocketmq關(guān)于tag的坑
- SpringBoot集成RocketMQ的使用示例
- 在SpringBoot中利用RocketMQ實現(xiàn)批量消息消費功能
- SpringBoot項目嵌入RocketMQ的實現(xiàn)示例
相關(guān)文章
Java中Double除保留后小數(shù)位的幾種方法(小結(jié))
這篇文章主要介紹了Java中Double保留后小數(shù)位的幾種方法,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2019-07-07
Mybatis-Plus集成Sharding-JDBC與Flyway實現(xiàn)多租戶分庫分表實戰(zhàn)
這篇文章主要為大家介紹了Mybatis-Plus集成Sharding-JDBC與Flyway實現(xiàn)多租戶分庫分表實戰(zhàn),有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪2023-11-11
IDEA 2021版新建Maven、TomCat工程的詳細教程
這篇文章主要介紹了IDEA 2021版新建Maven、TomCat工程,本文通過圖文并茂的形式給大家介紹的非常詳細,對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下2021-04-04
mybatis-plus分頁查詢?nèi)N方法小結(jié)
本文主要介紹了mybatis-plus分頁查詢?nèi)N方法,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2023-05-05
SpringBoot+aop實現(xiàn)主從數(shù)據(jù)庫的讀寫分離操作
讀寫分離的作用是為了緩解寫庫,也就是主庫的壓力,但一定要基于數(shù)據(jù)一致性的原則,就是保證主從庫之間的數(shù)據(jù)一定要一致,這篇文章給大家介紹SpringBoot+aop實現(xiàn)主從數(shù)據(jù)庫的讀寫分離操作,感興趣的朋友跟隨小編一起看看吧2024-03-03
java構(gòu)造器 默認構(gòu)造方法及參數(shù)化構(gòu)造方法
構(gòu)造器也叫構(gòu)造方法、構(gòu)造函數(shù),是一種特殊類型的方法,負責(zé)類中成員變量(域)的初始化。構(gòu)造器的用處是在創(chuàng)建對象時執(zhí)行初始化,當(dāng)創(chuàng)建一個對象時,系統(tǒng)會為這個對象的實例進行默認的初始化,下面文章將進入講解,需要的朋友可以參考下2021-10-10

