Celery 分布式任务队列入门指南:核心概念、特性与安装实践
【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery
Celery 是一个用 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 将消息投递给某个 worker;
- worker 执行任务并(可选地)把结果写回 result backend。
一个 Celery 系统可以由多个 worker 和多个 broker组成,从而支持高可用与水平扩展。Celery 本身用 Python 编写,但通信协议可以在任何语言中实现——除了 Python 之外,社区还提供了 Node.js(node-celery)、PHP 客户端、Go(gocelery、gopher-celery)、Rust(rusty-celery)等实现。语言互操作也可以通过另一种方式达成:暴露一个 HTTP 端点,然后由一个任务去请求它(即 webhook 模式)。
从源码看,这个"生产-投递-执行"的链路分散在多个模块中:任务定义与注册在 celery/app/task.py 与 celery/app/registry.py,消息发布与消费则依赖 kombu/amqp 连接 broker(应用封装见 celery/app/amqp.py),worker 侧的消费入口见 celery/worker/consumer。
运行 Celery 需要什么
版本要求
根据文档侧边栏的说明,Celery 5.x 系列支持的 Python 版本为Python 3.8、3.9、3.10、3.11、3.12、3.13,以及PyPy3.9+(v7.3.12+)。当前仓库是开发分支,celery/init.py 中标明的版本号为5.6.2(系列代号recovery)。
如果你运行的是更老版本的 Python,则需要搭配更老的 Celery:
| Python 版本 | 对应 Celery 版本 |
|---|---|
| Python 3.7 | Celery 5.2 或更早 |
| Python 3.6 | Celery 5.1 或更早 |
| Python 2.7 | Celery 4.x 系列 |
| Python 2.6 | Celery 3.1 系列或更早 |
| Python 2.5 | Celery 3.0 系列或更早 |
| Python 2.4 | Celery 2.2 系列或更早 |
另外需要说明:Celery 是一个资金极少的项目,官方不支持 Microsoft Windows,请勿针对该平台提交 issue。
消息传输(broker)
Celery 需要一个消息传输层来发送和接收消息。RabbitMQ 与 Redis 两种 broker 传输是功能完备(feature complete)的,此外还支持大量实验性方案,包括用于本地开发的 SQLite。Celery 可以运行在单台机器、多台机器,甚至跨数据中心部署。
各 broker 的详细配置请分别参考:
- docs/getting-started/backends-and-brokers/rabbitmq.rst
- docs/getting-started/backends-and-brokers/redis.rst
- docs/getting-started/backends-and-brokers/sqs.rst
- docs/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', broker='amqp://guest@localhost//') @app.task def hello(): return 'hello world'Celery(...)的第一个参数是当前模块名,用于自动生成任务名;broker关键字指定消息中间件的 URL。这里的amqp://guest@localhost//指向本地 RabbitMQ(RabbitMQ 也是默认选项)。
从实现上看,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.py |
| eventlet | 基于 eventlet 协程(green threads) | celery/concurrency/eventlet.py |
| gevent | 基于 gevent 协程 | celery/concurrency/gevent.py |
| thread | 多线程模型 | celery/concurrency/thread.py |
| solo | 单线程模型,常用于调试 | 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 目录:
- AMQP(RPC)、Redis
- Memcached(cache.py,支持 pylibmc 与纯 Python 的 pymemcache 两种驱动)
- SQLAlchemy(celery/backends/database)、Django ORM
- Apache Cassandra、Elasticsearch
- MongoDB、CouchDB、Couchbase、ArangoDB
- Amazon DynamoDB、Amazon S3
- Microsoft Azure Block Blob、Microsoft Azure Cosmos DB
- Google Cloud Storage
- 文件系统(File system)
Serialization(序列化)
- 序列化格式: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.py(Scheduler 类),调度表达式定义见 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.py,worker 的默认组件集合见 celery/worker/components.py,整体装配见 celery/worker/worker.py。
与 Web 框架集成
Celery 很容易与 Web 框架集成,其中一些框架甚至已有现成的集成包:
| 框架 | 集成包 |
|---|---|
| Pyramid | pyramid_celery |
| Pylons | celery-pylons |
| Flask | 不需要(可直接使用) |
| web2py | web2py-celery |
| Tornado | tornado-celery |
| Tryton | celery_tryton |
Django 用户请直接阅读 docs/django/first-steps-with-django.rst;Django 的接入(app 创建、shared_task等)实现在 celery/contrib/django/task.py。集成包并非必需,但它们能让开发更轻松,有时还会提供重要钩子——比如在fork(2)时关闭数据库连接。
仓库中还提供了可直接参考的示例工程:examples/app/myapp.py(单应用示例)与 examples/django/proj/celery.py(Django 项目中的 Celery 应用写法)。
安装 Celery
通过 pip 安装
Celery 发布在 Python Package Index(PyPI)上,可用标准 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 清单如下:
序列化器(Serializers)
| Bundle | 用途 |
|---|---|
celery[auth] | 使用auth安全序列化器 |
celery[msgpack] | 使用 msgpack 序列化器 |
celery[yaml] | 使用 yaml 序列化器 |
并发(Concurrency)
| Bundle | 用途 |
|---|---|
celery[eventlet] | 使用 eventlet 池 |
celery[gevent] | 使用 gevent 池 |
传输与后端(Transports and Backends)
| Bundle | 用途 |
|---|---|
celery[librabbitmq] | 使用 librabbitmq C 库 |
celery[redis] | 使用 Redis 作为消息传输或结果后端 |
celery[sqs] | 使用 Amazon SQS 作为消息传输(实验性) |
celery[tblib] | 使用task_remote_tracebacks特性 |
celery[memcache] | 使用 Memcached 作为结果后端(基于 pylibmc) |
celery[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 version)
Celery 开发版还依赖 kombu、amqp、billiard、vine 四个库的开发版本,可以通过 pip 安装各自的最新快照:
$ pip install https://github.com/celery/celery/zipball/main#egg=celery $ pip install https://github.com/celery/billiard/zipball/main#egg=billiard $ pip install https://github.com/celery/py-amqp/zipball/main#egg=amqp $ pip install https://github.com/celery/kombu/zipball/main#egg=kombu $ pip install https://github.com/celery/vine/zipball/main#egg=vine使用 git 方式安装开发版请参见 docs/contributing.rst。
完整的安装文档内容同时收录在 docs/includes/installation.txt 中,本文的安装章节即基于该文件展开。
快速跳转:按需查阅
官方文档以"我想……"为索引组织了一批高频入口,这里按仓库实际文件给出对应路径,方便按需查阅:
任务与结果
- 获取任务的返回值、内置任务状态、自定义任务状态、任务日志、最佳实践:见 docs/userguide/tasks.rst
- 调用任务(
delay/apply_async等):见 docs/userguide/calling.rst - 给一组任务添加回调(chord)、把任务拆成若干块(chunks):见 docs/userguide/canvas.rst
Worker 与运维
- 优化 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.rst
- API 参考:见 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),仅供参考