SpringBoot實(shí)現(xiàn)延遲消息的兩種方案及對(duì)象詳解
在日常業(yè)務(wù)開(kāi)發(fā)中,延遲消息是高頻剛需場(chǎng)景,幾乎所有中大型項(xiàng)目都會(huì)用到。
常見(jiàn)業(yè)務(wù)場(chǎng)景:
- 電商訂單:下單30分鐘未支付,自動(dòng)取消訂單
- 活動(dòng)營(yíng)銷(xiāo):優(yōu)惠券到期自動(dòng)失效、超時(shí)未領(lǐng)取提醒
- 社交場(chǎng)景:消息延時(shí)推送、超時(shí)未讀提醒
- 任務(wù)調(diào)度:延時(shí)重試失敗任務(wù)、定時(shí)回調(diào)業(yè)務(wù)接口
- 售后場(chǎng)景:超時(shí)未退貨自動(dòng)關(guān)閉售后單
很多同學(xué)第一反應(yīng)是用 定時(shí)任務(wù)輪詢(xún)數(shù)據(jù)庫(kù),但這種方式存在性能差、延遲不準(zhǔn)、數(shù)據(jù)庫(kù)壓力大、實(shí)時(shí)性低等致命問(wèn)題,完全不適合生產(chǎn)環(huán)境。
目前企業(yè)級(jí)主流實(shí)現(xiàn)方案只有兩種:
1. 死信隊(duì)列 TTL 延遲消息(RabbitMQ 原生,無(wú)需插件)
2. 延遲交換機(jī)插件延遲消息(生產(chǎn)首選、靈活度最高、本文重點(diǎn))
一、核心原理深度講解
1.1 RabbitMQ 原生短板
RabbitMQ 原生不支持真正的延遲消息,所有消息都是立即投遞、立即消費(fèi),沒(méi)有內(nèi)置定時(shí)延遲投遞機(jī)制。想要實(shí)現(xiàn)延遲效果,只能通過(guò)曲線(xiàn)方案實(shí)現(xiàn)。
1.2 兩種實(shí)現(xiàn)方案原理對(duì)比
方案一:死信隊(duì)列 TTL(原生無(wú)插件)
核心邏輯:消息先進(jìn)入普通隊(duì)列,設(shè)置過(guò)期時(shí)間(TTL),消息過(guò)期后無(wú)人消費(fèi),自動(dòng)變?yōu)樗佬?,被轉(zhuǎn)發(fā)到死信隊(duì)列,消費(fèi)者監(jiān)聽(tīng)死信隊(duì)列實(shí)現(xiàn)延遲消費(fèi)。
致命缺陷:
- 一個(gè)隊(duì)列只能設(shè)置統(tǒng)一過(guò)期時(shí)間,無(wú)法實(shí)現(xiàn)單條消息不同延遲
- 存在消息阻塞問(wèn)題:前面長(zhǎng)延遲消息未過(guò)期,后面短延遲消息會(huì)被卡住,延遲嚴(yán)重不準(zhǔn)
適用場(chǎng)景:全局統(tǒng)一延遲的簡(jiǎn)單業(yè)務(wù)(如所有訂單統(tǒng)一30分鐘超時(shí))
方案二:延遲交換機(jī)插件(x-delayed-message)? 生產(chǎn)首選
通過(guò)安裝官方延遲插件,RabbitMQ 會(huì)新增一種自定義交換機(jī)類(lèi)型:x-delayed-message。
核心原理:
1. 生產(chǎn)者發(fā)送消息時(shí),在消息頭攜帶 x-delay 延遲時(shí)間(毫秒)
2. 消息不會(huì)立即投遞到隊(duì)列,由延遲交換機(jī)內(nèi)部暫存
3. 等待指定延遲時(shí)間結(jié)束后,交換機(jī)自動(dòng)將消息路由到目標(biāo)隊(duì)列
4. 消費(fèi)者監(jiān)聽(tīng)隊(duì)列,完成延遲消費(fèi)
核心優(yōu)勢(shì):
- 單條消息獨(dú)立延遲,靈活度拉滿(mǎn)
- 無(wú)消息阻塞、時(shí)序精準(zhǔn)、延遲誤差極小
- 配置簡(jiǎn)單、代碼簡(jiǎn)潔、維護(hù)成本低
- 高并發(fā)場(chǎng)景性能穩(wěn)定,大廠(chǎng)普遍采用
二、前置環(huán)境準(zhǔn)備
2.1 安裝延遲消息插件
插件版本必須與當(dāng)前 RabbitMQ 版本完全一致,否則啟動(dòng)報(bào)錯(cuò)、功能失效。
插件下載地址:RabbitMQ 官方插件倉(cāng)庫(kù)
安裝步驟:
1. 將下載好的 rabbitmq_delayed_message_exchange-xxx.ez 放入 RabbitMQ plugins 目錄
2. 啟用插件:rabbitmq-plugins enable rabbitmq_delayed_message_exchange
3. 重啟 RabbitMQ 服務(wù):systemctl restart rabbitmq-server
4. 控制臺(tái)查看交換機(jī)類(lèi)型,出現(xiàn) x-delayed-message 即安裝成功
2.2 SpringBoot 項(xiàng)目依賴(lài)
SpringBoot 整合 RabbitMQ 核心依賴(lài),所有版本通用:
<!-- RabbitMQ AMQP 核心依賴(lài) -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>2.3 全局配置文件 application.yml
配置連接信息、消息確認(rèn)機(jī)制、重試機(jī)制,適配生產(chǎn)環(huán)境:
spring:
rabbitmq:
# 基礎(chǔ)連接配置
host:127.0.0.1
port:5672
username:guest
password:guest
virtual-host:/
# 開(kāi)啟生產(chǎn)者確認(rèn)
publisher-confirm-type:correlated
publisher-returns:true
# 消費(fèi)者手動(dòng)ACK
listener:
simple:
acknowledge-mode:manual
retry:
enabled:true
max-attempts: 3三、生產(chǎn)級(jí)完整代碼實(shí)現(xiàn)
3.1 RabbitMQ 延遲交換機(jī)配置類(lèi)
自定義延遲交換機(jī)、隊(duì)列、綁定關(guān)系,全部持久化,重啟不丟失配置:
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.DirectExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.rabbit.connection.ConnectionFactory;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
import java.util.Map;
/**
* 延遲消息隊(duì)列配置類(lèi)
* 基于 rabbitmq-delayed-message-exchange 插件實(shí)現(xiàn)
*/
@Configuration
publicclassDelayRabbitConfig {
// 延遲交換機(jī)名稱(chēng)
publicstaticfinalStringDELAY_EXCHANGE="business_delay_exchange";
// 延遲隊(duì)列名稱(chēng)
publicstaticfinalStringDELAY_QUEUE="business_delay_queue";
// 路由鍵
publicstaticfinalStringDELAY_ROUTING_KEY="business.delay.routing";
/**
* 構(gòu)建延遲交換機(jī)
* x-delayed-message:延遲交換機(jī)類(lèi)型
* x-delayed-type:轉(zhuǎn)發(fā)模式(direct/topic/fanout)
*/
@Bean
public DirectExchange delayExchange() {
Map<String, Object> args = newHashMap<>();
// 核心參數(shù):聲明為延遲交換機(jī)
args.put("x-delayed-type", "direct");
// 參數(shù):名稱(chēng)、持久化、不自動(dòng)刪除、自定義參數(shù)
returnnewDirectExchange(DELAY_EXCHANGE, true, false, args);
}
/**
* 延遲隊(duì)列(持久化)
*/
@Bean
public Queue delayQueue() {
returnnewQueue(DELAY_QUEUE, true);
}
/**
* 隊(duì)列與延遲交換機(jī)綁定
*/
@Bean
public Binding delayBinding(Queue delayQueue, DirectExchange delayExchange) {
return BindingBuilder.bind(delayQueue)
.to(delayExchange)
.with(DELAY_ROUTING_KEY);
}
/**
* 自定義RabbitTemplate,開(kāi)啟消息可靠投遞
*/
@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
RabbitTemplaterabbitTemplate=newRabbitTemplate(connectionFactory);
// 開(kāi)啟 mandatory,消息投遞失敗返回回調(diào)
rabbitTemplate.setMandatory(true);
return rabbitTemplate;
}
}3.2 延遲消息生產(chǎn)者(支持動(dòng)態(tài)自定義延遲時(shí)間)
核心亮點(diǎn):每條消息可單獨(dú)設(shè)置延遲時(shí)間,毫秒級(jí)單位,靈活適配不同業(yè)務(wù):
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
@RestController
publicclassDelayMsgProducer {
@Autowired
private RabbitTemplate rabbitTemplate;
/**
* 發(fā)送延遲消息接口
* @param message 消息內(nèi)容
* @param delayTime 延遲時(shí)間(單位:毫秒)
* @return 結(jié)果
*/
@GetMapping("/send/delay/message")
public String sendDelayMessage(@RequestParam String message,
@RequestParam Long delayTime) {
// 發(fā)送延遲消息
rabbitTemplate.convertAndSend(
DelayRabbitConfig.DELAY_EXCHANGE,
DelayRabbitConfig.DELAY_ROUTING_KEY,
message,
// 核心:設(shè)置單條消息延遲時(shí)間
msg -> {
msg.getMessageProperties().setHeader("x-delay", delayTime);
return msg;
}
);
return"延遲消息發(fā)送成功!預(yù)計(jì) " + delayTime / 1000 + " 秒后執(zhí)行,消息內(nèi)容:" + message;
}
}3.3 生產(chǎn)級(jí)消費(fèi)者(手動(dòng)ACK、異常重試、防消息丟失)
生產(chǎn)環(huán)境禁止自動(dòng)ACK,必須手動(dòng)確認(rèn),保證消息可靠投遞,失敗可重回隊(duì)列重試:
import com.rabbitmq.client.Channel;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
import java.io.IOException;
@Component
publicclassDelayMsgConsumer {
/**
* 監(jiān)聽(tīng)延遲隊(duì)列,手動(dòng)ACK模式
*/
@RabbitListener(queues = DelayRabbitConfig.DELAY_QUEUE)
publicvoidconsumeDelayMessage(String msg, Message message, Channel channel)throws IOException {
// 獲取消息唯一標(biāo)識(shí)
longdeliveryTag= message.getMessageProperties().getDeliveryTag();
try {
// 執(zhí)行業(yè)務(wù)邏輯
System.out.println("【延遲消息消費(fèi)成功】時(shí)間:" + System.currentTimeMillis() + ",消息內(nèi)容:" + msg);
// 手動(dòng)確認(rèn)消費(fèi)成功,刪除消息
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
// 消費(fèi)異常,拒絕消息,重回隊(duì)列重試
System.err.println("【延遲消息消費(fèi)失敗】異常信息:" + e.getMessage());
channel.basicNack(deliveryTag, false, true);
}
}
}四、接口測(cè)試
啟動(dòng)項(xiàng)目,訪(fǎng)問(wèn)接口,測(cè)試延遲效果:
測(cè)試地址:
http://localhost:8080/send/delay/message?message=訂單超時(shí)自動(dòng)取消&delayTime=10000
參數(shù)說(shuō)明:
- message:自定義消息內(nèi)容
- delayTime:延遲毫秒數(shù)(10000 = 10秒)
訪(fǎng)問(wèn)后,控制臺(tái)會(huì)在10秒后打印消費(fèi)日志,延遲效果精準(zhǔn)生效!
五、兩種延遲方案深度對(duì)比
實(shí)現(xiàn)方案 | 優(yōu)點(diǎn) | 缺點(diǎn) | 適用場(chǎng)景 |
死信隊(duì)列TTL | 原生支持、無(wú)需插件、零部署成本 | 單隊(duì)列統(tǒng)一延遲、消息阻塞、延遲不準(zhǔn)、不靈活 | 簡(jiǎn)單統(tǒng)一延遲業(yè)務(wù) |
延遲交換機(jī)插件 | 單消息獨(dú)立延遲、精準(zhǔn)無(wú)阻塞、靈活度高、代碼簡(jiǎn)潔 | 需安裝插件、重啟服務(wù) | 所有生產(chǎn)級(jí)延遲業(yè)務(wù)(推薦) |
六、注意事項(xiàng)
- 插件版本必須嚴(yán)格匹配 RabbitMQ 版本,版本不匹配直接失效
- 延遲時(shí)間單位是毫秒,千萬(wàn)不要傳秒,否則延遲嚴(yán)重偏差
- 超大延遲(超過(guò)3天)不建議使用,RabbitMQ重啟會(huì)丟失未執(zhí)行延遲消息
- 必須開(kāi)啟手動(dòng)ACK,禁止自動(dòng)ACK,防止消息丟失、業(yè)務(wù)未執(zhí)行但消息已刪除
- 延遲交換機(jī)必須配置
x-delayed-type參數(shù),否則無(wú)法生效 - 高并發(fā)場(chǎng)景建議配置消息重試、死信兜底,避免消息堆積
- 延遲消息不適合超高精準(zhǔn)定時(shí)任務(wù),毫秒級(jí)誤差可忽略,秒級(jí)完全精準(zhǔn)
七、總結(jié)
1、定時(shí)輪詢(xún)數(shù)據(jù)庫(kù)是最低效的延遲方案,生產(chǎn)環(huán)境直接淘汰;
2、死信隊(duì)列 TTL 適合簡(jiǎn)單統(tǒng)一延遲場(chǎng)景,局限性非常大;
3、延遲交換機(jī)插件方案靈活、精準(zhǔn)、穩(wěn)定,是目前企業(yè) SpringBoot 項(xiàng)目延遲消息的最優(yōu)解;
4、生產(chǎn)落地必須搭配 消息持久化、生產(chǎn)者確認(rèn)、消費(fèi)者手動(dòng)ACK、異常重試,保證消息可靠性。
以上就是SpringBoot實(shí)現(xiàn)延遲消息的兩種方案及對(duì)象詳解的詳細(xì)內(nèi)容,更多關(guān)于SpringBoot實(shí)現(xiàn)延遲消息的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!
相關(guān)文章
Quarkus的Spring擴(kuò)展快速改造Spring項(xiàng)目
這篇文章主要為大家介紹了Quarkus的Spring項(xiàng)目擴(kuò)展,帶大家快速改造Spring項(xiàng)目示例演繹,有需要的朋友可以借鑒參考下,希望能夠有所幫助2022-02-02
Java高并發(fā)編程之CAS實(shí)現(xiàn)無(wú)鎖隊(duì)列代碼實(shí)例
這篇文章主要介紹了Java高并發(fā)編程之CAS實(shí)現(xiàn)無(wú)鎖隊(duì)列代碼實(shí)例,在多線(xiàn)程操作中,我們通常會(huì)添加鎖來(lái)保證線(xiàn)程的安全,那么這樣勢(shì)必會(huì)影響程序的性能,那么為了解決這一問(wèn)題,于是就有了在無(wú)鎖操作的情況下依然能夠保證線(xiàn)程的安全,需要的朋友可以參考下2023-12-12
CompletableFuture創(chuàng)建及功能使用全面詳解
這篇文章主要為大家介紹了CompletableFuture創(chuàng)建及功能使用全面詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2023-07-07
Java實(shí)戰(zhàn)之仿天貓商城系統(tǒng)的實(shí)現(xiàn)
這篇文章主要介紹了如何利用Java制作一個(gè)基于SSM框架的迷你天貓商城系統(tǒng),文中采用的技術(shù)有JSP、Springboot、SpringMVC、Spring等,需要的可以參考一下2022-03-03
SpringBoot開(kāi)啟server:compression:enabled(Illegal characte
本文主要介紹了SpringBoot開(kāi)啟server:compression:enabled(Illegal character ((CTRL-CHAR, code 31)))的的問(wèn)題解決,具有一定的參考價(jià)值,感興趣的可以了解一下2025-03-03
Springboot+Redis實(shí)現(xiàn)API接口限流的示例代碼
本文主要介紹了Springboot+Redis實(shí)現(xiàn)API接口限流的示例代碼,文中通過(guò)示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2021-07-07

