網(wǎng)站首頁(yè) 編程語(yǔ)言 正文
日常需求開發(fā)過(guò)程中,不免會(huì)遇到需要通過(guò)代碼進(jìn)行異步處理的情況,比如批量發(fā)送郵件,批量發(fā)送短信,數(shù)據(jù)導(dǎo)入,為了減少用戶的等待,不希望一直菊花轉(zhuǎn)啊轉(zhuǎn),因此需要進(jìn)行異步處理,做法就是講要處理的數(shù)據(jù)添加到隊(duì)列當(dāng)中,然后按照排隊(duì)的先后順序進(jìn)行異步處理。
這個(gè)隊(duì)列,可以是專業(yè)的消息隊(duì)列,如 RocketMQ/RabbitMQ 等,一般項(xiàng)目中,如果只是為了進(jìn)行異步,未免有點(diǎn)殺雞用牛刀的意味。
也可以使用基于 JVM 內(nèi)存實(shí)現(xiàn)隊(duì)列,但是如果項(xiàng)目進(jìn)行了重啟,就會(huì)造成隊(duì)列數(shù)據(jù)丟失。
大部分的項(xiàng)目都會(huì)用到 Redis 中間件作為緩存使用,此時(shí)使用 Redis 的 list 結(jié)構(gòu)來(lái)實(shí)現(xiàn)隊(duì)列則是非常合適的選擇。
因此,本文主要講解基于 Redis 的方式實(shí)現(xiàn)異步隊(duì)列。
本文首發(fā)個(gè)人技術(shù)博客: https://nullpointer.pw/redis-block-queue.html
基于 Redis 的 list 實(shí)現(xiàn)隊(duì)列的方式也有多種,先說(shuō)第一種不推薦的方式,即使用LPUSH
生產(chǎn)消息,然后 while(true) 中通過(guò)RPOP
消費(fèi)消息,這種方式的確可以實(shí)現(xiàn),但是不斷代碼不斷的輪詢,勢(shì)必會(huì)消耗一些系統(tǒng)的資源。
第二種方式也是不推薦的方式,也是通過(guò) LPUSH
生產(chǎn)消息,然后通過(guò) BRPOP
進(jìn)行阻塞地等待并消費(fèi)消息,這種方式較第一種方式減少了無(wú)用的輪詢,降低系統(tǒng)資源的消耗,但是可能會(huì)存在隊(duì)列消息丟失的情況,如果取出了消息然后處理失敗,這個(gè)被取出的消息就將丟失。
第二種方式就是下文要介紹的方式,首先也是通過(guò) LPUSH
生產(chǎn)消息,然后通過(guò) BRPOPLPUSH
阻塞地等待 list 新消息到來(lái),有了新消息才開始消費(fèi),同時(shí)將消息備份到另外一個(gè) list 當(dāng)中,這種方式具備了第二種方式的優(yōu)點(diǎn),即減少了無(wú)用的輪詢,同時(shí)也對(duì)消息進(jìn)行了備份不會(huì)丟失數(shù)據(jù),如果處理成功,可以通過(guò) LREM
對(duì)備份的 list 中當(dāng)前的這條消息進(jìn)行刪除處理。這種方式實(shí)現(xiàn)方式可以參考 模式: 安全的隊(duì)列 .
Redis 基礎(chǔ)
# 將一個(gè)或多個(gè)值 value 插入到列表 key 的表頭 LPUSH key value [value …] # 阻塞式等待,將列表 source 中的最后一個(gè)元素 (尾元素) 彈出,并返回給客戶端。將 source 彈出的元素插入到列表 destination ,作為 destination 列表的的頭元素。超時(shí)參數(shù) timeout 接受一個(gè)以秒為單位的數(shù)字作為值。超時(shí)參數(shù)設(shè)為 0 表示阻塞時(shí)間可以無(wú)限期延長(zhǎng) (block indefinitely) 。 BRPOPLPUSH source destination timeout # 根據(jù)參數(shù) count 的值,移除列表中與參數(shù) value 相等的元素。 LREM key count value
代碼實(shí)現(xiàn)隊(duì)列消息生產(chǎn)者
筆者使用的是 Spring 相關(guān) API 實(shí)現(xiàn)對(duì) Redis 指令的調(diào)用。首先實(shí)現(xiàn)消息的生產(chǎn)代碼,封裝到一個(gè)工具類方法當(dāng)中。這里很簡(jiǎn)單,就是調(diào)用了 lpush 方法,將序列化的 key 和 value 添加到列表當(dāng)中去。
@Resource private RedisConnectionFactory connectionFactory; public void lPush(@Nonnull String key, @Nonnull String value) { RedisConnection connection = RedisConnectionUtils.getConnection(connectionFactory); try { byte[] byteKey = RedisSerializer.string().serialize(getKey(key)); byte[] byteValue = RedisSerializer.string().serialize(value); assert byteKey != null; connection.lPush(byteKey, byteValue); } finally { RedisConnectionUtils.releaseConnection(connection, connectionFactory); } }
代碼實(shí)現(xiàn)隊(duì)列消息消費(fèi)者
因?yàn)閷?shí)現(xiàn)隊(duì)列消費(fèi)消息的代碼比較多,不可能每個(gè)需要阻塞消費(fèi)的地方,對(duì)需要寫這一坨代碼,因此使用 Java8 的函數(shù)式接口實(shí)現(xiàn)方法的傳遞,同時(shí)阻塞式獲取消息代碼使用新線程去執(zhí)行。
有人看到以下代碼要吐槽了,不是說(shuō)不用 while(true) 嗎,怎么你這里面還是有,這里稍微解釋一下,因?yàn)?SpringBoot 一般會(huì)指定 timeout 的全局超時(shí)時(shí)間,即使 BRPOPLPUSH
設(shè)置了 0,即無(wú)限期,當(dāng)超出了 timeout 設(shè)置的值時(shí),就會(huì)拋出 QueryTimeoutException 異常導(dǎo)致線程退出,因此添加了 try/catch 對(duì)異常進(jìn)行捕獲并忽略,同時(shí)使用 while(true) 保證線程可以繼續(xù)執(zhí)行。
代碼中記錄了當(dāng)前消息處理結(jié)果,如果處理結(jié)果為成功,需要對(duì)備份隊(duì)列的當(dāng)前消息進(jìn)行刪除。
public void bRPopLPush(@Nonnull String key, Consumer<String> consumer) { CompletableFuture.runAsync(() -> { RedisConnection connection = RedisConnectionUtils.getConnection(connectionFactory); try { byte[] srcKey = RedisSerializer.string().serialize(getKey(key)); byte[] dstKey = RedisSerializer.string().serialize(getBackupKey(key)); assert srcKey != null; assert dstKey != null; while (true) { byte[] byteValue = new byte[0]; boolean success = false; try { byteValue = connection.bRPopLPush(0, srcKey, dstKey); if (byteValue != null && byteValue.length != 0) { consumer.accept(new String(byteValue)); success = true; } } catch (Exception ignored) { // 防止獲取 key 達(dá)到超時(shí)時(shí)間拋出 QueryTimeoutException 異常退出 } finally { if (success) { // 處理成功才刪除備份隊(duì)列的 key connection.lRem(dstKey, 1, byteValue); } } } } finally { RedisConnectionUtils.releaseConnection(connection, connectionFactory); } }); }
測(cè)試代碼
@Test public void testLPush() throws InterruptedException { String queueA = "queueA"; int i = 0; while (true) { String msg = "Hello-" + i++; redisBlockQueue.lPush(queueA, msg); System.out.println("lPush: " + msg); Thread.sleep(3000); } } @Test public void testBRPopLPush() { String queueA = "queueA"; redisBlockQueue.bRPopLPush(queueA, (val) -> { // 在這里處理具體的業(yè)務(wù)邏輯 System.out.println("val: " + val); }); // 防止 Junit 進(jìn)程退出 LockSupport.park(); }
項(xiàng)目使用方式
為了方便使用,我將其抽取為了一個(gè)工具類,使用時(shí)通過(guò) Spring 注入使用即可,
隊(duì)列消費(fèi)可以使用如下方式在項(xiàng)目啟動(dòng)的時(shí)候就進(jìn)行阻塞監(jiān)聽隊(duì)列,等待消費(fèi)
@Resource private RedisBlockQueue redisBlockQueue; @PostConstruct public void init() { redisBlockQueue.bRPopLPush(xx, (value) -> { //... }); }
本文完整代碼下載github 地址
原文鏈接:https://www.cnblogs.com/vcmq/p/13509269.html
相關(guān)推薦
- 2022-11-16 詳解C++中的左值,純右值和將亡值_C 語(yǔ)言
- 2022-11-06 SQL?Server?Reporting?Services?匿名登錄的問(wèn)題及解決方案_MsSql
- 2022-08-15 el-table表格怎么在表頭的某項(xiàng)添加提示信息
- 2022-08-19 python查看自己安裝的所有庫(kù)并導(dǎo)出的命令_python
- 2024-07-18 Spring Security之配置體系
- 2022-08-13 404究竟是什么意思呢?像404,200,503等數(shù)字究竟是什么東西
- 2022-12-15 C#入?yún)⑹褂靡妙愋鸵觬ef的原因解析_C#教程
- 2022-10-02 iOS簡(jiǎn)單抽屜效果的實(shí)現(xiàn)方法_IOS
- 最近更新
-
- window11 系統(tǒng)安裝 yarn
- 超詳細(xì)win安裝深度學(xué)習(xí)環(huán)境2025年最新版(
- Linux 中運(yùn)行的top命令 怎么退出?
- MySQL 中decimal 的用法? 存儲(chǔ)小
- get 、set 、toString 方法的使
- @Resource和 @Autowired注解
- Java基礎(chǔ)操作-- 運(yùn)算符,流程控制 Flo
- 1. Int 和Integer 的區(qū)別,Jav
- spring @retryable不生效的一種
- Spring Security之認(rèn)證信息的處理
- Spring Security之認(rèn)證過(guò)濾器
- Spring Security概述快速入門
- Spring Security之配置體系
- 【SpringBoot】SpringCache
- Spring Security之基于方法配置權(quán)
- redisson分布式鎖中waittime的設(shè)
- maven:解決release錯(cuò)誤:Artif
- restTemplate使用總結(jié)
- Spring Security之安全異常處理
- MybatisPlus優(yōu)雅實(shí)現(xiàn)加密?
- Spring ioc容器與Bean的生命周期。
- 【探索SpringCloud】服務(wù)發(fā)現(xiàn)-Nac
- Spring Security之基于HttpR
- Redis 底層數(shù)據(jù)結(jié)構(gòu)-簡(jiǎn)單動(dòng)態(tài)字符串(SD
- arthas操作spring被代理目標(biāo)對(duì)象命令
- Spring中的單例模式應(yīng)用詳解
- 聊聊消息隊(duì)列,發(fā)送消息的4種方式
- bootspring第三方資源配置管理
- GIT同步修改后的遠(yuǎn)程分支