Celery 完全指南

Celery 是 Python 中最成熟的分布式任务队列,用于: Flower 提供:任务实时监控、Worker 状态、任务历史、手动触发任务等。 同一个任务可能因网络问题被重复投递,确保多次执行结果相同: 只传 ID,不传整个对象(对象可能很大,且在任务执行时数据可能已变化): Celery 应用对象和任务定义通常存在循环导入风险,推荐用工厂模式或在 Django 中用 shared_task: Celery Beat 使用 celerybeat-schedule 文件存储调度状态,多个 Beat 进程同时运行会冲突。生产环境只能运行一个 Beat 实

分享

官方文档:https://docs.celeryq.dev/
最后更新:2026-03-29


1. 基础概念

Celery 是什么

Celery 是 Python 中最成熟的分布式任务队列,用于:

  • 异步任务:将耗时操作(发邮件、图片处理、报表生成)放到后台执行
  • 定时任务:周期性执行(Celery Beat)
  • 任务编排:串联(chain)、并发(group)、复杂 DAG 工作流

核心组件

Producer(生产者)
  → Broker(消息中间件:Redis / RabbitMQ)
    → Worker(消费者,执行任务)
      → Result Backend(结果存储:Redis / 数据库)
组件 推荐选择 说明
Broker Redis 生产首选,简单高效
Result Backend Redis 存储任务执行结果
Worker Celery Worker 进程 实际执行任务的进程
Beat Celery Beat 进程 定时触发任务的调度器

安装

pip install celery

# Redis 支持
pip install celery[redis]

# 监控工具
pip install flower

2. 快速开始

创建 Celery 应用

# src/celery_app.py
from celery import Celery

app = Celery(
    "myapp",
    broker="redis://localhost:6379/0",       # 任务队列
    backend="redis://localhost:6379/1",      # 结果存储
    include=["myapp.tasks.email", "myapp.tasks.report"],  # 自动发现任务模块
)

app.conf.update(
    task_serializer="json",
    result_serializer="json",
    accept_content=["json"],
    timezone="Asia/Shanghai",
    enable_utc=False,
    task_track_started=True,      # 记录任务开始状态
    result_expires=3600,          # 结果保存 1 小时
    worker_prefetch_multiplier=1, # 每次只预取 1 个任务(长任务推荐)
)

定义任务

# src/tasks/email.py
from celery_app import app
import time

@app.task
def send_welcome_email(user_id: int, email: str) -> bool:
    """发送欢迎邮件"""
    time.sleep(2)  # 模拟发送耗时
    print(f"发送邮件到 {email}")
    return True

@app.task(
    bind=True,                # bind=True 则第一个参数为 self(任务实例)
    max_retries=3,            # 最大重试次数
    default_retry_delay=60,   # 重试间隔(秒)
    time_limit=300,           # 硬时间限制(秒),超时 kill
    soft_time_limit=270,      # 软时间限制(秒),抛出 SoftTimeLimitExceeded
)
def send_report(self, report_id: int) -> dict:
    """生成报告,失败自动重试"""
    try:
        result = generate_report(report_id)
        return result
    except TemporaryError as exc:
        # 指数退避重试
        raise self.retry(exc=exc, countdown=2 ** self.request.retries)

启动 Worker

# 启动 Worker(-A 指定应用模块,-l 日志级别)
celery -A src.celery_app worker -l info

# 指定并发数(默认等于 CPU 核数)
celery -A src.celery_app worker -l info --concurrency=4

# 指定队列
celery -A src.celery_app worker -l info -Q default,email,high_priority

# 后台运行(守护进程)
celery -A src.celery_app worker -l info --detach

调用任务

# 异步调用(立即返回,任务放入队列)
result = send_welcome_email.delay(user_id=1, email="[email protected]")
print(result.id)  # 任务 ID

# 等待结果(会阻塞,生产环境一般不用)
output = result.get(timeout=10)
print(output)  # True

# apply_async:更多控制选项
result = send_report.apply_async(
    args=[42],
    countdown=10,         # 10 秒后执行
    expires=3600,         # 1 小时内未执行则丢弃
    queue="high_priority",
    task_id="custom-id-123",
)

3. 任务状态与结果

任务状态流转

PENDING → STARTED → SUCCESS
                  ↘ FAILURE
                  ↘ RETRY

查询任务状态

from celery.result import AsyncResult
from celery_app import app

def get_task_status(task_id: str) -> dict:
    result = AsyncResult(task_id, app=app)
    return {
        "task_id": task_id,
        "status": result.status,       # PENDING/STARTED/SUCCESS/FAILURE/RETRY
        "result": result.result if result.ready() else None,
        "traceback": result.traceback if result.failed() else None,
    }

更新任务进度(自定义状态)

@app.task(bind=True)
def process_large_file(self, file_id: int) -> dict:
    total = 1000
    for i in range(total):
        process_item(i)
        if i % 100 == 0:
            # 更新进度(可被前端轮询查询)
            self.update_state(
                state="PROGRESS",
                meta={"current": i, "total": total, "percent": i / total * 100},
            )
    return {"status": "done", "processed": total}

4. 任务编排

chain — 串行(前一个结果作为下一个输入)

from celery import chain

# 串行:download → process → upload
workflow = chain(
    download_file.s(url),      # .s() 创建签名(Signature)
    process_file.s(),          # 接收前一个任务的返回值作为第一个参数
    upload_result.s(bucket="output"),
)
result = workflow.delay()

group — 并发

from celery import group

# 并发处理多个 URL
job = group(
    fetch_url.s(url) for url in urls
)
result = job.apply_async()
outputs = result.get()  # 等待所有完成,返回结果列表

chord — 并发 + 汇总

from celery import chord

# 并发抓取,然后汇总
workflow = chord(
    group(fetch_url.s(url) for url in urls),
    aggregate_results.s(),   # 所有并发任务完成后执行汇总
)
result = workflow.delay()

混合编排

from celery import chain, group, chord

workflow = chain(
    validate_input.s(data),
    chord(
        group(process_chunk.s(chunk) for chunk in chunks),
        merge_results.s(),
    ),
    send_notification.s(email),
)

5. 定时任务(Celery Beat)

配置定时任务

# celery_app.py
from celery.schedules import crontab

app.conf.beat_schedule = {
    "每天早上 8 点发报告": {
        "task": "myapp.tasks.report.send_daily_report",
        "schedule": crontab(hour=8, minute=0),
        "args": (),
    },
    "每 5 分钟同步数据": {
        "task": "myapp.tasks.sync.sync_data",
        "schedule": 300,   # 秒数,等同于 timedelta(seconds=300)
    },
    "每周一清理日志": {
        "task": "myapp.tasks.maintenance.clean_logs",
        "schedule": crontab(hour=2, minute=0, day_of_week=1),  # 周一凌晨 2 点
    },
}

crontab 参数说明

参数 说明 示例
minute 分钟(0-59) "*/15" 每 15 分钟
hour 小时(0-23) "8,12,18" 8点、12点、18点
day_of_week 星期(0=周日,1=周一...6=周六) "1-5" 工作日
day_of_month 月中的日(1-31) "1" 每月 1 日
month_of_year 月(1-12) "*/3" 每季度

启动 Beat

# Beat 必须单独启动(不能和 worker 合并,除非开发环境)
celery -A src.celery_app beat -l info

# 开发环境:worker + beat 合并启动
celery -A src.celery_app worker -l info -B

6. 队列路由

# 定义多个队列,不同优先级/类型的任务路由到不同队列
app.conf.task_routes = {
    "myapp.tasks.email.*": {"queue": "email"},
    "myapp.tasks.report.*": {"queue": "low_priority"},
    "myapp.tasks.payment.*": {"queue": "high_priority"},
}

app.conf.task_default_queue = "default"

# 启动专门处理 email 队列的 worker
# celery -A src.celery_app worker -l info -Q email --concurrency=2

7. 与 FastAPI 集成

# src/main.py
from fastapi import FastAPI, BackgroundTasks
from celery.result import AsyncResult
from celery_app import app as celery_app
from tasks.email import send_welcome_email
from tasks.report import generate_report

api = FastAPI()

@api.post("/users/{user_id}/welcome-email")
async def trigger_welcome_email(user_id: int, email: str):
    task = send_welcome_email.delay(user_id=user_id, email=email)
    return {"task_id": task.id, "status": "queued"}

@api.get("/tasks/{task_id}")
async def get_task_result(task_id: str):
    result = AsyncResult(task_id, app=celery_app)
    if result.state == "PENDING":
        return {"task_id": task_id, "status": "pending"}
    elif result.state == "SUCCESS":
        return {"task_id": task_id, "status": "success", "result": result.result}
    elif result.state == "FAILURE":
        return {"task_id": task_id, "status": "failure", "error": str(result.result)}
    elif result.state == "PROGRESS":
        return {"task_id": task_id, "status": "progress", "meta": result.info}
    return {"task_id": task_id, "status": result.state}

8. 监控(Flower)

pip install flower

# 启动 Flower Web 监控界面(默认 http://localhost:5555)
celery -A src.celery_app flower

# 指定端口和认证
celery -A src.celery_app flower --port=5566 --basic_auth=admin:password

Flower 提供:任务实时监控、Worker 状态、任务历史、手动触发任务等。


9. 常用代码段

任务基类(统一日志和异常处理)

from celery import Task
import logging

logger = logging.getLogger(__name__)

class LoggedTask(Task):
    abstract = True  # 抽象任务,不会被注册为实际任务

    def on_success(self, retval, task_id, args, kwargs):
        logger.info("任务成功 task_id=%s retval=%s", task_id, retval)

    def on_failure(self, exc, task_id, args, kwargs, einfo):
        logger.error("任务失败 task_id=%s exc=%s", task_id, exc, exc_info=True)

    def on_retry(self, exc, task_id, args, kwargs, einfo):
        logger.warning("任务重试 task_id=%s exc=%s", task_id, exc)

@app.task(base=LoggedTask)
def my_task(x: int) -> int:
    return x * 2

防止重复执行(分布式锁)

import redis
from celery_app import app

redis_client = redis.Redis.from_url("redis://localhost:6379/0")

@app.task
def unique_task(resource_id: int):
    lock_key = f"task:lock:unique_task:{resource_id}"
    lock = redis_client.set(lock_key, "1", ex=600, nx=True)  # 10 分钟锁
    if not lock:
        return {"skipped": True, "reason": "任务正在执行中"}
    try:
        return do_work(resource_id)
    finally:
        redis_client.delete(lock_key)

10. 最佳实践

任务要幂等

同一个任务可能因网络问题被重复投递,确保多次执行结果相同:

@app.task
def create_invoice(order_id: int):
    # 先检查是否已创建,避免重复
    if Invoice.objects.filter(order_id=order_id).exists():
        return {"status": "already_exists"}
    return Invoice.create(order_id=order_id)

任务参数要轻量

只传 ID,不传整个对象(对象可能很大,且在任务执行时数据可能已变化):

# 不推荐
send_email.delay(user=user_obj)  # 序列化整个对象

# 推荐
send_email.delay(user_id=user_obj.id)  # 任务内部重新查询

任务超时保护

@app.task(time_limit=300, soft_time_limit=270)
def process_data(data_id: int):
    from celery.exceptions import SoftTimeLimitExceeded
    try:
        return expensive_operation(data_id)
    except SoftTimeLimitExceeded:
        # 软超时时做清理
        cleanup()
        raise

11. 踩坑与注意事项

循环导入

Celery 应用对象和任务定义通常存在循环导入风险,推荐用工厂模式或在 Django 中用 shared_task

# 在非 Django 项目中,用相对导入避免循环
from celery import shared_task  # Django 项目用 shared_task

@shared_task
def my_task():
    ...

beat 调度器文件锁

Celery Beat 使用 celerybeat-schedule 文件存储调度状态,多个 Beat 进程同时运行会冲突。生产环境只能运行一个 Beat 实例。容器化部署时用 DatabaseScheduler 代替文件调度器:

app.conf.beat_scheduler = "django_celery_beat.schedulers:DatabaseScheduler"  # Django
# 或使用 celery-redbeat(基于 Redis 的调度器,支持高可用)
pip install celery-redbeat
app.conf.beat_scheduler = "redbeat.RedBeatScheduler"

Worker 数量与任务类型匹配

  • I/O 密集型任务(HTTP 请求、数据库查询):使用 geventeventlet 并发模式,并发数可设较高(50-100)
  • CPU 密集型任务(图片处理、数据计算):使用默认 prefork 模式,并发数等于 CPU 核数
# gevent 模式(I/O 密集)
pip install gevent
celery -A app worker -P gevent --concurrency=100

最佳实践

任务必须是幂等的,不依赖任务执行次数:Celery 在 Worker 崩溃后可能重发任务(at-least-once 语义),任务执行两次不应产生副作用(如重复发邮件、重复扣款)。用数据库唯一约束或 Redis 去重保证幂等性。

不要在任务中共享可变状态(全局变量、类变量):Worker 使用进程池,进程间不共享内存;使用 gevent 模式时多个任务在同一进程的协程中执行,全局变量会被意外修改。所有状态应通过参数传入或从数据库读取。

为不同优先级和资源类型的任务配置独立队列:CPU 密集型任务(视频转码)和 I/O 密集型任务(发邮件)用不同队列和不同 Worker 配置,避免互相阻塞。

@app.task(queue='high_priority')
def send_urgent_email(): ...

@app.task(queue='low_priority')
def generate_report(): ...

设置 task_acks_late=True + task_reject_on_worker_lost=True 防止任务丢失:默认情况下 Worker 收到任务时立即 ack(确认),若处理中崩溃任务丢失。acks_late=True 在任务完成后才 ack,配合 reject_on_worker_lost=True 在 Worker 崩溃时将任务重新入队。

监控使用 Flower,不要只依赖日志flower 提供实时的 Task 状态、Worker 状态、队列积压可视化,是 Celery 的标准监控工具,比解析日志高效。


常见陷阱

陷阱:在任务中使用同步 ORM 调用阻塞 Worker

现象: 使用 gevent 或 eventlet 池时,任务执行极慢,Worker 频繁超时;CPU 使用率低但吞吐量差。

原因: gevent/eventlet 通过猴子补丁(monkey-patch)让标准库 IO 变成异步,但部分 ORM 或数据库驱动(如 psycopg2 非补丁版本)的 IO 不受影响,会阻塞整个 Worker 进程的协程调度。

解决: 使用与 gevent/eventlet 兼容的驱动(如 psycogreen),或切回 prefork 模式使用线程安全的同步驱动。

陷阱:Celery Beat 重复调度任务

现象: 定时任务被执行了多次,数据库中出现重复数据。

原因: 同时启动了多个 Celery Beat 进程(如在多个节点都启动了 Beat),每个 Beat 都会独立调度任务,导致任务重复发送。

解决: 整个集群只运行一个 Beat 进程。生产环境用 celery-redbeat 或将 Beat 部署为单点服务,不做水平扩展。

陷阱:任务参数包含不可序列化对象

现象: 提交任务时抛出 kombu.exceptions.EncodeError: cannot serialize ... object

原因: Celery 默认用 JSON 序列化任务参数,datetime(需转 .isoformat())、ORM 对象、文件句柄等不能直接 JSON 序列化。

解决: 只传基本类型(str/int/dict/list)作为任务参数;datetime 转字符串;ORM 对象传主键 ID,在任务内部重新查询。


参见

阅读更多

Web 安全基础

1. HTML 转义(服务端渲染必须): 2. CSP(Content Security Policy): 3. HttpOnly Cookie:防止 JS 读取会话 Cookie: 4. 前端框架防护: 攻击者在第三方网站构造一个表单,诱导已登录用户提交,浏览器会自动携带目标站的 Cookie。 触发条件: 1. 用户已登录目标网站(Cookie 有效) 2. 目标 API 仅凭 Cookie 识别用户身份 3. 请求来源未验证 1. CSRF Token(推荐): 2. SameSite Cookie: 3. 验证 Origin/Referer 头:

By yellowdog

HTTP 协议深度指南

HTTP(HyperText Transfer Protocol)是 Web 的基础传输协议,基于 TCP/IP,采用请求/响应模型。 相关文档:Web安全基础(/web-an-quan-ji-chu/) FastAPI完全指南(/fastapi-wan-quan-zhi-nan/) Nginx完全指南(/nginx-wan-quan-zhi-nan/) 幂等性:多次执行相同请求,服务器状态结果相同。PUT /users/1 多次执行结果一致;POST /users 每次创建新资源,非幂等。 浏览器直接从本地缓存读取,不向服务器发送请求。 缓存命中时,状

By yellowdog

系统设计基础

SLA 对照表: 选择建议:无状态服务(Web 层、API 层)优先水平扩展;数据库初期垂直扩展,达到瓶颈后考虑分库分表或读写分离。 缓存穿透(查询不存在的 key,每次都打到 DB): 缓存击穿(热点 key 过期,瞬间大量请求打到 DB): 缓存雪崩(大量 key 同时过期,或缓存服务宕机): 令牌桶 Python 实现: Redis 实现分布式限流(滑动窗口): URL 命名规则: Cursor 分页响应格式: 雪花算法结构(64 bit): 定义:分布式系统不能同时满足以下三个特性: 在分布式环境中 P 是必须保证的,所以实际是 CP vs AP

By yellowdog

算法思路与模板

二分查找要求序列有序,每次将搜索范围缩减一半,时间复杂度 O(log n)。 两个指针从两端向中间收缩,常用于有序数组。 滑动窗口维护一个满足条件的区间 left, right,right 不断向右扩张,条件不满足时收缩 left。 滑动窗口通用框架: 1. 确定"子问题":原问题可以分解为哪些规模更小的同类问题 2. 定义 dpi 或 dpij 的含义,要足够清晰 3. 推导状态转移方程 4. 确定初始状态(边界条件) 5. 确定计算顺序(确保依赖的子问题先计算) 每件物品最多选一次。dpj = 容量为 j 时的最大价值,逆序遍历容量防止重复选取。 每

By yellowdog