Redis+threading實(shí)現(xiàn)多線程消息隊(duì)列的使用示例
列表
lpush左插入、rpush右插入、lrange查詢集合
127.0.0.1:6379> lpush list v1 (integer) 1 127.0.0.1:6379> lpush list v2 (integer) 2 127.0.0.1:6379> lpush list v3 (integer) 3 127.0.0.1:6379> LRANGE list 0 -1 1) "v3" 2) "v2" 3) "v1" 127.0.0.1:6379> LRANGE list 0 1 1) "v3" 2) "v2" 127.0.0.1:6379> LRANGE list 0 0 1) "v3" 127.0.0.1:6379> rpush list rv0 (integer) 4 127.0.0.1:6379> lrange list 0 -1 1) "v3" 2) "v2" 3) "v1" 4) "rv0"
lpop左移除、rpop右移除
127.0.0.1:6379> lrange list 0 -1 1) "v3" 2) "v2" 3) "v1" 4) "rv0" 127.0.0.1:6379> lpop list "v3" 127.0.0.1:6379> lrange list 0 -1 1) "v2" 2) "v1" 3) "rv0" 127.0.0.1:6379> rpop list "rv0" 127.0.0.1:6379> lrange list 0 -1 1) "v2" 2) "v1"
lindex下標(biāo)查詢、llen長(zhǎng)度查詢
127.0.0.1:6379> lrange list 0 -1 1) "v4" 2) "v3" 3) "v2" 4) "v1" 127.0.0.1:6379> lindex list 1 "v3" 127.0.0.1:6379> lindex list 0 "v4" 127.0.0.1:6379> llen list (integer) 4
blpop、brpop
BRPOP 是 Redis 的一個(gè)阻塞式列表彈出命令,用于從指定的一個(gè)或多個(gè)列表中彈出最后一個(gè)元素。它和BLPOP 不同之處在于它是從列表的尾部彈出元素,而不是從頭部。
這種阻塞式彈出操作通常用于實(shí)現(xiàn)消息隊(duì)列。如果列表為空,就會(huì)阻塞等待直到有消息可供處理。timeout: 阻塞超時(shí)時(shí)間,如果所有指定的列表都為空,命令將阻塞直到有元素可彈出或超時(shí)。
# 連接到本地 Redis 服務(wù)器
r = redis.Redis(host='127.0.0.1', port=6379, decode_responses=True)
a = ['item1','item2','item3','item4','item5','item6']
# 將元素推入列表
r.rpush('my_queue',*a)
# 使用 blpop 彈出元素
result = r.blpop('my_queue', timeout=5) # 設(shè)置超時(shí)時(shí)間為 5 秒
['item1', 'item2', 'item3', 'item4', 'item5', 'item6']
['item2', 'item3', 'item4', 'item5', 'item6']
字符串
set、incr遞增、decr遞減
- 雖然輸入是int類型,但是set會(huì)自動(dòng)轉(zhuǎn)換為string
r.set('my_queue', 5)
# 對(duì) key 為 'my_queue' 的值執(zhí)行遞減操作
value = r.decr('my_queue')
# 獲取遞減后的值
print(f'New value: {value}')
new_value = r.incr('my_queue')
print(f'New value: {new_value}')
New value: 4
New value: 5
可以看到就算是賦值也已經(jīng)改變了my_queue鍵的值。
keys取鍵、get取值、delete
# 將元素推入列表
r.setnx('my_queue:ddd:count',123)
r.setnx('my_queue:aaa:count',456)
r.setnx('my_queue',789)
r.setnx('my_queue',101)
r.set('my_queue:kkk:count',159)
print(r.keys('*'))
keys = r.keys('my_queue:*:count')
# 打印匹配的鍵列表
print(keys)
value = r.get(keys[0])
print(value)
r.delete('my_queue:ddd:count')
print(r.keys('*'))
['my_queue:kkk:count', 'my_queue', 'my_queue:aaa:count', 'my_queue:ddd:count']
['my_queue:kkk:count', 'my_queue:aaa:count', 'my_queue:ddd:count']
159
['my_queue:kkk:count', 'my_queue', 'my_queue:aaa:count']
關(guān)于為什么delete了還能取到值(
delete只是刪除了redis隊(duì)列的鍵值對(duì),keys是已經(jīng)賦過值了所以不受影響。
r.setnx(f'rtp_task:{1}:{123}:count',123)
r.setnx(f'rtp_task:{2}:{456}:count',456)
r.setnx(f'rtp_task:{3}:{789}:count',789)
r.setnx(f'rtp_task:{4}:{101}:count',101)
r.setnx(f'rtp_task:{5}:{159}:count',159)
keys = r.keys('rtp_task:*:count')
r.delete('rtp_task:1:123:count')
print(r.keys("*"))
print(keys)
['rtp_task:4:101:count', 'rtp_task:2:456:count', 'rtp_task:3:789:count', 'rtp_task:5:159:count']
['rtp_task:4:101:count', 'rtp_task:2:456:count', 'rtp_task:3:789:count', 'rtp_task:1:123:count', 'rtp_task:5:159:count']
setnx
含義(setnx = SET if Not eXists):如果不存在,則set。
r.setnx('my_queue:ddd:count',123)
print(r.get('my_queue:ddd:count'))
123
r.setnx('my_queue:ddd:count',123)
r.setnx('my_queue:ddd:count',456)
print(r.get('my_queue:ddd:count'))
123
threading
Thread
創(chuàng)建線程
在創(chuàng)建線程時(shí),傳遞參數(shù)需要是一個(gè)可迭代的對(duì)象,如果只有一個(gè)參數(shù),需要在參數(shù)后面添加逗號(hào),以表示它是一個(gè)元組而不是一個(gè)單一的值。
如果寫成 args=(a,),它會(huì)被解釋為一個(gè)包含單一元素的元組,而如果你寫成 args=(a),它會(huì)被解釋為 args=a,這樣就不再是一個(gè)元組。
a = "this is message"
def iptest(message):
print(message)
t = threading.Thread(target=iptest, args=(a,))
start、join
a = "this is message"
def iptest(message):
print(message)
t = threading.Thread(target=iptest, args=(a,))
t.start()
this is message
join方法的作用是確保thread子線程執(zhí)行完畢后才能執(zhí)行下一個(gè)線程。
沒加join前:
def iptest():
print("message1\n")
for i in range(10):
# time.sleep() 函數(shù)推遲調(diào)用線程的運(yùn)行,可通過參數(shù)secs指秒數(shù),表示進(jìn)程掛起的時(shí)間。
time.sleep(0.1)
print('message2')
def main():
add_thread = threading.Thread(target=iptest, name="T2")
add_thread.start()
print("done")
if __name__ == '__main__':
main()
message1
donemessage2
加join后
message1
message2
done
消息隊(duì)列
import redis
import threading
import time
import json
# 連接到本地 Redis 服務(wù)器
r = redis.Redis(host='127.0.0.1', port=6379, decode_responses=True)
def producer(queue_name):
# 生產(chǎn)者線程,模擬向隊(duì)列中推送任務(wù)
for i in range(5):
message = {'task_id': i, 'data': f'Task {i}'}
r.rpush(queue_name, json.dumps(message))
time.sleep(1)
def consumer(queue_name, worker_id):
while True:
# 消費(fèi)者線程,使用 blpop 從隊(duì)列中阻塞獲取任務(wù)
message = r.blpop(queue_name, timeout=10)
if message:
task = json.loads(message[1])
print(f"Worker {worker_id} processing task: {task} \n")
# 模擬任務(wù)處理時(shí)間
time.sleep(2)
# 模擬任務(wù)處理完成后,更新任務(wù)計(jì)數(shù)
r.decr(f'rtp_task:{task["task_id"]}:count')
num_task = r.get(f'rtp_task:{task["task_id"]}:count')
print(f'taskid{task["task_id"]},num_task{num_task}')
if __name__ == '__main__':
# 設(shè)置初始任務(wù)計(jì)數(shù)
for i in range(5):
r.setnx(f'rtp_task:{i}:count', 3)
# 創(chuàng)建一個(gè)生產(chǎn)者線程
producer_thread = threading.Thread(target=producer, args=('product',))
producer_thread.start()
# 創(chuàng)建多個(gè)消費(fèi)者線程
num_consumers = 3
consumer_threads = []
for i in range(num_consumers):
consumer_thread = threading.Thread(target=consumer, args=('product', i + 1))
consumer_threads.append(consumer_thread)
consumer_thread.start()
# 等待生產(chǎn)者線程和消費(fèi)者線程完成
producer_thread.join()
for consumer_thread in consumer_threads:
consumer_thread.join()
Worker 2 processing task: {'task_id': 0, 'data': 'Task 0'}
Worker 1 processing task: {'task_id': 1, 'data': 'Task 1'}
Worker 3 processing task: {'task_id': 2, 'data': 'Task 2'}
taskid0,num_task2
taskid1,num_task2
Worker 2 processing task: {'task_id': 3, 'data': 'Task 3'}taskid2,num_task2
Worker 1 processing task: {'task_id': 4, 'data': 'Task 4'}taskid3,num_task2
taskid4,num_task2
到此這篇關(guān)于Redis+threading實(shí)現(xiàn)多線程消息隊(duì)列的使用示例的文章就介紹到這了,更多相關(guān)Redis threading多線程消息隊(duì)列內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- Redis消息隊(duì)列的三種實(shí)現(xiàn)方式
- Redis消息隊(duì)列、阻塞隊(duì)列、延時(shí)隊(duì)列的實(shí)現(xiàn)
- Redis使用ZSET實(shí)現(xiàn)消息隊(duì)列的項(xiàng)目實(shí)踐
- 一文弄懂Redis Stream消息隊(duì)列
- Redis使用ZSET實(shí)現(xiàn)消息隊(duì)列使用小結(jié)
- 詳解Redis Stream做消息隊(duì)列
- redis stream 實(shí)現(xiàn)消息隊(duì)列的實(shí)踐
- Redis?中使用?list,streams,pub/sub?幾種方式實(shí)現(xiàn)消息隊(duì)列的問題
- redis用list做消息隊(duì)列的實(shí)現(xiàn)示例
- Redis?使用?List?實(shí)現(xiàn)消息隊(duì)列的優(yōu)缺點(diǎn)
- golang實(shí)現(xiàn)redis的延時(shí)消息隊(duì)列功能示例
相關(guān)文章
redis.conf中使用requirepass不生效的原因及解決方法
本文主要介紹了如何啟用requirepass,以及啟用requirepass為什么不會(huì)生效,從代碼層面分析了不生效的原因,以及解決方法,需要的朋友可以參考下2023-07-07
Redis中5種數(shù)據(jù)結(jié)構(gòu)的使用場(chǎng)景介紹
這篇文章主要介紹了Redis中5種數(shù)據(jù)結(jié)構(gòu)的使用場(chǎng)景介紹,本文對(duì)Redis中的5種數(shù)據(jù)類型String、Hash、List、Set、Sorted Set做了講解,需要的朋友可以參考下2014-09-09
分布式使用Redis實(shí)現(xiàn)數(shù)據(jù)庫(kù)對(duì)象自增主鍵ID
本文介紹在分布式項(xiàng)目中使用Redis生成對(duì)象的自增主鍵ID,通過Redis的INCR等命令實(shí)現(xiàn)計(jì)數(shù)器功能,具有一定的參考價(jià)值,感興趣的可以了解一下2024-12-12
Redis集群模式和常用數(shù)據(jù)結(jié)構(gòu)詳解
Redis集群模式下的運(yùn)維指令主要用于集群的搭建、管理、監(jiān)控和維護(hù),講解了一些常用的Redis集群運(yùn)維指令,本文重點(diǎn)介紹了Redis集群模式和常用數(shù)據(jù)結(jié)構(gòu),需要的朋友可以參考下2024-03-03
Java實(shí)現(xiàn)多級(jí)緩存的方法詳解
對(duì)于高并發(fā)系統(tǒng)來說,有三個(gè)重要的機(jī)制來保障其高效運(yùn)行,它們分別是:緩存、限流和熔斷,所以本文就來和大家探討一下多級(jí)緩存的實(shí)現(xiàn)方法,希望對(duì)大家有所幫助2024-02-02
通過Redis實(shí)現(xiàn)Token黑名單機(jī)制的具體方案
如果沒有為Token提供主動(dòng)失效機(jī)制,一旦Token被簽發(fā),在過期之前將一直有效,存一些安全隱患,所以本文通過將已失效的Token存儲(chǔ)在Redis中,可以確保Token在被主動(dòng)注銷后無法繼續(xù)使用,需要的朋友可以參考下2025-11-11
redis字符串類型_動(dòng)力節(jié)點(diǎn)Java學(xué)院整理
這篇文章主要為大家詳細(xì)介紹了redis字符串類型的相關(guān)資料,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2017-08-08
Redis遍歷海量數(shù)據(jù)集的幾種實(shí)現(xiàn)方法
Redis作為一個(gè)高性能的鍵值存儲(chǔ)數(shù)據(jù)庫(kù),廣泛應(yīng)用于各種場(chǎng)景,包括緩存、消息隊(duì)列、排行榜,本文主要介紹了Redis遍歷海量數(shù)據(jù)集的幾種實(shí)現(xiàn)方法,文中通過示例代碼介紹的非常詳細(xì),需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2024-02-02

