> ## Content Index
> Fetch the complete content index at: https://blog.vercanti.com/llms.txt
> Use this file to discover other available public pages before exploring further.

# Celery 完全指南
- URL: https://blog.vercanti.com/celery-wan-quan-zhi-nan/
- Published: 2026-08-28T14:34:39.000Z
- Updated: 2026-08-28T14:57:01.000Z
- Description: Celery 是 Python 中最成熟的分布式任务队列，用于： Flower 提供：任务实时监控、Worker 状态、任务历史、手动触发任务等。 同一个任务可能因网络问题被重复投递，确保多次执行结果相同： 只传 ID，不传整个对象（对象可能很大，且在任务执行时数据可能已变化）： Celery 应用对象和任务定义通常存在循环导入风险，推荐用工厂模式或在 Django 中用 shared_task： Celery Beat 使用 celerybeat-schedule 文件存储调度状态，多个 Beat 进程同时运行会冲突。生产环境只能运行一个 Beat 实
- Author: yellowdog
- Tags: Python, 框架与库

> 官方文档：<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 进程   | 定时触发任务的调度器 |

### 安装

```bash
pip install celery

# Redis 支持
pip install celery[redis]

# 监控工具
pip install flower

```

---

## 2\. 快速开始

### 创建 Celery 应用

```python
# 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 个任务（长任务推荐）
)

```

### 定义任务

```python
# 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

```bash
# 启动 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

```

### 调用任务

```python
# 异步调用（立即返回，任务放入队列）
result = send_welcome_email.delay(user_id=1, email="alice@example.com")
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

```

### 查询任务状态

```python
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,
    }

```

### 更新任务进度（自定义状态）

```python
@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 — 串行（前一个结果作为下一个输入）

```python
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 — 并发

```python
from celery import group

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

```

### chord — 并发 + 汇总

```python
from celery import chord

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

```

### 混合编排

```python
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）

### 配置定时任务

```python
# 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

```bash
# Beat 必须单独启动（不能和 worker 合并，除非开发环境）
celery -A src.celery_app beat -l info

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

```

---

## 6\. 队列路由

```python
# 定义多个队列，不同优先级/类型的任务路由到不同队列
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 集成

```python
# 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）

```bash
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\. 常用代码段

### 任务基类（统一日志和异常处理）

```python
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

```

### 防止重复执行（分布式锁）

```python
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\. 最佳实践

### 任务要幂等

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

```python
@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，不传整个对象（对象可能很大，且在任务执行时数据可能已变化）：

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

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

```

### 任务超时保护

```python
@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`：

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

@shared_task
def my_task():
    ...

```

### beat 调度器文件锁

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

```python
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 请求、数据库查询）：使用 `gevent` 或 `eventlet` 并发模式，并发数可设较高（50-100）
- CPU 密集型任务（图片处理、数据计算）：使用默认 prefork 模式，并发数等于 CPU 核数

```bash
# 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 配置，避免互相阻塞。

```python
@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，在任务内部重新查询。

---

## 参见

- [FastAPI完全指南](https://blog.vercanti.com/fastapi-wan-quan-zhi-nan/)
- [asyncio异步编程完全指南](https://blog.vercanti.com/asyncio-yi-bu-bian-cheng-wan-quan-zhi-nan/)
- [Redis完全指南](https://blog.vercanti.com/redis-wan-quan-zhi-nan/)
- [RabbitMQ完全指南](https://blog.vercanti.com/rabbitmq-wan-quan-zhi-nan/)