Celery 是一款功能齐备的任务队列。它易于使用,让你不必先理解它所解决问题的全部复杂性,就能开始工作。它围绕最佳实践设计,使你的产品能够扩展、与其他语言集成,并提供在生产环境运行此类系统所需的工具与支持。
本教程介绍使用 Celery 最基础的知识,包括:
- 选择并安装消息传输组件(broker)。
- 安装 Celery 并创建第一个任务。
- 启动 worker 并调用任务。
- 跟踪任务在不同状态之间的变化,并检查返回值。
Celery 初看可能令人望而生畏,但不用担心,本教程会帮助你快速入门。我们有意保持内容简单,避免高级功能造成困惑。完成后,建议继续浏览其他文档。例如,后续步骤教程将展示 Celery 的能力。
选择消息代理
Celery 需要一个发送与接收消息的解决方案,通常是称为消息代理(message broker)的独立服务。有多种选择:
RabbitMQ
RabbitMQ 功能齐全、稳定、持久可靠,且易于安装,非常适合生产环境。使用 RabbitMQ 与 Celery 的详细说明见使用 RabbitMQ。
Ubuntu 或 Debian 用户可执行:
$ sudo apt-get install rabbitmq-server
若要通过 Docker 运行,则执行:
$ docker run -d -p 5672:5672 rabbitmq
命令完成后,消息代理就已在后台运行,准备为你传递消息:Starting rabbitmq-server: SUCCESS。
如果你不使用 Ubuntu 或 Debian,也不必担心。RabbitMQ 下载页面提供其他平台(包括 Microsoft Windows)的简单安装说明。
Redis
Redis 也功能齐全,但在进程突然终止或断电时,更容易发生数据丢失。详细说明见使用 Redis。如需通过 Docker 运行:
$ docker run -d -p 6379:6379 redis
其他消息代理
除上述选择外,还有其他实验性的传输实现,包括 Amazon SQS。完整列表见消息代理概览。
安装 Celery
Celery 发布于 Python Package Index(PyPI),因此可使用 pip 等标准 Python 工具安装:
$ pip install celery
应用
首先需要一个 Celery 实例,我们称其为 Celery 应用,简称 app。它是一切 Celery 操作的入口,例如创建任务和管理 worker,因此其他模块必须能够导入它。
本教程把所有内容放在一个模块中;较大项目应创建专用模块。创建 tasks.py:
from celery import Celery
app = Celery('tasks', broker='pyamqp://guest@localhost//')
@app.task
def add(x, y):
return x + y
Celery 的第一个参数是当前模块名。只有任务定义在 __main__ 模块中、需要自动生成任务名时,才需要它。
第二个参数是 broker 关键字参数,指定消息代理的 URL。这里使用 RabbitMQ,它也是默认选择。更多选择见上文“选择消息代理”。RabbitMQ 可使用 amqp://localhost;Redis 可使用 redis://localhost。
你已经定义了一个名为 add 的任务,返回两个数的和。
运行 Celery worker 服务
现在可以通过带 worker 参数执行程序来启动 worker:
$ celery -A tasks worker --loglevel=INFO
提示:如果 worker 未能启动,请参阅下面的“故障排查”。
生产环境通常需要让 worker 作为守护进程在后台运行。为此,使用平台提供的工具,或 supervisord 等工具;更多说明见守护进程化。
查看完整命令行选项:
$ celery worker --help
还有其他命令可用,也可查看帮助:
$ celery --help
调用任务
使用 delay() 方法调用任务。它是 apply_async() 的便捷快捷方式;后者提供更细致的执行控制。详见调用任务。
>>> from tasks import add
>>> add.delay(4, 4)
按教程执行后,先前启动的 worker 会处理该任务,你可以查看 worker 控制台输出验证。调用任务会返回 AsyncResult 实例,可用来检查状态、等待完成、获取返回值;任务失败时,也可获取异常和回溯。
默认不启用结果。如果要进行远程过程调用,或在数据库中跟踪任务结果,需要配置结果后端,下一节将介绍。
保存结果
要跟踪任务状态,Celery 必须把状态存储或发送到某个地方。内置结果后端包括 SQLAlchemy/Django ORM、MongoDB、Memcached、Redis、RPC(RabbitMQ/AMQP)等,你也可以自定义。
本例使用 rpc 结果后端,以临时消息形式回传状态。通过 Celery 的 backend 参数指定后端;使用配置模块时,也可通过 result_backend 设置指定。修改 tasks.py 中这一行,启用 rpc:// 后端:
app = Celery('tasks', backend='rpc://', broker='pyamqp://')
如果希望以 Redis 为结果后端,同时以 RabbitMQ 为消息代理,这是一个常见组合:
app = Celery('tasks', backend='redis://localhost', broker='pyamqp://')
更多说明见结果后端。配置好后端后,重启 worker,关闭当前 Python 会话,重新导入 tasks 模块,让变化生效。这次要保留调用任务返回的 AsyncResult 实例:
>>> from tasks import add # close and reopen to get updated 'app'
>>> result = add.delay(4, 4)
ready() 返回任务是否已处理完毕:
>>> result.ready()
False
也可以等待结果完成,但这很少使用,因为它会把异步调用变成同步调用:
>>> result.get(timeout=1)
8
如果任务抛出异常,get() 会重新抛出它;可以通过 propagate 参数改变行为:
>>> result.get(propagate=False)
任务抛出异常时,也可以读取原始回溯:
>>> result.traceback
注意:结果后端需要资源来存储和传递结果。为确保资源释放,调用任务后返回的每一个
AsyncResult实例,最终都必须调用get()或forget()。
完整结果对象参考见 celery.result。
配置
Celery 像家用设备一样,不需要很多配置就能运行。它有输入和输出:输入必须连接消息代理,输出可以选择连接结果后端。不过仔细看背面,你会发现一个盖板,打开后有许多滑块、旋钮和按钮——这就是配置。
默认配置足以满足多数用例,但还有许多选项,可以让 Celery 按照需求工作。了解可用选项有助于熟悉哪些方面可以配置,见配置与默认值参考。
可直接在 app 上设置,也可使用专门的配置模块。例如,修改 task_serializer 设置,配置任务载荷的默认序列化器:
app.conf.task_serializer = 'json'
同时配置多项设置时,可使用 update:
app.conf.update(
task_serializer='json',
accept_content=['json'], # Ignore other content
result_serializer='json',
timezone='Europe/Oslo',
enable_utc=True,
)
大型项目推荐使用专用配置模块。不建议硬编码周期任务间隔和任务路由选项,最好集中存放。对于库尤其如此,因为这允许用户控制任务行为。集中配置也方便系统管理员在系统发生故障时做简单调整。
调用 app.config_from_object(),告知实例使用哪个配置模块:
app.config_from_object('celeryconfig')
模块通常名为 celeryconfig,但可以使用任意模块名。上例要求名为 celeryconfig.py 的模块位于当前目录或 Python 搜索路径,能够被加载。示例如下:
celeryconfig.py
broker_url = 'pyamqp://'
result_backend = 'rpc://'
task_serializer = 'json'
result_serializer = 'json'
accept_content = ['json']
timezone = 'Europe/Oslo'
enable_utc = True
为验证配置文件正常且没有语法错误,可以尝试导入:
$ python -m celeryconfig
完整选项见配置与默认值。下面展示配置文件的作用:将行为异常的任务路由到专用队列:
celeryconfig.py
task_routes = {
'tasks.add': 'low-priority',
}
也可以不改变路由,而对任务限速,让此类型任务每分钟最多处理10个(10/m):
celeryconfig.py
task_annotations = {
'tasks.add': {'rate_limit': '10/m'}
}
以 RabbitMQ 或 Redis 为消息代理时,还可在运行期间让 worker 设置新的任务速率上限:
$ celery -A tasks control rate_limit tasks.add 10/m
worker@example.com: OK
new rate limit set successfully
任务路由见任务路由;注解见 task_annotations;远程控制命令与 worker 监控见监控与管理指南。
接下来读什么
故障排查
常见问题中也有故障排查章节。
worker 无法启动:权限错误
如果使用 Debian、Ubuntu 或其他基于 Debian 的发行版:原文记录 Debian 曾将特殊文件 /dev/shm 重命名为 /run/shm。一种简单的解决方式是创建符号链接:
# ln -s /run/shm /dev/shm
其他情形:如果提供 --pidfile、--logfile 或 --statedb 参数,必须确保它们指向的文件或目录对于启动 worker 的用户可读、可写。
结果后端不工作,或任务一直处于 PENDING 状态
所有任务默认都是 PENDING,因此这个状态叫“未知”或许更合适。发送任务时,Celery 不更新状态;任何没有历史记录的任务都被视为 pending——毕竟你只知道它的任务 ID。
- 确保任务未启用
ignore_result。启用后,worker 会跳过状态更新。 - 确保
task_ignore_result设置未启用。 - 确保没有旧 worker 仍在运行。很容易误启动多个 worker,所以启动新 worker 前应正确关闭旧的。如果旧 worker 使用的结果后端与你预期不同,它可能仍在运行并接走任务。可以将
--pidfile设置为绝对路径,避免这种情况。 - 确保客户端使用正确的后端。如果客户端和 worker 配置了不同后端,你就无法收到结果。检查配置:
>>> result = task.delay()
>>> print(result.backend)
来源:Celery 入门第一步。作者:Ask Solem 与贡献者。文档版本:Celery stable 5.6(网页标题为5.6.3)。页面页脚版权:Copyright © 2009–2023, Ask Solem & contributors。版权页署名为 Copyright © 2009–2016, Ask Solem;保留所有权利。Celery 文档按 Creative Commons Attribution-ShareAlike 4.0 International 许可提供,可分享、改编并用于商业用途,但须署名;修改后的作品仅可按相同或兼容许可分发。本中文改编采用该相同许可。Celery 软件另按 BSD 3-Clause 许可提供。











暂无评论内容