
Celery 4.3 版本变更解析新结果后端、任务优先级体系与可靠性修复全指南【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery本指南以仓库内 docs/history/changelog-4.3.rst 为骨架逐条解析 Celery 4.3.x 系列4.3.0 RC1 / RC2 / 4.3.0 / 4.3.1的变更内容并结合当前仓库源码celery/app/defaults.py、celery/app/control.py、celery/signals.py、celery/worker/strategy.py 等说明每个特性背后的实现与配置方式。读完本文你将掌握 4.3 引入的新配置项任务优先级继承、结果接受格式、chord join 超时等、新增的四个结果后端ArangoDB、CosmosDB、Azure Block Blob、S3、Pidbox 正则/通配符广播、task_received信号以及多个关键缺陷修复的来龙去脉可直接据此评估升级与迁移成本。版本总览与发布信息Celery 4.3 是一个功能与修复并重的次要版本发布日期与负责人如下摘自 changelog-4.3.rst版本发布日期发布者4.3.02019-03-31Omer Katz4.3.0 RC22019-03-03Omer Katz4.3.0 RC12019-02-20Omer Katz4.3.12020-09-10Omer Katz其中 4.3.1 为纯依赖修复版本唯一变更是把 vine 依赖版本限制在 5.0.0 以下- Limit vine version to be below 5.0.0避免上游破坏性变更影响 Celery 的 Promise 回调机制。4.3 版本同时正式宣告 Python 3.7 支持并将 Python 3.6 以下的 CLI 崩溃修复为可正常运行对 RabbitMQ 则结束了 2.x 的官方支持HA 头调整为兼容 RabbitMQ 3.x。新增配置项任务确认、优先级与结果存储4.3 系列引入了多个全局配置项多数现在仍保留在 celery/app/defaults.py 的NAMESPACES结构中。下面逐项说明默认值与语义。acks_on_failure_or_timeout应用级任务确认策略此前acks_on_failure_or_timeout仅是任务级选项4.3.0 将其补全为应用级配置并在 4.3.0 RC1 中新增了同名任务设置。该选项的核心价值在于 SQS 场景当任务失败或超时仍然确认ack消息时死信队列dead letter queue将无法工作。开启原生 SQS 消息生命周期后重试依赖消息重新投递redeliveries超过定义次数后消息转入死信队列。当前源码中该配置位于task命名空间defaults.pytaskNamespace( ... acks_on_failure_or_timeoutOption( True, typebool, deprecate_by6.0, remove_by7.0, alttask_acks_on_failure and task_acks_on_timeout, ), acks_on_failureOption(None, typebool), acks_on_timeoutOption(None, typebool), ... )配置方式app.conf.task_acks_on_failure_or_timeout False # 默认 True值得注意从源码可以看到该选项已在 6.0 版本被弃用、计划 7.0 移除替代方案是拆分为task_acks_on_failure与task_acks_on_timeout两个独立选项。升级到新版本时应优先使用拆分后的配置。task_default_priority任务默认优先级4.3.0 新增允许在应用级别设置任务默认优先级当某个具体任务未显式指定优先级时该值生效。源码中默认值为Nonedefaults.pydefault_priorityOption(None, typestring),app.conf.task_default_priority 5task_inherit_parent_priority链式任务优先级继承4.3.0 新增默认Falsedefaults.py。设置为True后Celery 任务会继承前一个链接任务link的优先级使整条链chain保持一致的优先级语义。官方示例c celery.chain( add.s(2), # priorityNone add.s(3).set(priority5), # priority5 add.s(4), # priority5继承 add.s(5).set(priority3), # priority3 add.s(6), # priority3继承 )父子任务场景同样生效app.task(bindTrue) def child_task(self): pass app.task(bindTrue) def parent_task(self): child_task.delay() # child_task 也会获得 priority5 parent_task.apply_async(args[], priority5)result_extended在结果后端存储扩展任务属性4.3.0 RC1 引入默认Falsedefaults.py。置为True后任务结果后端会额外存储任务的属性如参数、任务名、时间戳等便于后续查询与分析app.conf.result_extended True同时celery/app/task.py 的Task.update_state在 4.3.0 RC1 起接受关键字参数允许向结果后端传递额外字段这些字段默认不被使用但自定义结果后端可以据此决定结果如何存储。result_accept_content结果后端的可接受内容类型4.3.0 RC1 引入默认Nonedefaults.py取值类型为列表。其背景是签名消息使用特殊的auth序列化器但结果后端通常仍希望保持json以避免加密内容写入结果存储。因此该选项独立于全局accept_content专门控制结果后端接受的反序列化内容类型app.conf.result_accept_content [json, msgpack]在 celery/backends/base.py 中可以看到其优先级逻辑显式传入的accept参数 conf.result_accept_contentconf.accept_content。result_chord_join_timeoutchord join 超时可配置化此前 celery.result.GroupResult.join 在 chord 场景的等待超时固定为 3 秒。4.3.0 将其提取为配置项result_chord_join_timeout当前默认值仍为3.0秒defaults.pychord_join_timeoutOption(3.0, typefloat),app.conf.result_chord_join_timeout 10.0 # 秒event_exchange 与 control_exchange同 vhost 多应用隔离4.3.0 RC1 新增允许在同一个 RabbitMQ vhost 上运行多个相互独立的 Celery 应用。事件交换默认名为celeryevdefaults.py控制Pidbox交换默认名为celerydefaults.pyapp.conf.event_exchange myapp_events # 独立事件交换 app.conf.control_exchange myapp_control # 独立 Pidbox 交换新增信号task_received4.3.0 在 celery/signals.py 中新增了task_received信号task_received Signal( nametask_received, providing_args{request} )该信号在 worker 收到任务消息、进行过期/撤销检查之后、进入限流与执行队列之前触发。其触发点在 celery/worker/strategy.pysignals.task_received.send(senderconsumer, requestreq)监听方式from celery.signals import task_received task_received.connect def on_task_received(sender, request, **kwargs): print(f收到任务 {request.id}任务名 {request.name})request即 celery/worker/request.py 中的Request实例可访问任务 ID、任务名、参数等属性。该信号适合在任务入队瞬间做监控、审计或指标采集。远程控制广播支持正则与通配符匹配多 worker4.3.0 最实用的运维特性之一Control的broadcast/inspect现在支持用正则表达式regex或通配符glob模式一次性匹配多个 Pidbox从而批量 inspect 或 ping worker。实现位于 celery/app/control.pyControl._prepare在收到各节点回复后若提供了pattern则用match(node, pattern, matcher)过滤节点def __init__(self, destinationNone, timeout1.0, callbackNone, connectionNone, appNone, limitNone, patternNone, matcherNone): ... self.pattern pattern self.matcher matcher使用示例——匹配所有以worker-开头的节点并 pingapp.control.ping(patternrworker-*) app.control.ping(patternrworker-\d, matcherregex)注意matcher默认是 glob 语义传regex才使用正则匹配参见_request中patternself.pattern, matcherself.matcher的透传。这让大规模集群的批量巡检、批量撤销revoke、批量开启事件enable_events等远程操作从“逐节点列举”升级为“按模式匹配”。命名空间包支持与任务自动发现4.3.0 增加了对 PEP 420 隐式命名空间包namespace packages的支持允许从命名空间包中加载任务模块。相关能力在 celery/app/base.py 的autodiscover_tasks中体现。4.3.0 RC1 还有一个相关增强当autodiscover_tasks的related_name参数为None时会直接尝试导入包本身而不是导入包.tasks模块。源码 docstring 明确说明base.pyrelated_name (Optional[str]): The name of the module to find. Defaults to tasks. IfNonewill only try to import the package, i.e. look for module.示例# 尝试导入 foo 包本身而非 foo.tasks app.autodiscover_tasks([foo], related_nameNone)这为任务分布在命名空间包根__init__.py中的项目提供了便捷入口。新增结果后端四个 PaaS 后端与一个数据库后端4.3 系列一口气新增了四个云 PaaS 结果后端和一个数据库后端分别位于 celery/backends/ 目录ArangoDB 后端4.3.0 RC2基于 ArangoDB 的图/文档数据库实现于 celery/backends/arangodb.py。配置前缀arangodb_backend_settings对应 defaults.py 中的arangodb_backend_settings字典选项。Azure Block Blob 后端4.3.0 RC1基于 azure-storage 库将结果存入 Azure Blob Storage实现于 celery/backends/azureblockblob.py。仓库保留了端到端测试设置环境变量AZUREBLOCKBLOB_URLazureblockblob://{ConnectionString}连接串可在 Azure Portal 存储账户的 Access Keys 面板找到即可激活验证。默认容器名为celery并支持azureblockblob_retry_initial_backoff_sec默认 2、azureblockblob_retry_max_attempts默认 3等重试参数defaults.py。CosmosDB 后端4.3.0 RC1基于 pydocumentdb 库将结果存入 Azure CosmosDB实现于 celery/backends/cosmosdbsql.py。默认数据库名celerydb、集合名celerycoldefaults.py。S3 后端4.3.0 RC1将结果存入 AWS S3实现于 celery/backends/s3.py。配套的s3_access_key_id、s3_secret_access_key、s3_bucket、s3_base_path、s3_endpoint_url、s3_region等选项见 defaults.py可指向任意 S3 兼容对象存储。以 S3 为例的最小配置app.conf.result_backend s3://mybucket/prefix/ app.conf.s3_access_key_id AKIA... app.conf.s3_secret_access_key ... app.conf.s3_region us-east-1序列化与安全msgpack 替换、auth 序列化器重写4.3.0 将废弃的msgpack-python包替换为msgpack升级时需同步调整依赖声明。4.3.0 RC1 则对auth序列化器做了彻底重写此前“horribly broken”从 pyOpenSSL 迁移到 cryptography 库签名消息的安全性得到实质提升。配合新增的result_accept_content见上文可以做到“消息走 auth 签名、结果后端仍存 json 明文”的灵活组合。Canvas 与任务执行修复4.3 修复了大量 canvas工作流原语与任务执行相关问题这些修复大多在当前仓库中仍可验证其行为celery.chain.apply不再忽略关键字参数4.3.0 RC1应用链时会正确透传 kwargs。ResultSet 不再缓存 join 结果4.3.0 RC1此前ResultSet.get会向结果缓存写入数据一旦某个结果包含异常join 会意外失败现在移除了该缓存。chord task_always_eager 支持4.3.0 RC1task_always_eagerTrue时 chord 也能正确执行。chord 内嵌 group 中的子 chord/链4.3.0形如chord(group([chain(dummy.si(), chord(group([...]), ...)), ...]), callback)的复杂画布能正确执行最后一个任务官方回归用例见 t/unit/tasks/test_canvas.py。复杂画布的错误回调不再抛 AttributeError4.3.0含错误回调error callback的深度嵌套画布可正常构建。类class-based任务错误回调处理修复4.3.0 RC1。eager 任务往返序列化修复4.3.0修复了 eager 任务序列化时 serializer 选择不一致的问题。现在的选择顺序为apply_async的serializer参数 → Producer 的 serializer → 应用级task_serializer配置默认json。错误处理器允许非注册任务4.3.0link_error可以指向当前 worker 未注册的任务签名例如from celery import Signature Signature( bar, args[foo], link_errorSignature(msg.err, queuemsg) ).apply_async()MaxRetriesExceededError 携带任务参数4.3.0 RC1重试次数用尽时抛出的 celery/exceptions.py 异常会保留任务参数便于排查。稳定性与并发修复Asynpool 关闭 socket 死锁修复4.3.0 RC1此前 celery/concurrency/asynpool.py 的AsynPool关闭 socket 时只移除 hub 中的队列写入者writer而未移除读取者reader导致文件描述符死锁、worker 逐渐停止接收新任务现在在同一轮循环迭代中同时关闭 reader 与 writer。心跳连接断开后重试4.3.0 RC1此前 worker 会持续向已断开的连接写入事件分发器不断向出站缓冲追加消息造成内存泄漏现在连接死亡后会重试。Redis 结果后端 ResultConsumer 竞态修复4.3.0celery/backends/redis.py 的ResultConsumer不再假设start先于drain_events调用修复了 Gevent worker 池下的竞态条件。AsyncResult 回调内存泄漏修复4.3.0在 vine 1.2.0 提供完整 Python 2 weakref 支持后重新为 AsyncResult 回调 Promise 引入绑定方法的弱引用。worker 解码错误不再崩溃4.3.0 RC1v2 协议下消费端遇到 kombu.exceptions.DecodeError 时优雅处理而非使 worker 崩溃。Redis 结果后端 result_expiresNone 崩溃修复4.3.0result_expires设为None永久保留结果时不再崩溃该配置默认值为timedelta(days1)defaults.py。Celery Beat 与调度修复微秒级调度4.3.0 RC1beat调度时间戳现在正确处理微秒。时区感知的时间戳计算4.3.0 RC1计算时间戳时正确考虑时区。Scheduler.schedules_equal空值容错4.3.0 RC1celery/beat.py 中任一参数为None时不再异常。Windows 10 陈旧 PID 文件清理4.3.0 RC1beat 在SystemExit时移除过期 PID 文件解决 Windows 上 beat 无法启动的问题。兼容性与依赖变更汇总Python新增 Python 3.7 支持修复 Python 3.6 下 CLI 崩溃此前因使用 Python 3.6 才引入的ModuleNotFoundError异常导致。Django放弃 Django 1.11 支持移除旧版 djcelery loaderloader 在 sys.path 中改为前置当前工作目录保证项目目录优先于系统模块Request向后端store_result传入Context而非self修复 django-celery-results 下撤销任务导致结果数据损坏的问题相关实现见 celery/backends/base.py。依赖下限提升Kombu 提升至 4.3RC1再到 4.4RC2Billiard 提升至 3.6py-redis 提升至 3.2.0因早期版本存在多个影响 Celery 的 bugeventlet 提升至 0.24.1vine 限制在 5.0.04.3.1。Riak 后端向用户警告 Python 3.7 下可能的不兼容问题。Cython 支持支持 Cython 化的 Celery 任务。Redis 与 RabbitMQ SSLRedis broker/结果后端的 SSL 支持获得多项修复与增强。RabbitMQHA 头调整以兼容 RabbitMQ 3.x同时结束对 RabbitMQ 2.x 的官方支持。命令行与部署改进celery update错误处理改进4.3.0 RC1命令执行出错时提供更清晰的反馈。应用加载失败提示4.3.0此前 app 无法加载时 CLI 直接抛异常现在打印明确的错误信息。celery report增加内核版本4.3.0 RC1环境诊断报告新增内核版本字段与其他平台信息一并输出。init.d 服务停止脚本修复4.3.0 RC1修复 extra/generic-init.d/ 中celeryd/celerybeat的 stop 逻辑。升级建议与小结从 4.3 系列变更可以看到 Celery 在“云原生结果后端”“细粒度任务确认与优先级控制”“运维可观测性”三个方向的明显投入。升级到 4.3 时建议重点关注SQS 用户利用task_acks_on_failure_or_timeoutFalse配合原生死信队列同时注意该选项在后续版本中的弃用与拆分task_acks_on_failure/task_acks_on_timeout。序列化依赖msgpack-python→msgpack的替换需要同步修改requirements/中的依赖声明auth序列化器改用 cryptography 后若曾依赖 pyOpenSSL 需调整。批量运维立即利用 Pidbox 的正则/通配符广播能力减少大规模集群巡检脚本的复杂度。结果后端选型新增的 S3、Azure Block Blob、CosmosDB、ArangoDB 后端为“低成本、可扩展、免运维”的结果存储提供了更丰富的选项result_accept_content与result_extended则让结果存储的安全边界与信息维度都可按需定制。以上所有配置项的默认值均可在 celery/app/defaults.py 的NAMESPACES中逐一核对相关行为也可在 t/unit/ 下的单元测试中进一步验证。【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考