news 2026/9/3 21:05:00

Python消息队列:Celery上手

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Python消息队列:Celery上手

在上一篇的博文中, 实现了一个异步任务情景, 会在调用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 + y
1 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 于苏州

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/3 21:03:35

beego启动流程

从 import beego 到端口可服务的完整链路:包 init → 默认应用构建 → Run 执行。1. 两个阶段阶段一:包初始化(import 即发生)web 包 init → NewHttpSever() → BeeApp 就绪config 包 init → 读 conf/app.conf → BConfig 就绪阶段二:web.Run()(main 里调用)initBeforeHTTPRu…

作者头像 李华
网站建设 2026/9/3 21:02:30

新能源汽车VCU整车控制器开发源码、原理图与PCB设计全解析

简介:这份VCU整套开发资料面向电动汽车控制系统的工程师、嵌入式开发者及车辆工程专业学生,提供从底层驱动、应用层逻辑到硬件设计的完整闭环。压缩包共477个文件,约95.78MB,主要包含C/C源码(.c/.h/.cpp)、…

作者头像 李华
网站建设 2026/9/3 20:52:27

BQ79616+BQ79600底层驱动开发:帧格式、CRC校验与踩坑实战

简介:BQ79616与BQ79600是TI推出的高精度锂离子电池监控芯片,广泛用于电动车、储能设备等BMS电芯电压采集场景。该驱动源码面向BMS嵌入式开发工程师,覆盖芯片初始化、I2C/SPI通信、菊花链地址管理、中断响应、均衡控制及故障诊断等关键环节&am…

作者头像 李华
网站建设 2026/9/3 20:46:07

天正图纸自动转换为标准AutoCAD图元工具

产品概述 以前在设计院时这个需求比较烦人,批量处理天正图纸时,不得不用天正软件自带的命令一张张处理,后来开发的插件终于实现了用脚本一张张处理,但是仍然不支持并发。现在总算弄成了一个CAD插件,实现了批量处理。 …

作者头像 李华
网站建设 2026/9/3 20:45:25

FunASR Docker部署实战:搭建2Pass实时语音识别服务并完成WebSocket测试

最近需要在本地部署一个语音识别服务,用于后续的音频转文字功能。经过对比后选择了阿里开源的 FunASR。 本文记录一次完整的 FunASR 部署过程,包括: Docker 部署 FunASR Runtime配置 Paraformer 离线模型配置 Online 实时模型配置 VAD、标点…

作者头像 李华