列:Celery上手)
在上一篇的博文中, 實(shí)現(xiàn)了一個(gè)異步任務(wù)情景, 會(huì)在調(diào)用web服務(wù)后馬上返回結(jié)果, 而后臺(tái)會(huì)接著執(zhí)行這個(gè)任務(wù), 這是我工作里的一個(gè)實(shí)際需求, 在我費(fèi)盡周折把這個(gè)功能編寫完成后, 我才知曉有一個(gè)現(xiàn)成的工具能夠達(dá)成這個(gè)功能, 這便是今天要學(xué)習(xí)的。這是一個(gè)有著這般特性的框架, 它簡(jiǎn)單, 靈活, 可靠, 屬于分布式任務(wù)執(zhí)行框架范疇, 它能夠支持大量任務(wù)的并發(fā)執(zhí)行情況。該框架采用典型生產(chǎn)者與消費(fèi)者模型。生產(chǎn)者將任務(wù)提交到任務(wù)隊(duì)列里。眾多消費(fèi)者從任務(wù)隊(duì)列當(dāng)中獲取任務(wù)去執(zhí)行。有這樣一種設(shè)計(jì)模式被稱作生產(chǎn)者和消費(fèi)者模型, 在該模式下, 生產(chǎn)者是任務(wù)的發(fā)布者, 消費(fèi)者是任務(wù)的獲取者, 生產(chǎn)者與消費(fèi)者不存在直接關(guān)聯(lián), 他們之間的交流借助中間人來達(dá)成, 這個(gè)中間人也被叫做消息隊(duì)列。處于這個(gè)進(jìn)程里, 生產(chǎn)者如同懸賞榜上之張貼告示者那般, 把任務(wù)投放至消息隊(duì)列當(dāng)中, 任務(wù)于任務(wù)隊(duì)列里逐一執(zhí)行完畢后, 把結(jié)果傳送給消費(fèi)者, 于生產(chǎn)環(huán)境下, 任務(wù)隊(duì)列通常借助Redis予以實(shí)現(xiàn)。實(shí)際場(chǎng)景的實(shí)際場(chǎng)景在日常生活中經(jīng)常出現(xiàn)舉例說, 在Web應(yīng)用里, 當(dāng)用戶引發(fā)了一個(gè)得長(zhǎng)時(shí)間開展的操作之時(shí)高計(jì)算或者高IO等會(huì)形成阻塞的任務(wù)類型, 能夠?qū)⑵洚?dāng)作任務(wù)給予異步去執(zhí)行, 執(zhí)行完畢之后再返還給用戶。這個(gè)時(shí)間段用戶無需等待。有著這樣一種情況, 對(duì)于身為用戶的他而言, 在其點(diǎn)擊了執(zhí)行按鈕之后, 便得到了一個(gè)任務(wù)ID, 至于程序呢, 是在后臺(tái)方面執(zhí)行的, 接下來, 用戶所只需要做的僅僅是等上一段時(shí)間, 通過這個(gè)任務(wù)ID去拿到任務(wù)執(zhí)行之后的結(jié)果便可。還有一個(gè)場(chǎng)景是定時(shí)任務(wù)例如需要定時(shí)向一些地址發(fā)布郵件。在著手代碼之初, 我察覺到了極大的問題, 于我的設(shè)備之上運(yùn)行程序之際, 一直出現(xiàn)報(bào)錯(cuò)。: not to起初, 我方才覺得那當(dāng)屬代碼邏輯之問題, 而最終經(jīng)查找發(fā)覺乃是其最新版本當(dāng)下并不予以支許了, 于此情形能夠采用WSL或者借助 -A --poolsolo -l info來開展執(zhí)行操作。若采用如此這般的后者方式, 那就表明了是以單線程模式來對(duì)代碼予以執(zhí)行的。最簡(jiǎn)單的案例先是一個(gè)最為簡(jiǎn)單的案例, 我們存在一個(gè)進(jìn)行計(jì)算的程序, 它承擔(dān)著把輸入的兩個(gè)數(shù)字加起來的職責(zé), 得以獲取結(jié)果為了去模擬具備高計(jì)算量的程序, 我們于計(jì)算之際添加上秒。目前, 我們期望用戶在運(yùn)行該程序時(shí)間段, 程序不會(huì)因sleep長(zhǎng)達(dá)兩秒而出現(xiàn)阻塞狀況, 而是能夠于后臺(tái)開展執(zhí)行操作。此一過程, 我們將其放置于隊(duì)列當(dāng)中。想要達(dá)成這個(gè)目標(biāo), 我們需要去實(shí)現(xiàn)以下幾個(gè)方面的內(nèi)容:我們得去達(dá)成本地的一個(gè)消息代理的實(shí)現(xiàn), 就像Redis那樣, 我們要去實(shí)現(xiàn)一個(gè)生產(chǎn)者程序, 其職責(zé)是生成任務(wù)給消息代理發(fā)送過去, 還要去實(shí)現(xiàn)一個(gè)消費(fèi)者程序, 在負(fù)責(zé)從消息隊(duì)列那兒接收任務(wù)后去執(zhí)行它們。在這個(gè)過程中生產(chǎn)者不負(fù)責(zé)執(zhí)行程序只負(fù)責(zé)發(fā)布任務(wù)。以下是代碼的實(shí)現(xiàn)一開始, 我們借助在本地的6379端口來開啟redis服務(wù), 在這里就不再詳細(xì)敘述了。消費(fèi)者程序我們命名為tasks.py。1 2 3 4 5 6 7 8 9 10 11 12import time from celery import Celery broker redis://127.0.0.1:6379 backend redis://127.0.0.1:6379/0 app Celery(my_task, brokerbroker, backendbackend) app.task def add(x, y): time.sleep(2) # 模擬耗時(shí)操作 return x y消費(fèi)者程序當(dāng)中, 定義了消息代理, 其是用redis實(shí)現(xiàn)的, 還定義了結(jié)果后端, 這也是用redis實(shí)現(xiàn)的。按其名稱含義來說, 其中一個(gè)是用來連接消息隊(duì)列的, 另一個(gè)是用來存儲(chǔ)結(jié)果的。在起始點(diǎn), 創(chuàng)建了一個(gè)實(shí)例, 它的稱謂是。其中, app.task屬于一個(gè)裝飾器范疇, 該裝飾器會(huì)把被其修飾的函數(shù)登記成為任務(wù)。進(jìn)而使得這個(gè)函數(shù)能夠以異步方式來進(jìn)行調(diào)用了。生產(chǎn)者程序命名為.py負(fù)責(zé)發(fā)布任務(wù)。1 2 3 4 5from tasks import add # 異步任務(wù) add.delay(2, 8) print(hello world)生產(chǎn)者里頭, 最先導(dǎo)入了歸消費(fèi)者所有的add函數(shù), add函數(shù)經(jīng)app.task進(jìn)行包裝, 搖身一變成了一個(gè)任務(wù), 到了這時(shí)候我們憑借delay方法從而能異步執(zhí)行它, 且傳進(jìn)兩個(gè)參數(shù)2, 8。執(zhí)行異步任務(wù)之際, 程序并非會(huì)干等著兩秒來返回結(jié)果, 而是即刻去執(zhí)行下面的print(hello world), 并且add的結(jié)果會(huì)于后臺(tái)開展計(jì)算然后返回。如何執(zhí)行他們呢首先需要在命令行執(zhí)行1celery -A tasks worker --poolsolo -l info正在開啟一個(gè)用于監(jiān)聽隊(duì)列, 而執(zhí)行任務(wù)的工作進(jìn)程。-A所代表的應(yīng)用模塊名源自tasks.per, 用以表明要開啟工作進(jìn)程, 進(jìn)而示意日志的級(jí)別。啟動(dòng)后能看到成功連接的日志于是乎, 于此之際, 我們于另外的一個(gè)命令行那兒去執(zhí)行.py。緊接著, 命令行便會(huì)即刻返回hello world。當(dāng)此之時(shí)呀, 程序?qū)?huì)就在后臺(tái)進(jìn)行執(zhí)行, 能夠在進(jìn)程的后臺(tái)部位看到接收以及執(zhí)行的結(jié)果。這樣就實(shí)現(xiàn)了一個(gè)最簡(jiǎn)單的用例。app.task裝飾器將程序包裝成實(shí)例的那個(gè), 是app.task這個(gè)裝飾器, 這里面存在幾個(gè)需要留意的要點(diǎn)。1 2 3app.task(bindTrue) def add(self, x, y): print(self.request.id)此時(shí)程序的第一個(gè)參數(shù)必須是任務(wù)實(shí)例不然拿不到任務(wù)id。1 2 3app.task(nametasks.add) # 不顯式設(shè)置的話也為task.add def add(x, y): return x y1 2 3 4 5 6 7app.task(bindTrue) def send_twitter_status(self, oauth, tweet): try: twitter Twitter(oauth) twitter.update_status(tweet) except (Twitter.FailWhaleError, Twitter.LoginError) as exc: raise self.retry(excexc)或者一種更方便的方法1 2 3 4app.task(autoretry_for(FailWhaleError,), retry_kwargs{max_retries: 5}) def refresh_timeline(user): return twitter.refresh_timeline(user)Delay方法所提供的delay方法, 是一個(gè)屬于異步執(zhí)行的接口, 它是對(duì)另外一個(gè)接口進(jìn)行的封裝。在執(zhí)行之后, 它們會(huì)返回一個(gè)實(shí)例, 這個(gè)實(shí)例的作用是用來跟蹤任務(wù)的狀態(tài), 也就是專門用來存儲(chǔ)這個(gè)的。結(jié)果的獲取我們能夠于上面所提及的代碼之中直接獲取結(jié)果, 以及與任務(wù)相關(guān)聯(lián)的信息, 情況如下:1 2 3 4 5 6 7 8 9 10from tasks import add # 異步任務(wù) res add.delay(2, 8) print(hello world) res.get(timeout1) # 10如果出現(xiàn)報(bào)錯(cuò)會(huì)將調(diào)用棧返回 res.id # 獲取任務(wù)id res.get(propagateFalse) # 10但是不返回報(bào)錯(cuò)信息 res.state # 任務(wù)狀態(tài)包含PENDING/STARTED/SUCCESS/FAILURE等在這兒直接獲取結(jié)果, 事實(shí)上有點(diǎn)類似順序執(zhí)行情況。要是拿到了任務(wù)id, 那需要靠再一個(gè)不同模樣的服務(wù)去查看相應(yīng)任務(wù)狀態(tài)該怎么操作呢?1 2 3from tasks import app # 先導(dǎo)入Celery實(shí)例 res app.AsyncResult(given-task-id) # 這時(shí)候就可以和上面一樣獲取任務(wù)結(jié)果了構(gòu)建鏈與同樣, 亦援手鏈?zhǔn)秸{(diào)用。設(shè)若需求于一項(xiàng)任務(wù)回返之后調(diào)用另外一項(xiàng)任務(wù)。于此便牽扯到簽名。所謂簽名指的乃是把一項(xiàng)任務(wù)的實(shí)行選項(xiàng)跟參數(shù)予以打包, 諸如:1 2 3add.signature((2, 2), countdown10) # 為add任務(wù)增加了22的參數(shù)和倒計(jì)時(shí)10秒的執(zhí)行選項(xiàng) add.s(2, 2) # 簡(jiǎn)寫對(duì)于上面這個(gè)簽名也可以直接執(zhí)行1 2 3s1 add.s(2, 2) res s1.delay() res.get()如果使用鏈的話是這樣的1 2 3 4 5from celery import chain from tasks import add, multiply # (4 4) * 8 chain(add.s(4,4) | multiply.s(8))().get()路由支持路由也就是根據(jù)名稱將結(jié)果發(fā)到不同隊(duì)列1 2 3 4 5app.conf.update( task_routes { tasks.add: {queue: add_queue}, }, )在執(zhí)行時(shí)在方法中加入queue參數(shù)1 2from tasks import add add.apply_async((2, 2), queueadd_queue)并在執(zhí)行時(shí)使用-Q來選擇隊(duì)列1celery -A tasks worker -Q add_queue讀取配置文件處于上面提及的程序里, 和的配置是書寫于程序之中的, 不過呢, 它同樣能夠被寫成配置文件, 要運(yùn)用 app 的方式去加載配置。必須留意, 配置文件得跟啟動(dòng)文件放置于同一個(gè)路徑之下。舉例來說:在項(xiàng)目路徑下創(chuàng)建.py內(nèi)容為1 2 3 4 5 6 7 8 9 10from datetime import timedelta from celery.schedules import crontab broker_url redis://127.0.0.1:6379 # 指定 Broker result_backend redis://127.0.0.1:6379/0 # 指定 Backend broker_connection_retry_on_startup True imports ( # 指定導(dǎo)入的任務(wù)模塊 tasks, )相應(yīng)的tasks.py也要修改一下修改后內(nèi)容如下1 2 3 4 5 6 7 8 9 10import time from celery import Celery app Celery(demo) # Celery實(shí)例的名稱 app.config_from_object(celery_config) app.task def add(x, y): time.sleep(2) # 模擬耗時(shí)操作 return x y最初的時(shí)候, 定義出來的地址以及 app 都是寫在 task.py 這個(gè)文件當(dāng)中的, 然而現(xiàn)如今, 僅僅只要在 task.py 里面直接加載配置文件就行了。2024/5/26 于蘇州