聊聊通過(guò)celery_one避免Celery定時(shí)任務(wù)重復(fù)執(zhí)行的問(wèn)題
在使用Celery統(tǒng)計(jì)每日訪問(wèn)數(shù)量的時(shí)候,發(fā)現(xiàn)一個(gè)任務(wù)會(huì)同時(shí)執(zhí)行兩次,發(fā)現(xiàn)同一時(shí)間內(nèi)(1s內(nèi))竟然同時(shí)發(fā)送了兩次任務(wù),也就是同時(shí)產(chǎn)生了兩個(gè)worker,造成統(tǒng)計(jì)兩次,一直找不到原因。
參考:http://m.fzitv.net/article/226849.htm
有人使用 Redis 實(shí)現(xiàn)了分布式鎖,然后也有人使用了 Celery Once。
Celery Once 也是利用 Redis 加鎖來(lái)實(shí)現(xiàn), Celery Once 在 Task 類基礎(chǔ)上實(shí)現(xiàn)了 QueueOnce 類,該類提供了任務(wù)去重的功能,所以在使用時(shí),我們自己實(shí)現(xiàn)的方法需要將 QueueOnce 設(shè)置為 base
@task(base=QueueOnce, once={'graceful': True})
后面的 once 參數(shù)表示,在遇到重復(fù)方法時(shí)的處理方式,默認(rèn) graceful 為 False,那樣 Celery 會(huì)拋出 AlreadyQueued 異常,手動(dòng)設(shè)置為 True,則靜默處理。
另外如果要手動(dòng)設(shè)置任務(wù)的 key,可以指定 keys 參數(shù)
@celery.task(base=QueueOnce, once={'keys': ['a']})
def slow_add(a, b):
sleep(30)
return a + b
解決步驟
Celery One允許你將Celery任務(wù)排隊(duì),防止多次執(zhí)行
安裝
pip install -U celery_once
要求,需要Celery4.0,老版本可能運(yùn)行,但不是官方支持的。
使用celery_once,tasks需要繼承一個(gè)名為QueueOnce的抽象base tasks
Once安裝完成后,需要配置一些關(guān)于ONCE的選項(xiàng)在Celery配置中
from celery import Celery
from celery_once import QueueOnce
from time import sleep
celery = Celery('tasks', broker='amqp://guest@localhost//')
# 一般之前的配置沒(méi)有這個(gè),需要添加上
celery.conf.ONCE = {
'backend': 'celery_once.backends.Redis',
'settings': {
'url': 'redis://localhost:6379/0',
'default_timeout': 60 * 60
}
}
# 在原本沒(méi)有參數(shù)的里面加上base
@celery.task(base=QueueOnce)
def slow_task():
sleep(30)
return "Done!"
要確定配置,需要取決于使用哪個(gè)backend進(jìn)行鎖定,查看Backends
在后端,這將覆蓋apply_async和delay。它不影響直接調(diào)用任務(wù)。
在運(yùn)行任務(wù)時(shí),celery_once檢查是否沒(méi)有鎖定(針對(duì)Redis鍵)。否則,任務(wù)將正常運(yùn)行。一旦任務(wù)完成(或由于異常而結(jié)束),鎖將被清除。如果在任務(wù)完成之前嘗試再次運(yùn)行該任務(wù),將會(huì)引發(fā)AlreadyQueued異常。
example.delay(10)
example.delay(10)
Traceback (most recent call last):
..
AlreadyQueued()
result = example.apply_async(args=(10))
result = example.apply_async(args=(10))
Traceback (most recent call last):
..
AlreadyQueued()
graceful:如果在任務(wù)的選項(xiàng)中設(shè)置了once={'graceful': True},或者在運(yùn)行時(shí)設(shè)置了apply_async,則任務(wù)可以返回None,而不是引發(fā)AlreadyQueued異常。
from celery_once import AlreadyQueued
# Either catch the exception,
try:
example.delay(10)
except AlreadyQueued:
pass
# Or, handle it gracefully at run time.
result = example.apply(args=(10), once={'graceful': True})
# or by default.
@celery.task(base=QueueOnce, once={'graceful': True})
def slow_task():
sleep(30)
return "Done!"
其他功能請(qǐng)?jiān)L問(wèn):https://pypi.org/project/celery_once/
到此這篇關(guān)于通過(guò)celery_one避免Celery定時(shí)任務(wù)重復(fù)執(zhí)行的文章就介紹到這了,更多相關(guān)Celery定時(shí)任務(wù)重復(fù)執(zhí)行內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Python Pytest裝飾器@pytest.mark.parametrize詳解
本文主要介紹了Python Pytest裝飾器@pytest.mark.parametrize詳解,文中通過(guò)示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2021-08-08
python 還原梯度下降算法實(shí)現(xiàn)一維線性回歸
這篇文章主要介紹了python 還原梯度下降算法實(shí)現(xiàn)一維線性回歸,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧2020-10-10
一鍵搞定python連接mysql驅(qū)動(dòng)有關(guān)問(wèn)題(windows版本)
這篇文章主要介紹了對(duì)于mysql驅(qū)動(dòng)問(wèn)題折騰了一下午,現(xiàn)共享出解決方案,需要的朋友可以參考下2016-04-04
Python調(diào)用微信公眾平臺(tái)接口操作示例
這篇文章主要介紹了Python調(diào)用微信公眾平臺(tái)接口操作,結(jié)合具體實(shí)例形式分析了Python針對(duì)微信接口數(shù)據(jù)傳輸?shù)南嚓P(guān)操作技巧,需要的朋友可以參考下2017-07-07
Python實(shí)現(xiàn)線程池之線程安全隊(duì)列
這篇文章主要為大家詳細(xì)介紹了Python實(shí)現(xiàn)線程池之線程安全隊(duì)列,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2022-05-05
python密碼學(xué)簡(jiǎn)單替代密碼解密及測(cè)試教程
這篇文章主要介紹了python密碼學(xué)簡(jiǎn)單替代密碼解密及測(cè)試教程,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2022-05-05
詳解Django中 render() 函數(shù)的使用方法
這篇文章主要介紹了Django中 render() 函數(shù)的使用方法,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2021-04-04

