Celery 2.2 版本变更全解析:Kombu 替换、task.request 上下文与事件驱动的演进里程碑
【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery
本文基于 Celery 仓库中的官方变更历史文档 docs/history/changelog-2.2.rst,系统梳理 Celery 2.2 系列(2.2.0~2.2.8)的核心变更。作为"分布式任务队列"的经典版本,2.2 完成了消息库从 Carrot 到 Kombu 的迁移、引入task.request上下文、Eventlet 并发池、进程自动伸缩与远程调试等影响深远的能力。阅读本文,你将掌握 2.2 版本引入的关键概念(task.request、celery.task装饰器、BROKER_TRANSPORT_OPTIONS、CELERY_TASK_PUBLISH_RETRY等)的用法与背后的实现原理,并能对照当前仓库源码验证这些特性的真实演化轨迹。
版本总览
Celery 2.2 系列共发布 9 个版本,时间跨度从 2011 年 2 月的 2.2.0 到 2011 年 11 月的 2.2.8,全部由 Ask Solem 发布:
| 版本 | 发布日期 | 定位 |
|---|---|---|
| 2.2.0 | 2011-02-01 | 里程碑版本:Kombu 替换 Carrot、task.request上下文、Eventlet 支持 |
| 2.2.1 | 2011-02-02 | 修复 Eventlet 内存泄漏等问题 |
| 2.2.2 | 2011-02-03 | 修复celerybeat调度读取、eta与retry组合等问题 |
| 2.2.3 | 2011-02-12 | 修复路由回归、日志格式化等问题 |
| 2.2.4 | 2011-02-19 | 修复结果轮询、SQLAlchemy 后端date_done回归 |
| 2.2.5 | 2011-03-28 | 日志轮转(WatchedFileHandler)、BROKER_TRANSPORT_OPTIONS等新特性 |
| 2.2.6 | 2011-04-15 | 依赖 Kombu 1.1.2、修复 Python 2.5 兼容 |
| 2.2.7 | 2011-06-13 | 新增日志信号、Redis 后端兼容 2.4.4 |
| 2.2.8 | 2011-11-25 | 安全修复 CELERYSA-0001(--uid/--gid权限问题) |
2.2.0:架构级的三大里程碑
1. Carrot 正式退役,Kombu 成为新的消息层
2.2.0 最重要的事件是Carrot 被 Kombu 彻底替换。文档明确说明,Kombu 是 Python 的"下一代消息库",修复了 Carrot 中难以在不破坏向后兼容的前提下修复的多个缺陷,并带来:
- 虚拟传输(virtual transports)的一流支持:Redis、Django ORM、SQLAlchemy、Beanstalk、MongoDB、CouchDB 以及内存传输;
- 一致的错误处理与内省能力;
- 优雅处理连接与信道错误,确保操作可靠执行;
- 消息压缩(
:mod:zlib、:mod:bz2或自定义压缩方案)。
这意味着ghettoq不再需要——其功能已默认内置。虚拟传输对 exchange(direct 与 topic)的支持更加完善,其中 Redis 传输甚至支持 fanout exchange,从而能够执行 worker 远程控制命令。
这一架构变更延续至今:当前仓库中消息发布逻辑仍在 celery/app/amqp.py 中通过 Kombu 完成,例如任务发布的重试策略就读取task_publish_retry与task_publish_retry_policy配置(见 celery/app/amqp.py)。
2. 魔法关键字参数(Magic Keyword Arguments)进入弃用轨道
旧式任务会接收task_id、delivery_info、task_retries等隐式关键字参数,文档指出这些"魔法关键字参数"导致了大量问题:与装饰器配合时的怪异行为、以及使用者无意识下的关键字命名冲突。
弃用路径分为三步:
celery.decorators模块被弃用,装饰器迁移至celery.task;celery.task中的装饰器默认禁用关键字参数;- 文档中的全部示例改用
celery.task。
旧风格(启用魔法关键字参数):
from celery.decorators import task @task() def add(x, y, **kwargs): print('In task %s' % kwargs['task_id']) return x + y新风格(禁用魔法关键字参数):
from celery.task import task @task() def add(x, y): print('In task %s' % add.request.id) return x + y任务还可以通过设置task.accept_magic_kwargs属性决定是否接受魔法关键字参数。弃用时间线为:2.2 发出PendingDeprecationWarning,2.4 升级为DeprecationWarning,4.0 移除celery.decorators模块且accept_magic_kwargs不再生效。在当代代码中,该属性仍保留在 celery/app/utils.py 的应用构造参数列表中,见证了这一演化。
3. task.request:线程本地的请求上下文
魔法关键字参数的替代方案是task.request上下文。它基于线程本地存储(thread-local storage)保存与当前请求相关的状态,可变且可扩展——你可以向其中添加自定义属性,且只会被当前任务请求看到。
文档给出的映射表如下(即"旧参数 → 新写法"):
| 魔法关键字参数 | 替换为 |
|---|---|
kwargs['task_id'] | self.request.id |
kwargs['delivery_info'] | self.request.delivery_info |
kwargs['task_retries'] | self.request.retries |
kwargs['logfile'] | self.request.logfile |
kwargs['loglevel'] | self.request.loglevel |
kwargs['task_is_eager'] | self.request.is_eager |
| 新增 | self.request.args |
| 新增 | self.request.kwargs |
以下方法会自动使用当前上下文,无需再手动传递kwargs:
task.retrytask.get_loggertask.update_state
2.2.5 还补充了"Task.request 上下文现在总是被初始化,确保直接调用任务函数也能正常工作",2.2.4 则新增了request.taskset(当前 taskset id)。时至今日,task.request仍是 Celery 任务上下文的标准入口。
2.2.0 的并发、调试与运维新能力
Eventlet 并发池:I/O 密集型任务的福音
2.2.0 引入Eventlet 支持,并明确说明"这对 I/O 密集型任务是好消息"。切换池实现有两种方式:
$ celery worker --pool=eventlet或全局设置CELERYD_POOL,值可以是类的完整名称,也可以是别名:processes、eventlet、gevent。
当时的 gevent 池仍属实验性质,缺少 ETA 任务调度能力(该能力在 2.2.5 中补上,但要求CELERY_DISABLE_RATE_LIMITS=True)。与事件池生命周期相关的信号eventlet_pool_started、eventlet_pool_preshutdown、eventlet_pool_postshutdown、eventlet_pool_apply至今仍定义在 celery/signals.py,并在 celery/concurrency/eventlet.py 的池启动与关闭流程中实际发送。
Worker 进程自动伸缩(Autoscaling)
:option:--autoscale <celery worker --autoscale>选项可配置子进程的上下限:
--autoscale=AUTOSCALE Enable autoscaling by providing max_concurrency,min_concurrency. Example: --autoscale=10,3 (always keep 3 processes, but grow to 10 if necessary).即--autoscale=10,3表示始终保留 3 个进程,必要时增长到 10 个。当前仓库的 celery/worker/autoscale.py 中Autoscaler类(celery/worker/autoscale.py)正是这一特性的实现,支持通过update(max=, min=)动态调整上下限。
任务远程调试:celery.contrib.rdb
celery.contrib.rdb是pdb的扩展,让没有终端访问权限的进程也能远程调试:
from celery.contrib import rdb from celery.task import task @task() def add(x, y): result = x + y # 设置断点 rdb.set_trace() return resultset_trace在当前代码位置设置断点并创建一个可 telnet 的 socket。多个进程可同时启动调试器,因此端口不是固定的:调试器从基础端口(默认 6900)开始搜索可用端口,基础端口可通过环境变量CELERY_RDB_PORT修改;默认只允许本机访问,如需外部访问需设置CELERY_RDB_HOST。
worker 遇到断点时会输出:
[INFO/MainProcess] Received task: tasks.add[d7261c71-4962-47e5-b342-2448bedd20e8] [WARNING/PoolWorker-1] Remote Debugger:6900: Please telnet 127.0.0.1 6900. Type `exit` in session to continue. [2011-01-18 14:25:44,119: WARNING/PoolWorker-1] Remote Debugger:6900: Waiting for client...telnet 进入后即获得pdb交互 shell:
$ telnet localhost 6900 Connected to localhost. Escape character is '^]'. > /opt/devel/demoapp/tasks.py(128)add() -> return result (Pdb)输入help可查看命令列表。这一实现至今仍保留在 celery/contrib/rdb.py 中:Rdb(Pdb)类负责建立 socket 并接受客户端连接(celery/contrib/rdb.py),端口搜索逻辑会结合 worker 进程名中的编号做端口偏移(celery/contrib/rdb.py),使得同一主机上多个 worker 进程的调试端口互不冲突;set_trace则在 celery/contrib/rdb.py 定义。
事件体系重构:topic 交换与多监视器
事件从 direct 交换改为瞬时(transient)topic 交换。这意味着:
- 事件在有消费者之前不会被存储,消费者停止后事件立即消失;
- 可以同时运行多个监视器;
- 事件的路由键就是事件类型(如
worker.started、worker.heartbeat、task.succeeded),消费者可按类型过滤; - 每个消费者创建唯一队列,效果上相当于广播交换。
由此带来新的可能性:worker 可以监听其他 worker 的事件来感知"邻居",甚至在它们宕机时重启它们(或用于任务/自动伸缩优化)。
注意:事件交换名从"celeryevent"改名为"celeryev"以避免与旧版本冲突。如需删除旧交换,可执行:
$ camqadm exchange.delete celeryeventCELERYD_EVENT_EXCHANGE、CELERYD_EVENT_ROUTING_KEY、CELERYD_EVENT_EXCHANGE_TYPE设置不再使用。同时所有 worker 事件新增三个字段:sw_ident(worker 软件名,如"py-celery")、sw_ver(软件版本)、sw_sys(操作系统)。新增CELERY_SEND_TASK_SENT_EVENT设置,开启后每条任务都会发送事件,使监视器能在 worker 接收任务之前就跟踪到任务。task-started事件现在还会携带接受任务的子进程 PID。
2.2.0 的配置与 CLI 演进
Worker 无配置启动与命令行内联配置
worker 现在不需要配置文件即可启动,配置可以直接写在命令行上,位于最后一个参数之后、以两个短横线分隔:
$ celery worker -l info -I tasks -- broker.host=localhost broker.vhost=/app同时配置对象现在是原始配置的别名,运行时对原始配置的修改会即时反映到 Celery。celery.conf被弃用,修改celery.conf.ALWAYS_EAGER不再生效;默认配置移入celery.app.defaults模块,所有配置项及其类型均可内省。配置文件与加载器也可在命令行指定:
$ celery worker --config=celeryconfig.py --loader=myloader.Loader消息发布重试与消息压缩
新增任务消息发布重试能力,应对连接丢失或失败场景。默认关闭,可通过CELERY_TASK_PUBLISH_RETRY启用,并用CELERY_TASK_PUBLISH_RETRY_POLICY微调策略。Task.apply_async同时新增retry与retry_policy关键字参数(使用retry参数需要手动管理 publisher/连接)。该配置延续至今,当前默认值已改为开启,定义在 celery/app/defaults.py。
新增消息压缩支持:通过CELERY_MESSAGE_COMPRESSION设置或apply_async的compression参数(也可由路由器设置)。
远程终止任务:revoke + terminate
:control:revoke远程控制命令现在支持terminate参数,可远程终止正在处理任务的 worker 进程。默认信号为TERM,可通过signal参数指定(signal模块中任何信号的英文大写名称)。终止任务同时会撤销它:
>>> from celery.task.control import revoke >>> revoke(task_id, terminate=True) >>> revoke(task_id, terminate=True, signal='KILL') >>> revoke(task_id, terminate=True, signal='SIGKILL')结果等待与轮询增强
TaskSetResult.join_native:后端优化版join(),利用后端批量获取多个结果的能力(当时仅 AMQP 后端支持,Memcached 与 Redis 计划后续支持);TaskSetResult.join与AsyncResult.wait改进:两者新增interval关键字参数(默认 0.5 秒)控制轮询间隔;result.wait()新增propagate参数,设为False时错误以返回值形式返回而非抛出;- 文档特别警告:使用数据库结果后端时应降低轮询频率,频繁轮询会导致数据库高负载。
其他值得关注的变更
- 远程控制命令:新增
active_queues(返回 worker 当前消费的队列声明); celery multi与守护化:worker 的内建守护化支持(celery multi)不再视为实验特性,进入生产可用阶段;celerybeat与celeryev新增--detach守护化选项;- 信号:新增
beat_init(celerybeat启动时派发,发送者为celery.beat.Service实例)与beat_embedded_init(内嵌启动celerybeat时额外派发),两者仍定义于 celery/signals.py; - SIGUSR1 线程栈转储:worker 收到
SIGUSR1时记录所有线程的堆栈(CPython 2.4、Windows 或 Jython 上不可用); - 安全:移除
celery.task.RemoteExecuteTask及dmap、dmap_async、execute_remote(低危:用 pickle 执行任意代码在消息代理被非法访问时是潜在安全隐患);stats命令不再传输 broker 密码(低危); - 模块重构:
celery.worker.listener更名为celery.worker.consumer,CarrotListener更名为Consumer;celery.task.schedules弃用,改用celery.schedules;celery.execute.apply_async/apply/delay_task弃用;远程控制命令改由kombu.pidbox(通用进程邮箱)提供; - 测试覆盖:新增大量单元测试,总覆盖率 95%。
2.2.1~2.2.4:修复潮
这几个版本以修复为主,重点如下:
- 2.2.1:修复 Eventlet 池内存泄漏(Issue #308);恢复被意外移除的弃用函数
celery.execute.delay_task;celeryd_detach对不存在的用户/组名给出可读错误;日志错误时对 unicode 解码错误的更智能处理。 - 2.2.2:修复
celerybeat无法正确读取调度导致CELERYBEAT_SCHEDULE条目不被调度;eta参数现在可与task.retry一起使用(此前会被countdown覆盖);错误日志重新包含exc_info;守护化教程修正--time-limit 300→--time-limit=300笔误。 - 2.2.3:
Task.retry支持max_retries参数覆盖默认值;修复Task.exchange与Task.routing_key不再生效的回归;multiprocessing.cpu_count在不支持平台可能抛NotImplementedError的处理;远程控制命令active_queues现在能反映运行时新增的队列,且返回的 exchange 键改为完整的 exchange 声明字典;修复celery worker -Q删除未用队列声明导致路由失败的问题——队列不再被移除,而是通过app.amqp.queues.consume_from()作为消费列表;celeryctl支持inspect active_queues。 - 2.2.4:修复 2.2.3 破坏的错误日志(traceback 不再被记录);AMQP 结果后端在队列中存在多条结果消息时轮询失败的问题;
TaskSet.apply_async()/TaskSet.apply()支持taskset_id关键字参数(Issue #331),当前 taskset id 可通过request.taskset获取(Issue #329);SQLAlchemy 结果后端恢复被意外移除的date_done(Issue #325),并为Task.id与TaskSet.taskset_id增加唯一约束(需要重建表)。
2.2.5:日志轮转与传输选项
2.2.5 在修复之外带来几项实用新特性:
- 日志轮转:通过
WatchedFileHandler支持外部工具(如logrotate.d)轮转日志(Issue #321)。其原理是文件被重命名或删除后重新打开文件——这正是当前仓库 celery/utils/log.py 中日志处理链路的处理方式。 BROKER_TRANSPORT_OPTIONS:新设置,用于向特定 broker 传输传递额外参数。- gevent 支持 ETA 任务(仍需
CELERY_DISABLE_RATE_LIMITS=True)。 - Eventlet 新增四个信号:
eventlet_pool_started、eventlet_pool_preshutdown、eventlet_pool_postshutdown、eventlet_pool_apply。 TaskSet.apply/TaskSet.apply_async接受可选的taskset_id参数;taskset_id 进入 Task 请求上下文;SQLAlchemy 结果后端为 taskset_id 与 task_id 增加唯一约束(需重建表)。- worker 广播命令返回的请求信息中包含
worker_pid;Task.after_return现在总是在结果写入之后调用;移除了未使用的AsyncResult.uuid属性。 - 性能优化:速率限制在无任务时不再 sleep,而是等待"任务已接收"条件变量;multiprocessing.Pool 在标记
WorkerLostError前等待 10 秒(给结果处理器机会取回已发布的结果),关闭时 ResultHandler 在 5 秒后超时退出;prefetch count 超过 short 上限 65535 时暂时禁用、低于上限后重新启用(Issue #359)。 - 内部模块
celery.worker.controllers更名为celery.worker.mediator;celery.contrib.batches为批处理任务设置 loglevel/logfile 使task.get_logger可用(Issue #357);cursesmon修复未绑定局部变量错误(Issue #303);app.config_from_object/config_from_envvar对所有 loader 生效;结果后端名称未知时给出用户友好错误(Issue #349);Cassandra 结果后端适配最新pycassa。
2.2.6~2.2.8:兼容性与安全收尾
- 2.2.6:依赖 Kombu 1.1.2;明确排除
python-dateutil2.x(仅支持 Python 3),误装者可降级:pip install -U python-dateutil==1.5.0(或easy_install -U python-dateutil==1.5.0)。修复:WatchedFileHandler破坏 Python 2.5 支持(Issue #367);任务显式设置名称时不再使用app.main;Python 2.5 下邮件发送因版本检测 bug 失效(Issue #378);Beat.ScheduleEntry新增可覆写的_default_now方法以改变last_run_at默认值;进程清理中的错误不再传播(改为记录日志),避免干扰任务结果发布(Issue #365);Djangoshell_plus下任务定义失效(Issue #366);AsyncResult.get恢复接受interval与propagate参数;修复 worker 在socket.error时不退出的 bug。 - 2.2.7:新增信号
after_setup_logger与after_setup_task_logger,可在 Celery 完成日志配置后增强日志配置(当前定义于 celery/signals.py,携带logger、loglevel、logfile、format、colorize参数);Redis 结果后端兼容 Redis 2.4.4;multi的--gid选项现在正确生效;worker 重试时误用 traceback 的 repr 而非字符串表示;App.config_from_object现在加载模块而非模块的属性;修复对象日志输出<Unrepresentable: ...>的问题。 - 2.2.8(安全修复):CELERYSA-0001——
celery multi、celeryd_detach、celery beat、celery events使用--uid/--gid参数时,守护进程设置的是有效 ID(effective id)而非真实 ID(real id),导致权限没有真正降级,之后可能重新获取超级用户权限。完整的漏洞说明见仓库内的 docs/sec/CELERYSA-0001.txt。
从 2.2 看 Celery 的演进主线
对照当前仓库源码,2.2 时代确立的许多设计至今仍是 Celery 的骨架:
- Kombu 作为统一消息层:虚拟传输、连接优雅重试、消息压缩的能力延续至今;
task.request上下文:从替代魔法关键字参数出发,演化为如今任务执行时携带id、args、kwargs、retries、delivery_info等信息的标准上下文对象;- 信号体系:
after_setup_logger、beat_init、eventlet_pool_*等信号至今保留在 celery/signals.py 中并继续被 celery/concurrency/eventlet.py 等模块使用; - 配置内省与命令行内联配置:
celery.app.defaults中的配置声明方式成为后续所有配置项的基础; - 自动伸缩、远程调试、远程终止任务:分别对应今天 celery/worker/autoscale.py、celery/contrib/rdb.py 与
revoke(terminate=True)的成熟实现。
对于希望理解 Celery 设计脉络的读者,2.2 版本变更史是一份浓缩的架构决策记录:它回答了"为什么任务不再使用魔法参数"、"为什么选择 Kombu"、"事件系统为何采用 topic 交换"等根本性问题。若需完整了解每个版本的条目细节,可直接阅读仓库中的 docs/history/changelog-2.2.rst 原文,以及相邻版本的 docs/history/changelog-2.1.rst 与 docs/history/changelog-2.3.rst 以对比演进。
【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考