SpringBoot配置RocketMQ的詳細(xì)過程
- 引入maven
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
</dependency>- 配置yml
# rocketmq 配置項(xiàng),對(duì)應(yīng) RocketMQProperties 配置類
rocketmq:
name-server: 127.0.0.1:9876 # RocketMQ Namesrv
# Producer 配置項(xiàng)
producer:
group: demo-producer-group # 生產(chǎn)者分組
send-message-timeout: 3000 # 發(fā)送消息超時(shí)時(shí)間,單位:毫秒。默認(rèn)為 3000 。
compress-message-body-threshold: 4096 # 消息壓縮閥值,當(dāng)消息體的大小超過該閥值后,進(jìn)行消息壓縮。默認(rèn)為 4 * 1024B
max-message-size: 4194304 # 消息體的最大允許大小。。默認(rèn)為 4 * 1024 * 1024B
retry-times-when-send-failed: 2 # 同步發(fā)送消息時(shí),失敗重試次數(shù)。默認(rèn)為 2 次。
retry-times-when-send-async-failed: 2 # 異步發(fā)送消息時(shí),失敗重試次數(shù)。默認(rèn)為 2 次。
retry-next-server: false # 發(fā)送消息給 Broker 時(shí),如果發(fā)送失敗,是否重試另外一臺(tái) Broker 。默認(rèn)為 false
access-key: # Access Key ,可閱讀 https://github.com/apache/rocketmq/blob/master/docs/cn/acl/user_guide.md 文檔
secret-key: # Secret Key
enable-msg-trace: true # 是否開啟消息軌跡功能。默認(rèn)為 true 開啟??砷喿x https://github.com/apache/rocketmq/blob/master/docs/cn/msg_trace/user_guide.md 文檔
customized-trace-topic: RMQ_SYS_TRACE_TOPIC # 自定義消息軌跡的 Topic 。默認(rèn)為 RMQ_SYS_TRACE_TOPIC 。
# Consumer 配置項(xiàng)
consumer:
listeners: # 配置某個(gè)消費(fèi)分組,是否監(jiān)聽指定 Topic 。結(jié)構(gòu)為 Map<消費(fèi)者分組, <Topic, Boolean>> 。默認(rèn)情況下,不配置表示監(jiān)聽。
test-consumer-group:
topic1: false # 關(guān)閉 test-consumer-group 對(duì) topic1 的監(jiān)聽消費(fèi)- 配置變量
/**
* 消息隊(duì)列相關(guān)常亮配置,包括group、topic、tag
**/
public class MqTopicConstant {
/**
* 示例消息隊(duì)列,topic1個(gè)
*/
public static final String DEMO_TOPIC = "test-top-1";
/**
* 注冊(cè)tag
*/
public static final String DEMO_TAG_REGISTERED = "registered";
/**
* 修改tag
*/
public static final String DEMO_TAG_MODIFY = "modify";
}- 創(chuàng)建Service
import com.alibaba.fastjson.JSON;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.client.producer.SendStatus;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.util.List;
@Component
public class RocketMQService {
private static final Logger log = LoggerFactory.getLogger(RocketMQService.class);
@Resource
private RocketMQTemplate template;
/**
* 發(fā)送普通消息
*
* @param topic topic
* @param message 消息體
*/
public void sendMessage(String topic, Object message) {
this.template.convertAndSend(topic, message);
log.info("普通消息發(fā)送完成:message = {}", message);
}
/**
* 發(fā)送同步消息
*
* @param topic topic
* @param message 消息體
*/
public void syncSendMessage(String topic, Object message) {
SendResult sendResult = this.template.syncSend(topic, message);
log.info("同步發(fā)送消息完成:message = {}, sendResult = {}", message, sendResult);
}
/**
* 發(fā)送同步消息
*
* @param topic topic
* @param message 消息體
*/
public SendResult syncSendMessageR(String topic, Object message) {
SendResult sendResult = this.template.syncSend(topic, message);
log.info("同步發(fā)送消息完成:message = {}, sendResult = {}", message, sendResult);
SendStatus sendStatus = sendResult.getSendStatus();
log.info("狀態(tài)打印 : {}" , sendStatus);
//SEND_OK
return sendResult;
}
/**
* 發(fā)送異步消息
*
* @param topic topic
* @param message 消息體
*/
public void asyncSendMessage(String topic, final Object message) {
this.template.asyncSend(topic, message, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
log.info("異步消息發(fā)送成功, SendStatus = {}", sendResult.getSendStatus());
// log.info("異步消息發(fā)送成功,message = {}, SendStatus = {}", message, sendResult.getSendStatus());
}
@Override
public void onException(Throwable e) {
log.info("異步消息發(fā)送異常,exception = {}", e.getMessage());
}
});
}
/**
* 發(fā)送單向消息
*
* @param topic topic
* @param message 消息體
*/
public void sendOneWayMessage(String topic, Object message) {
this.template.sendOneWay(topic, message);
log.info("單向發(fā)送消息完成:message = {}", message);
}
/**
* 同步發(fā)送批量消息
*
* @param topic topic
* @param messageList 消息集合
* @param timeout 超時(shí)時(shí)間(毫秒)
*/
public void syncSendMessages(String topic, List<Message<?>> messageList, long timeout) {
this.template.syncSend(topic, messageList, timeout);
log.info("同步發(fā)送批量消息完成:message = {}", JSON.toJSONString(messageList));
}
/**
* 發(fā)送攜帶 tag 的消息(過濾消息)
*
* @param topic topic,RocketMQTemplate將 topic 和 tag 合二為一了,底層會(huì)進(jìn)行
* 拆分再組裝。只要在指定 topic 時(shí)跟上 {:tags} 就可以指定tag
* 例如 test-topic:tagA
* @param message 消息體
*/
public void syncSendMessageWithTag(String topic, Object message) {
this.template.syncSend(topic, message);
log.info("發(fā)送帶 tag 的消息完成:message = {}", message);
}
/**
* 同步發(fā)送延時(shí)消息
*
* @param topic topic
* @param message 消息體
* @param timeout 超時(shí)
* @param delayLevel 延時(shí)等級(jí):現(xiàn)在RocketMq并不支持任意時(shí)間的延時(shí),需要設(shè)置幾個(gè)固定的延時(shí)等級(jí),
* 從1s到2h分別對(duì)應(yīng)著等級(jí) 1 到 18,消息消費(fèi)失敗會(huì)進(jìn)入延時(shí)消息隊(duì)列
* "1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h";
*/
public void syncSendDelay(String topic, Object message, long timeout, int delayLevel) {
this.template.syncSend(topic, MessageBuilder.withPayload(message).build(), timeout, delayLevel);
log.info("已同步發(fā)送延時(shí)消息 message = {}", message);
}
/**
* 異步發(fā)送延時(shí)消息
*
* @param topic topic
* @param message 消息對(duì)象
* @param timeout 超時(shí)時(shí)間
* @param delayLevel 延時(shí)等級(jí)
*/
public void asyncSendDelay(String topic, final Object message, long timeout, int delayLevel) {
this.template.asyncSend(topic, MessageBuilder.withPayload(message).build(), new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
log.info("異步發(fā)送延時(shí)消息成功,message = {}", message);
}
@Override
public void onException(Throwable throwable) {
log.error("異步發(fā)送延時(shí)消息發(fā)生異常,exception = {}", throwable.getMessage());
}
}, timeout, delayLevel);
log.info("已異步發(fā)送延時(shí)消息 message = {}", message);
}
/**
* 發(fā)送事務(wù)消息
*
* @param topic topic,RocketMQTemplate將 topic 和 tag 合二為一了,底層會(huì)進(jìn)行
* 拆分再組裝。只要在指定 topic 時(shí)跟上 {:tags} 就可以指定tag
* 例如 test-topic:tagA
* @param message 消息對(duì)象
* @param arg 傳給事務(wù)監(jiān)聽器的參數(shù)(可以作為事務(wù)處理的唯一ID,來驗(yàn)證本地事務(wù))
*/
public void sendMessageInTransaction(String topic, final Object message ,final Object arg) {
TransactionSendResult res = this.template.sendMessageInTransaction(topic, MessageBuilder.withPayload(message).build(), arg);
if (res.getLocalTransactionState().equals(LocalTransactionState.COMMIT_MESSAGE) && res.getSendStatus().equals(SendStatus.SEND_OK)) {
log.info("【生產(chǎn)者】事物消息發(fā)送成功;成功結(jié)果:{}", res);
} else {
log.info("【生產(chǎn)者】事務(wù)發(fā)送失敗:失敗原因:{}", res);
}
}
}- 事務(wù)監(jiān)聽器
import lombok.extern.slf4j.Slf4j;
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;
@Slf4j
@Component
@RocketMQTransactionListener
public class TranscationRocketListener implements RocketMQLocalTransactionListener {
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message message, Object arg) {
// 獲取第三個(gè)參數(shù)(用戶自定義參數(shù))
Integer customArg = (Integer) arg;
log.info("執(zhí)行本地事務(wù),自定義參數(shù): {}", customArg);
String tag = String.valueOf(message.getHeaders().get("rocketmq_TAGS"));
log.info("這里是校驗(yàn)TAG: {}" , tag );
//RocketMQLocalTransactionState.COMMIT
//RocketMQLocalTransactionState.ROLLBACK
//RocketMQLocalTransactionState.UNKNOWN
return RocketMQLocalTransactionState.UNKNOWN;
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message message) {
log.info("檢查本地交易: {}", message);
return RocketMQLocalTransactionState.COMMIT;
}
}狀態(tài)解釋
- RocketMQLocalTransactionState.COMMIT: 含義: 表示本地事務(wù)已成功執(zhí)行,允許提交消息。這意味著消息將對(duì)消費(fèi)者可見,可以被正常消費(fèi)。 使用場景: 當(dāng)您的業(yè)務(wù)邏輯成功執(zhí)行且希望該消息能夠被下游系統(tǒng)處理時(shí),應(yīng)返回此狀態(tài)。 - RocketMQLocalTransactionState.ROLLBACK: 含義: 表示本地事務(wù)執(zhí)行失敗,要求回滾消息。即該消息不會(huì)被發(fā)送給任何消費(fèi)者。 使用場景: 如果您的業(yè)務(wù)邏輯執(zhí)行過程中遇到錯(cuò)誤或異常情況,不希望該消息影響下游系統(tǒng),則應(yīng)回滾事務(wù),返回此狀態(tài)。 - RocketMQLocalTransactionState.UNKNOWN: 含義: 表示當(dāng)前無法確定事務(wù)的狀態(tài),可能是因?yàn)榫W(wǎng)絡(luò)問題或其他原因?qū)е聲簳r(shí)無法判斷。RocketMQ 會(huì)定期調(diào)用 checkLocalTransaction 方法來檢查事務(wù)的狀態(tài)。 使用場景: 當(dāng)您不確定事務(wù)是否成功完成時(shí)(例如,遠(yuǎn)程服務(wù)調(diào)用超時(shí)),可以返回此狀態(tài)。RocketMQ 將通過 checkLocalTransaction 方法嘗試再次確認(rèn)事務(wù)狀態(tài)。
- 測(cè)試發(fā)送消息
@Autowired
private RocketMQService rocketMQService;
## 發(fā)送事務(wù)消息
rocketMQService.sendMessageInTransaction(GpsConstants.DEDUCTION_MESSP_TOPIC, "這是測(cè)試消息");
## 消費(fèi)事務(wù)消息
@Service
@RocketMQMessageListener(topic = GpsConstants.DEDUCTION_MESSP_TOPIC
, consumerGroup = GpsConstants.DEDUCTION_MESSP_GROUP)
public class deductionTopic implements RocketMQListener<String> {
@Override
public void onMessage(String msg) {
System.out.println("msg : " + msg);
//這里就會(huì)打印msg : 這是測(cè)試消息
}
}提示 : 同一個(gè)topic,不同的consumerGroup都會(huì)消費(fèi),根據(jù)自己的業(yè)務(wù)指定不同的 consumerGroup 處理不同的業(yè)務(wù),如果不需要,則一個(gè)topic只能用一次
到此這篇關(guān)于SpringBoot配置RocketMQ的詳細(xì)過程的文章就介紹到這了,更多相關(guān)SpringBoot配置RocketMQ內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Spring?Boot實(shí)現(xiàn)MyBatis動(dòng)態(tài)創(chuàng)建表的操作語句
這篇文章主要介紹了Spring?Boot實(shí)現(xiàn)MyBatis動(dòng)態(tài)創(chuàng)建表,MyBatis提供了動(dòng)態(tài)SQL,我們可以通過動(dòng)態(tài)SQL,傳入表名等信息然組裝成建表和操作語句,本文通過案例講解展示我們的設(shè)計(jì)思路,需要的朋友可以參考下2024-01-01
Java getRealPath("/")與getContextPath()區(qū)別詳細(xì)分析
這篇文章主要介紹了Java getRealPath("/")與getContextPath()區(qū)別詳細(xì)分析,本篇文章通過簡要的案例,講解了該項(xiàng)技術(shù)的了解與使用,以下就是詳細(xì)內(nèi)容,需要的朋友可以參考下2021-08-08
java使用Hutool工具庫實(shí)現(xiàn)圖形驗(yàn)證碼
本文介紹了java使用Hutool工具庫實(shí)現(xiàn)圖形驗(yàn)證碼的方法,線段干擾驗(yàn)證碼、圓圈干擾驗(yàn)證碼和擾亂干擾驗(yàn)證碼,通過Hutool簡化了驗(yàn)證碼的生成與驗(yàn)證過程,提升了開發(fā)效率,文章詳細(xì)描述了線段干擾驗(yàn)證碼的登錄驗(yàn)證步驟,并提供了前端和后端代碼示例2026-05-05
Java調(diào)用Coze?API詳細(xì)實(shí)現(xiàn)步驟
Coze是一款由字節(jié)跳動(dòng)推出的低代碼AI開發(fā)平臺(tái),讓非專業(yè)開發(fā)者也能輕松創(chuàng)建智能應(yīng)用,這篇文章主要介紹了Java調(diào)用Coze?API的相關(guān)資料,文中通過代碼介紹的非常詳細(xì),需要的朋友可以參考下2026-01-01
Java SpringMVC 異常處理SimpleMappingExceptionResolver類詳解
這篇文章主要介紹了SpringMVC 異常處理SimpleMappingExceptionResolver類詳解,本篇文章通過簡要的案例,講解了該項(xiàng)技術(shù)的了解與使用,以下就是詳細(xì)內(nèi)容,需要的朋友可以參考下2021-09-09
SpringAOP中基于注解實(shí)現(xiàn)通用日志打印方法詳解
這篇文章主要介紹了SpringAOP中基于注解實(shí)現(xiàn)通用日志打印方法詳解,在日常開發(fā)中,項(xiàng)目里日志是必不可少的,一般有業(yè)務(wù)日志,數(shù)據(jù)庫日志,異常日志等,主要用于幫助程序猿后期排查一些生產(chǎn)中的bug,需要的朋友可以參考下2023-12-12
Java的MyBatis框架中實(shí)現(xiàn)多表連接查詢和查詢結(jié)果分頁
這篇文章主要介紹了Java的MyBatis框架中實(shí)現(xiàn)多表連接查詢和查詢結(jié)果分頁,借助MyBatis框架中帶有的動(dòng)態(tài)SQL查詢功能可以比普通SQL查詢做到更多,需要的朋友可以參考下2016-04-04

