RocketMQ的四種常用消息隊(duì)列及代碼演示
消息隊(duì)列
普通消息隊(duì)列
普通消息隊(duì)列是最基本的一種消息隊(duì)列,可以按照先進(jìn)先出(FIFO)的順序存儲(chǔ)消息,并且可以被多個(gè)消費(fèi)者同時(shí)消費(fèi)。可以通過(guò)在生產(chǎn)者端指定主題名稱(chēng)和標(biāo)簽來(lái)創(chuàng)建普通消息隊(duì)列。
順序消息隊(duì)列
順序消息隊(duì)列可以保證相同主題和相同消息鍵的消息按照嚴(yán)格的順序被消費(fèi),例如可以用于訂單等需要保證處理順序的場(chǎng)景。可以通過(guò)在創(chuàng)建普通消息隊(duì)列時(shí)指定MessageQueueSelector對(duì)象和鍵來(lái)創(chuàng)建順序消息隊(duì)列。
延遲消息隊(duì)列
延遲消息隊(duì)列是一種可以在指定時(shí)間后被消費(fèi)的消息隊(duì)列??梢栽谏a(chǎn)者端指定消息發(fā)送的時(shí)間戳和延遲時(shí)間,RocketMQ會(huì)根據(jù)這些信息將消息存儲(chǔ)到延遲消息隊(duì)列中,并在指定的時(shí)間后投遞消息到普通消息隊(duì)列中。
事務(wù)消息隊(duì)列
事務(wù)消息隊(duì)列是一種可以保證消息投遞的事務(wù)性消息隊(duì)列。在生產(chǎn)者端發(fā)送事務(wù)消息時(shí),會(huì)先向RocketMQ發(fā)送一條預(yù)提交消息,然后在本地事務(wù)執(zhí)行成功后再提交或回滾事務(wù)。如果提交事務(wù),則RocketMQ會(huì)將消息投遞到消費(fèi)者,否則將不會(huì)投遞該消息??梢酝ㄟ^(guò)在創(chuàng)建事務(wù)消息隊(duì)列時(shí)指定本地事務(wù)執(zhí)行器來(lái)創(chuàng)建事務(wù)消息隊(duì)列。
除此之外,RocketMQ還支持多主題(Topic)、多消息生產(chǎn)者(Producer)和多消費(fèi)者組(Consumer Group)的概念,可以為不同的業(yè)務(wù)場(chǎng)景創(chuàng)建不同的消息隊(duì)列。
代碼演示
普通消息隊(duì)列
@Service
public class MyProducer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
public void sendMessage(String message) {
rocketMQTemplate.convertAndSend("myTopic", message);
}
}順序消息隊(duì)列
@Service
public class MyProducer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
public void sendOrderMessage(String message, int orderId) {
rocketMQTemplate.setMessageQueueSelector(new OrderMessageQueueSelector());
rocketMQTemplate.convertAndSend("myTopic", message, orderId);
}
}
class OrderMessageQueueSelector implements MessageQueueSelector {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message message, Object orderId) {
int index = (int) orderId % mqs.size();
return mqs.get(index);
}
}延遲消息隊(duì)列
@Service
public class MyProducer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
public void sendDelayMessage(String message, long delayTime) {
rocketMQTemplate.syncSend("myTopic", MessageBuilder.withPayload(message)
.build(), 3000, 2, delayTime);
}
}事務(wù)消息隊(duì)列
@Service
public class MyTransactionListener implements RocketMQLocalTransactionListener {
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message message, Object o) {
// 執(zhí)行本地事務(wù)
// 如果本地事務(wù)執(zhí)行成功,則返回RocketMQLocalTransactionState.COMMIT
// 如果本地事務(wù)執(zhí)行失敗,則返回RocketMQLocalTransactionState.ROLLBACK
return RocketMQLocalTransactionState.UNKNOWN;
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message message) {
// 檢查本地事務(wù)狀態(tài)
// 如果本地事務(wù)執(zhí)行成功,則返回RocketMQLocalTransactionState.COMMIT
// 如果本地事務(wù)執(zhí)行失敗,則返回RocketMQLocalTransactionState.ROLLBACK
// 如果本地事務(wù)狀態(tài)未知,則返回RocketMQLocalTransactionState.UNKNOWN
return RocketMQLocalTransactionState.UNKNOWN;
}
}
@Service
public class MyProducer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Autowired
private MyTransactionListener transactionListener;
public void sendTransactionMessage(String message) {
rocketMQTemplate.setTransactionListener(transactionListener);
rocketMQTemplate.sendMessageInTransaction("myTransactionGroup", "myTopic",
MessageBuilder.withpayload(message).build(), null);
}
}到此這篇關(guān)于RocketMQ的四種常用消息隊(duì)列及代碼演示的文章就介紹到這了,更多相關(guān)RocketMQ常用消息隊(duì)列內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
SpringCloud+nacos部署在多ip環(huán)境下統(tǒng)一nacos服務(wù)注冊(cè)ip(親測(cè)有效)
在部署SpringCoud項(xiàng)目的時(shí)候分服務(wù)器部署注冊(cè)同一個(gè)nacos服務(wù),但是在服務(wù)器有多個(gè)ip存在的同時(shí)(內(nèi)外網(wǎng)),就會(huì)出現(xiàn)注冊(cè)服務(wù)ip不同的問(wèn)題,導(dǎo)致一些接口無(wú)法連接訪問(wèn),經(jīng)過(guò)多次排查終于找到問(wèn)題并找到解決方法,需要的朋友可以參考下2023-04-04
Java重寫(xiě)equals及hashcode方法流程解析
這篇文章主要介紹了Java重寫(xiě)equals及hashcode方法流程解析,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2020-04-04
springboot發(fā)送郵件功能的實(shí)現(xiàn)代碼
發(fā)郵件是一個(gè)很常見(jiàn)的功能,在java中實(shí)現(xiàn)需要依靠JavaMailSender這個(gè)接口,今天通過(guò)本文給大家分享springboot發(fā)送郵件功能的實(shí)現(xiàn)代碼,感興趣的朋友跟隨小編一起看看吧2021-07-07
Apache?Commons?BeanUtils:?JavaBean操作方法
這篇文章主要介紹了Apache?Commons?BeanUtils:?JavaBean操作的藝術(shù),有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2023-12-12
idea項(xiàng)目debug模式啟動(dòng),斷點(diǎn)失效,斷點(diǎn)紅點(diǎn)內(nèi)無(wú)對(duì)勾問(wèn)題及解決
這篇文章主要介紹了idea項(xiàng)目debug模式啟動(dòng),斷點(diǎn)失效,斷點(diǎn)紅點(diǎn)內(nèi)無(wú)對(duì)勾問(wèn)題及解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2023-10-10
SpringBoot解決循環(huán)調(diào)用問(wèn)題
作者在將SpringBoot從1.5版本升級(jí)至2.6版本,并遷移至阿里云上運(yùn)行后,遇到了循環(huán)調(diào)用問(wèn)題,在Jetty容器中運(yùn)行沒(méi)問(wèn)題,但在Tomcat容器中就出現(xiàn)了循環(huán)引用問(wèn)題,原因是SpringBoot 2.6不鼓勵(lì)循環(huán)引用,暴露出該問(wèn)題,作者提供了兩種解決思路2024-10-10

