Python 基础体系 · 第 87/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。
Celery 任务队列:Broker、Worker、ACK、重试、定时和幂等
Celery 解决的问题是:把一次函数调用转换成一条消息,再由独立进程异步执行这条消息。
在同步 Web 请求中,调用函数通常意味着:
HTTP 请求进程
└── 直接执行函数
如果函数需要访问外部 API、生成报表、发送邮件或处理大量数据,请求进程会被长时间占用。Celery 将流程拆成:
Web 进程 ──发布任务消息──> Broker ──投递任务消息──> Worker
│ │
└───────────────查询任务状态/结果 <────────────┘
其中:
- Broker:任务消息的中转服务;
- Worker:真正执行任务的 Celery 进程;
- ACK:Worker 向 Broker 确认“这条消息已经被接收并处理到某个阶段”;
- 重试:任务失败后重新安排一次执行;
- 定时:由
celery beat按时间规则发布任务; - 幂等:同一任务执行一次或多次,最终业务结果都符合预期。
Celery 当前稳定文档对应 5.6 系列;本文示例使用 Python 3.14 语法,但 Celery、Broker 驱动和 Python 3.14 的具体兼容矩阵仍应以实际安装时的官方发布说明为准。Python 3.14 官方文档当前为 3.14.7。(docs.celeryq.dev)
一、先建立正确模型:Celery 传递的是消息,不是函数调用
1. 任务调用实际上经历了两个阶段
下面的代码看起来像是在调用函数:
result = send_email.delay(
user_id=42,
template="welcome",
)
但 delay() 并不会在当前进程中执行 send_email()。它只是将任务名称、参数和任务 ID 序列化成一条消息,发布到 Broker。
可以把它形式化为:
调用阶段:
f(args, kwargs)
↓
publish(task_name, task_id, args, kwargs)
↓
返回 AsyncResult(task_id)
执行阶段:
Worker 收到消息
↓
根据 task_name 查找已注册任务
↓
调用 f(args, kwargs)
Celery 任务消息通常包含任务名称、任务 ID、位置参数、关键字参数,以及重试次数、ETA 等元数据。Broker 负责在生产者和消费者之间路由消息。(docs.celeryq.dev)
因此,下面两种代码的语义完全不同:
# 同步调用:当前进程执行
send_email(user_id=42, template="welcome")
# 异步调用:当前进程只发布消息
send_email.delay(user_id=42, template="welcome")
delay() 返回的是任务结果的引用,而不是任务函数的返回值:
async_result = send_email.delay(user_id=42)
print(async_result.id)
# 例如:'9d8c...'
print(async_result.ready())
# 通常为 False
如果需要保存任务状态或返回值,必须配置 Result Backend。Broker 与 Result Backend 是两个不同的角色:
| 组件 | 主要职责 | 典型数据 |
|---|---|---|
| Broker | 传递待执行消息 | 队列中的任务消息 |
| Worker | 消费并执行消息 | 任务执行过程 |
| Result Backend | 保存执行状态和结果 | PENDING、STARTED、SUCCESS、FAILURE |
| Beat | 按时间发布任务 | 周期性任务消息 |
不要把“任务已经发布到 Broker”误认为“任务已经执行成功”。
二、Broker:任务消息的中转站
1. Broker 保存的是“待处理消息”
假设 Web 进程发布了任务:
generate_report.delay(report_id=1001)
消息进入 Broker 后,典型状态是:
READY
↓ Worker 读取
UNACKNOWLEDGED
↓ Worker 完成并确认
REMOVED / ACKED
在消息被 ACK 之前,Broker 通常仍认为这条消息没有完成。Celery 文档明确指出,任务消息在被 Worker 确认前不会从队列中移除;如果 Worker 在此期间失效,消息可能重新投递。(docs.celeryq.dev)
Broker 不是数据库意义上的业务事实存储。它解决的是:
“现在有没有一个消费者可以接收这条消息?”
而不是:
“这个订单最终是否已经支付?”
订单支付状态应由业务数据库保存,不能只依赖 Celery 的任务状态。
2. RabbitMQ、Redis 和 Kafka 不是同一类抽象
Celery 常用 RabbitMQ 或 Redis 作为 Broker。当前 Celery Worker 文档列出的 Broker 支持包括 AMQP 和 Redis。(docs.celeryq.dev)
RabbitMQ
RabbitMQ 基于 AMQP 模型,通常包含:
Producer
↓
Exchange ──根据 routing key──> Queue
↓
Consumer
Celery 使用的消息会经过 Exchange、Queue 和路由键。Exchange 决定消息进入哪些队列,Queue 保存待消费消息,Worker 作为 Consumer 从 Queue 获取消息。
适合关注以下语义的场景:
- 明确的队列和路由;
- ACK、拒绝、重新入队;
- 持久化消息;
- 死信交换器;
- 多种消费者绑定同一消息拓扑。
Redis
Redis 可以作为 Celery Broker,但它不是 AMQP Broker。Redis 传输层使用自己的队列和确认机制,部分语义通过可见性超时实现。
需要特别注意:任务执行时间超过 Redis 的 visibility_timeout 时,原消息可能重新变得可见,而原 Worker 仍在执行。于是可能出现:
t0 Worker A 取出任务
t1 Worker A 仍在执行
t2 visibility_timeout 到期
t3 Worker B 又取到同一任务
t4 A 和 B 同时执行
这不是 Celery 任务代码“自动执行了两次”,而是 Broker 认为未确认消息已经超时,可以重新投递。长时间运行的任务必须让可见性超时覆盖最长执行时间,或者使用更适合该语义的 Broker。
Kafka
Kafka 的核心抽象是追加日志、分区和消费者组,而 Celery 常见的抽象是任务队列和任务 ACK。两者都能传递消息,但消费语义不同:
| 特性 | Celery + RabbitMQ/Redis | Kafka |
|---|---|---|
| 消费模型 | 任务被 Worker 取走并确认 | 消费者提交 offset |
| 消息保留 | 通常直到确认或过期 | 按保留策略保存 |
| 重放 | 通常不是主要能力 | 原生支持按 offset 重放 |
| 路由 | 队列、Exchange、routing key | Topic、Partition |
| 典型用途 | 异步任务、延迟执行、重试 | 事件流、日志流、数据管道 |
因此,“把 Celery Broker 换成 Kafka”不是简单替换连接字符串的问题。需要重新确认消息确认、分区顺序、消费组、重放和任务重复执行的设计。
三、Worker:注册任务、取消息、执行任务
1. Worker 启动时必须拥有任务代码
定义任务:
# app/tasks.py
from celery import Celery
celery_app = Celery(
"demo",
broker="redis://localhost:6379/0",
backend="redis://localhost:6379/1",
)
@celery_app.task
def add(x: int, y: int) -> int:
return x + y
启动 Worker:
celery -A app.tasks:celery_app worker --loglevel=INFO
启动参数的含义是:
-A app.tasks:celery_app:找到 Celery 应用对象;worker:启动任务消费者;--loglevel=INFO:输出普通任务生命周期日志。
另一个 Python 进程发布任务:
# app/producer.py
from app.tasks import add
result = add.delay(2, 3)
print(result.id)
Worker 收到消息后,会根据消息中的任务名查找任务注册表:
消息中的任务名:app.tasks.add
↓
Worker 注册表
↓
找到 add 函数
↓
执行 add(2, 3)
如果 Worker 没有导入该任务,常见结果不是“任务执行失败”,而是类似“未注册任务”的错误。原因是消息已经到达了 Worker,但 Worker 不知道该任务名称对应哪段代码。
2. 并发不是“一个 Worker 只能执行一个任务”
Worker 通常会创建多个执行单元。对于 CPU 密集型任务,常见模型是多进程;对于 I/O 密集型任务,也可以根据任务特征选择线程或协程相关池,但不同池对信号、超时和第三方库的支持不同。
假设:
Worker 并发数 = 4
队列中任务数 = 10
理想化状态为:
正在执行:4 个
尚未取出:6 个
但实际 Worker 还可能提前预取消息。预取意味着 Worker 一次从 Broker 获取多于当前正在执行数量的消息,目的是减少网络等待。
预取会影响:
- 某个 Worker 是否“囤积”大量任务;
- 长任务和短任务是否公平;
- Worker 崩溃后有多少消息重新投递;
- 内存占用;
- 任务延迟的观测结果。
Celery 文档还特别说明,ETA/countdown 任务会被 Worker 取入内存,由内部计时器等待执行,因此它们不完全受普通进程预取窗口的约束;Celery 5.6 提供 worker_eta_task_limit 限制内存中等待执行的 ETA 任务数量。(docs.celeryq.dev)
3. 不要在 Web 请求中等待任务结果
下面的写法会抵消异步任务的价值:
result = add.delay(2, 3)
value = result.get()
HTTP 请求进程仍然要等待 Worker 执行完成。更合理的接口通常返回任务 ID:
from fastapi import FastAPI
from app.tasks import add
api = FastAPI()
@api.post("/tasks/add")
def submit_add(x: int, y: int):
result = add.delay(x, y)
return {
"task_id": result.id,
"status": "submitted",
}
客户端随后查询:
@api.get("/tasks/{task_id}")
def task_status(task_id: str):
result = add.AsyncResult(task_id)
response = {
"task_id": task_id,
"status": result.status,
}
if result.successful():
response["result"] = result.result
if result.failed():
response["error"] = str(result.result)
return response
FastAPI 是 ASGI 应用框架;Celery Worker 是独立的任务执行进程。它们之间不是通过 ASGI 生命周期自动连接,而是通过 Broker 消息通信。ASGI 规定的是异步服务器与应用之间的调用接口,不会替代 Celery 的任务投递、ACK 或重试机制。(docs.python.org)
四、ACK:消息确认不等于业务成功
1. ACK 的精确定义
ACK,Acknowledgement,表示消费者向 Broker 发送确认:
“我已经处理了这条消息,你可以把它从未确认消息集合中移除。”
它不天然表示:
“业务数据库已经提交。”
“外部 API 已经成功。”
“任务结果永远不会丢失。”
ACK 只是消息生命周期中的确认信号。业务成功需要由任务代码和业务存储共同定义。
2. 默认的提前 ACK
Celery 默认通常采用提前 ACK:
1. Worker 从 Broker 取出消息
2. Worker 立即 ACK
3. Worker 开始执行任务
4. Worker 执行失败
如果第 3 步之后 Worker 崩溃,Broker 已经收到 ACK,不会再投递原消息:
Broker:消息已确认
Worker:执行到一半崩溃
结果:任务可能丢失
提前 ACK 的优点是:已经开始执行的任务不会因为 Worker 崩溃而自动重复。缺点是:任务可能在真正完成前丢失。
3. acks_late=True:执行完成后 ACK
对于可以安全重复执行的任务,可以使用晚 ACK:
@celery_app.task(
acks_late=True,
)
def rebuild_cache(user_id: int):
return rebuild_user_cache(user_id)
理想路径变成:
1. Worker 从 Broker 取出消息
2. Worker 执行任务
3. 任务成功返回
4. Worker ACK
如果第 2 步中 Worker 崩溃:
1. 消息还没有 ACK
2. Broker 重新投递
3. 其他 Worker 可能再次执行
所以 acks_late=True 不是“保证不丢消息”,而是把风险从:
执行中崩溃 → 任务丢失
移动为:
执行中崩溃 → 任务重复执行
Celery 官方要求,使用晚 ACK 的任务应当具备幂等性,因为 Worker 可能在执行中崩溃后再次执行同一任务。(docs.celeryq.dev)
4. 晚 ACK 也不是绝对的“崩溃必重试”
Celery 对执行子进程异常退出、被信号终止等情况,默认可能仍然 ACK 消息。这样做是为了避免段错误、管理员主动杀进程、OOM 等场景形成无限重试循环。
如果确实希望 Worker 子进程丢失时重新入队,可以配置:
celery_app.conf.task_reject_on_worker_lost = True
但这可能造成消息循环:
任务执行 → Worker 被杀 → 重新入队 → 再次执行 → 再次被杀
因此,启用它前必须同时具备:
- 幂等任务;
- 有限重试;
- 对 OOM 或确定性崩溃的根因处理;
- 失败任务隔离或死信处理。
Celery 文档明确警告,task_reject_on_worker_lost 可能引起消息循环。(docs.celeryq.dev)
五、ACK、异常和任务状态的完整路径
考虑一个晚 ACK 任务:
@celery_app.task(
bind=True,
acks_late=True,
)
def charge_order(self, order_id: int):
charge_payment(order_id)
return "charged"
情况一:任务成功
Broker: READY
↓ consume
Worker: 执行中
↓ 返回 "charged"
Worker: ACK
↓
Broker: 删除消息
Backend: SUCCESS
情况二:普通异常
Broker: READY
↓ consume
Worker: 执行中
↓ 抛出异常
Worker: 记录 FAILURE
Worker: 根据 ACK/拒绝配置处理消息
默认配置下,晚 ACK 任务失败后通常会 ACK,而不是自动重新入队。要重试,必须显式调用 retry() 或配置自动重试。
情况三:Worker 在业务副作用后崩溃
1. 调用支付服务成功
2. Worker 尚未 ACK
3. Worker 崩溃
4. Broker 重新投递
5. 同一订单再次调用支付服务
这是最危险的情况:消息至少一次投递与外部副作用之间存在窗口。
解决方案不是简单地把 acks_late 改回 False,因为那会把风险变成任务丢失。正确做法是让业务副作用本身具备幂等能力,例如向支付服务传递稳定的幂等键:
payment_key = f"order:{order_id}:payment"
payment_client.charge(
order_id=order_id,
idempotency_key=payment_key,
)
六、重试:重新发布任务,而不是回到当前 Python 调用栈
1. 手动重试
Celery 的 self.retry() 通常用于可恢复错误:
from celery import Celery
celery_app = Celery(
"demo",
broker="redis://localhost:6379/0",
backend="redis://localhost:6379/1",
)
@celery_app.task(
bind=True,
max_retries=5,
default_retry_delay=30,
)
def fetch_remote_data(self, resource_id: str):
try:
return remote_client.fetch(resource_id)
except TemporaryNetworkError as exc:
raise self.retry(exc=exc)
bind=True 使第一个参数成为任务实例 self,从而调用 self.retry()。
retry() 的关键语义是:
- 记录当前任务进入
RETRY状态; - 发布一条新的任务消息;
- 通常使用同一个任务 ID;
- 通过异常控制流立即离开当前任务函数。
因此下面的代码不会继续执行:
@celery_app.task(bind=True)
def example(self):
print("before")
raise self.retry(countdown=10)
print("after") # 不会执行
Celery 文档说明,retry() 会抛出内部的重试异常;它不是普通业务错误,Worker 会据此记录正确的重试状态。(docs.celeryq.dev)
2. 可重试错误与不可重试错误必须区分
不是所有异常都适合重试:
@celery_app.task(bind=True, max_retries=4)
def process_invoice(self, invoice_id: int):
try:
invoice = load_invoice(invoice_id)
return render_invoice(invoice)
except TemporaryNetworkError as exc:
raise self.retry(exc=exc, countdown=10)
适合重试的例子:
- 网络连接短暂失败;
- 外部服务返回 502、503;
- 数据库暂时不可用;
- 限流错误,且服务提供了稍后重试的信号。
不适合盲目重试的例子:
- 参数格式错误;
- 订单不存在;
- 权限不足;
- 业务状态已经非法;
- 程序逻辑错误;
- 外部 API 返回明确的永久失败。
如果把永久错误配置为无限重试,队列会不断产生相同任务,最终掩盖真正问题。
3. 指数退避和抖动
设第 次重试的基础等待时间为:
其中:
- :第几次重试;
- :基础延迟;
- :最大延迟;
- :本次重试前的等待时间。
例如 ,最大值为 60 秒:
| 重试次数 | 基础延迟 |
|---|---|
| 1 | 3 秒 |
| 2 | 6 秒 |
| 3 | 12 秒 |
| 4 | 24 秒 |
| 5 | 48 秒 |
| 6 | 60 秒 |
如果所有任务同时失败并按照相同延迟重试,会产生“惊群”:
10:00:00 1000 个任务同时失败
10:00:03 1000 个任务同时重试
10:00:03 外部服务再次被打满
抖动会把等待时间从固定值变成一个区间内的随机值:
Celery 的自动重试支持 retry_backoff、retry_backoff_max 和 retry_jitter;文档说明抖动默认启用,并将退避值作为随机等待上限。(docs.celeryq.dev)
示例:
@celery_app.task(
autoretry_for=(TemporaryNetworkError,),
retry_backoff=True,
retry_backoff_max=60,
retry_jitter=True,
max_retries=5,
)
def sync_profile(user_id: int):
return remote_client.sync_profile(user_id)
这里的边界是:自动重试只能根据异常类型做粗粒度判断。若不同 HTTP 状态码代表不同策略,手动 retry() 往往更清楚。
4. 重试不是 ACK 的替代品
重试与 ACK 处理的是不同问题:
| 机制 | 解决的问题 |
|---|---|
| ACK | Broker 是否可以移除当前消息 |
retry() |
是否发布下一次任务尝试 |
acks_late |
Worker 崩溃时当前消息是否可能重新投递 |
| 幂等 | 重复执行时业务结果是否安全 |
例如:
任务调用 self.retry()
↓
当前尝试被标记为 RETRY
↓
发布下一条重试消息
如果任务在调用外部服务后进程崩溃,代码根本没有机会调用 retry();这时是否再次执行由 ACK 和 Broker 的重新投递机制决定。因此生产系统必须同时设计“显式异常重试”和“进程崩溃重投”。
七、幂等:重复执行时仍然得到正确业务结果
1. 幂等不是“函数返回值一样”这么简单
数学上的幂等通常表示:
在业务系统中,更准确的表达是:
再次执行同一业务操作:
其中:
- :业务系统当前状态;
- :同一业务操作的幂等键;
- :第一次成功执行后的状态。
例如:
# 幂等:把缓存设置为某个确定值
cache.set("user:42:name", "Alice")
# 非幂等:每执行一次都增加余额
account.balance += 100
set 重复执行仍然是 "Alice";+= 100 执行两次会变成增加 200。
2. 常见幂等实现:数据库唯一键
假设每次订单扣款都有唯一业务键:
CREATE TABLE payment_attempt (
id BIGINT PRIMARY KEY,
order_id BIGINT NOT NULL,
idempotency_key VARCHAR(128) NOT NULL UNIQUE,
status VARCHAR(32) NOT NULL,
provider_payment_id VARCHAR(128),
created_at TIMESTAMP NOT NULL
);
任务执行流程:
@celery_app.task(
bind=True,
acks_late=True,
)
def charge_order(self, order_id: int, idempotency_key: str):
with db.transaction():
attempt = db.query_one(
"""
SELECT id, status, provider_payment_id
FROM payment_attempt
WHERE idempotency_key = %s
FOR UPDATE
""",
[idempotency_key],
)
if attempt is not None:
if attempt["status"] == "SUCCEEDED":
return {
"status": "already_succeeded",
"provider_payment_id": attempt["provider_payment_id"],
}
if attempt["status"] == "PROCESSING":
return {"status": "already_processing"}
db.execute(
"""
INSERT INTO payment_attempt
(order_id, idempotency_key, status, created_at)
VALUES (%s, %s, 'PROCESSING', CURRENT_TIMESTAMP)
ON CONFLICT (idempotency_key) DO NOTHING
""",
[order_id, idempotency_key],
)
# 外部服务也必须接收同一个幂等键
payment = provider.charge(
order_id=order_id,
idempotency_key=idempotency_key,
)
with db.transaction():
db.execute(
"""
UPDATE payment_attempt
SET status = 'SUCCEEDED',
provider_payment_id = %s
WHERE idempotency_key = %s
""",
[payment.id, idempotency_key],
)
return {"status": "succeeded", "provider_payment_id": payment.id}
这里有两个不同层次的保护:
- 数据库唯一约束防止本地重复创建扣款尝试;
- 外部支付服务的幂等键防止网络超时后重复扣款。
只在本地数据库中加锁并不能保证外部服务不重复收费,因为事务锁无法覆盖外部 HTTP 请求。
3. “先写数据库,再发布任务”的竞态
下面的流程存在竞态:
def create_article():
article = db.insert_article(...)
expand_article.delay(article.id)
db.commit()
如果任务在 db.commit() 前开始执行,Worker 可能查不到文章。
更严重的情况是:
1. 数据库写入尚未提交
2. 任务消息已经发布
3. 当前事务回滚
4. Worker 执行任务
5. Worker 永远找不到该数据
事务型应用应当在事务提交后再发送任务。Celery 文档以 Django 为例提供了 delay_on_commit(),它在事务成功提交后才发布任务;该方法自 Celery 5.4 添加。(docs.celeryq.dev)
非 Django 项目可以使用本地事务提交回调,或者采用 Outbox Pattern:
同一个数据库事务内:
写入业务表
写入 outbox_event 表
↓ 提交成功
独立发布器读取 outbox_event
↓
发布 Celery 任务
↓
标记 outbox_event 已发布
Outbox 发布器自身也可能重复发布,因此任务仍然必须幂等。Outbox 解决的是“业务写入成功但消息未发布”的一致性窗口,不会自动提供 exactly-once 执行。
4. 幂等与并发不是一回事
即使任务幂等,也可能发生并发覆盖:
Worker A:读取余额 100
Worker B:读取余额 100
Worker A:写入 110
Worker B:写入 110
如果业务要求同一订单同时只能有一个处理者,需要额外的并发控制:
- 数据库行锁;
- 乐观锁版本号;
- Redis 分布式锁;
- 唯一约束;
- 状态机条件更新。
例如使用条件更新:
UPDATE orders
SET status = 'PROCESSING',
version = version + 1
WHERE id = ?
AND status = 'PAID'
AND version = ?;
如果影响行数为 0,说明其他 Worker 已经抢先修改,当前任务不能继续处理。
八、定时任务:Beat 负责“何时发布”,Worker 负责“何时执行”
1. celery beat 不是 Worker
Celery Beat 是调度器。它只负责按照 schedule 发布任务消息:
Beat
└── 到达调度时间
└── 发布 refresh_cache.delay()
└── Broker
└── Worker 执行
启动方式:
celery -A app.tasks:celery_app beat --loglevel=INFO
celery -A app.tasks:celery_app worker --loglevel=INFO
生产环境中,同一套周期任务通常只运行一个 Beat 实例。若两个 Beat 同时运行,可能各自认为“现在应该发送一次”,于是产生两条相同任务消息。Celery 官方文档要求确保同一个 schedule 同时只有一个调度器。(docs.celeryq.dev)
2. beat_schedule 示例
from celery.schedules import crontab
celery_app.conf.timezone = "Asia/Shanghai"
celery_app.conf.beat_schedule = {
"refresh-cache-every-5-minutes": {
"task": "app.tasks.refresh_cache",
"schedule": 300.0,
"args": (),
},
"close-expired-orders-every-minute": {
"task": "app.tasks.close_expired_orders",
"schedule": crontab(minute="*"),
"args": (),
},
}
间隔调度和 Crontab 调度的时间语义不同:
# 每 300 秒一次,通常相对于调度器运行和上次调度时间
schedule = 300.0
# 每小时的第 0 分钟触发
schedule = crontab(minute=0)
Celery 的周期调度默认使用 UTC,也可以通过 timezone 配置时区。周期任务如果上一次尚未完成,下一次仍可能被发布,因此可能重叠执行。(docs.celeryq.dev)
3. “每五分钟执行”不等于“前一次完成五分钟后执行”
设任务执行时间为 8 分钟,周期为 5 分钟:
12:00 Beat 发布任务 A
12:05 Beat 发布任务 B
12:08 A 完成
12:10 Beat 发布任务 C
如果任务不能并发运行,需要锁或状态条件:
@celery_app.task
def close_expired_orders():
lock = distributed_lock("close-expired-orders", ttl=240)
if not lock.acquire():
return {"status": "skipped", "reason": "already_running"}
try:
return close_orders()
finally:
lock.release()
锁的 TTL 必须覆盖任务执行时间,并考虑 Worker 崩溃后的自动过期。若任务可能超过固定 TTL,需要可续租锁,或者改用数据库状态机。
4. 时间语义中的三个陷阱
时区不一致
如果 Beat 使用 Asia/Shanghai,Worker 使用 UTC,消息中的 ETA 和日志时间可能造成误判。建议:
- 系统内部统一保存 UTC 时间;
- 对外展示时转换为用户时区;
- 明确 Beat 的
timezone; - 日志同时输出时区或 UTC 偏移。
夏令时重复或缺失
某些时区在夏令时切换时会出现:
- 某个本地时间不存在;
- 某个本地时间出现两次。
定时任务如果要求严格业务时间,不能只写一句“每天 02:30 执行”,还必须定义重复和缺失时刻如何处理。
任务积压
Beat 只负责发布,不关心 Worker 是否已经积压。如果 Worker 执行速度小于发布速度,队列长度会持续增长:
其中:
- :任务进入速率;
- :任务完成速率;
- :队列积压量。
当 时,积压会增长。周期任务尤其容易产生这种问题,因为它会持续发布任务。
九、ETA、countdown 和 Beat 的边界
Celery 支持一次性的延迟任务:
send_reminder.apply_async(
args=[order_id],
countdown=600,
)
表示任务最早在 600 秒后执行。
也可以使用绝对时间:
from datetime import datetime, timedelta, timezone
eta = datetime.now(timezone.utc) + timedelta(minutes=10)
send_reminder.apply_async(
args=[order_id],
eta=eta,
)
countdown 是相对延迟,eta 是具体时间点。它们都描述“一次任务何时可以开始”,不是周期调度。
长时间延迟任务不应无限堆积在 Worker 内存中。Celery 文档指出,ETA/countdown 任务会被 Worker 预取并放入内存等待执行;大量此类任务可能造成内存压力。(docs.celeryq.dev)
因此可以这样区分:
| 需求 | 合适机制 |
|---|---|
| 10 分钟后发送一次提醒 | countdown 或 eta |
| 每天凌晨执行清理 | Beat + Crontab |
| 未来几个月的大量用户提醒 | 数据库记录到期时间 + 扫描器 |
| 需要可取消、可审计的预约任务 | 数据库调度模型 + Worker |
十、从 HTTP 请求到任务完成的端到端示例
项目结构:
project/
├── app/
│ ├── __init__.py
│ ├── celery_app.py
│ ├── tasks.py
│ └── api.py
1. 创建 Celery 应用
# app/celery_app.py
from celery import Celery
celery_app = Celery(
"project",
broker="redis://localhost:6379/0",
backend="redis://localhost:6379/1",
)
celery_app.conf.update(
task_track_started=True,
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="Asia/Shanghai",
enable_utc=True,
)
accept_content=["json"] 可以限制 Worker 接受的序列化格式。不要为了方便在不可信输入环境中启用任意对象反序列化;任务消息本质上是跨进程输入,序列化格式与 Broker 访问权限都属于安全边界。
2. 编写可重试、可观测任务
# app/tasks.py
from celery.exceptions import SoftTimeLimitExceeded
from .celery_app import celery_app
class TemporaryRemoteError(Exception):
pass
@celery_app.task(
bind=True,
acks_late=True,
autoretry_for=(TemporaryRemoteError,),
retry_backoff=True,
retry_backoff_max=60,
retry_jitter=True,
max_retries=5,
soft_time_limit=120,
)
def build_report(self, report_id: int) -> dict:
try:
report = load_report(report_id)
if report is None:
# 数据不存在通常是永久错误,不应自动重试
return {
"status": "skipped",
"reason": "report_not_found",
}
output_uri = render_report(report)
save_report_result(
report_id=report_id,
output_uri=output_uri,
idempotency_key=f"report:{report_id}",
)
return {
"status": "success",
"report_id": report_id,
"output_uri": output_uri,
}
except TemporaryRemoteError:
raise
except SoftTimeLimitExceeded:
mark_report_failed(report_id, reason="soft_time_limit")
raise
这个任务仍然要求 save_report_result() 幂等。因为 render_report() 成功后、数据库更新前,如果 Worker 崩溃,晚 ACK 会使任务再次执行。
3. FastAPI 发布任务
# app/api.py
from fastapi import FastAPI, HTTPException
from .tasks import build_report
api = FastAPI()
@api.post("/reports/{report_id}/build", status_code=202)
def submit_report(report_id: int):
if not report_exists(report_id):
raise HTTPException(status_code=404, detail="report not found")
result = build_report.delay(report_id)
return {
"task_id": result.id,
"report_id": report_id,
"status": "PENDING",
}
@api.get("/tasks/{task_id}")
def get_task(task_id: str):
result = build_report.AsyncResult(task_id)
response = {
"task_id": task_id,
"status": result.status,
}
if result.status == "SUCCESS":
response["result"] = result.result
elif result.status == "FAILURE":
response["error"] = str(result.result)
return response
接口返回 202 Accepted 的含义是:请求已经接受,但任务还没有完成。不能把 202 当成业务成功。
启动:
uvicorn app.api:api --host 0.0.0.0 --port 8000
celery -A app.celery_app:celery_app worker --loglevel=INFO
请求:
curl -X POST http://localhost:8000/reports/42/build
可能返回:
{
"task_id": "2c9d8a1e-...",
"report_id": 42,
"status": "PENDING"
}
查询:
curl http://localhost:8000/tasks/2c9d8a1e-...
可能经历:
{"task_id": "...", "status": "PENDING"}
{"task_id": "...", "status": "STARTED"}
{
"task_id": "...",
"status": "SUCCESS",
"result": {
"status": "success",
"report_id": 42,
"output_uri": "s3://bucket/reports/42.pdf"
}
}
STARTED 不是所有配置下都会默认记录;示例中通过 task_track_started=True 开启。Celery 的任务状态包括 PENDING、STARTED、SUCCESS、FAILURE 和 RETRY 等,状态数据由 Result Backend 保存。(docs.celeryq.dev)
十一、任务状态、消息状态和业务状态必须分开
下面三种状态经常被错误地混为一谈:
Celery 任务状态
PENDING → STARTED → SUCCESS
↘ FAILURE
↘ RETRY
它回答:
“这一次 Celery 任务尝试处于什么状态?”
Broker 消息状态
READY → DELIVERED/UNACKED → ACKED
↘ REDELIVERED
它回答:
“这条消息是否还可能再次交付?”
业务状态
CREATED → PROCESSING → SUCCEEDED
↘ FAILED
它回答:
“订单、报表或支付业务处于什么状态?”
一个任务完全可能出现:
Celery:SUCCESS
Broker:消息已 ACK
业务:外部文件实际不存在
也可能出现:
Celery:FAILURE
业务:外部支付已经成功
后一种情况常发生在外部副作用已经完成,但 Worker 在保存结果前抛出异常。业务系统必须以自己的业务记录和外部系统查询结果为准,而不能只看 Celery 的 SUCCESS 或 FAILURE。
十二、失败路径:从故障反推配置
1. Web 进程发布前崩溃
数据库事务成功
Web 进程在 publish 前崩溃
结果:业务数据存在,但任务消息不存在
适合使用:
- 事务提交回调;
- Outbox 表;
- 后台补偿扫描。
2. Broker 发布成功,HTTP 响应返回前崩溃
消息已经进入 Broker
HTTP 响应尚未返回
客户端重试请求
结果:可能发布两条任务
解决方式是为业务请求生成稳定的幂等键,并让任务或发布记录按该键去重。仅仅依赖 Celery 自动生成的随机 task ID,无法识别两次业务请求是否代表同一个操作。
3. Worker 取消息后崩溃
- 提前 ACK:任务可能丢失;
- 晚 ACK:任务可能重新投递;
- 晚 ACK + 幂等:通常可以接受重复执行;
- 晚 ACK + 非幂等:可能产生重复扣款、重复发货等严重后果。
4. 任务异常后没有重试
常见原因:
- 没有调用
self.retry(); - 异常没有匹配
autoretry_for; - 异常属于永久错误;
- 已达到
max_retries; - 任务失败后被 ACK,而团队误以为 ACK 会触发重试。
5. 任务无限重试
常见原因:
max_retries=None;- 把所有
Exception都配置为自动重试; - 外部服务永久返回 400,但代码仍然重试;
- 没有死信或人工处理路径。
十三、诊断方法:先判断消息在哪个阶段丢失
遇到“任务没有执行”时,不要直接重启 Worker。按以下路径定位:
第一步:确认生产者是否成功发布
记录:
业务请求 ID
业务幂等键
Celery task ID
任务名称
目标队列
发布时间
如果没有 task ID,问题可能发生在发布前或发布阶段。
第二步:确认 Worker 是否注册任务
检查 Worker 启动日志中是否加载了对应模块,确认任务名称完全一致:
app.tasks.build_report
以下差异都会导致注册失败:
app.tasks.build_report
project.app.tasks.build_report
第三步:观察任务状态变化
状态序列可以帮助定位:
PENDING 长期不变
可能表示:
- Worker 未启动;
- 任务进入了错误队列;
- Worker 没有订阅该队列;
- Result Backend 不可用;
- 任务 ID 不存在。
STARTED 长期不变
可能表示:
- 任务卡死;
- 外部 I/O 无超时;
- Worker 进程资源耗尽;
- 任务正在等待锁。
RETRY 反复出现
说明:
- 异常确实触发了重试;
- 需要检查退避时间、异常类型和最大重试次数;
- 还要确认外部依赖是否持续不可用。
第四步:检查 Broker 队列和未确认消息
需要区分:
队列中没有消息
与:
消息被 Worker 取走但未 ACK
前者可能是已经消费、发布失败或路由错误;后者可能是任务正在执行、Worker 卡死或 Broker 连接异常。
第五步:检查 Result Backend 的局限
Result Backend 不是可靠业务审计日志:
- 结果可能过期;
- 某些后端通过轮询读取状态;
- 消息型结果后端可能不适合多个进程同时等待同一结果;
- 只保存任务状态,不一定保存完整业务事实。
Celery 文档指出,不同 Result Backend 有不同限制;例如数据库后端轮询会带来成本,消息型后端的结果消息默认可能是非持久的。(docs.celeryq.dev)
十四、常见误解
误解一:Celery 提供 exactly-once 执行
通常不能这样假设。
在网络断开、Worker 崩溃、ACK 丢失、Broker 重投或生产者重试的情况下,同一业务任务可能执行多次。工程上更现实的目标是:
至少一次投递 + 幂等业务处理 + 有限重试 + 可补偿
误解二:任务 ID 就是幂等键
Celery 自动生成的任务 ID 标识一次消息任务,但客户端重试一次业务请求时,可能生成新的任务 ID:
第一次请求 → task_id=A
HTTP 超时后重试 → task_id=B
如果两次请求实际上代表同一个业务操作,A 和 B 并不能自动被识别为重复。幂等键应由业务语义生成,例如:
order:1001:payment
user:42:monthly-report:2026-09
误解三:重试会撤销之前的副作用
不会。
下面的代码在外部服务成功后抛异常:
try:
provider.send_email(...)
except TemporaryRemoteError as exc:
raise self.retry(exc=exc)
如果 provider.send_email() 实际已经成功,但响应在网络中丢失,任务重试可能再次发送邮件。必须使用外部服务的幂等键,或者建立本地发送记录。
误解四:Beat 保证任务只执行一次
Beat 只负责发布周期消息,不保证:
- 只有一个 Beat 实例;
- 任务不会重复发布;
- 任务不会重叠执行;
- Worker 不会因 Broker 重投而重复执行。
周期任务的“只处理一次”需要由锁、唯一约束或业务状态机保证。
误解五:acks_late=True 总是更可靠
它提高了 Worker 崩溃后重新执行的可能性,但也增加重复副作用的概率。任务是否使用晚 ACK,取决于任务能否安全重复执行,而不是取决于“晚 ACK 看起来更先进”。
十五、一个可操作的设计判断顺序
设计一个 Celery 任务时,按以下顺序回答问题:
1. 任务是否允许重复执行?
如果不能,先改造业务操作:
稳定幂等键
↓
唯一约束或幂等记录
↓
外部调用使用同一幂等键
↓
重复调用返回已有结果
不要先配置重试,再考虑幂等。
2. 哪些错误可以恢复?
建立明确分类:
可恢复:
网络超时、503、限流
不可恢复:
参数错误、权限错误、资源不存在
需人工处理:
外部状态不确定、金额不一致、数据损坏
3. Worker 崩溃时,选择丢失还是重复?
- 任务是纯计算且结果可重新生成:可以倾向晚 ACK;
- 任务会产生不可逆外部副作用:必须先有幂等设计;
- 任务不能重复且不能丢失:需要业务级状态机、外部幂等协议或人工补偿,单靠 ACK 无法解决。
4. 任务是否可能超时?
为外部调用设置连接超时和读取超时:
http_client.get(
url,
connect_timeout=5,
read_timeout=30,
)
Celery 的任务时间限制只能终止或中断执行,不能撤销已经发送给外部服务的请求。超时后是否重试,仍然要根据副作用是否已经发生来判断。
5. 周期任务是否可能重叠?
如果可能,设计:
分布式锁
或
数据库条件更新
或
按时间窗口唯一约束
例如日报任务的幂等键可以是:
daily-report:2026-09-01
这样即使 Beat 发布两次,也只生成一份业务结果。
结语:把 Celery 当作“至少一次消息驱动的执行系统”
Celery 的核心不是把函数放到后台,而是建立一条跨进程、跨网络的执行链:
业务请求
↓
发布任务消息
↓
Broker 暂存和投递
↓
Worker 获取消息
↓
ACK 或重新投递
↓
任务成功、失败或重试
↓
业务状态持久化
其中最重要的因果关系是:
ACK 决定消息是否还可能被 Broker 投递
重试决定是否主动安排下一次尝试
幂等决定重复执行是否会破坏业务
Beat 决定何时发布周期消息
Worker 决定如何并发执行任务
Result Backend 只记录任务状态和结果,不替代业务数据库
因此,可靠的 Celery 任务通常不是“配置了几项参数”的结果,而是下面这个组合:
当任务可能重复、可能延迟、可能在副作用后崩溃,并且这些情况都能通过业务状态和幂等键恢复时,Celery 才真正适合承载生产任务队列。
系列导航与关联阅读
- 系列入口:Python 完整学习路线:从语言模型、并发到 Web、数据、AI 与生产交付
- 上一篇:Python 与 Redis:连接池、Pipeline、事务、缓存和分布式锁
- 下一篇:Python 消息系统:RabbitMQ、Kafka、消费语义、重试和死信
- 延伸:Python 定时任务:时间语义、调度器、重复执行、锁和补偿
官方资料
本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。

评论
0 条讨论