Celery 入门第一步

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 许可提供。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容