Celery 分布式任务队列入门指南:核心概念、特性与安装实践
Celery 分布式任务队列入门指南核心概念、特性与安装实践【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celeryCelery 是一个用 Python 编写的分布式任务队列用于把耗时、异步或需要跨机器执行的工作从应用主流程中剥离出来交给常驻的 worker 进程处理。本指南基于仓库中的 docs/getting-started/introduction.rst 展开覆盖任务队列的基本原理、运行环境要求、核心能力全景与安装方式并结合仓库源码celery/app/base.py、celery/init.py 等补充实现层面的细节。读完本文你将理解 Celery 的架构角色知道如何选择 broker 与 result backend掌握 pip 安装与各功能 bundle 的用法并了解官方推荐的学习路径。什么是任务队列Task Queue任务队列是在线程或机器之间分发工作的一种机制。一个任务队列的输入是一个被称为任务task的工作单元而专门的 worker 进程会持续监控任务队列等待新的工作到来并执行。Celery 通过消息进行通信通常借助一个 broker消息中间件在客户端与 worker 之间做中介客户端client将一条消息任务放入队列broker 将消息投递给某个 workerworker 执行任务并可选地把结果写回 result backend。一个 Celery 系统可以由多个 worker 和多个 broker组成从而支持高可用与水平扩展。Celery 本身用 Python 编写但通信协议可以在任何语言中实现——除了 Python 之外社区还提供了 Node.jsnode-celery、PHP 客户端、Gogocelery、gopher-celery、Rustrusty-celery等实现。语言互操作也可以通过另一种方式达成暴露一个 HTTP 端点然后由一个任务去请求它即 webhook 模式。从源码看这个生产-投递-执行的链路分散在多个模块中任务定义与注册在 celery/app/task.py 与 celery/app/registry.py消息发布与消费则依赖 kombu/amqp 连接 broker应用封装见 celery/app/amqp.pyworker 侧的消费入口见 celery/worker/consumer。运行 Celery 需要什么版本要求根据文档侧边栏的说明Celery 5.x 系列支持的 Python 版本为Python 3.8、3.9、3.10、3.11、3.12、3.13以及PyPy3.9v7.3.12。当前仓库是开发分支celery/init.py 中标明的版本号为5.6.2系列代号recovery。如果你运行的是更老版本的 Python则需要搭配更老的 CeleryPython 版本对应 Celery 版本Python 3.7Celery 5.2 或更早Python 3.6Celery 5.1 或更早Python 2.7Celery 4.x 系列Python 2.6Celery 3.1 系列或更早Python 2.5Celery 3.0 系列或更早Python 2.4Celery 2.2 系列或更早另外需要说明Celery 是一个资金极少的项目官方不支持 Microsoft Windows请勿针对该平台提交 issue。消息传输brokerCelery 需要一个消息传输层来发送和接收消息。RabbitMQ 与 Redis 两种 broker 传输是功能完备feature complete的此外还支持大量实验性方案包括用于本地开发的 SQLite。Celery 可以运行在单台机器、多台机器甚至跨数据中心部署。各 broker 的详细配置请分别参考docs/getting-started/backends-and-brokers/rabbitmq.rstdocs/getting-started/backends-and-brokers/redis.rstdocs/getting-started/backends-and-brokers/sqs.rstdocs/getting-started/backends-and-brokers/index.rst开始学习如果你是第一次接触 Celery或者从 3.1 之前的版本升级过来建议按顺序阅读两篇入门教程docs/getting-started/first-steps-with-celery.rst选择并安装 broker、安装 Celery、创建第一个任务、启动 worker、调用任务、查看任务状态与返回值、配置任务序列化与路由docs/getting-started/next-steps.rst展示 Celery 的进阶能力应用、任务、画布工作流、路由、周期任务、监控、安全等。Celery 的核心特质文档用四个关键词概括 Celery 的设计取向Simple简单Celery 易于使用和维护不需要配置文件即可跑起来。以下是最简单的 Celery 应用from celery import Celery app Celery(hello, brokeramqp://guestlocalhost//) app.task def hello(): return hello worldCelery(...)的第一个参数是当前模块名用于自动生成任务名broker关键字指定消息中间件的 URL。这里的amqp://guestlocalhost//指向本地 RabbitMQRabbitMQ 也是默认选项。从实现上看app Celery(...)创建的是 celery/app/base.py 中的Celery类实例它是你在 Celery 中想做的一切的入口点创建任务、管理 worker、访问配置app.conf、获取结果等。同时 celery/init.py 通过local.recreate_module实现了懒加载——from celery import Celery不会立刻导入整个库而是在首次访问属性时才加载对应子模块这让库的启动更快。Highly Available高可用worker 和客户端在连接丢失或失败时都会自动重试部分 broker 还通过 Primary/Primary 或 Primary/Replica 复制提供 HA 能力。这意味着你可以放心地部署多个 worker 与多个 broker 实例单点故障不会让任务系统整体瘫痪。Fast快速项目文档声明单个 Celery 进程每分钟可以处理数百万个任务往返延迟低于毫秒级在使用 RabbitMQ、librabbitmq 与优化配置的前提下。这一性能目标也解释了为什么仓库同时维护了多种消息传输与并发池实现便于针对不同场景做取舍。Flexible灵活Celery 的几乎每个部分都可以被扩展或单独使用自定义池pool实现、序列化器、压缩方案、日志、调度器、消费者、生产者、broker 传输等等。这种可插拔设计在源码目录结构中体现得很直观——celery/concurrency、celery/backends、celery/loaders、celery/schedules.py 都是相对独立、可替换的组件。Celery 支持什么Broker、并发、结果存储与序列化Brokers消息中间件RabbitMQ、Redis功能完备Amazon SQS以及其他实验性传输完整清单见 docs/getting-started/backends-and-brokers/index.rst。Concurrency并发模型Celery 提供多种 worker 并发实现对应源码见 celery/concurrency 目录并发模型说明源码prefork基于 multiprocessing 的多进程模型默认选项celery/concurrency/prefork.pyeventlet基于 eventlet 协程green threadscelery/concurrency/eventlet.pygevent基于 gevent 协程celery/concurrency/gevent.pythread多线程模型celery/concurrency/thread.pysolo单线程模型常用于调试celery/concurrency/solo.py值得注意的实现细节eventlet/gevent 需要尽早完成 monkey-patch。celery/init.py 中的maybe_patch_concurrency会在解析命令行参数如-P eventlet/--pool gevent时在导入任何其他内容之前执行 monkey patch并预先实例化对应的并发池实现。Result Stores结果存储后端Celery 支持把任务状态与返回值存储到多种后端源码实现全部位于 celery/backends 目录AMQPRPC、RedisMemcachedcache.py支持 pylibmc 与纯 Python 的 pymemcache 两种驱动SQLAlchemycelery/backends/database、Django ORMApache Cassandra、ElasticsearchMongoDB、CouchDB、Couchbase、ArangoDBAmazon DynamoDB、Amazon S3Microsoft Azure Block Blob、Microsoft Azure Cosmos DBGoogle Cloud Storage文件系统File systemSerialization序列化序列化格式pickle、json、yaml、msgpack压缩方案zlib、bzip2加密消息签名见 celery/security 模块提供证书、密钥与签名序列化支持。核心特性一览Monitoring监控worker 会持续发出监控事件流内置和外部工具可以用它实时了解集群正在做什么。深入内容见 docs/userguide/monitoring.rst。监控事件的接收与解析实现在 celery/events 模块事件状态模型见 celery/events/state.py。Work-flows工作流借助一组强大的原语——官方称之为 canvas——可以组合出简单到复杂的工作流包括**分组group、链式chain、分块chunking**等。核心实现在 celery/canvas.py使用教程见 docs/userguide/canvas.rst。canvas 中chord、chunks、group、chain、signature等符号从 celery/init.py 的__all__直接导出。Time Rate Limits时间与速率限制你可以控制每秒/每分钟/每小时能执行多少个任务或一个任务允许运行多长时间这些限制可以设为全局默认值、针对特定 worker、或针对单个任务类型。参见 docs/userguide/workers.rst 中关于时间限制与速率限制的章节。Scheduling调度可以用秒数或datetime 对象指定任务的执行时间也可以用周期任务处理重复事件支持简单的interval表达式以及支持分钟、小时、星期几、月内第几天、年内第几月的Crontab 表达式。周期任务调度器由celery beat负责实现见 celery/beat.pyScheduler 类调度表达式定义见 celery/schedules.py使用教程见 docs/userguide/periodic-tasks.rst。Resource Leak Protection资源泄漏防护--max-tasks-per-child选项用于应对用户任务造成的资源泄漏如内存或文件描述符——这类问题往往超出你的控制范围。该选项让 worker 的子进程在执行指定数量的任务后被回收重建。详见 docs/userguide/workers.rst 中--max-tasks-per-child一节。相关命令选项在celery worker --help中有完整列表。User Components用户自定义组件每个 worker 组件都可以定制用户还可以定义额外组件。worker 是通过 bootsteps 构建起来的——bootsteps 是一个依赖图dependency graph允许对 worker 内部机制进行细粒度控制。框架实现在 celery/bootsteps.pyworker 的默认组件集合见 celery/worker/components.py整体装配见 celery/worker/worker.py。与 Web 框架集成Celery 很容易与 Web 框架集成其中一些框架甚至已有现成的集成包框架集成包Pyramidpyramid_celeryPylonscelery-pylonsFlask不需要可直接使用web2pyweb2py-celeryTornadotornado-celeryTrytoncelery_trytonDjango 用户请直接阅读 docs/django/first-steps-with-django.rstDjango 的接入app 创建、shared_task等实现在 celery/contrib/django/task.py。集成包并非必需但它们能让开发更轻松有时还会提供重要钩子——比如在fork(2)时关闭数据库连接。仓库中还提供了可直接参考的示例工程examples/app/myapp.py单应用示例与 examples/django/proj/celery.pyDjango 项目中的 Celery 应用写法。安装 Celery通过 pip 安装Celery 发布在 Python Package IndexPyPI上可用标准 Python 工具安装$ pip install -U CeleryBundles功能包Celery 定义了一组bundles用于一次性安装 Celery 以及某个功能所需的依赖。可以在 requirements 文件或 pip 命令行中用方括号指定多个 bundle 用逗号分隔$ pip install celery[librabbitmq] $ pip install celery[librabbitmq,redis,auth,msgpack]各 bundle 的依赖定义可在仓库 requirements/extras 目录下逐一核对。可用 bundle 清单如下序列化器SerializersBundle用途celery[auth]使用auth安全序列化器celery[msgpack]使用 msgpack 序列化器celery[yaml]使用 yaml 序列化器并发ConcurrencyBundle用途celery[eventlet]使用 eventlet 池celery[gevent]使用 gevent 池传输与后端Transports and BackendsBundle用途celery[librabbitmq]使用 librabbitmq C 库celery[redis]使用 Redis 作为消息传输或结果后端celery[sqs]使用 Amazon SQS 作为消息传输实验性celery[tblib]使用task_remote_tracebacks特性celery[memcache]使用 Memcached 作为结果后端基于 pylibmccelery[pymemcache]使用 Memcached 作为结果后端纯 Python 实现celery[cassandra]使用 Apache Cassandra/Astra DB 作为结果后端DataStax 驱动celery[couchbase]使用 Couchbase 作为结果后端celery[arangodb]使用 ArangoDB 作为结果后端celery[elasticsearch]使用 Elasticsearch 作为结果后端celery[riak]使用 Riak 作为结果后端celery[dynamodb]使用 AWS DynamoDB 作为结果后端celery[zookeeper]使用 Zookeeper 作为消息传输celery[sqlalchemy]使用 SQLAlchemy 作为结果后端受支持celery[pyro]使用 Pyro4 消息传输实验性celery[slmq]使用 SoftLayer Message Queue 传输实验性celery[consul]使用 Consul.io KV 存储作为消息传输或结果后端实验性celery[django]指定 Django 支持所需的最低版本仅作参考通常不建议写进依赖celery[gcs]使用 Google Cloud Storage 作为结果后端实验性celery[gcpubsub]使用 Google Cloud Pub/Sub 作为消息传输实验性从源码安装从 PyPI 下载最新版本https://pypi.org/project/celery/后解压安装$ tar xvfz celery-0.0.0.tar.gz $ cd celery-0.0.0 $ python setup.py build # python setup.py install最后一条命令在未使用虚拟环境时必须以特权用户执行。使用开发版本development versionCelery 开发版还依赖 kombu、amqp、billiard、vine 四个库的开发版本可以通过 pip 安装各自的最新快照$ pip install https://github.com/celery/celery/zipball/main#eggcelery $ pip install https://github.com/celery/billiard/zipball/main#eggbilliard $ pip install https://github.com/celery/py-amqp/zipball/main#eggamqp $ pip install https://github.com/celery/kombu/zipball/main#eggkombu $ pip install https://github.com/celery/vine/zipball/main#eggvine使用 git 方式安装开发版请参见 docs/contributing.rst。完整的安装文档内容同时收录在 docs/includes/installation.txt 中本文的安装章节即基于该文件展开。快速跳转按需查阅官方文档以我想……为索引组织了一批高频入口这里按仓库实际文件给出对应路径方便按需查阅任务与结果获取任务的返回值、内置任务状态、自定义任务状态、任务日志、最佳实践见 docs/userguide/tasks.rst调用任务delay/apply_async等见 docs/userguide/calling.rst给一组任务添加回调chord、把任务拆成若干块chunks见 docs/userguide/canvas.rstWorker 与运维优化 worker见 docs/userguide/optimizing.rst查看运行中的 worker、清空所有消息purge、检查 worker 状态、注册的任务列表、迁移任务到新 broker见 docs/userguide/monitoring.rst运行时修改 worker 队列、编写自定义远程控制命令见 docs/userguide/workers.rst 与 docs/userguide/routing.rst任务重试、跟踪任务开始时间、获取当前任务 ID 与投递队列信息见 docs/userguide/tasks.rst配置与进阶全部配置项参考见 docs/userguide/configuration.rst应用app的概念与创建见 docs/userguide/application.rst事件消息类型列表见 docs/reference/celery.events.rst安全见 docs/userguide/security.rst守护进程化daemonizing见 docs/userguide/daemonizing.rst信号signals见 docs/userguide/signals.rst常见问题见 docs/faq.rstAPI 参考见 docs/reference/index.rst参与贡献见 docs/contributing.rst结语Celery 是一个开箱即用的分布式任务队列最简单的应用只有几行代码但通过 broker、结果后端、并发池、canvas 工作流、周期调度与 bootstep 组件体系它可以支撑从单机脚本到跨数据中心集群的各种规模。上手时建议从 docs/getting-started/first-steps-with-celery.rst 起步完成第一个任务的创建、调用与结果获取再按需深入本文列出的各专题文档安装时优先使用pip install Celery并按实际用到的功能选择对应的 bundle 组合。【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考