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