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 保存执行状态和结果 PENDINGSTARTEDSUCCESSFAILURE
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() 的关键语义是:

  1. 记录当前任务进入 RETRY 状态;
  2. 发布一条新的任务消息;
  3. 通常使用同一个任务 ID;
  4. 通过异常控制流立即离开当前任务函数。

因此下面的代码不会继续执行:

@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. 指数退避和抖动

设第 nn 次重试的基础等待时间为:

dn=min(dmax,b×2n1)d_n = \min(d_{\max}, b \times 2^{n-1})

其中:

  • nn:第几次重试;
  • bb:基础延迟;
  • dmaxd_{\max}:最大延迟;
  • dnd_n:本次重试前的等待时间。

例如 b=3b=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  外部服务再次被打满

抖动会把等待时间从固定值变成一个区间内的随机值:

DnU(0,dn)D_n \sim U(0, d_n)

Celery 的自动重试支持 retry_backoffretry_backoff_maxretry_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. 幂等不是“函数返回值一样”这么简单

数学上的幂等通常表示:

f(f(x))=f(x)f(f(x)) = f(x)

在业务系统中,更准确的表达是:

Sexecute(k)SS \xrightarrow{execute(k)} S'

再次执行同一业务操作:

Sexecute(k)SS' \xrightarrow{execute(k)} S'

其中:

  • SS:业务系统当前状态;
  • kk:同一业务操作的幂等键;
  • SS':第一次成功执行后的状态。

例如:

# 幂等:把缓存设置为某个确定值
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}

这里有两个不同层次的保护:

  1. 数据库唯一约束防止本地重复创建扣款尝试;
  2. 外部支付服务的幂等键防止网络超时后重复扣款。

只在本地数据库中加锁并不能保证外部服务不重复收费,因为事务锁无法覆盖外部 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 执行速度小于发布速度,队列长度会持续增长:

ΔQ=λinλout\Delta Q = \lambda_{in} - \lambda_{out}

其中:

  • λin\lambda_{in}:任务进入速率;
  • λout\lambda_{out}:任务完成速率;
  • QQ:队列积压量。

λin>λout\lambda_{in} > \lambda_{out} 时,积压会增长。周期任务尤其容易产生这种问题,因为它会持续发布任务。


九、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 分钟后发送一次提醒 countdowneta
每天凌晨执行清理 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 的任务状态包括 PENDINGSTARTEDSUCCESSFAILURERETRY 等,状态数据由 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 的 SUCCESSFAILURE


十二、失败路径:从故障反推配置

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 任务通常不是“配置了几项参数”的结果,而是下面这个组合:

可靠异步处理=明确 Broker 语义+正确 ACK 策略+有限且分类的重试+稳定幂等键+可持久化业务状态+失败后的补偿路径\text{可靠异步处理} = \text{明确 Broker 语义} + \text{正确 ACK 策略} + \text{有限且分类的重试} + \text{稳定幂等键} + \text{可持久化业务状态} + \text{失败后的补偿路径}

当任务可能重复、可能延迟、可能在副作用后崩溃,并且这些情况都能通过业务状态和幂等键恢复时,Celery 才真正适合承载生产任务队列。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。