Redisson延遲隊(duì)列實(shí)現(xiàn)原理
簡(jiǎn)介
上次發(fā)文是在5個(gè)月前,講了一篇redisson分布式鎖的實(shí)現(xiàn)原理,這次講講延遲隊(duì)列的實(shí)現(xiàn)原理。
API的使用
blockingQueue = redissonClient.getBlockingQueue(name); delayQueue = redissonClient.getDelayedQueue(blockingQueue);
我們可以看到是先獲取了一個(gè)阻塞隊(duì)列然后將其包裝為一個(gè)延遲隊(duì)列。
核心實(shí)現(xiàn)
一個(gè)延遲隊(duì)列會(huì)在Redisson內(nèi)部維護(hù)的channel和數(shù)據(jù)類型,外界無感知,它實(shí)際在內(nèi)部維護(hù)了以下4個(gè)數(shù)據(jù)結(jié)構(gòu):
- redisson_delay_queue_timeout:{name},sorted set數(shù)據(jù)類型,存放所有延遲任務(wù),按照延遲任務(wù)的到期時(shí)間戳(提交任務(wù)時(shí)的時(shí)間戳 + 延遲時(shí)間)來排序的,所以列表的最前面的第一個(gè)元素就是整個(gè)延遲隊(duì)列中最早要被執(zhí)行的任務(wù),這個(gè)概念很重要
- redisson_delay_queue:{name},list數(shù)據(jù)類型,核心過程用不上,后文會(huì)討論他的作用
- {name},list數(shù)據(jù)類型,被稱為目標(biāo)隊(duì)列,這個(gè)里面存放的任務(wù)都是已經(jīng)到了延遲時(shí)間的,可以被消費(fèi)者獲取的任務(wù),所以上面demo中的RBlockingQueue的take方法是從這個(gè)目標(biāo)隊(duì)列中獲取到任務(wù)的
- redisson_delay_queue_channel:{name},是一個(gè)channel,用來通知客戶端開啟一個(gè)延遲任務(wù)
初始化過程
首先當(dāng)調(diào)用redissonClient.getDelayedQueue(blockingQueue)時(shí)候,其實(shí)是new了一個(gè)RedissonDelayedQueue,我們看一下他的構(gòu)造方法。
protected RedissonDelayedQueue(QueueTransferService queueTransferService, Codec codec, final CommandAsyncExecutor commandExecutor, String name) {
super(codec, commandExecutor, name);
// ... 初始化相關(guān)隊(duì)列名稱(channelName, queueName, timeoutSetName)...
// 1. 創(chuàng)建QueueTransferTask子類實(shí)例(這里是匿名實(shí)現(xiàn)類,QueueTransferTask是抽象類)
QueueTransferTask task = new (commandExecutor.getConnectionManager()) {
@Override
protected RFuture<Long> pushTaskAsync() {
// 核心Lua腳本,負(fù)責(zé)轉(zhuǎn)移到期任務(wù)
return commandExecutor.evalWriteAsync(getRawName(), LongCodec.INSTANCE, RedisCommands.EVAL_LONG,
"local expiredValues = redis.call('zrangebyscore', KEYS[2], 0, ARGV[1], 'limit', 0, ARGV[2]); " +
"if #expiredValues > 0 then " +
"for i, v in ipairs(expiredValues) do " +
"local randomId, value = struct.unpack('dLc0', v);" +
"redis.call('rpush', KEYS[1], value);" + // 轉(zhuǎn)移到目標(biāo)隊(duì)列
"redis.call('lrem', KEYS[3], 1, v);" + // 從備份隊(duì)列移除
"end; " +
"redis.call('zrem', KEYS[2], unpack(expiredValues));" + // 從延遲隊(duì)列移除
"end; " +
// 返回下一個(gè)任務(wù)的到期時(shí)間
"local v = redis.call('zrange', KEYS[2], 0, 0, 'WITHSCORES'); " +
"if v[1] ~= nil then " +
"return v[2]; " +
"end " +
"return nil;",
Arrays.asList(getRawName(), timeoutSetName, queueName),
System.currentTimeMillis(), 100);
}
@Override
protected RTopic getTopic() {
return RedissonTopic.createRaw(LongCodec.INSTANCE, commandExecutor, channelName);
}
};
// 2. 啟動(dòng)任務(wù)
queueTransferService.schedule(queueName, task);
}
可以看到客戶端會(huì)創(chuàng)建一個(gè)延遲任務(wù)QueueTransferTask。這個(gè)延遲任務(wù)是redisson內(nèi)部維護(hù)的,這個(gè)延遲任務(wù)會(huì)向Redis Server發(fā)送一段lua腳本,Redis執(zhí)行l(wèi)ua腳本中的命令,并且是原子性的。這段lua腳本主要干了兩件事:
- 將到了延遲時(shí)間的任務(wù)從Zsetredisson_delay_queue_timeout:{name}中移除,存到List類型的{name}這個(gè)目標(biāo)隊(duì)列
- 獲取到Zsetredisson_delay_queue_timeout:{name}中目前最早到過期時(shí)間的延遲任務(wù)的到期時(shí)間戳,并返回給客戶端。(這里埋下伏筆,思考為什么需要這個(gè)到期時(shí)間戳)
最后調(diào)用了schedule,我們?cè)賮砜匆幌逻@個(gè)的源碼。
public class QueueTransferService {
private final Map<String, QueueTransferTask> tasks = new ConcurrentHashMap();
//...
public void schedule(String name, QueueTransferTask task) {
this.tasks.compute(name, (k, t) -> {
if (t == null) {
task.start();
return task;
} else {
t.incUsage();
return t;
}
});
}
//....
}
可以看到,執(zhí)行邏輯是先從記錄所有task的Map中獲取是否有同名的,如果有就增加計(jì)數(shù),說明該任務(wù)已經(jīng)被初始化過監(jiān)聽器了直接跳過。如果沒有就調(diào)用task的start,為task設(shè)置監(jiān)聽器。
在前文提到這個(gè)task其實(shí)是QueueTransferTask的子類,這里的start其實(shí)調(diào)用的是父類QueueTransferTask的start方法。
public abstract class QueueTransferTask {
//.....
//開啟這個(gè)延遲任務(wù),初始化2個(gè)監(jiān)聽器
public void start() {
RTopic schedulerTopic = this.getTopic();
this.statusListenerId = schedulerTopic.addListener(new BaseStatusListener() {
//當(dāng)連接時(shí)候調(diào)用一次pushTask,后文會(huì)提及作用
public void onSubscribe(String channel) {
QueueTransferTask.this.pushTask();
}
});
this.messageListenerId = schedulerTopic.addListener(Long.class, new MessageListener<Long>() {
//當(dāng)從監(jiān)聽的channel中,監(jiān)聽到消息時(shí)候會(huì)調(diào)用此方法,這里后文會(huì)提及作用。
public void onMessage(CharSequence channel, Long startTime) {
QueueTransferTask.this.scheduleTask(startTime);
}
});
}
private void pushTask() {
if (this.usage != 0) {
RFuture<Long> startTimeFuture = this.pushTaskAsync();
//這里res就是剛才返回的最早過期時(shí)間戳
startTimeFuture.whenComplete((res, e) -> {
//有異常就重試
if (e != null) {
if (!this.serviceManager.isShuttingDown(e)) {
log.error(e.getMessage(), e);
this.scheduleTask(System.currentTimeMillis() + 5000L);
}
} else {
if (res != null) {
//將時(shí)間戳傳遞給scheduleTask進(jìn)行調(diào)度任務(wù)
this.scheduleTask(res);
}
}
});
}
}
//真正去規(guī)劃做任務(wù)的邏輯
private void scheduleTask(Long startTime) {
if (this.usage != 0) {
if (startTime != null) {
TimeoutTask oldTimeout = (TimeoutTask)this.lastTimeout.get();
if (oldTimeout != null) {
oldTimeout.getTask().cancel();
}
long delay = startTime - System.currentTimeMillis();
if (delay > 10L) {
//使用netty時(shí)間輪,安排任務(wù)
Timeout timeout = this.serviceManager.newTimeout(new TimerTask() {
public void run(Timeout timeout) throws Exception {
//執(zhí)行剛才匿名實(shí)現(xiàn)的task的pushTaskAsync方法
QueueTransferTask.this.pushTask();
TimeoutTask currentTimeout = (TimeoutTask)QueueTransferTask.this.lastTimeout.get();
if (currentTimeout != null && currentTimeout.getTask() == timeout) {
QueueTransferTask.this.lastTimeout.compareAndSet(currentTimeout, (Object)null);
}
}
}, delay, TimeUnit.MILLISECONDS);
this.lastTimeout.compareAndSet(oldTimeout, new TimeoutTask(startTime, timeout));
} else {
this.pushTask();
}
}
}
}
}
start方法主要是設(shè)置了指定主題(主題名:redisson_delay_queue_channel:{name})兩個(gè)發(fā)布訂閱的監(jiān)聽器。
pushTask方法會(huì)調(diào)用pushTaskAsync方法,即執(zhí)行一次前文提到的lua腳本。
當(dāng)指定主題有新訂閱時(shí)調(diào)用 pushTask() 方法
- BaseStatusListener:這是一個(gè)連接狀態(tài)監(jiān)聽器。當(dāng)客戶端與 Redis Server 成功建立訂閱連接時(shí),也就是說項(xiàng)目啟動(dòng)的時(shí)候會(huì)執(zhí)行一次客戶端延遲任務(wù),它的 onSubscribe方法會(huì)被觸發(fā),并立即執(zhí)行一次 pushTask()。
- 這是因?yàn)轫?xiàng)目在重啟時(shí),由于沒有客戶端延遲任務(wù)的執(zhí)行,可能會(huì)出現(xiàn)redisson_delay_queue_timeout:{name}隊(duì)列中有到期但是沒有被放到目標(biāo)隊(duì)列的可能,重啟就執(zhí)行一次就是為了保證到期的數(shù)據(jù)能被及時(shí)放到目標(biāo)隊(duì)列中。
當(dāng)指定主題有新消息時(shí)調(diào)用 scheduleTask(startTime) 方法。它的作用是:
計(jì)算延遲時(shí)間:接著,它用傳入的 startTime(即下一個(gè)任務(wù)的到期時(shí)間戳)減去當(dāng)前時(shí)間戳,得出需要延遲的時(shí)間 delay。
安排新任務(wù):然后,根據(jù)計(jì)算出的延遲時(shí)間做出決策:
- 如果延遲時(shí)間 > 10毫秒:說明任務(wù)還需要等待一段時(shí)間才到期。此時(shí),它會(huì)利用 Netty 的 HashedWheelTimer(時(shí)間輪) 提交一個(gè)一次性的定時(shí)任務(wù)。這個(gè)定時(shí)任務(wù)設(shè)定的延遲時(shí)間就是剛剛計(jì)算出的 delay毫秒。時(shí)間輪到點(diǎn)后,會(huì)觸發(fā)執(zhí)行 pushTask()方法,從而進(jìn)行任務(wù)轉(zhuǎn)移 。
- 如果延遲時(shí)間 <= 10毫秒:說明任務(wù)已經(jīng)到期,或者即將在瞬間到期(10毫秒被認(rèn)為是可立即執(zhí)行的閾值)。此時(shí),它會(huì)立即執(zhí)行 pushTask()方法,而不是等待定時(shí)器,以保證任務(wù)能被盡快處理 。
取消舊任務(wù):首先,它會(huì)檢查是否存在之前已經(jīng)安排但尚未執(zhí)行的定時(shí)任務(wù)(oldTimeout)。如果存在,就將其取消。這是為了確??偸菆?zhí)行最新的、最早到期的任務(wù),避免舊的任務(wù)干擾調(diào)度 (這一點(diǎn)和channel有關(guān),后文會(huì)解釋)
回收伏筆,這也就解釋了為什么要將這個(gè)最早快要到過期時(shí)間的時(shí)間戳返回來。
可以認(rèn)為QueueTransferTask不是一個(gè)死板的鬧鐘到點(diǎn)了即使沒有什么任務(wù)也會(huì)吵你,而是一個(gè)靈活的秘書,到達(dá)重要時(shí)刻就會(huì)提醒你參加重要活動(dòng)。
比如上述所說的將最快過期的時(shí)間戳返回給客戶端,客戶端通過scheduleTask使用這個(gè)時(shí)間戳開啟一個(gè)時(shí)間輪,讓客戶端阻塞到達(dá)這個(gè)時(shí)間戳,一旦到達(dá)這個(gè)時(shí)間戳說明redisson_delay_queue_timeout:{name}中上面說的最早到過期時(shí)間的任務(wù)已經(jīng)到期了,于是客戶端開始執(zhí)行l(wèi)ua腳本操作,及時(shí)將到了延遲時(shí)間的任務(wù)放到目標(biāo)隊(duì)列中。然后再次發(fā)布剩余的延遲任務(wù)中最早到期的任務(wù)到期時(shí)間戳到channe中,如此循環(huán)往復(fù),一直運(yùn)行下去,保證redisson_delay_queue_timeout:{name}中到期的數(shù)據(jù)能及時(shí)放到目標(biāo)隊(duì)列中。
所以,上述說了一大堆的主要的作用就是保證到了延遲時(shí)間的任務(wù)能夠及時(shí)被放到目標(biāo)隊(duì)列。
發(fā)送延遲任務(wù)
生產(chǎn)者在提交任務(wù)的時(shí)候調(diào)用delayQueue.offer時(shí)候翻看源碼最后會(huì)調(diào)用到offerAsync方法。
其實(shí)是將任務(wù)放入到了Zset類型的redisson_delay_queue_timeout:{name}中,分?jǐn)?shù)就是提交任務(wù)的時(shí)間戳+延遲時(shí)間,就是延遲任務(wù)的到期時(shí)間戳。這段lua腳本會(huì)將數(shù)據(jù)原子性的放入Zset的redisson_delay_queue_timeout:{name}并放一份一樣的到List的redisson_delay_queue:{name},然后發(fā)送消息到channel中,帶上延遲的時(shí)間戳。
@Override
public RFuture<Void> offerAsync(V e, long delay, TimeUnit timeUnit) {
if (delay < 0) {
throw new IllegalArgumentException("Delay can't be negative");
}
long delayInMs = timeUnit.toMillis(delay);
long timeout = System.currentTimeMillis() + delayInMs;
long randomId = ThreadLocalRandom.current().nextLong();
return commandExecutor.evalWriteNoRetryAsync(getRawName(), codec, RedisCommands.EVAL_VOID,
//把 `timeout`(過期時(shí)間戳)、`randomId`(隨機(jī)ID)和 `encode(e)`(任務(wù)真實(shí)數(shù)據(jù))打包成了一個(gè)二進(jìn)制的 `value`。
//即使你添加了兩個(gè)一模一樣的任務(wù)(內(nèi)容一樣、時(shí)間一樣),因?yàn)橛?`randomId`,打包后的 `value` 也是不一樣的。這就允許隊(duì)列里存在重復(fù)的任務(wù)。
"local value = struct.pack('dLc0', tonumber(ARGV[2]), string.len(ARGV[3]), ARGV[3]);"
+ "redis.call('zadd', KEYS[2], ARGV[1], value);"
//存入redisson_delay_queue:{name},注意這里并不是存入結(jié)果List!!
//用于保證 `RDelayedQueue` 這個(gè)對(duì)象本身有數(shù)據(jù)可查(比如你調(diào)用 `contains()` 或 `iterator()` 時(shí),就是查這里)。
+ "redis.call('rpush', KEYS[3], value);"
// 取出 ZSet 里排第一(最早要執(zhí)行)的任務(wù)
+ "local v = redis.call('zrange', KEYS[2], 0, 0); "
//判斷:剛才新加的這個(gè)任務(wù)(value),是不是就是那個(gè)排第一的(v[1])?
+ "if v[1] == value then "
//如果是,說明新任務(wù)插隊(duì)成功,成為了最早的任務(wù)!發(fā)送 PUBLISH 通知客戶端:"嘿,最早的時(shí)間變了,快調(diào)整鬧鐘!"
+ "redis.call('publish', KEYS[4], ARGV[1]); "
+ "end;",
Arrays.<Object>asList(getRawName(), timeoutSetName, queueName, channelName),
timeout, randomId, encode(e));
}
為什么要發(fā)送這個(gè)channel?取消任務(wù)的伏筆回收。比如此時(shí)有一個(gè)任務(wù)是1小時(shí)后才需要執(zhí)行,客戶端會(huì)阻塞到1小時(shí)之后。那此時(shí)來了一個(gè)任務(wù)是10秒后執(zhí)行,可是客戶端已經(jīng)在阻塞了咋辦?所以就需要這個(gè)channel。前文提到一旦監(jiān)聽器監(jiān)聽到消息時(shí)候就會(huì)把這個(gè)channel帶的時(shí)間戳發(fā)送給scheduleTask去執(zhí)行。此時(shí)發(fā)現(xiàn)下一次的時(shí)間戳久于當(dāng)前這個(gè)目標(biāo)時(shí)間戳。那么就會(huì)取消掉這個(gè)任務(wù)??偨Y(jié):就是它就是用于處理 “插隊(duì)” 情況。如果沒有這段代碼:消費(fèi)者會(huì)傻傻地睡到原定時(shí)間(比如10:00),導(dǎo)致那個(gè)新插入隊(duì) 09:00 的新任務(wù)被推遲整整一小時(shí)才執(zhí)行。有了這段代碼,消費(fèi)者能實(shí)時(shí)感知到“任務(wù)變了”,動(dòng)態(tài)調(diào)整自己的睡眠時(shí)間,保證延遲任務(wù)的準(zhǔn)時(shí)性。
至此,講清楚了4個(gè)數(shù)據(jù)結(jié)構(gòu)他們各自的本職工作。
總結(jié)
來看看整體的運(yùn)行原理圖

在翻看源碼時(shí)候發(fā)現(xiàn)本文講解的RDelayedQueue 已經(jīng)被廢棄了,官方更推薦RReliableQueue 基于 Redis Stream 或者更復(fù)雜的 Lua 結(jié)構(gòu)來實(shí)現(xiàn),引入了消息隊(duì)列必須具備的 ACK(確認(rèn))機(jī)制。消費(fèi)者拿到消息,但 Redis 不會(huì)立刻刪除它,而是把它標(biāo)記為“處理中。如果你的業(yè)務(wù)涉及訂單超時(shí)取消、支付回調(diào)等關(guān)鍵鏈路,強(qiáng)烈建議遷移到 RReliableQueue,因?yàn)榕f版在極端并發(fā)或宕機(jī)場(chǎng)景下確實(shí)存在丟單風(fēng)險(xiǎn)。后續(xù)會(huì)找時(shí)間出一篇RReliableQueue的實(shí)現(xiàn)原理
到此這篇關(guān)于Redisson延遲隊(duì)列實(shí)現(xiàn)原理的文章就介紹到這了,更多相關(guān)Redisson延遲隊(duì)列內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- Redisson延遲隊(duì)列實(shí)現(xiàn)訂單關(guān)閉的操作方法
- Spring?Boot集成Redisson實(shí)現(xiàn)延遲隊(duì)列
- Spring Boot 項(xiàng)目集成 Redisson 實(shí)現(xiàn)延遲隊(duì)列的詳細(xì)過程
- SpringBoot中Redisson延遲隊(duì)列的示例
- Redisson延遲隊(duì)列執(zhí)行流程源碼解析
- 分布式利器redis及redisson的延遲隊(duì)列實(shí)踐
- SpringBoot集成Redisson實(shí)現(xiàn)延遲隊(duì)列的場(chǎng)景分析
相關(guān)文章
虛擬線程在Spring?Boot中的正確使用方式及最佳實(shí)踐
虛擬線程是Java19引入的一項(xiàng)新特性,它屬于Project Loom項(xiàng)目的一部分,與傳統(tǒng)的線程不同,虛擬線程并不是由操作系統(tǒng)直接管理,而是由Java虛擬機(jī)控制,這篇文章主要介紹了虛擬線程在Spring?Boot中的正確使用方式及最佳實(shí)踐的相關(guān)資料,需要的朋友可以參考下2026-03-03
idea配置檢查XML中SQL語法及書寫sql語句智能提示的方法
idea連接了數(shù)據(jù)庫,也可以執(zhí)行SQL查到數(shù)據(jù),但是無法識(shí)別sql語句中的表導(dǎo)致沒有提示,下面這篇文章主要給大家介紹了關(guān)于idea配置檢查XML中SQL語法及書寫sql語句智能提示的相關(guān)資料,需要的朋友可以參考下2023-03-03
Maven?項(xiàng)目用Assembly打包可執(zhí)行jar包的方法
這篇文章主要介紹了Maven?項(xiàng)目用Assembly打包可執(zhí)行jar包的方法,該方法只可打包非spring項(xiàng)目的可執(zhí)行jar包,需要的朋友可以參考下2023-03-03
Java調(diào)用Vue前端頁面生成PDF文件的完整代碼
在Vue應(yīng)用中,將頁面導(dǎo)出為PDF文件通常涉及到前端技術(shù)的組合,這篇文章主要介紹了Java調(diào)用Vue前端頁面生成PDF文件的相關(guān)資料,文中通過代碼介紹的非常詳細(xì),需要的朋友可以參考下2025-10-10
java插入排序和希爾排序?qū)崿F(xiàn)思路及代碼
這篇文章主要介紹了插入排序和希爾排序兩種排序算法,文章通過代碼示例和圖解詳細(xì)介紹了這兩種排序算法的實(shí)現(xiàn)過程和原理,需要的朋友可以參考下2025-03-03
java公眾平臺(tái)通用接口工具類HttpConnectUtil實(shí)例代碼
下面小編就為大家分享一篇java公眾平臺(tái)通用接口工具類HttpConnectUtil實(shí)例代碼,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過來看看吧2018-01-01
攔截器獲取request的值之后,Controller拿不到值的解決
這篇文章主要介紹了攔截器獲取request的值之后,Controller拿不到值的解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2021-10-10
修改及反編譯可運(yùn)行Jar包實(shí)現(xiàn)過程詳解
這篇文章主要介紹了如何修改及反編譯可運(yùn)行Jar包,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2020-09-09
SpringBoot快速搭建RESTful應(yīng)用的流程步驟
在現(xiàn)代 Web 應(yīng)用開發(fā)中,RESTful API 是一種流行的設(shè)計(jì)范式,Spring Boot 提供了一套簡(jiǎn)潔的注解和約定,使得開發(fā)者能夠輕松地創(chuàng)建 RESTful 服務(wù),本課程將通過一個(gè)簡(jiǎn)單的實(shí)例,展示如何使用SpringBoot快速搭建RESTful應(yīng)用2025-07-07

