Celery Next Steps 实战指南:从最小示例到任务编排、路由与远程控制
2026/9/19 19:43:45 网站建设 项目流程

Celery Next Steps 实战指南:从最小示例到任务编排、路由与远程控制

【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery

本文是 Celery 官方入门文档的续篇:在 First Steps with Celery 搭建起最小示例之后,本篇将带你系统掌握在生产项目中集成 Celery 的完整路径——从工程化项目布局、启动与管理 worker、--app参数解析规则,到调用 API、任务状态机、基于 signature 的 Canvas 工作流编排、队列路由、远程控制、时区与吞吐量优化。读完本文,你将能独立设计一个结构清晰、可路由、可监控的 Celery 任务系统,并理解每个环节背后的源码级原理。

在应用中使用 Celery

项目布局(Project Layout)

入门指南刻意保持最小化,而在真实应用中,通常需要把 Celery 集成进一个可扩展的包结构。文档推荐的标准布局如下:

src/ proj/__init__.py /celery.py /tasks.py

核心思想是:在proj/celery.py中创建并配置唯一的 Celery 应用(app)实例,然后在proj/tasks.py中定义任务。仓库中的完整可运行示例位于 examples/next-steps/,其中还包含一份 setup.py,演示了如何将任务打包为可分发到 PyPI 或私有包索引的 Python 包(install_requires=['celery>=5.0'])。

proj/celery.py:创建 app 实例

示例模块 examples/next-steps/proj/celery.py 内容如下:

from celery import Celery app = Celery('proj', broker='amqp://', backend='rpc://', include=['proj.tasks']) # Optional configuration, see the application user guide. app.conf.update( result_expires=3600, ) if __name__ == '__main__': app.start()

在这个模块中创建的Celery实例(即所谓的app)是整个项目集成 Celery 的入口——项目内任何地方使用 Celery 时,只需from proj.celery import app。三个核心构造参数含义如下:

  • broker:指定消息代理(broker)的 URL。示例使用amqp://(本地 RabbitMQ 默认地址)。更多选择见 backends-and-brokers。
  • backend:指定结果后端(result backend),用于跟踪任务状态与返回值。Celery 默认禁用结果存储(因为没有一种后端适合所有场景),这里使用rpc://后端只是为了演示如何获取结果。如果业务不需要结果,更明智的做法是保持禁用;也可以为单个任务通过@task(ignore_result=True)关闭结果存储。各后端的优劣权衡可参考 backends-and-brokers。
  • include:worker 启动时要导入的模块列表。必须把任务模块加进来,worker 才能发现并注册这些任务。

app.conf.update(result_expires=3600)是可选的全局配置调用,这里把任务结果的有效期设置为 3600 秒;app.start()分支则允许通过python -m proj.celery方式直接以该模块为程序入口启动 worker。

proj/tasks.py:定义任务

示例任务模块 examples/next-steps/proj/tasks.py:

from .celery import app @app.task def add(x, y): return x + y @app.task def mul(x, y): return x * y @app.task def xsum(numbers): return sum(numbers)

@app.task装饰器将普通函数注册为 Celery 任务,worker 会基于这些注册信息在接收消息时找到对应的执行函数。

启动 Worker

proj的上级目录(按上面的布局即src)中启动 worker:

$ celery -A proj worker -l INFO

启动成功后会出现 banner 和日志信息:

--------------- celery@halcyon.local v4.0 (latentcall) --- ***** ----- -- ******* ---- [Configuration] - *** --- * --- . broker: amqp://guest@localhost:5672// - ** ---------- . app: __main__:0x1012d8590 - ** ---------- . concurrency: 8 (processes) - ** ---------- . events: OFF (enable -E to monitor this worker) - ** ---------- - *** --- * --- [Queues] -- ******* ---- . celery: exchange:celery(direct) binding:celery --- ***** ----- [2012-06-08 16:23:51,078: WARNING/MainProcess] celery@halcyon.local has started.

逐项解读 banner 中的关键信息:

  • broker:即celery模块中broker参数指定的 URL,也可用命令行-b选项覆盖。
  • concurrency:prefork 模式下用于并发处理任务的进程数。当所有进程都在忙碌时,新任务必须等待。默认值为机器 CPU 数(含核心数),可用celery worker -c自定义。没有普适的推荐值,若任务大多为 I/O 密集型可以尝试调大;实验经验表明超过 CPU 数两倍往往无效甚至降低性能。除默认的 prefork 池外,Celery 还支持 Eventlet、Gevent 以及单线程模式(见 Concurrency)。
  • events:是否发送 worker 内部动作的监控事件消息,供celery events、Flower 等监控程序使用(可用-E开启,详见 Monitoring and Management)。
  • Queues:worker 消费的队列列表。worker 可同时消费多个队列,用于消息路由、服务质量(QoS)、关注点分离与优先级控制(详见 Routing Guide)。

通过--help可查看完整命令行参数清单:

$ celery worker --help

更详尽的参数说明见 Workers Guide。

停止 Worker

前台运行时直接按Control-c即可停止。worker 支持的完整信号列表同样见 Workers Guide。

后台运行:celery multi

生产环境需要后台运行 worker,完整方案见 daemonizing 教程。守护脚本底层依赖celery multi命令来启动一个或多个后台 worker:

$ celery multi start w1 -A proj -l INFO celery multi v4.0.0 (latentcall) > Starting nodes... > w1.halcyon.local: OK

重启:

$ celery multi restart w1 -A proj -l INFO celery multi v4.0.0 (latentcall) > Stopping nodes... > w1.halcyon.local: TERM -> 64024 > Waiting for 1 node..... > w1.halcyon.local: OK > Restarting node w1.halcyon.local: OK celery multi v4.0.0 (latentcall) > Stopping nodes... > w1.halcyon.local: TERM -> 64052

停止:

$ celery multi stop w1 -A proj -l INFO

stop是异步命令,不会等待 worker 真正退出;如需确保当前正在执行的任务全部完成后再退出,应使用stopwait

$ celery multi stopwait w1 -A proj -l INFO

注意celery multi不保存任何 worker 信息,重启时必须传入相同的命令行参数;停止时只需保持 pidfile 和 logfile 参数一致即可。

默认情况下 pid 和日志文件会创建在当前目录。为防止多个 worker 相互叠加启动,建议把它们放到专用目录:

$ mkdir -p /var/run/celery $ mkdir -p /var/log/celery $ celery multi start w1 -A proj -l INFO --pidfile=/var/run/celery/%n.pid \ --logfile=/var/log/celery/%n%I.log

multi还支持一次启动多个 worker,并为不同 worker 指定不同参数,例如:

$ celery multi start 10 -A proj -l INFO -Q:1-3 images,video -Q:4,5 data \ -Q default -L:4,5 debug

该命令会启动 10 个 worker:编号 1-3 消费images,video队列且日志级别为默认 INFO,编号 4、5 消费data队列并开启 debug 日志,其余消费默认队列。更多示例可参考 API 参考中的 celery.bin.multi 模块。其底层实现位于 celery/apps/multi.py,例如start(L411)、restart(L436)、stop/stopwait(L448-L452)等方法分别对应各子命令的执行逻辑。

关于 --app 参数

--app参数指定要使用的 Celery app 实例,形式为module.path:attribute,但也支持快捷形式:只给包名时,Celery 会按如下顺序查找 app 实例(对应实现见 celery/app/utils.py 的 find_app)。

--app=proj为例:

  1. 名为proj.app的属性;或
  2. 名为proj.celery的属性;或
  3. proj模块中值为 Celery 应用实例的任意属性。

若以上均未找到,则尝试proj.celery子模块:

  1. 名为proj.celery.app的属性;或
  2. 名为proj.celery.celery的属性;或
  3. proj.celery模块中值为 Celery 应用实例的任意属性。

这套查找机制与文档惯例保持一致:单模块项目用proj:app,较大项目用proj.celery:app

调用任务

使用delay方法即可异步调用任务:

>>> from proj.tasks import add >>> add.delay(2, 2)

delay实际是apply_async的星号参数快捷方式(见 celery/app/task.py 中 Task.delay,其实现为return self.apply_async(args, kwargs)):

>>> add.apply_async((2, 2))

apply_async允许指定执行选项,比如执行时间(countdown)和发送到的队列:

>>> add.apply_async((2, 2), queue='lopri', countdown=10)

上面的示例会把任务发送到名为lopri的队列,并且任务最早在消息发出 10 秒后才执行。apply_async支持的完整选项在 celery/app/task.py 中有详细 docstring,常用的还包括:

  • eta:任务执行的绝对时间(datetime),与countdown二选一;
  • expires:任务过期时间(秒数或datetime),过期后不再执行;
  • priority:任务优先级(0-9),具体语义与 broker 相关;
  • time_limit/soft_time_limit:覆盖默认的(软)时间限制;
  • serializer/compression:消息序列化方式与压缩方式;
  • link/link_error:任务成功/失败后触发的签名(后续 Canvas 部分会用到)。

直接调用任务则会在当前进程内同步执行,不发送任何消息:

>>> add(2, 2) 4

这对应Task.__call__(celery/app/task.py#L506-L513),它会压入请求上下文后直接执行self.run(*args, **kwargs)

delayapply_async与直接调用(__call__)三者共同构成 Celery 的 Calling API,signature 也复用这套 API。更详细的讲解见 Calling 用户指南。

任务结果与状态

每次任务调用都会获得一个唯一标识(UUID),即任务 id。delayapply_async返回AsyncResult实例,可用于跟踪任务执行状态——但前提是配置了结果后端(见 task result backends)。

结果默认禁用,因为没有一种后端适合所有应用;对很多任务而言保留返回值意义不大,因此这是合理的默认值。需要注意:结果后端用于监控任务和 worker,监控依靠独立的事件消息(见 Monitoring)。

配置结果后端后,可以取回任务返回值:

>>> res = add.delay(2, 2) >>> res.get(timeout=1) 4

通过id属性获取任务 id:

>>> res.id d6b3aea2-fb9b-4ebc-8da4-848818db9114

任务抛出异常时,可以检查异常与 traceback——默认情况下res.get()会传播任何错误:

>>> res = add.delay(2, '2') >>> res.get(timeout=1)
Traceback (most recent call last): File "<stdin>", line 1, in <module> ... TypeError: unsupported operand type(s) for +: 'int' and 'str'

如果不想让异常向上传播,可传入propagate

>>> res.get(propagate=False) TypeError("unsupported operand type(s) for +: 'int' and 'str'")

此时返回的是被抛出的异常实例本身,因此判断任务成败需要用结果实例上的对应方法:

>>> res.failed() True >>> res.successful() False

其依据是任务的state(状态):

>>> res.state 'FAILURE'

任务状态机

任务同一时刻只能处于一个状态,但可以沿多个状态推进。典型任务的阶段为:

PENDING -> STARTED -> SUCCESS

STARTED是特殊状态,只有启用task_track_started设置、或为任务设置@task(track_started=True)时才会记录。而PENDING实际上并非真实记录的状态,而是任何未知任务 id 的默认状态:

>>> from proj.celery import app >>> res = app.AsyncResult('this-id-does-not-exist') >>> res.state 'PENDING'

任务被重试时状态流转会更复杂。以重试两次的任务为例:

PENDING -> STARTED -> RETRY -> STARTED -> RETRY -> STARTED -> SUCCESS

所有预定义状态的常量定义在 celery/states.py:PENDINGRECEIVEDSTARTEDSUCCESSFAILUREREVOKEDRETRY。更完整的任务状态说明见 tasks 用户指南中的 Task States 一节。

调用任务的细节可继续阅读 Calling Guide。

Canvas:设计工作流

delay足以应对大部分场景,但有时需要把一次任务调用的 signature 传给另一个进程或作为参数传给其他函数——为此 Celery 引入了signature(签名)。

签名把单次任务调用的参数和执行选项包装起来,使其可以传给函数,甚至可以序列化后通过网络传输。在 celery/canvas.py 中,Signature本质上是dict的子类,docstring 明确指出它"包装单次任务调用的参数和执行选项,用作 group 等结构中的组成部分,或把任务作为回调传递"。

add任务创建参数(2, 2)、countdown 为 10 秒的签名:

>>> add.signature((2, 2), countdown=10) tasks.add(2, 2)

星号参数快捷方式:

>>> add.s(2, 2) tasks.add(2, 2)

签名也支持 Calling API

签名实例同样具备delayapply_async方法,区别在于签名本身可能已指定了参数签名。add接收两个参数,指定两个参数的签名就是完整签名:

>>> s1 = add.s(2, 2) >>> res = s1.delay() >>> res.get() 4

也可以构造不完整的签名,即partials(部分签名):

# incomplete partial: add(?, 2) >>> s2 = add.s(2)

s2现在是一个还缺一个参数的 partial 签名,调用时补足即可:

# resolves the partial: add(8, 2) >>> res = s2.delay(8) >>> res.get() 10

这里传入的 8 被前插到已有参数 2 之前,形成完整的add(8, 2)。关键字参数也可在之后追加,新参数与已有关键字参数合并且新值优先

>>> s3 = add.s(2, 2, debug=True) >>> s3.delay(debug=False) # debug is now False.

综上,签名支持 Calling API 意味着:

  • sig.apply_async(args=(), kwargs={}, **options):以可选的部分位置参数、部分关键字参数和部分执行选项调用签名;
  • sig.delay(*args, **kwargs)apply_async的星号参数版本,参数前插到签名已有参数之前,关键字参数与已有键合并。

原语(Primitives)

这些签名能组合出什么样的工作流?这就要引出 Canvas 的六大原语:

  • group:并行调用一组任务(celery/canvas.py#L1540)
  • chain:串行链接任务(celery/canvas.py#L1370)
  • chord:带回调的 group
  • map/starmap(celery/canvas.py#L1458-L1474)
  • chunks:将参数列表分块处理(celery/canvas.py#L1485)

这些原语本身就是签名对象,因此可以任意组合成复杂的工作流。

注意:以下示例都要获取结果,因此需要先配置结果后端——上文示例项目已通过backend参数完成配置。

Groups:并行调用

group并行调用一组任务,返回一个特殊的结果实例,可整体查看结果并按顺序取回返回值:

>>> from celery import group >>> from proj.tasks import add >>> group(add.s(i, i) for i in range(10))().get() [0, 2, 4, 6, 8, 10, 12, 14, 16, 18]

部分 group(调用时再补参):

>>> g = group(add.s(i) for i in range(10)) >>> g(10).get() [10, 11, 12, 13, 14, 15, 16, 17, 18, 19]
Chains:串行链接

任务可以串联:前一个任务返回后自动调用下一个:

>>> from celery import chain >>> from proj.tasks import add, mul # (4 + 4) * 8 >>> chain(add.s(4, 4) | mul.s(8))().get() 64

部分 chain:

>>> # (? + 4) * 8 >>> g = chain(add.s(4) | mul.s(8)) >>> g(4).get() 64

chain 也可以直接用管道符|书写:

>>> (add.s(4, 4) | mul.s(8))().get() 64
Chords:带回调的 group

chord 是带回调的 group:

>>> from celery import chord >>> from proj.tasks import add, xsum >>> chord((add.s(i, i) for i in range(10)), xsum.s())().get() 90

group 链到另一个任务时会自动转换为 chord:

>>> (group(add.s(i, i) for i in range(10)) | xsum.s())().get() 90

由于所有原语都是签名类型,它们几乎可以任意组合,例如:

>>> upload_document.s(file) | group(apply_filter.s() for filter in filters)

更多工作流设计详见 Canvas 用户指南。

路由

Celery 支持 AMQP 提供的全部路由能力,也支持把消息发往命名队列的简单路由。

:setting:task_routes设置可以按任务名路由,并把路由规则集中在一处:

app.conf.update( task_routes = { 'proj.tasks.add': {'queue': 'hipri'}, }, )

也可以在运行时通过apply_asyncqueue参数指定队列:

>>> from proj.tasks import add >>> add.apply_async((2, 2), queue='hipri')

然后让 worker 用-Q选项消费该队列:

$ celery -A proj worker -Q hipri

可以用逗号分隔列表指定多个队列。例如让 worker 同时消费默认队列和hipri队列(默认队列因历史原因名为celery):

$ celery -A proj worker -Q hipri,celery

队列顺序无关紧要,worker 会平等对待所有队列。想充分利用 AMQP 的完整路由能力,参见 Routing Guide。

远程控制

使用 RabbitMQ(AMQP)、Redis 或 Qpid 作为 broker 时,可以在运行时控制和检查 worker。例如查看 worker 当前正在处理的任务:

$ celery -A proj inspect active

其实现依赖广播消息机制,因此集群中每个worker 都会收到所有远程控制命令。可以用--destination选项指定一个或多个要响应的 worker(逗号分隔的主机名列表):

$ celery -A proj inspect active --destination=celery@example.com

不指定 destination 时,所有 worker 都会执行并回复。

celery inspect命令只返回 worker 内部的信息和统计,不改变任何东西。查看全部 inspect 子命令:

$ celery -A proj inspect --help

celery control命令则包含真正在运行时改变 worker 行为的命令:

$ celery -A proj control --help

例如强制 worker 开启事件消息(用于监控任务和 worker):

$ celery -A proj control enable_events

开启事件后,可以启动事件转储器观察 worker 行为:

$ celery -A proj events --dump

或启动 curses 交互界面:

$ celery -A proj events

监控结束后再关闭事件:

$ celery -A proj control disable_events

celery status同样基于远程控制命令,显示集群中在线的 worker 列表:

$ celery -A proj status

更多命令与监控内容见 Monitoring Guide。

时区

Celery 内部及消息中的所有时间都使用 UTC 时区。worker 收到消息时(例如设置了 countdown 的任务)会把 UTC 时间转换为本地时间。如需使用与系统时区不同的时区,通过timezone设置配置:

app.conf.timezone = 'Europe/London'

优化

默认配置并非为吞吐量而优化:它默认在"大量短任务"与"少量长任务"之间取中间路线,即吞吐量与公平调度之间的折中。如果有严格的公平调度需求或希望针对吞吐量优化,请阅读 Optimizing Guide。

下一步做什么

读完本文后,建议继续阅读 User Guide 全面掌握 Celery 的完整功能与最佳实践;需要查阅 API 细节时,可随时参考 API Reference。本指南对应的可运行示例代码均位于仓库 examples/next-steps/,可对照本文逐步实践。

【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询