Celery Next Steps 实战指南:从最小示例到任务编排、路由与远程控制 📅 发布时间:2026/9/19 19:43:49 👁 浏览次数: 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[celery5.0]。proj/celery.py创建 app 实例示例模块 examples/next-steps/proj/celery.py 内容如下from celery import Celery app Celery(proj, brokeramqp://, backendrpc://, include[proj.tasks]) # Optional configuration, see the application user guide. app.conf.update( result_expires3600, ) 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_resultTrue)关闭结果存储。各后端的优劣权衡可参考 backends-and-brokers。includeworker 启动时要导入的模块列表。必须把任务模块加进来worker 才能发现并注册这些任务。app.conf.update(result_expires3600)是可选的全局配置调用这里把任务结果的有效期设置为 3600 秒app.start()分支则允许通过python -m proj.celery方式直接以该模块为程序入口启动 worker。proj/tasks.py定义任务示例任务模块 examples/next-steps/proj/tasks.pyfrom .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 和日志信息--------------- celeryhalcyon.local v4.0 (latentcall) --- ***** ----- -- ******* ---- [Configuration] - *** --- * --- . broker: amqp://guestlocalhost: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] celeryhalcyon.local has started.逐项解读 banner 中的关键信息broker即celery模块中broker参数指定的 URL也可用命令行-b选项覆盖。concurrencyprefork 模式下用于并发处理任务的进程数。当所有进程都在忙碌时新任务必须等待。默认值为机器 CPU 数含核心数可用celery worker -c自定义。没有普适的推荐值若任务大多为 I/O 密集型可以尝试调大实验经验表明超过 CPU 数两倍往往无效甚至降低性能。除默认的 prefork 池外Celery 还支持 Eventlet、Gevent 以及单线程模式见 Concurrency。events是否发送 worker 内部动作的监控事件消息供celery events、Flower 等监控程序使用可用-E开启详见 Monitoring and Management。Queuesworker 消费的队列列表。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例如startL411、restartL436、stop/stopwaitL448-L452等方法分别对应各子命令的执行逻辑。关于 --app 参数--app参数指定要使用的 Celery app 实例形式为module.path:attribute但也支持快捷形式只给包名时Celery 会按如下顺序查找 app 实例对应实现见 celery/app/utils.py 的 find_app。以--appproj为例名为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), queuelopri, countdown10)上面的示例会把任务发送到名为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 APIsignature 也复用这套 API。更详细的讲解见 Calling 用户指南。任务结果与状态每次任务调用都会获得一个唯一标识UUID即任务 id。delay和apply_async返回AsyncResult实例可用于跟踪任务执行状态——但前提是配置了结果后端见 task result backends。结果默认禁用因为没有一种后端适合所有应用对很多任务而言保留返回值意义不大因此这是合理的默认值。需要注意结果后端不用于监控任务和 worker监控依靠独立的事件消息见 Monitoring。配置结果后端后可以取回任务返回值 res add.delay(2, 2) res.get(timeout1) 4通过id属性获取任务 id res.id d6b3aea2-fb9b-4ebc-8da4-848818db9114任务抛出异常时可以检查异常与 traceback——默认情况下res.get()会传播任何错误 res add.delay(2, 2) res.get(timeout1)Traceback (most recent call last): File stdin, line 1, in module ... TypeError: unsupported operand type(s) for : int and str如果不想让异常向上传播可传入propagate res.get(propagateFalse) 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_startedTrue)时才会记录。而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.pyPENDING、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), countdown10) 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, debugTrue) s3.delay(debugFalse) # debug is now False.综上签名支持 Calling API 意味着sig.apply_async(args(), kwargs{}, **options)以可选的部分位置参数、部分关键字参数和部分执行选项调用签名sig.delay(*args, **kwargs)apply_async的星号参数版本参数前插到签名已有参数之前关键字参数与已有键合并。原语Primitives这些签名能组合出什么样的工作流这就要引出 Canvas 的六大原语group并行调用一组任务celery/canvas.py#L1540chain串行链接任务celery/canvas.py#L1370chord带回调的 groupmap/starmapcelery/canvas.py#L1458-L1474chunks将参数列表分块处理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带回调的 groupchord 是带回调的 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), queuehipri)然后让 worker 用-Q选项消费该队列$ celery -A proj worker -Q hipri可以用逗号分隔列表指定多个队列。例如让 worker 同时消费默认队列和hipri队列默认队列因历史原因名为celery$ celery -A proj worker -Q hipri,celery队列顺序无关紧要worker 会平等对待所有队列。想充分利用 AMQP 的完整路由能力参见 Routing Guide。远程控制使用 RabbitMQAMQP、Redis 或 Qpid 作为 broker 时可以在运行时控制和检查 worker。例如查看 worker 当前正在处理的任务$ celery -A proj inspect active其实现依赖广播消息机制因此集群中每个worker 都会收到所有远程控制命令。可以用--destination选项指定一个或多个要响应的 worker逗号分隔的主机名列表$ celery -A proj inspect active --destinationceleryexample.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),仅供参考