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