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 请求、数据库查询):使用
gevent或eventlet并发模式,并发数可设较高(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,在任务内部重新查询。