Celery 4.1 系列版本变更深度解读:修复项、新特性与升级注意事项
【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery
本文以 Celery 官方 4.1.x 系列变更记录(docs/history/changelog-4.1.rst)为核心骨架,结合当前仓库源码逐项剖析 4.1.0 与 4.1.1 两个版本中的关键修复、行为变化与新增能力,帮助使用 4.1 系列或计划从 3.x/4.0 升级的开发者快速定位需要关注的点,并给出可落地的升级与验证建议。读完本文,你将掌握 4.1 系列中 Worker 信号、Redis 后端 SSL、任务同步子任务限制、Beat 调度堆更新等核心机制的实现位置与使用方式。
一、版本总览:4.1 系列改了什么
Celery 4.1 系列包含两个版本:
- 4.1.0(2017-07-25 发布):本系列的功能与修复主版本;
- 4.1.1(2018-05-21 发布):紧急修复版本,涉及 Kombu 模块更名带来的兼容性问题,官方在变更记录中明确提示"请尽快升级或将 Kombu 固定到 4.1.0"。
4.1 系列最重要的架构前提是依赖Kombu 4.1.0(见 requirements/default.txt 中对应依赖声明),这是 4.0 系列以来消息层重构的延续。在 4.1 中,Kombu 将原本名为async的模块更名为asynchronous,这是 4.1.1 中被标记为Breaking Change的变更——若你的代码直接from kombu.async import ...,在 Kombu 4.1.x 后续版本中会失败,因此官方强烈建议升级或固定版本。
从变更记录看,4.1 系列的改动集中在以下几个域:
| 领域 | 典型变更 |
|---|---|
| App / 配置 | 恢复CELERY_SEND_EVENTS兼容、修复 ETA/重试/过期任务、Broadcast 队列恢复 |
| Worker | 新增worker_shutting_down信号、gevent 嵌入场景关闭修复、启动时不再过早消费队列 |
| Canvas | 修复链中替换任务、chord 空 header、链去重、链内任务反序列化 |
| Results | Redis 后端 SSL、Elasticsearch 字段复用与键序列化、MongoDB 二进制编码、DynamoDB 后端 |
| Beat | crontab 反序列化修复、调度堆更新机制、pidfile默认值 |
| Task | 类方法定义任务、disable_sync_subtasks、Python 3.6 支持、kwargs 处理 |
| Utils | maybe_make_aware、任务参数处理修复 |
二、Worker 生命周期:新增worker_shutting_down信号
4.1 为 Worker 新增了worker_shutting_down信号(Issue #3998),用于在 Worker 收到关闭请求(软关闭Warm或强制关闭Cold)时、正式停止消费之前发出通知。它比既有的worker_shutdown更早触发,适合做优雅退出前的资源回收或状态上报。
信号定义位于 celery/signals.py:
worker_shutting_down = Signal(name='worker_shutting_down')发送点在 celery/apps/worker.py 的关闭信号处理器中,仅在主进程(MainProcess)触发,并携带关闭方式与退出码:
def _handle_request(*args): with in_sighandler(): from celery.worker import state if current_process()._name == 'MainProcess': if callback: callback(worker) if verbose: safe_say(f'worker: {how} shutdown (MainProcess)', sys.__stdout__) signals.worker_shutting_down.send( sender=worker.hostname, sig=sig, how=how, exitcode=exitcode, ) setattr(state, {'Warm': 'should_stop', 'Cold': 'should_terminate'}[how], exitcode)可以看到该信号携带sig(触发信号)、how(Warm/Cold)、exitcode三个参数。使用方式:
from celery.signals import worker_shutting_down @worker_shutting_down.connect def on_worker_shutting_down(sender, sig, how, exitcode, **kwargs): # 在此处执行资源清理、优雅下线注册等逻辑 print(f"worker {sender} shutting down via {how}")同样在本系列中,App 层还修复了before_task_publish信号不携带properties的问题(Issue #4035),发布时业务自定义的属性现在可以被该信号读取到。与之配套,celery/signals.py 中before_task_publish的providing_args明确包含body, exchange, routing_key, headers, properties, declare, retry_policy。
三、任务执行修复:ETA、重试、过期与同步子任务
4.1 修复了一批任务调度与执行层面的问题,其中几个与线上稳定性强相关:
3.1 ETA 与时区
- 修复配置定义时区时的任务 ETA 问题(Issue #3867 / #3753):当
timezone在配置中定义后,带 ETA 的任务可能被错误调度; - 修复重试任务带过期时间(expiration)的问题(Issue #3790 / #3734):重试时任务过期时间被错误计算。
这两类问题都与时间处理相关,建议在 4.1 上重新跑一遍带eta、expires和重试策略的用例。
3.2disable_sync_subtasks:默认禁止任务内同步等待子任务
这是 4.1 引入的行为约束:默认禁止任务内部再同步调用并等待子任务结果,以避免任务池被自我阻塞耗尽(经典死锁场景:同一 Worker 内任务 A 同步等待任务 B,而 B 排在 A 之后执行)。
实现位于 celery/result.py,AsyncResult.get()在disable_sync_subtasks=True(默认值)时调用assert_will_not_block():
if disable_sync_subtasks: assert_will_not_block()其文档明确警告:"Disable tasks to wait for sub tasks, this is the default configuration. CAUTION do not enable this unless you must."(默认禁止,除非必要否则不要开启)。assert_will_not_block通过判断调用线程是否为 Worker 任务执行线程来抛出RuntimeError,防止死锁。若你的业务确需同步等待,可以显式传入disable_sync_subtasks=False,但必须自行评估阻塞风险。
同时,本系列还修复了同步apply阻塞执行时请求上下文缺少hostname的问题(Issue #3716 / #3735),并新增了hostname属性,保证阻塞执行模式下日志与请求上下文完整。
3.3 类方法定义任务
4.1 允许类方法(classmethod)定义任务(Issue #3952 / #3863)。在此之前任务多定义在模块级函数或类实例上;现在形如:
class MyClass: @classmethod def process(cls, x): # 需要访问 cls 的任务逻辑 ...也可以通过app.task包装。相关实现见 celery/app/task.py 中任务注册与run解析逻辑,测试用例可参考 t/unit/tasks/test_tasks.py。
3.4 Python 3.6 与 kwargs 处理
- 正式支持 Python 3.6(Issue #3904 等),并修复 Python 3 下任务带关键字参数时的处理问题(Issue #3657 / #3678);
- 协议兼容:Task 协议 1 中缺失的
*args/**kwargs在协议 2 中返回空值(Issue #3687),避免旧消息反序列化时报错。
四、Canvas 工作流修复:链、组、chord 与替换任务
Canvas(任务工作流编排)在 4.1 中获得大量修复,涉及chain、group、chord与replace的组合场景:
- 替换任务后的链顺序(Issue #3730):
replace后链中后续任务顺序被纠正; - 任务被替换为 group 后不完成(Issue #3725 / #3731):group 替换场景下链的完成状态修复;
- chord 空 header 抛
IndexError(Issue #3847):修复 chord 头部为空列表时的崩溃;Lookup task only if list has items即对应此修复; - chord 中链去重(Issue #3771 / #3779):避免同一链在 chord 中被重复展开执行;
- 链中任务全部反序列化(Issue #4015):链在发送前对其包含的所有任务做反序列化校验,尽早暴露序列化错误。
这些修复说明 4.1 重点巩固了"链 + 组 + chord + replace"组合工作流的正确性。若你在用复杂 canvas 结构,建议对照 t/unit/tasks/test_canvas.py 与 t/unit/test_canvas.py 中的用例回归验证。
五、结果后端(Results)增强:Redis SSL、Elasticsearch、DynamoDB
5.1 Redis 后端 SSL 支持
4.1 为 Redis 结果后端新增SSL 选项(Issue #3830 / #3831)。核心实现在 celery/backends/redis.py:redis_backend_use_ssl必须是一个包含ssl_cert_reqs、ssl_ca_certs、ssl_certfile、ssl_keyfile键的字典(与 broker 的broker_use_ssl一致):
ssl = _get('redis_backend_use_ssl') if ssl: self.connparams.update(ssl) self.connparams['connection_class'] = self.connection_class_sslssl_cert_reqs支持CERT_REQUIRED、CERT_OPTIONAL、CERT_NONE,字符串值会被转换为对应的ssl常量并做合法性校验(celery/backends/redis.py)。配置示例:
# celery 配置 result_backend = 'redis://redis.example.com:6379/0' redis_backend_use_ssl = { 'ssl_cert_reqs': 'CERT_REQUIRED', 'ssl_ca_certs': '/path/to/ca.crt', 'ssl_certfile': '/path/to/client.crt', 'ssl_keyfile': '/path/to/client.key', }同时支持在 URL 查询串中携带ssl_cert_reqs等参数(celery/backends/redis.py),并使用rediss://协议;源码注释特别强调:"A rediss:// URL must have parameter ssl_cert_reqs and this must be set to something valid",且CERT_NONE时 Celery 会给出安全提示(模块头注释),实际生产请务必配置为CERT_REQUIRED。
5.2 Elasticsearch 后端多项修复
- 键序列化修复(Issue #3924):结果键序列化方式修正;
- 文档 ID 序列化修复:确保文档 ID 类型稳定;
- 字段复用(Issue #3708):同一任务的多次结果写入不再每次生成新字段,而是复用既有字段,避免索引膨胀;
- 后端选项设置支持(Issue #3736 关联):可自定义 Elasticsearch 后端连接选项。
实现见 celery/backends/elasticsearch.py。
5.3 MongoDB 与 DynamoDB
- MongoDB:修复二进制编码(binary encodings)场景下的集成问题(Issue #3575),见 celery/backends/mongodb.py;
- DynamoDB:新增 AWS DynamoDB 作为结果后端的能力(Issue #3736),相关文档见 docs/internals/reference/celery.backends.dynamodb.rst,实现见 celery/backends/dynamodb.py。
5.4 其他后端修复
- Unicode 异常消息(Issue #3858 / #3903):任务抛出的异常允许包含 Unicode 消息;
- Flower REST API 兼容:修复 Celery 在使用 Flower REST API 时的事件状态问题,保证
Task.as_dict()在信息不完整时也能工作(相关事件状态逻辑见 celery/events/state.py)。
六、Beat 调度器:crontab 反序列化与堆更新机制
4.1 对 Beat 调度器做了两处值得关注的增强:
6.1 crontab 的 pickle 还原修复
修复了 pickledcrontab调度在反序列化后无法正确恢复的问题(Issue #3826 / #3827):celery.schedule.crontab的__reduce__被修正,保证跨进程/持久化场景下调度表还原正确。实现见 celery/schedules.py 中crontab类的 reduce 相关方法。
6.2 调度堆(heap)透明更新
新增透明的调度堆更新方法(Issue #3721):当周期任务在运行时被动态增删改时,Beat 的调度堆能随之刷新,而不是停留在启动时的快照。核心逻辑在 celery/beat.py:
def populate_heap(self, event_t=event_t, heapify=heapq.heapify): """Populate the heap with the data contained in the schedule.""" priority = 5 self._heap = [] for entry in self.schedule.values(): is_due, next_call_delay = entry.is_due() self._heap.append(event_t( self._when(entry, 0 if is_due else next_call_delay) or 0, priority, entry )) heapify(self._heap)而调度器在tick()中检测到调度表变化(schedules_equal比较新旧调度)时自动重新populate_heap()(celery/beat.py),这正是"Populate heap when periodic tasks are changed"的实现。同时 4.1 还为celery beat的--pidfile指定了默认值(Issue #3722),避免未指定时产生歧义。另外,Scheduler.schedule现在返回调度字典的浅拷贝,防止外部误修改污染内部状态。
七、配置兼容:CELERY_SEND_EVENTS与 Broadcast 队列
CELERY_SEND_EVENTS兼容恢复(Issue #3997):3.1.x 用户习惯使用CELERY_SEND_EVENTS控制事件发送,4.0 中该配置曾被改为CELERYD_SEND_EVENTS;4.1 恢复对CELERY_SEND_EVENTS的支持,降低升级迁移成本;- Broadcast 队列恢复(Issue #3934):恢复 Broadcast(广播)队列的行为,使
celery control broadcast类操作(如远程shutdown、rate_limit广播)在 4.1 中正常工作; - 日志 Formatter 增强(Issue #3994):使
id、name总能通过extra从logging.Formatter访问,便于自定义日志格式时输出任务 ID 与任务名。
八、Platforms、系统集成与周边修复
8.1 信号支持检测返回布尔值
Platforms模块中检查某个信号是否被支持的方法(如SIGKILL在 Windows 上的可用性检测)现在始终返回布尔值(Issue #3962),避免返回None造成的隐式真值误判,相关实现见 celery/platforms.py。
8.2 Systemd 配置 loglevel 恢复
修复 systemd 配置中ExecStart丢失loglevel的问题(Issue #4023)。仓库中 extra/systemd/celery.service 的启动命令如下:
ExecStart=/bin/sh -c '${CELERY_BIN} -A $CELERY_APP multi start $CELERYD_NODES \ --pidfile="${CELERYD_PID_FILE}" \ --logfile="${CELERYD_LOG_FILE}" \ --loglevel="${CELERYD_LOG_LEVEL}" $CELERYD_OPTS'--loglevel="${CELERYD_LOG_LEVEL}"会从环境变量恢复日志级别,配套的环境变量定义见 extra/systemd/celery.conf。
8.3 gevent 嵌入场景的关闭修复
修复 Consumer 在嵌入 gevent 应用时无法正确关闭的问题(Issue #3745 / #3746):嵌入模式下 Worker 关闭流程与 gevent 的事件循环配合被修正。
8.4 启动时不再过早消费队列
修复 Worker 启动时在就绪(ready)之前就开始消费队列的问题(Issue #3620):确保所有启动步骤(包括连接建立、初始化)完成后再开始拉取消息,避免启动阶段丢消息。
8.5maybe_make_aware修复
utils.time.maybe_make_aware在传入的 datetime 已是 aware(带时区)时不再重复修改(Issue #3849 / #3850),避免时区被意外覆盖。实现见 celery/utils/time.py。
九、升级到 4.1 的实操建议
基于上述变更,从 3.x/4.0 升级到 4.1 时建议按以下顺序处理:
- 固定依赖版本:若暂不升级到 4.1.1,请务必将 Kombu 固定到
4.1.0(4.1.1 的 Breaking Change 针对 Kombu 模块更名);若直接采用 4.1.1,确保业务代码不直接引用kombu.async模块名(新名为kombu.asynchronous); - 回归测试 ETA/过期/重试:重点覆盖带
eta、expires的任务与autoretry_for场景; - 检查同步子任务调用:代码中
result.get()/result.wait()若出现在任务函数内部,需要显式传disable_sync_subtasks=False或重构为异步回调/链式任务; - 验证 canvas 组合:跑一遍
chain、group、chord、replace的组合用例(参考 t/unit/tasks/test_canvas.py); - 安全配置 Redis SSL:若 Redis 后端走公网/跨网段,配置
redis_backend_use_ssl并将ssl_cert_reqs设为CERT_REQUIRED; - Python 版本:4.1 支持到 Python 3.6,低于 3.6 的版本建议评估是否继续使用 4.0 或升级 Python;
- 利用新信号:需要优雅下线逻辑的团队,可在
worker_shutting_down中注册清理动作。
完整版本差异还可对照 docs/history/changelog-4.0.rst 与 docs/history/whatsnew-5.0.rst 了解后续演进;4.2 的变更概览见 docs/history/whatsnew-4.2.rst。
十、小结
Celery 4.1 系列是一个典型的"稳定性加固 + 能力补全"版本:它修复了 ETA/时区、重试过期、canvas 组合、结果后端序列化等一批直接影响生产可靠性的缺陷,同时引入worker_shutting_down信号、Redis 结果后端 SSL、DynamoDB 后端、类方法任务定义与disable_sync_subtasks防死锁机制。对于仍在 3.x/4.0 上运行、又暂未计划升级到 5.x 的团队,4.1.1(配合 Kombu 版本固定)是一个值得认真评估的稳定目标版本。
【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考