Celery 的核心价值不是”把函数放到后台运行”这么简单,而是把工作封装成消息,让生产者与执行者解耦,并允许多个 Worker 跨进程、跨机器并行处理。
它适合解决这些问题:
- Web 请求不应该等待的耗时工作,例如发邮件、生成报表、处理图片。
- 可以并行拆分的批量任务。
- 失败后需要自动重试的外部服务调用。
- 定时或周期性任务。
- 需要分队列、分机器、限速执行的后台作业。
它不直接解决:
- 业务操作的幂等性与数据一致性。
- “绝对只执行一次”。分布式系统中通常应按可能重复执行来设计。
- 任意长时间的精确延时调度。
- 事务提交与任务发布之间的原子性。
1. 总体架构
flowchart LR
subgraph Producers[任务生产者]
Web[Web / API]
CLI[脚本 / CLI]
Beat[Celery Beat]
end
subgraph Messaging[消息系统]
Broker[(Broker
RabbitMQ / Redis)]
end
subgraph Execution[执行系统]
W1[Worker A]
W2[Worker B]
Pool1[执行池
prefork / thread / gevent]
Pool2[执行池]
W1 --> Pool1
W2 --> Pool2
end
Backend[(Result Backend
Redis / DB / rpc)]
Monitor[监控与控制
events / inspect / Flower]
Web -->|发布任务消息| Broker
CLI -->|发布任务消息| Broker
Beat -->|到期后发布| Broker
Broker -->|投递| W1
Broker -->|投递| W2
Pool1 -->|状态 / 结果| Backend
Pool2 -->|状态 / 结果| Backend
W1 -.事件.-> Monitor
W2 -.事件.-> Monitor
1.1 Celery Application
Celery(...) 创建应用对象,它是整个系统的入口,持有:
- Broker 和 Backend 配置。
- 任务注册表。
- 序列化、路由、重试、时区等配置。
- Worker、Beat、Canvas 和控制命令需要的上下文。
1.2 Task
Task 是可以通过消息调用的工作单元。@app.task 会把普通函数转换成 Task 对象,并以唯一任务名注册。
消息里通常不会携带 Python 函数本身,只携带任务名与参数。Worker 必须导入相同任务代码,才能从注册表找到并执行它。
1.3 Broker
Broker 是消息中介:生产者把任务消息发给它,Worker 从队列取消息。
- RabbitMQ:专门的消息代理,功能完整,适合重视消息语义和持久性的生产环境。
- Redis:可同时充当 Broker 和 Backend,部署简单,适合快速上手和大量小消息;必须考虑内存、持久化与
visibility_timeout。
Broker 不负责执行 Python 函数,也不等同于 Result Backend。
1.4 Worker
Worker 是常驻消费者。它连接 Broker,接收消息、找到 Task、提交到执行池,最后确认或拒绝消息。
Worker 主进程通常负责连接、调度、心跳和控制;实际任务由 prefork 子进程或其他并发池执行。
1.5 Result Backend
Backend 保存或传递任务状态、返回值、异常和 traceback。它是可选的:
- 不需要查询结果时可以不配置。
- 配置后可以使用
AsyncResult查询状态和结果。 - 不需要结果的任务应设置
ignore_result=True,减少存储开销。
1.6 Celery Beat
Beat 是调度器,不是执行器。它判断某个周期任务是否到期,到期后把普通任务消息发给 Broker,仍由 Worker 执行。
同一份 schedule 同一时间只能运行一个 Beat,否则会重复投递。
1.7 Kombu 与并发池
Celery 借助 Kombu 统一访问 RabbitMQ、Redis、SQS 等传输。Redis 中的 unacked 等内部结构属于具体 Kombu transport 的实现,不是所有 Broker 的共同数据模型。
prefork 进程池由 Celery/Billiard 管理。其他并发方式包括 thread、gevent、eventlet 和 solo;不同池不一定保留全部特性。
2. 一次任务经历了什么
sequenceDiagram
autonumber
participant C as Client
participant T as Task.apply_async
participant A as Celery App
participant K as Kombu Producer
participant B as Broker
participant W as Worker Consumer
participant P as Execution Pool
participant R as Result Backend
C->>T: add.delay(2, 3)
T->>A: send_task(name, args, kwargs)
A->>A: 生成 task_id、计算路由、构造协议消息
A->>K: publish(body, headers, exchange, routing_key)
K->>B: 发送消息
B-->>C: 发布成功不代表执行成功
T-->>C: AsyncResult(task_id)
B->>W: 投递消息
W->>W: 解析任务名并查注册表
W->>P: 提交 Request
P->>P: 执行 Task.run
alt 成功
P->>R: SUCCESS + return value
else Retry
P->>B: 发布新的重试消息
P->>R: RETRY + exception
else 失败
P->>R: FAILURE + exception + traceback
end
W->>B: 按确认策略 ACK / Reject
有三个容易混淆的”成功”:
- 客户端成功调用
delay():只表示拿到了一个异步结果对象。 - 消息成功发布到 Broker:只表示 Broker 接收了消息。
- Worker 成功执行任务:此时 Backend 才可能记录
SUCCESS。
这三个环节各自的确认机制、失败模式与排查方法见 8.1。
3. 创建一个可以运行的项目
下面使用 Redis 同时充当 Broker 和 Backend,便于学习。生产环境可以把 Broker 换为 RabbitMQ。
3.1 环境要求
Celery 5.6 支持 Python 3.9–3.13;这是最后支持 Python 3.9 的 Celery 系列。Celery 官方不支持 Microsoft Windows,Windows 开发环境建议使用 WSL、Linux 容器或 Linux 虚拟机。
安装:
1 | python -m venv .venv |
启动 Redis:
1 | docker run --name celery-redis -p 6379:6379 -d redis:7 |
3.2 项目结构
1 | celery-demo/ |
celery_demo/celery_app.py:
1 | import os |
celery_demo/tasks.py:
1 | from .celery_app import app |
run_task.py:
1 | from celery_demo.tasks import add |
3.3 启动 Worker
在项目根目录运行:
1 | celery -A celery_demo.celery_app:app worker --loglevel=INFO |
启动日志的 [tasks] 区域应包含:
1 | celery_demo.tasks.add |
另开终端运行:
1 | python run_task.py |
预期结果为 5。
3.4 RabbitMQ 配置
启动 RabbitMQ:
1 | docker run --name celery-rabbitmq -p 5672:5672 -d rabbitmq:4 |
把 Broker URL 改为:
1 | pyamqp://guest@localhost// |
Backend 可以继续使用 Redis。RabbitMQ 主要是 Broker;rpc:// 虽可作为 Backend,但结果通过每个客户端的临时队列返回,语义与持久化 Backend 不同。
4. Task 的注册与调用原理
4.1 @app.task 做了什么
flowchart TD
F[普通 Python 函数 add] --> D[@app.task]
D --> N[生成唯一任务名
celery_demo.tasks.add]
N --> C[动态创建 Task 子类]
C --> R[把原函数设为 run]
R --> I[实例化 Task]
I --> G[写入 app.tasks 注册表]
G --> B[绑定 Application 配置]
因此 Task 更接近”全局注册对象”,而不是每次调用都新建的函数包装器。不要把请求级可变状态保存在 Task 实例属性中。
4.2 三种调用方式
1 | # 当前进程同步执行,不发送消息 |
delay() 不能设置 countdown、eta、expires、queue、priority 等选项。
4.3 AsyncResult 不是后台线程
AsyncResult 只是任务 ID 与 Backend 的查询句柄:
1 | result = add.delay(2, 3) |
注意:
get()会阻塞当前调用方。- 在 Web 请求中频繁
get()会抵消异步处理的价值。 - 不要在一个 Celery Task 中用
.delay().get()等待另一个任务,Worker 池耗尽时可能死锁;应使用 Canvas。 - Backend 会消耗资源。对每个需要结果的
AsyncResult,最终应调用get()或forget();不需要结果就设置ignore_result=True。
5. Celery 消息协议与发布过程
任务消息不是”序列化后的函数”,而是”任务名称 + 参数 + 执行元数据”。协议 v2 可抽象成:
flowchart LR
subgraph Headers[Headers]
H1[task / id]
H2[eta / expires]
H3[retries / timelimit]
H4[root_id / parent_id]
H5[group / origin / stamps]
end
subgraph Body[Body]
B1[args]
B2[kwargs]
B3[callbacks / errbacks]
B4[chain / chord]
end
Headers --> S[Serializer
默认 JSON]
Body --> S
S --> K[Kombu Producer]
K --> Q[Exchange / Routing key / Queue]
发布的大致步骤:
- 校验参数形状。
- 生成或接受自定义
task_id。 - 将
countdown转换为绝对eta。 - 根据
task_routes或调用参数选择 queue、exchange、routing key。 - 生成消息 headers 与 body。
- 用 JSON 等 serializer 编码。
- 通过 Kombu Producer 发布到具体 Broker。
- 返回
AsyncResult(task_id)。
默认优先使用 JSON。不要为了方便随意开放 pickle:反序列化不可信 pickle 数据可以执行任意代码。
参数也应保持小而稳定。传数据库对象时,通常只传主键,让 Worker 自己重新查询;不要发送巨大的二进制数据或依赖 Python 进程内对象。
6. Worker 内部结构
flowchart TD
B[(Broker)] --> C[Consumer]
C --> Decode[解析消息协议]
Decode --> Lookup{任务名已注册?}
Lookup -- 否 --> Reject[Unknown task / Reject]
Lookup -- 是 --> Strategy[Task Strategy]
Strategy --> Check{撤销或过期?}
Check -- 是 --> Skip[跳过并更新状态]
Check -- 否 --> ETA{有 ETA?}
ETA -- 是 --> Timer[Worker Timer 等待]
Timer --> Reserve[标记 reserved]
ETA -- 否 --> Rate[限速 / QoS / prefetch]
Rate --> Reserve
Reserve --> Request[Request 生命周期对象]
Request --> Pool[执行池]
Pool --> Trace[build_tracer / trace_task]
Trace --> User[Task.run 用户代码]
User --> Outcome{执行结果}
Outcome -->|成功| Success[SUCCESS]
Outcome -->|Retry| Retry[RETRY + 再发布]
Outcome -->|异常| Failure[FAILURE]
Success --> Backend[(Backend)]
Retry --> Backend
Failure --> Backend
6.1 Consumer
Consumer 从消息 header 读取任务名,在 Worker 的任务策略表中查找对应 Task。Worker 没有导入任务模块时,会出现 Received unregistered task。
6.2 Strategy 与 Request
Strategy 负责把 Broker 消息转换成 Request,检查 ETA、撤销、过期和 rate limit,然后提交给 Worker Controller。
Request 保存本次调用的 ID、参数、重试次数、投递信息、确认回调、时间限制等,并处理成功、失败、超时和 Worker 子进程丢失。
6.3 Execution Pool
默认 prefork 使用多个子进程执行任务:
- 适合 CPU 密集任务,并提供较完整的 Celery 特性。
--concurrency默认与 CPU 数量有关,但最佳值必须压测。- I/O 密集任务也可以考虑 thread、gevent 或 eventlet,但切换池之前应确认 soft timeout、信号和第三方库兼容性。
solo单线程执行,常用于调试;运行任务时也会阻塞控制命令。
6.4 Trace
真正包围用户代码的是 trace 层。它负责:
- 建立
Task.request上下文。 - 发送 prerun、postrun、success、failure、retry 等信号。
- 捕获
Retry、Ignore、Reject和普通异常。 - 调用 callback、errback、chain、chord。
- 把最终状态和结果写入 Backend。
- 清理 traceback 引用,避免长期 Worker 的内存泄漏。
6.5 多个 Worker:并行抢占,不是串行
多个 Worker(无论同一台机器还是不同机器)连接同一个 Broker、消费同一个队列时,它们之间是并行抢占关系,而不是串行接力:
flowchart LR
Q[(同一个队列)] -->|投递| W1[Worker A]
Q -->|投递| W2[Worker B]
Q -->|投递| W3[Worker C]
W1 --> T1[执行任务 1]
W2 --> T2[执行任务 2]
W3 --> T3[执行任务 3]
这是经典的**竞争消费者(Competing Consumers)**模型:
- 每条消息只投递给一个 Worker:Broker 保证一条任务消息在正常情况下只被一个消费者接收,不会被多个 Worker 重复处理。
- 不同 Worker 同时执行不同任务:多个 Worker 并行从队列取消息,谁先取到谁执行,因此整体是并行处理;总吞吐 ≈ Worker 数量 × 每个 Worker 的并发槽位。
- 没有严格的全局顺序:并行抢占意味着”先发布”不等于”先完成”。完成顺序取决于任务耗时、Worker 负载和 prefetch 情况,不要依赖默认队列保证顺序。
- 单个任务仍是串行的:一个任务只由一个 Worker 完整执行,不会被拆分到多个 Worker 并行;Worker 之间也不存在”等上一个任务完成”的约束。
需要串行或保序时,不能靠多 Worker 默认队列,具体方案见 6.6。
6.6 如何保证任务顺序
顺序需求先分清范围,再选方案。默认的并行抢占只保证”每条消息只执行一次”,不保证执行顺序。
方案一:全局严格 FIFO(所有任务按发布顺序执行)
代价是放弃并行度:
1 | celery -A celery_demo.celery_app:app worker -Q order -n order@%h --concurrency=1 |
- 一个专用队列只由一个 Worker 消费,且
--concurrency=1(或solo池)。 - 队列 FIFO + 单消费者单槽位 ⇒ 严格按序执行。
- 吞吐被压到串行水平,只适合”数量少但顺序敏感”的任务,例如流水线式对账。
方案二:按业务 key 保序(最常见)
同一用户/订单的多个事件不能乱序,但不同用户之间可以并行:
1 | # 发布时按 key 选队列:同 key 永远进同一队列 |
1 | # 每个队列配一个 --concurrency=1 的专用 Worker |
- 同 key 始终进入同一队列、由同一 Worker 串行处理 ⇒ 同 key 有序;不同 key 落在不同 Worker ⇒ 跨 key 并行。
- RabbitMQ 场景可以配合 consistent hash exchange,避免手写取模。
方案三:版本号校验(容忍乱序到达)
消息带递增序号或版本,执行时与数据库中的 last_version 比较:
1 |
|
- 乱序到达的旧事件被忽略,新事件先到也可以暂存等待缺口补齐。
- 配合”重试补缺口”(
self.retry(countdown=...))可以应对短暂乱序,但要有最大重试次数防死循环。 - 适合事件流、状态机这类”最终一致即可、可按版本推进”的业务。
方案四:工作流内部分步保序
一次业务的多步依赖用 Canvas 表达,而不是自己串:
1 | workflow = chain( |
chain 上一步完成后再发布下一步,天然保序;避免在任务内 .delay().get() 等待(会占死 Worker 槽位)。
顺序会被哪些机制破坏
- 失败重试(
retry/autoretry)会把消息重新入队,可能插到其他任务前面。 acks_late=True时 Worker 崩溃导致的重新投递,会打乱顺序。- Redis 的
visibility_timeout到期后消息重新入队。 - 多个 Beat 实例重复投递同一批周期任务。
结论:顺序与并行度天然冲突,不存在”既有全局限序又有高吞吐”的免费方案。先把必须保序的范围缩到最小(通常是单个业务实体),再用”按 key 分片 + 单并发消费”或”版本号校验”解决,而不是要求整个系统全局有序。
7. 任务状态与 Result Backend
stateDiagram-v2
[*] --> PENDING: 已发布或 Backend 不认识此 ID
PENDING --> STARTED: track_started=True
PENDING --> SUCCESS: 执行成功
PENDING --> FAILURE: 执行失败
PENDING --> RETRY: 请求重试
STARTED --> SUCCESS
STARTED --> FAILURE
STARTED --> RETRY
RETRY --> STARTED
RETRY --> SUCCESS
RETRY --> FAILURE
PENDING --> REVOKED: 被撤销
STARTED --> REVOKED: 终止/撤销路径
SUCCESS --> [*]
FAILURE --> [*]
REVOKED --> [*]
7.1 状态含义
PENDING:等待执行,或者 Backend 完全不知道这个 ID。不能仅凭它判断任务仍在队列。STARTED:Worker 已开始执行;默认不上报,需开启task_track_started。SUCCESS:成功,并可能保存返回值。FAILURE:失败,通常保存异常与 traceback。RETRY:已安排重试。REVOKED:已撤销。
7.2 Backend 如何写结果
执行成功后,trace 层调用 Backend 的 mark_as_done();失败和重试分别进入 mark_as_failure() 与 mark_as_retry(),最终由 store_result() 编码并写入具体存储。
如果没有 Backend,Celery 使用禁用 Backend:任务照常执行,但无法查询持久状态或结果。
7.3 Backend 选择
- Redis:读取快,适合短期结果;关注内存、过期和持久化。
- SQLAlchemy/Django DB:便于审计和长期保存,但频繁轮询数据库成本较高。
rpc://:通过消息返回状态,结果通常只有发起任务的客户端能消费,且默认是瞬态消息。
没有一种 Backend 适合所有场景。只把真正需要查询的结果写入 Backend,并配置合理的 result_expires。
8. 最重要的可靠性问题:ACK 与幂等
8.1 一条消息的三个确认点:发送、存储、消费
“客户端正常发送 → Broker 正常存储 → Worker 正常消费”是三个独立的环节,各有各的确认机制和失败模式,不能混为一谈。
flowchart LR
C[客户端] -->|1. 发布| B[Broker]
B -->|2. 存储| B
B -->|3. 投递 + ACK| W[Worker]
第一步:客户端 → Broker(发布确认)
delay()/apply_async()是同步发布:Kombu 建立连接并发送消息。连接失败或 Broker 拒绝时delay()会抛异常(如OperationalError),不会静默丢失。- 发布重试由
task_publish_retry(默认开启)与task_publish_retry_policy控制重试次数和退避。 - RabbitMQ 的 publisher confirms:开启
confirm_publish后,Broker 会回传确认,表示消息已进入 Broker。注意 confirm 只代表”消息被接收”,不代表”已写入磁盘”,更不代表”会被执行”。 - Redis 发布:消息
LPUSH进 list 即返回成功;是否落盘取决于 Redis 自身的 RDB/AOF 配置。
发布端完整配置(RabbitMQ 场景):
1 | app.conf.update( |
开启 confirms 后,delay() 只有在 Broker 确认接收后才返回;确认超时或 Broker 拒绝时会抛异常,配合 task_publish_retry 重试。但confirm 成功 ≠ 已写入磁盘 ≠ 会被执行。
第二步:Broker 存储
- RabbitMQ:队列
durable=True且消息 persistent(delivery_mode=2)时,Broker 重启后消息不丢;否则重启即丢。 - Redis:消息存在 list 中,持久化依赖 Redis 配置;同时 Redis transport 用
unacked结构跟踪”已取走但未确认”的消息(见第三步)。 - 消息在队列里不代表会被执行:可能没有 Worker 消费该队列,也可能已经
expires过期。
Broker 存储端配置(RabbitMQ 场景):
1 | app.conf.update( |
队列 durable + 消息 persistent 组合才能保证”Broker 重启后消息还在”。只声明 durable 队列而消息是 transient,重启后消息仍会丢。
确认”存储成功”:publisher confirms、查看队列长度(RabbitMQ 管理台 / XLEN / inspect)、对队列积压设置告警。
第三步:Broker → Worker(消费确认,即 ACK)
ACK 是 AMQP 协议层面的确认。RabbitMQ 和 Redis 两套实现:
- RabbitMQ(AMQP 原生):Worker 处理完消息后调用
basic_ack,Broker 才删除消息。确认时机由acks_late决定:- 默认
acks_late=False:收到消息执行前就 ACK。执行中崩溃 → 消息已确认,不会重投 → 可能丢失。 acks_late=True:执行成功后才 ACK。执行中崩溃 → 未确认 → Broker 重新投递 → 可能重复执行。
- 默认
- Redis(Kombu transport 模拟):Worker 取消息时把消息从 list 移到
unacked集合并记录时间戳;ACK 时从unacked删除。若超过visibility_timeout(默认 1 小时)仍未确认,消息会被恢复重新入队,其他 Worker 可再次消费。因此 Redis 场景下即使acks_late=False,Worker 崩溃也可能因超时重投。
确认”消费成功”:任务状态写入 Backend(SUCCESS)、Worker 日志、inspect active / reserved、Flower 监控。
三个环节的失败点对照
| 环节 | 确认机制 | 典型失败 |
|---|---|---|
| 客户端发布 | delay() 不抛异常 / publisher confirms |
连接断开、Broker 拒绝;发布重试耗尽 |
| Broker 存储 | durable 队列 + persistent 消息 | Broker 重启丢未持久化消息;队列无人消费;消息过期 |
| Worker 消费 | ACK(acks_late 决定时机) |
提前 ACK 后崩溃 → 丢失;延后 ACK 崩溃 → 重投重复 |
端到端实践建议
- 发布端:捕获发布异常并补偿;开启 confirms;不要把”发布成功”当”执行成功”。
- Broker:RabbitMQ 用 durable + persistent;Redis 开启持久化并调好
visibility_timeout。 - Worker:要求不丢选
acks_late=True+ 幂等;要求不重复则靠唯一键/幂等表兜底。 - 验收链路:发布不抛错 → 队列长度符合预期 → Worker 执行日志 → Backend 状态为
SUCCESS。
8.2 默认是提前确认
Celery 默认 acks_late=False,Worker 在执行前确认消息。这可以避免一个已经开始的任务因 Worker 异常而自动重复执行,但 Worker 在确认后、完成前崩溃时,任务可能丢失。
设置 acks_late=True 后,成功执行完才确认;Worker 中途故障时消息可能重新投递,因此任务可能执行多次。
flowchart TD
M[Worker 收到任务] --> Late{acks_late?}
Late -- 否 --> EarlyAck[执行前 ACK]
EarlyAck --> Run1[执行任务]
Run1 --> Crash1{中途崩溃?}
Crash1 -- 是 --> Lost[通常不会重投
可能丢失]
Crash1 -- 否 --> Done1[完成]
Late -- 是 --> Run2[先执行任务]
Run2 --> Crash2{中途崩溃?}
Crash2 -- 否 --> Ack2[执行后 ACK]
Ack2 --> Done2[完成]
Crash2 -- 是 --> LostPolicy{reject_on_worker_lost
及 Broker 行为}
LostPolicy --> Redeliver[可能重新投递]
Redeliver --> Duplicate[可能重复执行]
即使启用 acks_late,子进程被信号终止时,Celery 默认仍可能确认消息。如果确实要在 Worker 子进程丢失时重新入队,可研究 task_reject_on_worker_lost;官方警告它可能制造无限消息循环。
8.3 什么是幂等任务
同一参数执行一次和多次,业务最终结果相同,且不会产生非预期副作用。
非幂等例子:
1 |
|
Worker 可能已经扣款,但在 ACK 前崩溃;重投后再次扣款。
一种常见防重方式是使用稳定的业务键:
1 |
|
真正可靠的实现还需要数据库唯一约束、事务边界和外部服务的幂等键共同配合,不能只靠一行 Python 判断。
8.4 交付语义的现实理解
flowchart LR
AtMost[提前 ACK 倾向] -->|故障窗口| Loss[可能丢失]
AtLeast[延后 ACK 倾向] -->|重投窗口| Dup[可能重复]
Dup --> Idempotency[业务幂等 / 去重]
Loss --> Reconcile[补偿 / 对账 / 重建任务]
不要宣称 Celery 为业务操作提供”恰好一次”。应根据任务价值选择:
- 可以接受偶尔丢失、不允许重复:可能保留默认提前确认。
- 不允许丢失、可以通过幂等处理重复:考虑
acks_late=True。 - 涉及钱、库存、权益:还要有唯一约束、事务、对账与补偿机制。
9. 重试、超时和失败控制
9.1 手动重试
1 | import requests |
self.retry() 的实现不是在当前栈重新调用函数,而是:
sequenceDiagram
participant Task as 当前 Task
participant Retry as self.retry
participant Broker
participant Trace
participant Backend
Task->>Retry: 捕获临时异常
Retry->>Retry: retries + 1,复制当前 Signature
Retry->>Broker: apply_async 再发布
Retry-->>Trace: 抛出特殊 Retry 异常
Trace->>Backend: 写入 RETRY
Broker-->>Task: 到期后再次投递
重试通常沿用同一个任务 ID,因此同一个 AsyncResult 可以观察其状态变化。
9.2 自动重试与指数退避
1 | import requests |
指数退避减轻下游故障时的压力;jitter 让大量任务不要在同一瞬间重试。
不要对所有异常盲目永久重试:参数错误、权限错误、数据不存在等永久错误重试也不会成功。
9.3 I/O 超时与 Celery Time Limit
优先给每个 I/O 操作设置库级超时,例如 HTTP 的连接超时与读取超时。Celery 的:
soft_time_limit:尝试在任务中抛出SoftTimeLimitExceeded,允许清理。time_limit:硬限制,可能终止执行进程。
硬终止会影响 ACK、事务和资源清理,不应代替正常的网络与数据库超时。
9.4 任务过期
1 | send_notification.apply_async(args=(user_id,), expires=60) |
expires 表示超过时间后不再执行,适合已经失去业务价值的任务。它不是”最多运行 60 秒”;运行时长限制应使用 time limit。
10. ETA、Countdown 与周期任务
10.1 短延时
1 | add.apply_async(args=(2, 3), countdown=10) |
countdown=10 表示最早在 10 秒后执行,不保证恰好在第 10 秒执行;队列拥堵和 Worker 负载都会造成延迟。
sequenceDiagram
participant Client
participant Broker
participant Worker
participant Timer as Worker Timer
participant Pool
Client->>Broker: 发布 eta/countdown 任务
Broker->>Worker: Worker 提前取走消息
Worker->>Timer: 保存在内存直到 ETA
Note over Worker,Timer: 尚未开始执行,通常也尚未确认
Timer->>Pool: 到时提交执行
大量远期 ETA/countdown 任务会占用 Worker 内存。使用 Redis 时,超过 visibility_timeout 还可能重新投递;RabbitMQ 也有 consumer acknowledgement timeout。官方建议 ETA/countdown 主要用于数分钟内的短延时。
10.2 周期任务
在 celery_app.py 中添加:
1 | from celery.schedules import crontab |
任务定义:
1 |
|
分别启动:
1 | celery -A celery_demo.celery_app:app worker -l INFO |
Beat 原理:
flowchart TD
Load[加载 beat_schedule] --> Heap[按下次触发时间建立最小堆]
Heap --> Tick[tick 检查堆顶]
Tick --> Due{是否到期?}
Due -- 否 --> Sleep[返回建议休眠时间]
Sleep --> Tick
Due -- 是 --> Reserve[推进该条目的下次运行时间]
Reserve --> Publish[apply_async 发布普通任务]
Publish --> Broker[(Broker)]
Broker --> Worker[Worker 执行]
Reserve --> Heap
Beat 不检查上一次任务是否完成。若执行时间超过周期,任务会重叠;需要业务锁或其他单实例机制。
11. 路由与多队列
把不同性质的任务放在同一队列,会产生”短任务被长任务堵住”的问题。
flowchart LR
Producer[生产者] --> Router{task_routes}
Router -->|emails.*| EmailQ[(email queue)]
Router -->|reports.*| ReportQ[(report queue)]
Router -->|default| DefaultQ[(default queue)]
EmailQ --> EmailW[高并发 I/O Worker]
ReportQ --> ReportW[低并发重任务 Worker]
DefaultQ --> DefaultW[通用 Worker]
配置:
1 | app.conf.task_routes = { |
启动指定队列的 Worker:
1 | celery -A celery_demo.celery_app:app worker -l INFO -Q email -n email@%h |
也可以在单次调用时指定:
1 | generate_report.apply_async(args=(report_id,), queue="reports") |
分队列的价值:资源隔离、独立扩容、独立并发数和故障隔离。它不是权限隔离;安全仍依赖网络、Broker 账户和访问控制。
12. Canvas:组合任务而不是同步等待
12.1 Signature
1 | sig = add.s(2, 3) |
Signature 是可序列化的任务调用描述,包含任务、参数和执行选项。.si() 创建 immutable signature,不接收上一步结果。
12.2 Chain
1 | from celery import chain |
flowchart LR
A[add 2,3] -->|结果 5| B[add 5,10]
B -->|结果 15| R[最终结果]
12.3 Group
1 | from celery import group |
flowchart LR
Start[Group] --> T1[add 0,0]
Start --> T2[add 1,1]
Start --> T3[...]
Start --> T4[add 9,9]
12.4 Chord
Chord 是”一组并行任务全部完成后执行回调”:
1 | from celery import chord |
flowchart LR
S[Chord Header] --> A[Task A]
S --> B[Task B]
S --> C[Task C]
A --> Barrier{全部完成?}
B --> Barrier
C --> Barrier
Barrier -->|是| Callback[Chord Body / Callback]
Chord 依赖 Result Backend 协调组结果,并非所有 Backend 能力完全相同。
12.5 为什么不用 .get() 串联
坏的设计:
1 |
|
这会让当前 Worker 槽位阻塞等待其他 Worker 槽位。池中所有槽位都这样等待时会死锁。使用 chain 会在上一步完成后发布下一步,不占用等待中的执行槽位。
13. Prefetch、并发与公平性
Worker 可以提前保留多条消息。prefetch 上限通常与并发数及 worker_prefetch_multiplier 有关。
flowchart LR
Q[(Queue)] -->|prefetch| W1[Worker A
已保留多条长任务]
Q -->|prefetch| W2[Worker B]
W1 --> P11[Pool slot 1]
W1 --> P12[Pool slot 2]
W1 --> Wait[本地等待的 reserved tasks]
过度 prefetch 会让任务集中到少数 Worker,尤其影响长任务公平性。需要结合任务时长、并发池和 Broker 进行压测。
常用控制:
1 | app.conf.update( |
worker_prefetch_multiplier=1常用于长任务,减少每个进程预留的额外任务,但不是通用性能最优值。worker_max_tasks_per_child定期替换子进程,缓解第三方库资源泄漏。worker_max_memory_per_child在子进程超过内存阈值后替换它。
不要照抄并发数。CPU 密集、I/O 密集、任务耗时分布和外部服务容量都会改变最佳配置。
14. 监控、检查与控制
14.1 基础命令
1 | celery -A celery_demo.celery_app:app status |
RabbitMQ 和 Redis 支持事件监控与远程控制;并非所有 Broker transport 都支持。
14.2 事件与 Flower
启动事件:
1 | celery -A celery_demo.celery_app:app worker -l INFO -E |
Flower 是常见的 Web 监控工具,但监控界面不能代替指标和告警。至少应观察:
- 队列长度和最老任务等待时间。
- 发布失败率、执行成功率和重试率。
- 任务运行时间分位数。
- Worker 在线数、并发槽位和心跳。
- Broker 连接、内存、磁盘与持久化健康度。
- Backend 容量、错误和结果过期情况。
14.3 Revoke 不是安全取消
1 | celery -A celery_demo.celery_app:app control revoke TASK_ID |
Revoke 主要让 Worker 跳过尚未执行的任务。terminate=True 会杀死正在执行任务的进程,是管理员处理卡死任务的最后手段;目标进程可能已开始处理其他任务,因此官方明确不建议在程序逻辑中自动使用。
15. Worker 关闭与发布上线
优先使用 TERM 触发 warm shutdown:Worker 停止接收工作,并等待当前任务完成。不要把 kill -9 当正常部署方式。
flowchart TD
Signal[收到信号] --> Kind{信号 / 阶段}
Kind -->|TERM 或第一次 Ctrl+C| Warm[Warm shutdown
等待运行中任务完成]
Kind -->|QUIT| Cold[Cold shutdown
停止运行中任务]
Warm --> Done[正常退出]
Cold --> Soft{启用 soft shutdown?}
Soft -- 是 --> Grace[等待有限时间后终止]
Soft -- 否 --> Stop[立即终止]
Celery 5.5 起提供可配置的 soft shutdown。使用 ETA 任务时还要关注 worker_enable_soft_shutdown_on_idle,否则 Worker 看似空闲但内存中仍有 ETA 任务,冷关闭可能造成任务风险。
生产环境应使用 systemd、容器编排器或其他进程监督系统管理 Worker,并确保:
- 停止宽限期大于正常任务完成时间。
- readiness 与 liveness 检查不会制造重启循环。
- Worker 进程以非 root 账户运行。
- 发布新版本时,新 Worker 已加载与消息任务名兼容的代码。
16. 安全原则
Celery 的 Broker 是系统信任边界。能够向队列写消息的主体,可能让 Worker 调用已注册任务。
必须做到:
- Broker 不暴露在公网,使用防火墙、私网和最小权限账户。
- 生产环境使用 TLS,并正确校验证书。
- 使用 JSON 等安全 serializer;不要接受不可信 pickle。
- 配置
accept_content=["json"],限制 Worker 可反序列化的内容。 - 不在任务参数、日志、
argsrepr和kwargsrepr中泄露密码或令牌。 - Worker 使用低权限系统用户,并隔离网络、文件和云权限。
- 对任务参数做业务鉴权与校验;”任务已注册”不代表调用者有权执行。
- 定期更新 Celery、Kombu、Broker 和依赖。
序列化签名可以验证消息来源与完整性,但不提供加密;敏感数据仍应通过 TLS 和应用层控制保护。
17. 常见故障排查
17.1 Received unregistered task
原因:Worker 没有导入任务模块,或生产者和 Worker 使用的任务名/版本不同。
检查:
1 | celery -A celery_demo.celery_app:app inspect registered |
确认 include、autodiscover、包路径和 Worker 启动目录正确,并在修改代码后重启 Worker。
17.2 一直是 PENDING
依次检查:
- 是否配置 Result Backend。
- 生产者和 Worker 是否使用同一个 Backend URL 与数据库编号。
- Worker 是否实际收到任务。
- 任务是否设置
ignore_result=True。 - Backend 是否可连接、结果是否已过期。
- Task ID 是否真实存在;未知 ID 也显示
PENDING。
17.3 delay() 很慢或报连接错误
delay() 本身需要连接 Broker 并发布消息。检查 DNS、端口、TLS、凭据、连接池和 Broker 健康度。任务发布重试与任务执行重试是两套机制,不要混淆。
17.4 任务重复执行
检查:
- 是否启用了
acks_late。 - Worker 是否在执行中崩溃或丢失 Broker 连接。
- Redis
visibility_timeout是否小于任务或 ETA 时长。 - Beat 是否启动了多个实例。
- 调用方是否因发布结果不确定而重复发布。
- 业务是否有唯一键和幂等处理。
17.5 任务丢失
检查:
- 是否在任务开始前 ACK,随后 Worker 崩溃。
- Broker 是否启用持久化,队列和消息是否 durable。
- 是否使用
kill -9、冷关闭或过短的容器停止宽限期。 - 消息是否过期、被 revoke、路由到无人消费的队列。
- 发布方是否把”调用
delay()没抛错”误当成业务完成。
17.6 Worker 内存持续增长
检查:
- 大量远期 ETA/countdown 任务。
- 过大的消息或结果。
- 任务代码和第三方库泄漏。
- traceback、全局缓存和数据库连接未释放。
- prefetch 过高。
可以结合 worker_max_tasks_per_child、worker_max_memory_per_child 缓解,但仍应找到根因。
17.7 队列有任务但处理很慢
检查 active、reserved、scheduled 和 stats;区分是 Worker 数量不足、长任务阻塞、下游限流、prefetch 不公平、任务重试风暴,还是某个队列根本没有 Worker 消费。
18. 一份较稳妥的基础配置
以下是起点,不是所有系统的最终答案:
1 | app.conf.update( |
任务级可靠性应逐个设置,而不是全局盲目开启:
1 |
|
19. 从开发到生产的检查清单
任务设计
- 参数小、可 JSON 序列化,只传必要 ID。
- 明确任务是否幂等,重复执行有什么后果。
- 永久错误与临时错误分开处理。
- 所有网络和数据库 I/O 有超时。
- 不需要结果的任务使用
ignore_result=True。 - 不在 Task 内同步等待另一个 Task。
消息与可靠性
- 根据丢失与重复风险选择 ACK 策略。
-
acks_late=True的任务有数据库或外部服务级幂等保护。 - 配置发布失败处理;不要把发布成功等同于任务成功。
- Redis 的
visibility_timeout与最长任务/ETA 相匹配。 - 远期调度不大量使用 ETA/countdown。
Worker 与部署
- 不同资源类型的任务拆分队列和 Worker。
- 并发数和 prefetch 通过压测确定。
- 使用 warm shutdown 和足够的停止宽限期。
- Worker 由进程监督系统管理,以非 root 账户运行。
- Beat 同一 schedule 只有一个实例。
监控
- 监控队列长度、等待时间、运行时间、失败率和重试率。
- 监控 Worker 心跳、Broker 与 Backend 容量。
- 对任务长期堆积和重试风暴设置告警。
- 有补偿、重放或对账工具,而不只依赖日志。
20. 源码阅读地图
flowchart TD
Base[celery/app/base.py
Application、装饰器、任务注册、send_task]
Task[celery/app/task.py
Task、delay、apply_async、retry]
AMQP[celery/app/amqp.py
消息协议、路由、publish]
Kombu[Kombu
具体 Broker transport]
Consumer[celery/worker/consumer/consumer.py
消息消费、任务名查找]
Strategy[celery/worker/strategy.py
ETA、限速、Request 构造]
Request[celery/worker/request.py
执行池、ACK、失败处理]
Trace[celery/app/trace.py
执行包装、状态、回调、信号]
Backend[celery/backends/base.py
结果状态统一接口]
Beat[celery/beat.py
周期调度最小堆与发布]
Base --> Task --> AMQP --> Kombu --> Consumer --> Strategy --> Request --> Trace --> Backend
Beat --> AMQP
建议按图中的顺序读源码。先理解一条普通任务链路,再阅读 Canvas、具体 Backend 和具体 Kombu transport,否则容易陷入实现细节。