与结果后端(Backend)选型与配置指南)
Celery 消息代理Broker与结果后端Backend选型与配置指南【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celeryCelery 支持多种消息传输Message Transport与结果存储Result Store方案。本指南以 Celery 官方文档 docs/getting-started/backends-and-brokers/index.rst 为主体系统梳理 RabbitMQ、Redis、Amazon SQS、Kafka、Google Cloud Pub/Sub 等 Broker 的选型对比、安装方式、核心配置与注意事项并结合仓库源码说明各传输在底层的工作原理。读完本文你将能够根据可靠性、监控能力、运维成本等维度为项目选择合适的 Broker 与 Result Backend并正确完成配置与调优。Broker 与 Backend两个容易混淆的概念在 Celery 的架构中有两个承担不同职责的存储/传输组件Broker消息代理 / Message Transport负责传输任务消息。生产者如 Web 应用把任务投递到 BrokerWorker 从 Broker 消费任务并执行。Backend结果后端 / Result Store负责存储任务的执行状态与返回值。启用后你可以通过AsyncResult查询任务结果。同一个组件可以同时担任两种角色也可以只担任一种。官方文档在 index.rst 的 Summaries 一节明确说明Redis可以同时作为 Broker 和 BackendRabbitMQ主要作为 Broker也可通过rpc://后端存储结果Amazon SQS、Kafka、Google Cloud Pub/Sub仅作为 BrokerSQLAlchemy仅作为 Backend用于对接 MySQL、PostgreSQL、SQLite 等关系型数据库。该结论与仓库目录结构一致Redis 后端实现在 celery/backends/redis.pyRPC 后端在 celery/backends/rpc.pySQLAlchemy 后端位于 celery/backends/database/而各 Broker 的传输层由 Kombu 提供requirements/extras/目录下的 extras 依赖包可佐证。Broker 选型总览官方对比表index.rst 给出了各传输的官方对比表这是选型最直接的依据NameStatusMonitoringRemote ControlRabbitMQStableYesYesRedisStableYesYesAmazon SQSStableNoNoZookeeperExperimentalNoNoKafkaExperimentalNoNoGC PubSubExperimentalYesYes对表中各列的理解官方文档给出了明确解释Status状态Experimental 状态的 Broker 可能可以正常工作但它们没有专职维护者使用时需自行评估风险。Monitoring监控缺少监控支持意味着该传输没有实现事件events机制因此 Flower、celery events、celerymon等基于事件流的监控工具将无法工作。这也是 SQS、Kafka 无法使用事件监控的根本原因。Remote Control远程控制指在运行时使用celery inspect和celery control命令以及其他基于远程控制 API 的工具检查和管理工作者的能力。SQS 等传输由于缺乏该支持无法执行远程 worker 管理操作。各传输的详细配置文档见 RabbitMQ、Redis、Amazon SQS、Kafka、Google Cloud Pub/Sub。RabbitMQ默认消息代理RabbitMQ 是 Celery 的默认 Broker不需要安装任何额外依赖只需在配置中给出 Broker 地址即可broker_url amqp://myuser:mypasswordlocalhost:5672/myvhost安装 RabbitMQ 服务端安装方式参考 RabbitMQ 官方下载与安装文档macOS 可用 Homebrew 的brew install rabbitmq安装安装后需将/usr/local/sbin加入PATH。需要注意的是如果你使用 DHCP 随机分配主机名必须永久配置主机名因为 RabbitMQ 使用主机名进行节点间通信——若主机名以 IP 开头如23.10.112.31.comcast.netRabbitMQ 会尝试使用rabbit23这样的非法节点名$ sudo scutil --set HostName myhost.local并在/etc/hosts中添加对应解析127.0.0.1 localhost myhost myhost.local启动与停止服务时务必使用rabbitmqctl而非kill$ sudo rabbitmq-server # 前台启动 $ sudo rabbitmq-server -detached # 后台启动注意只有一个连字符 $ sudo rabbitmqctl stop # 优雅停止创建用户、虚拟主机与权限使用 Celery 前需要创建 RabbitMQ 用户、虚拟主机vhost并授权$ sudo rabbitmqctl add_user myuser mypassword $ sudo rabbitmqctl add_vhost myvhost $ sudo rabbitmqctl set_user_tags myuser mytag $ sudo rabbitmqctl set_permissions -p myvhost myuser .* .* .*将myuser、mypassword、myvhost替换为实际值。完整访问控制细节可参考 RabbitMQ 官方 Admin Guide 与 Access Control 文档。Quorum Queues5.5 版本起Celery 支持 RabbitMQ 的 Quorum Queues只需通过x-queue-type头声明队列类型为quorumfrom kombu import Queue task_queues [Queue(my-queue, queue_arguments{x-queue-type: quorum})] broker_transport_options {confirm_publish: True}也可以不显式声明队列而是修改默认队列类型并通过task_routes路由task_default_queue_type quorum task_default_exchange_type topic task_default_queue my-queue broker_transport_options {confirm_publish: True} task_routes { *: { routing_key: my-queue, }, }Celery 会通过worker_detect_quorum_queues设置自动检测是否使用了 Quorum Queues该检测逻辑实现在 celery/utils/quorum_queues.py 的detect_quorum_queues函数中当 driver 类型为amqp时遍历app.amqp.queues检查各队列的queue_arguments中是否含有x-queue-type: quorum。官方建议保持该自动检测默认开启。使用 Quorum Queues 的限制Quorum Queues 要求禁用全局 QoS这意味着 per-channel QoS 变为静态部分 Celery 功能将不可用自动伸缩Autoscaling依赖在进程创建/终止时动态调整 prefetch count因此在检测到 Quorum Queues 时不会生效worker_enable_prefetch_count_reduction设置即使设为True也会成为空操作ETA/Countdown 任务在到达执行时间前会阻塞 worker无法再通过提高 prefetch count 获取其他任务。为解决 ETA/Countdown 的调度问题Celery 在检测到 Quorum Queues 时会自动启用 Native Delayed Delivery原生延迟投递利用 RabbitMQ 自身的机制来调度延迟任务。该设计借鉴自 NServiceBus 的 delayed delivery 实现。默认情况下 Native Delayed Delivery 的延迟队列本身也是 Quorum Queues如需改为经典队列可设置broker_native_delayed_delivery_queue_type为classic。Redis既当代理又当结果后端Redis 是官方文档重点推荐且最灵活的方案之一既能作为 Broker也能作为 Result Backend且两者可同时使用同一个 Redis 实例。安装Redis 支持需要额外依赖可通过celery[redis]bundle 一次安装$ pip install -U celery[redis]对应依赖定义在 requirements/extras/redis.txt内容为kombu[redis]即传输层由 Kombu 的 redis transport 提供。配置 Brokerapp.conf.broker_url redis://localhost:6379/0URL 格式为redis://:passwordhostname:port/db_number除 scheme 外的所有字段都是可选的默认主机为localhost、端口6379、数据库0。此外还支持三种特殊形式凭证提供器credential providerredis://hostname:port/db_number?credential_providermymodule.myfile.myclassUnix socket 连接redissocket:///path/to/redis.sock如需指定数据库号可附加virtual_host参数redissocket:///path/to/redis.sock?virtual_hostdb_numberRedis Sentinelapp.conf.broker_url sentinel://localhost:26379;sentinel://localhost:26380;sentinel://localhost:26381 app.conf.broker_transport_options { master_name: cluster1 }还可以通过sentinel_kwargs向 Sentinel 客户端传递额外参数如密码app.conf.broker_transport_options { sentinel_kwargs: { password: password } }Visibility Timeout可见性超时Visibility timeout 定义 worker 确认任务前需要等待的秒数超时后消息会被重新投递给其他 worker。默认值为1 小时可通过broker_transport_options调整app.conf.broker_transport_options {visibility_timeout: 3600} # 1 hour.配置结果后端如果希望把任务状态与返回值也存到 Redisapp.conf.result_backend redis://localhost:6379/0使用 Sentinel 时需要为结果后端单独指定master_nameapp.conf.result_backend_transport_options {master_name: mymaster}结果后端还支持三个高级选项global_keyprefix全局 key 前缀默认无前缀。当多个应用共享同一个 Redis 数据库时可以给所有结果 key 加前缀以避免冲突app.conf.result_backend_transport_options { global_keyprefix: my_prefix_ }retry_policy连接超时重试策略app.conf.result_backend_transport_options { retry_policy: { timeout: 5.0 } }additional_connection_errors附加连接错误某些 Redis 代理或云厂商会抛出 Celery 无法识别的自定义异常。为了能让这些异常被自动重试可在result_backend_transport_options下配置additional_connection_errors点号分隔的导入字符串与异常类均支持该特性在本文档中标记为 5.7 版本引入app.conf.result_backend_transport_options { additional_connection_errors: ( my_redis_proxy.CustomConnectionError, ), }Kombu 在 Redis 上的存储模型理解 key 结构Redis 本身没有队列queue和交换机exchange的原生概念因此 Kombu 的 redis transport 在普通 Redis 数据结构之上模拟了这些概念。理解底层的 key 结构有助于调试、共享 Redis 实例或配置持久化队列是 list每个队列是一个 Redis 列表key 即队列名默认celery。投递任务是对该列表执行LPUSHworker 则通过阻塞式BRPOP消费因此每个队列内部是 FIFO 顺序。普通任务消息不使用 Pub/Sub。启用消息优先级时每个队列会按优先级档位拆分为多个 list如celery、celery\x06\x163、……。绑定关系是 set每个 exchange 的路由表是一个 Redis setkey 名为_kombu.binding.exchange name如_kombu.binding.celery成员记录哪些队列以何种 routing key 绑定到该 exchange。未确认消息保存在 hash 中Redis 没有消息确认机制Kombu 用unackedhashdelivery tag → 消息和unacked_indexsorted set记录接收时间来模拟。确认任务会同时删除这两处条目由unacked_mutexkey 保护的周期性清理会把超过 visibility timeout 的条目重新推回队列 list供其他 worker 消费。广播使用 Pub/Subfanout exchange如celery.pidbox用于远程控制命令是唯一使用 Redis Pub/Sub 的场景。频道名带数据库号前缀如/0.celery.pidbox因此同一服务器上不同数据库的应用不会收到彼此的广播。Pub/Sub 消息只投递给当前在线的订阅者且不会持久化。上述所有 key 默认都没有前缀。如果自己的应用恰好也使用celery或unacked这样的 key就会与 Broker 冲突。最安全的做法是使用独立 Redis 实例或独立数据库号如果必须共享数据库可以用global_keyprefix给 Broker 的所有 key 加前缀app.conf.broker_transport_options {global_keyprefix: myapp:}结果后端的结果存储使用独立的celery-task-meta-*key前缀配置见上文redis-result-backend-global-keyprefix。另外务必注意不要对 Broker 所在的数据库执行FLUSHDB也不要在服务器上执行FLUSHALL——它们会一次性删除所有已排队任务、未确认任务和绑定关系无论 Celery 是否在运行FLUSHALL还会清空服务器上的所有数据库。持久化任务消息就是普通的 Redis key因此 Redis 重启后任务能否存活完全取决于服务器持久化配置默认的 RDB 快照可能丢失自上次快照以来发布的任务开启 AOFappendonly yes并配合appendfsync everysec可将丢失窗口缩小到最多 1 秒还需防止内存压力下 Redis 驱逐 key见下文 Caveats。Serverless 方案文档还提及两种 serverless/替代型 RedisUpstash提供 serverless Redis 数据库服务采用最终一致性模型与多级存储架构支持持久化存储适合微服务或希望降低运维成本的环境官方提供与 Celery 集成的示例。Dragonfly声称是 Redis 的 drop-in 替代品旨在充分利用现代云硬件、降低成本和提升性能。以上均为第三方服务/产品的官方宣传定位请以各厂商文档为准。Redis 的 CaveatsVisibility timeout 与 ETA/countdown/retry 冲突如果任务未在 visibility timeout 内被确认会被重新投递并再次执行。对于执行时间超过 visibility timeout 的 ETA/countdown/retry 任务这会形成执行→超时→重投→再执行的死循环。缓解办法是调大 visibility timeout 以匹配最长 ETA但这并不被推荐Celery 会在 worker 关闭时重投消息过长的 timeout 只会延迟丢失任务如断电或 worker 被强制终止时的重投。重要提示如果需要调度更远未来的任务请优先考虑数据库支持的计划任务方案——周期性任务periodic tasks不受 visibility timeout 影响因为该概念与 ETA/countdown 相互独立。如需调大 timeout必须同时设置以下三个同名配置项缺一不可app.conf.broker_transport_options {visibility_timeout: 43200} app.conf.result_backend_transport_options {visibility_timeout: 43200} app.conf.visibility_timeout 43200值为表示秒数的整数。注意如果多个应用共享同一个 Broker 但设置不同将取最短的值未设置时也会发送默认值参与比较。软关闭Soft Shutdown在关闭过程中开启task_acks_late的 worker 会尝试把未确认消息重新入队但若 worker 被强制终止冷关闭可能来不及重投消息要等到 visibility timeout 过后才能被再次消费。visibility timeout 越大这种延迟越严重。软关闭引入了冷关闭之前的一段限时温关闭阶段显著提高重投成功率。启用方式是把worker_soft_shutdown_timeout设为大于 0 的浮点秒数若同时通过环境变量将REMAP_SIGTERM配置为SIGQUITworker 收到TERM信号以及QUIT信号时会先进入软关闭。Key 驱逐Redis 在特定情况下可能驱逐 key。若遇到如下错误InconsistencyError: Probably the key (_kombu.binding.celery) has been removed from the Redis database.应在 redis 配置文件中设置maxmemory以及maxmemory-policy为noeviction或allkeys-lru防止键被驱逐。Group 结果排序result_chord_ordered4.4.6 及更早版本使用无序列表存储 group 的结果对象可能导致结果返回顺序与任务定义顺序不一致4.4.7 引入可选修复5.0 起默认开启opt-out。该行为由result_chord_ordered控制# Specifying this for workers running Celery 4.4.6 or earlier has no effect app.conf.result_backend_transport_options { result_chord_ordered: True # or False }由于这是运行时行为变更共享同一 Redis 后端的 worker 必须保持一致若集群中仍存在 4.4.6 及更早版本 worker则 5.0 worker 必须显式设为False若集群全为 4.4.7 及以后版本建议全部设为True以便将来迁移。迁移行为会中断 Redis 中现有结果若由迁移后的 worker 执行下游任务可能造成破坏请提前规划。Amazon SQS完全托管的云原生 Broker如果项目已深度集成 AWS 并熟悉 SQSSQS 是很好的 Broker 选择极度可扩展、完全托管任务委派方式与 RabbitMQ 类似但缺少 RabbitMQ 的部分特性例如 worker 远程控制命令。安装$ pip install celery[sqs]依赖定义见 requirements/extras/sqs.txt包含boto31.26.143、urllib31.26.16、kombu[sqs]5.5.0等。配置Broker URL 格式为sqs://aws_access_key_id:aws_secret_access_key必须保留末尾的符号并应对密码进行 URL 编码以保证正确解析from kombu.utils.url import safequote aws_access_key safequote(ABCDEFGHIJKLMNOPQRST) aws_secret_key safequote(ZYXK7NiynG/TogH8NjP9nlE73sq3) broker_url sqs://{aws_access_key}:{aws_secret_key}.format( aws_access_keyaws_access_key, aws_secret_keyaws_secret_key, )安全警告不要在 Django 的debugTrue下使用这种 URL 内嵌凭证的方式——调试模式下 Django 会展示环境变量SQS URL 中的 AWS 访问密钥可能暴露到互联网。生产环境请关闭 debug 模式或改用下面的方式。凭证也可以通过环境变量AWS_ACCESS_KEY_ID与AWS_SECRET_ACCESS_KEY提供此时 Broker URL 只需写sqs://如果实例使用 IAM 角色同样只需sqs://Kombu 会尝试从实例元数据获取访问令牌。关键选项通过 broker_transport_options 配置Region默认us-east-1可配置为其他区域broker_transport_options {region: eu-west-1}Visibility Timeout默认30 分钟注意与 Redis 的 1 小时不同。该选项用于创建 SQS 队列时设置使用预定义队列predefined_queues时无效broker_transport_options {visibility_timeout: 3600} # 1 hour.Polling Interval轮询间隔两次无结果轮询之间的休眠秒数支持 int 或 float默认1 秒。轮询越频繁成本越高增大轮询间隔可以省钱broker_transport_options {polling_interval: 0.3}过于频繁的轮询可能造成 busy loop 并大量占用 CPU如果需要亚毫秒级延迟请改用 RabbitMQ 或 Redis。Long Polling长轮询默认启用ReceiveMessage的WaitTimeSeconds为 10 秒可配置合法值 0~20broker_transport_options {wait_time_seconds: 15}注意新创建的队列包括 Celery 创建的其 Receive Message Wait Time 队列属性默认是 0。Queue Prefix队列名前缀默认不加前缀避免与其他服务冲突可加前缀broker_transport_options {queue_name_prefix: celery-}预定义队列Predefined Queues如果希望 Celery 使用 AWS 中预先定义好的队列且绝不列出、创建或删除 SQS 队列可通过predefined_queues传入队列名到 URL 的映射broker_transport_options { predefined_queues: { my-q: { url: https://ap-southeast-2.queue.amazonaws.com/123456/my-q, access_key_id: xxx, secret_access_key: xxx, } } }重要警告在predefined_queues中不要对access_key_id和secret_access_key做 URL 编码不要使用safequote——URL 编码只应用于 Broker URL 中的凭证。若在predefined_queues中使用编码后的凭证会导致签名不匹配错误The request signature we calculated does not match the signature you provided.。Broker URL 与预定义队列组合的正确写法import os from kombu.utils.url import safequote from celery import Celery # Raw credentials from environment AWS_ACCESS_KEY_ID os.getenv(AWS_ACCESS_KEY_ID) AWS_SECRET_ACCESS_KEY os.getenv(AWS_SECRET_ACCESS_KEY) # URL-encode ONLY for broker URL aws_access_key_encoded safequote(AWS_ACCESS_KEY_ID) aws_secret_key_encoded safequote(AWS_SECRET_ACCESS_KEY) # Use encoded credentials in broker URL broker_url fsqs://{aws_access_key_encoded}:{aws_secret_key_encoded} celery_app Celery(tasks, brokerbroker_url) celery_app.conf.broker_transport_options { region: us-east-1, predefined_queues: { my-queue: { url: https://sqs.us-east-1.amazonaws.com/123456/my-queue, # Use RAW credentials here (NOT encoded) access_key_id: AWS_ACCESS_KEY_ID, secret_access_key: AWS_SECRET_ACCESS_KEY, }, }, }使用predefined_queues时visibility timeout 应在 AWS 控制台直接配置在 SQS 队列上而非通过上面的visibility_timeout选项。Back-off Policy退避策略退避策略利用 SQS 的 visibility timeout 机制按重试次数逐步拉大两次重试之间的间隔。重试次数由 SQS 通过消息属性ApproximateReceiveCount管理无需用户额外干预broker_transport_options { predefined_queues: { my-q: { url: https://ap-southeast-2.queue.amazonaws.com/123456/my-q, access_key_id: xxx, secret_access_key: xxx, backoff_policy: {1: 10, 2: 20, 3: 40, 4: 80, 5: 320, 6: 640}, backoff_tasks: [svc.tasks.tasks.task1] } } }backoff_policy是重试次数 → 重试间隔秒数即 SQS visibility timeout的字典backoff_tasks是应用该策略的任务名列表。以上配置的效果AttemptDelay2nd attempt20 seconds3rd attempt40 seconds4th attempt80 seconds5th attempt320 seconds6th attempt640 secondsSTS Token 认证支持通过sts_role_arn和sts_token_timeout使用 AWS STS 认证sts_role_arn是被假设的 IAM 角色 ARNsts_token_timeout是令牌超时默认也是最小值900 秒超时后自动创建新令牌broker_transport_options { predefined_queues: { my-q: { url: https://ap-southeast-2.queue.amazonaws.com/123456/my-q, access_key_id: xxx, secret_access_key: xxx, backoff_policy: {1: 10, 2: 20, 3: 40, 4: 80, 5: 320, 6: 640}, backoff_tasks: [svc.tasks.tasks.task1] } }, sts_role_arn: arn:aws:iam::xxx:role/STSTest, # optional sts_token_timeout: 900 # optional }SQS 的 Caveats任务未在visibility_timeout内确认会被重新投递执行执行时间超过 timeout 的 ETA/countdown/retry 任务会反复重执行。AWS 目前支持的最大 visibility timeout 为12 小时43200 秒broker_transport_options {visibility_timeout: 43200}周期任务不受影响与 ETA/countdown 相互独立。注意visibility timeout 调大只会延迟丢失任务在断电或 worker 被强杀后的重投。SQS不支持 worker 远程控制命令表格中 Remote Control 为 No。SQS不支持事件events因此不能与celery events、celerymon或 Django Admin monitor 配合使用。使用 FIFO 队列时发布消息可能需要设置MessageGroupId与MessageDeduplicationId等附加属性可通过apply_async的关键字参数传递task.apply_async( MessageGroupIdYourMessageGroupId, MessageDeduplicationIdYourMessageDeduplicationId, )与 Redis 相同的软关闭soft shutdown建议同样适用强杀 worker 会推迟未确认任务的重投可设置worker_soft_shutdown_timeout大于 0 引入限时温关闭阶段配合环境变量REMAP_SIGTERMSIGQUIT时worker 收到TERM/QUIT信号会先进入软关闭。结果后端警告AWS 家族里目前没有内置的 SQS 结果后端切勿将amqp结果后端与 SQS 一起使用——它会为每个任务创建一个队列且不会被回收导致不必要的费用。Kafka实验性 BrokerKafka 作为 Broker 目前处于Experimental状态没有专职维护者且从仓库文档看目前只能使用单个 worker相关限制见 Kombu 上游 issue 讨论。配置在celeryconfig.py中的典型配置以 Confluent Kafka 为例使用 SASL/SCRAM 认证import os task_serializer json broker_transport_options { # allow_create_topics: True, } broker_connection_retry_on_startup True # For using SQLAlchemy as the backend # result_backend dbpostgresql://postgres:examplelocalhost/postgres broker_transport_options.update({ security_protocol: SASL_SSL, sasl_mechanism: SCRAM-SHA-512, }) sasl_username os.environ[SASL_USERNAME] sasl_password os.environ[SASL_PASSWORD] broker_url fconfluentkafka://{sasl_username}:{sasl_password}broker:9094 broker_transport_options.update({ kafka_admin_config: { sasl.username: sasl_username, sasl.password: sasl_password, }, kafka_common_config: { sasl.username: sasl_username, sasl.password: sasl_password, security.protocol: SASL_SSL, sasl.mechanism: SCRAM-SHA-512, bootstrap.servers: broker:9094, } })其中allow_create_topics仅在 topic 尚不存在时需要tasks.py中的用法与普通 Celery 应用无异from celery import Celery app Celery(tasks) app.config_from_object(celeryconfig) app.task def add(x, y): return x y认证方面SASL 用户名和密码通过环境变量传入见上方示例。工作方式与限制Celery 队列会被路由到 Kafkatopic例如队列名为add_queueKafka 中就会创建/使用名为add_queue的 topic。对于支持的后端canvas 的典型机制chain、group、chord通常可以工作。主要限制当前使用 Kafka 作为 Broker 时只支持单个 worker详见 Kombu 上游 issue #1785 的讨论大规模集群场景需谨慎评估。Google Cloud Pub/Sub实验性 BrokerGoogle Cloud Pub/Sub 是5.5 版本起引入的 Broker。如果项目已深度集成 Google Cloud 并熟悉 Pub/Sub它是一个很好的选择极度可扩展、完全托管任务委派方式与 RabbitMQ 类似且支持监控与远程控制见对比表。安装与配置$ pip install celery[gcpubsub]依赖定义见 requirements/extras/gcpubsub.txtkombu[gcpubsub]5.5.0。Broker URL 需包含 Google 项目 ID且必须加上projects/前缀gcpubsub://projects/project-id登录凭证使用环境中配置的常规 GCP 凭证。关键选项Resource expiry资源过期默认配置追求开箱即用、成本可控pubsub 消息与订阅在24 小时后过期可通过expiration_seconds调整expiration_seconds 86400Ack Deadline Seconds确认截止定义 Pub/Sub 基础设施等待 worker 确认任务的秒数超时后消息被重新投递给其他 worker。默认240 秒worker 会自动为所有待处理消息续期broker_transport_options {ack_deadline_seconds: 60} # 1 minute.Polling Interval轮询间隔默认0.1 秒。不过这不意味着 worker 会每 0.1 秒轰炸 Pub/Sub API——无消息时 worker 会被 Pub/Sub API 的阻塞调用挂起直到有新消息或 10 秒超时才会返回broker_transport_options {polling_interval: 0.3}过于频繁的轮询会形成 busy loop 并占用 CPU需要亚毫秒级延迟时请改用 RabbitMQ 或 Redis。Queue Prefix队列名前缀默认 Celery 会给队列名加kombu-前缀可通过queue_name_prefix调整broker_transport_options {queue_name_prefix: kombu-}Pool start method进程池启动方式Pub/Sub 驱动基于 gRPC会启动后台 C 线程而 C 线程不能安全地 fork。使用默认的fork启动方式时子进程可能继承损坏的 gRPC/TLS 状态——例如在task_time_limit硬杀后替代子进程可能挂起。因此建议将 prefork 池子进程的启动方式改为spawn让每个子进程获得干净的解释器worker_pool_start_method spawn结果后端Google Cloud StorageGCS是存放结果的合适候选可参考仓库中 GCS 后端文档celery/backends/gcs.py 与对应测试 t/unit/backends/test_gcs.py。Caveats使用 Flower 监控时需要--inspect-timeout10选项才能正确检测 worker 状态。空闲订阅无排队消息配置为 24 小时后自动移除以降低成本。排队与未确认消息同样在 24 小时后自动清理。通道队列大小只是近似值Pub/Sub API 不提供获取订阅内消息精确数量的方法。孤儿 topic没有订阅的 topic不会被自动删除GCP 对每个项目有 1 万个 topic 的硬限制建议定期手动清理孤儿 topic。最大消息大小限制为10MB作为变通方案可将大消息存入 GCS 后端再把 GCS URL 传给任务。结果后端Result Backend选型除了上文已详述的 Redis 后端外官方文档 Summaries 部分还重点提到了另外几种后端组合RabbitMQ 的 rpc:// 后端RabbitMQ 可以通过rpc://后端存储结果。该后端会为每个客户端创建独立的临时队列。适合轻量、短生命周期结果的场景注意与 SQS 的兼容性问题见上文警告。SQLAlchemy 后端SQLAlchemy 后端允许 Celery 对接 MySQL、PostgreSQL、SQLite 等关系型数据库其实现位于 celery/backends/database/__init__.py、models.py、session.py完整选项见配置文档 docs/userguide/configuration.rst 中关于 database result backend 的章节。它是用 SQL 数据库作为结果后端的官方途径适合对结果持久化有更高要求的场景。常见组合建议官方文档给出了一条非常实用的经验法则RabbitMQ作为 Broker Redis作为 Backend是极常见的组合。如果对结果存储的长期持久性有更严格的保障需求可考虑PostgreSQL 或 MySQL通过 SQLAlchemyCassandra或自定义后端。选型决策要点总结结合官方对比表与各传输文档选型时可从以下几个维度权衡监控与运维需求需要 Flower /celery events等事件监控、以及celery inspect/celery control远程控制应优先 RabbitMQ、Redis 或 GC PubSub三者均支持 Monitoring 与 Remote ControlSQS、Kafka 则不具备这些能力。云厂商绑定与托管成本深度使用 AWS 选 SQS深度使用 GCP 选 Pub/Sub两者都极度可扩展、完全托管想要自建则选 RabbitMQ 或 Redis。消息体大小与吞吐RabbitMQ 处理大消息优于 Redis而 Redis 在大量小消息快速涌入时表现良好但大消息可能阻塞系统。高并发快速涌入场景下除非 RabbitMQ 已大规模部署否则可考虑 Redis 或 SQS。结果持久化要求Redis 作为后端超快但受内存限制且需要自行设计持久化对结果持久性要求高时改用 PostgreSQL/MySQLSQLAlchemy、Cassandra 或自定义后端。实验性传输的风险Kafka当前仅支持单 worker、Zookeeper、GC PubSub 均标注为 Experimental可能可用但没有专职维护者生产环境需额外评估。以上所有配置均以当前仓库Celery 5.6.2 开发分支见 celery/init.py随附文档与源码为准各传输更完整的配置项列表可继续阅读 docs/userguide/configuration.rst 中的 broker 与 backend 章节以及仓库中对应的后端实现与单元测试如 t/unit/backends/test_redis.py、t/unit/backends/test_rpc.py。【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考