在上一篇的博文中, 实现了一个异步任务情景, 会在调用web服务后马上返回结果, 而后台会接着执行这个任务, 这是我工作里的一个实际需求, 在我费尽周折把这个功能编写完成后, 我才知晓有一个现成的工具能够达成这个功能, 这便是今天要学习的。
这是一个有着这般特性的框架, 它简单, 灵活, 可靠, 属于分布式任务执行框架范畴, 它能够支持大量任务的并发执行情况。该框架采用典型生产者与消费者模型。生产者将任务提交到任务队列里。众多消费者从任务队列当中获取任务去执行。
有这样一种设计模式被称作生产者和消费者模型, 在该模式下, 生产者是任务的发布者, 消费者是任务的获取者, 生产者与消费者不存在直接关联, 他们之间的交流借助中间人来达成, 这个中间人也被叫做消息队列。
处于这个进程里, 生产者如同悬赏榜上之张贴告示者那般, 把任务投放至消息队列当中, 任务于任务队列里逐一执行完毕后, 把结果传送给消费者, 于生产环境下, 任务队列通常借助Redis予以实现。
实际场景
的实际场景在日常生活中经常出现:
举例说, 在Web应用里, 当用户引发了一个得长时间开展的操作之时(高计算或者高IO等会形成阻塞的任务类型), 能够将其当作任务给予异步去执行, 执行完毕之后再返还给用户。这个时间段用户无需等待。
有着这样一种情况, 对于身为用户的他而言, 在其点击了执行按钮之后, 便得到了一个任务ID, 至于程序呢, 是在后台方面执行的, 接下来, 用户所只需要做的仅仅是等上一段时间, 通过这个任务ID去拿到任务执行之后的结果便可。
还有一个场景是定时任务:例如需要定时向一些地址发布邮件。
在着手代码之初, 我察觉到了极大的问题, 于我的设备之上运行程序之际, 一直出现报错。
: not to
起初, 我方才觉得那当属代码逻辑之问题, 而最终经查找发觉乃是其最新版本当下并不予以支许了, 于此情形能够采用WSL或者借助 -A --pool=solo -l info来开展执行操作。若采用如此这般的后者方式, 那就表明了是以单线程模式来对代码予以执行的。
最简单的案例
先是一个最为简单的案例, 我们存在一个进行计算的程序, 它承担着把输入的两个数字加起来的职责, 得以获取结果,为了去模拟具备高计算量的程序, 我们于计算之际添加上秒。
目前, 我们期望用户在运行该程序时间段, 程序不会因sleep长达两秒而出现阻塞状况, 而是能够于后台开展执行操作。此一过程, 我们将其放置于队列当中。想要达成这个目标, 我们需要去实现以下几个方面的内容:
我们得去达成本地的一个消息代理的实现, 就像Redis那样, 我们要去实现一个生产者程序, 其职责是生成任务给消息代理发送过去, 还要去实现一个消费者程序, 在负责从消息队列那儿接收任务后去执行它们。
在这个过程中,生产者不负责执行程序,只负责发布任务。
以下是代码的实现:
一开始, 我们借助在本地的6379端口来开启redis服务, 在这里就不再详细叙述了。
消费者程序,我们命名为tasks.py。
1 2 3 4 5 6 7 8 9 10 11 12
import time from celery import Celery broker = 'redis://127.0.0.1:6379' backend = 'redis://127.0.0.1:6379/0' app = Celery('my_task', broker=broker, backend=backend) @app.task def add(x, y): time.sleep(2) # 模拟耗时操作 return x + y消费者程序当中, 定义了消息代理, 其是用redis实现的, 还定义了结果后端, 这也是用redis实现的。按其名称含义来说, 其中一个是用来连接消息队列的, 另一个是用来存储结果的。
在起始点, 创建了一个实例, 它的称谓是。其中, @app.task属于一个装饰器范畴, 该装饰器会把被其修饰的函数登记成为任务。进而使得这个函数能够以异步方式来进行调用了。
生产者程序,命名为.py,负责发布任务。
1 2 3 4 5
from tasks import add # 异步任务 add.delay(2, 8) print('hello world')生产者里头, 最先导入了归消费者所有的add函数, add函数经@app.task进行包装, 摇身一变成了一个任务, 到了这时候,我们凭借delay方法从而能异步执行它, 且传进两个参数2, 8。
执行异步任务之际, 程序并非会干等着两秒来返回结果, 而是即刻去执行下面的print('hello world'), 并且add的结果会于后台开展计算然后返回。
如何执行他们呢?首先需要在命令行执行:
1
celery -A tasks worker --pool=solo -l info正在开启一个用于监听队列, 而执行任务的工作进程。-A所代表的应用模块名源自tasks.per, 用以表明要开启工作进程, 进而示意日志的级别。
启动后能看到成功连接的日志:
于是乎, 于此之际, 我们于另外的一个命令行那儿去执行.py。紧接着, 命令行便会即刻返回hello world。当此之时呀, 程序将会就在后台进行执行, 能够在进程的后台部位看到接收以及执行的结果。
这样就实现了一个最简单的用例。
app.task装饰器
将程序包装成实例的那个, 是@app.task这个装饰器, 这里面存在几个需要留意的要点。
1 2 3
@app.task(bind=True) def add(self, x, y): print(self.request.id)此时程序的第一个参数必须是任务实例,不然拿不到任务id。
1 2 3
@app.task(name='tasks.add') # 不显式设置的话也为task.add def add(x, y): return x + y1 2 3 4 5 6 7
@app.task(bind=True) 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(exc=exc)或者一种更方便的方法:
1 2 3 4
@app.task(autoretry_for=(FailWhaleError,), retry_kwargs={'max_retries': 5}) def refresh_timeline(user): return twitter.refresh_timeline(user)Delay方法
所提供的delay方法, 是一个属于异步执行的接口, 它是对另外一个接口进行的封装。在执行之后, 它们会返回一个实例, 这个实例的作用是用来跟踪任务的状态, 也就是专门用来存储这个的。
结果的获取
我们能够于上面所提及的代码之中直接获取结果, 以及与任务相关联的信息, 情况如下:
1 2 3 4 5 6 7 8 9 10
from tasks import add # 异步任务 res = add.delay(2, 8) print('hello world') res.get(timeout=1) # 10,如果出现报错会将调用栈返回 res.id # 获取任务id res.get(propagate=False) # 10,但是不返回报错信息 res.state # 任务状态,包含PENDING/STARTED/SUCCESS/FAILURE等在这儿直接获取结果, 事实上有点类似顺序执行情况。要是拿到了任务id, 那需要靠再一个不同模样的服务去查看相应任务状态该怎么操作呢?
1 2 3
from tasks import app # 先导入Celery实例 res = app.AsyncResult('given-task-id') # 这时候就可以和上面一样获取任务结果了构建链
与同样, 亦援手链式调用。设若需求于一项任务回返之后调用另外一项任务。于此便牵扯到签名。所谓签名指的乃是把一项任务的实行选项跟参数予以打包, 诸如:
1 2 3
add.signature((2, 2), countdown=10) # 为add任务增加了2,2的参数,和倒计时10秒的执行选项 add.s(2, 2) # 简写对于上面这个签名,也可以直接执行:
1 2 3
s1 = add.s(2, 2) res = s1.delay() res.get()如果使用链的话是这样的:
1 2 3 4 5
from celery import chain from tasks import add, multiply # (4 + 4) * 8 chain(add.s(4,4) | multiply.s(8))().get()路由
支持路由,也就是根据名称将结果发到不同队列:
1 2 3 4 5
app.conf.update( task_routes = { 'tasks.add': {'queue': 'add_queue'}, }, )在执行时,在方法中加入queue参数:
1 2
from tasks import add add.apply_async((2, 2), queue='add_queue')并在执行时使用-Q来选择队列:
1
celery -A tasks worker -Q add_queue读取配置文件
处于上面提及的程序里, 和的配置是书写于程序之中的, 不过呢, 它同样能够被写成配置文件, 要运用 app 的方式去加载配置。必须留意, 配置文件得跟启动文件放置于同一个路径之下。举例来说:
在项目路径下创建.py,内容为:
1 2 3 4 5 6 7 8 9 10
from 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也要修改一下,修改后内容如下:
1 2 3 4 5 6 7 8 9 10
import 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 这个文件当中的, 然而现如今, 仅仅只要在 task.py 里面直接加载配置文件就行了。
2024/5/26 于苏州