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 INFOstop是异步命令,不会等待 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.logmulti还支持一次启动多个 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为例:
- 名为
proj.app的属性;或 - 名为
proj.celery的属性;或 proj模块中值为 Celery 应用实例的任意属性。
若以上均未找到,则尝试proj.celery子模块:
- 名为
proj.celery.app的属性;或 - 名为
proj.celery.celery的属性;或 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)。
delay、apply_async与直接调用(__call__)三者共同构成 Celery 的 Calling API,signature 也复用这套 API。更详细的讲解见 Calling 用户指南。
任务结果与状态
每次任务调用都会获得一个唯一标识(UUID),即任务 id。delay和apply_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 -> SUCCESSSTARTED是特殊状态,只有启用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:PENDING、RECEIVED、STARTED、SUCCESS、FAILURE、REVOKED、RETRY。更完整的任务状态说明见 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
签名实例同样具备delay和apply_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:带回调的 groupmap/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() 64chain 也可以直接用管道符|书写:
>>> (add.s(4, 4) | mul.s(8))().get() 64Chords:带回调的 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() 90group 链到另一个任务时会自动转换为 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_async的queue参数指定队列:
>>> 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 --helpcelery 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_eventscelery 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),仅供参考