Web 应用里总有些操作不能让用户盯着浏览器等:发送一封邮件、生成一份 PDF 报表、批量处理图片、同步第三方 API。如果这些活儿都塞在主请求里,接口分分钟超时,用户体验也会崩。Celery 是 Python 生态里最受欢迎的分布式任务队列之一,这篇文章从安装、任务定义、Worker 调度、定时任务到生产踩坑,带你跑通一条完整链路。
全文代码可直接复制运行,示例基于 Celery 5.4 与 Redis 7,操作系统以 Linux 为主,Windows 用户注意文中标注的特殊参数即可。
一、为什么需要任务队列
先看一下没有任务队列时的典型场景。用户点击"导出月度报表",后端开始查库、聚合、渲染 Excel,最后返回下载链接。整个过程可能持续几十秒,HTTP 连接一直挂着,中间只要网络抖动一下,用户就要重新点。更麻烦的是,这种同步请求会长期占用 Web 服务器的工作线程,高并发时很容易把连接池打满。
任务队列的思路很简单:把耗时的活包装成一个任务,扔进队列里立刻返回"任务已提交",再由后台 Worker 异步消费。主请求在毫秒级结束,真正的重活交给 Worker 慢慢做。除了削峰,任务队列还能做:
- 解耦:业务系统与具体执行逻辑分离,任务代码可以独立迭代;
- 削峰填谷:突发流量先排队,Worker 按自身能力处理,避免数据库被打挂;
- 定时调度:配合 Celery Beat 实现类 Cron 的周期任务,不用依赖系统 crontab;
- 失败重试:任务失败自动按策略重试,避免手动补偿,也更容易追踪;
- 分布式扩展:Worker 可以部署在多台机器上,按负载水平扩容。
Python 生态里可选方案不少:Celery、RQ、Huey、Dramatiq。Celery 功能最全、社区最大、插件最多,但配置也稍复杂。如果你只是偶尔发几封邮件,RQ 可能更轻;如果要做企业级调度、路由、监控,Celery 仍是首选。
二、Celery 核心架构
上图是 Celery 最常见的部署形态,包含四个核心角色:
- Producer(生产者):你的 Web 应用,调用
task.delay()或task.apply_async()把任务发到 Broker; - Broker(消息中间件):Redis 或 RabbitMQ,负责存储和分发任务消息;
- Worker(消费者):实际执行任务的进程,可以水平扩展;
- Backend(结果后端):Redis、数据库等,用于存储任务结果和状态。
Beat 是独立进程,负责按时间表把定时任务投递到 Broker,本身不执行任务。任务执行仍由 Worker 完成。Producer 和 Worker 之间不直接通信,所有协调都通过 Broker 完成,这也是 Celery 可以水平扩展的关键。
三、Broker 选型:Redis 还是 RabbitMQ
Broker 的选择直接影响可靠性、吞吐量和运维成本。下面是两款最常用 Broker 的对比:
| 维度 | Redis | RabbitMQ |
|---|---|---|
| 部署难度 | 极简单,一条命令启动 | 稍复杂,需 Erlang 环境 |
| 吞吐量 | 非常高,适合短时任务 | 高,支持更复杂路由 |
| 消息持久化 | RDB/AOF,非专业消息持久化 | 原生支持队列和消息持久化 |
| 优先级队列 | 支持有限 | 原生支持 |
| 路由能力 | 基本 | Exchange + Routing Key 非常灵活 |
| 运维监控 | Redis 自身监控 | 自带 Management UI |
| 适用场景 | 中小规模、任务可丢失可重跑 | 金融、电商等对可靠性要求高的场景 |
个人建议:刚起步用 Redis 足够;任务量上来、对可靠性要求变高后,切到 RabbitMQ。切换成本不高,主要改 broker URL 即可。
四、最小可运行示例
4.1 安装依赖
这里用 Redis 做 Broker 和 Backend,安装 redis 客户端与 Celery:
python -m venv venv
source venv/bin/activate # Windows 用 venv\Scripts\activate
pip install celery[redis] redis
4.2 创建 Celery 应用
新建 celery_app.py:
from celery import Celery
app = Celery(
'demo',
broker='redis://localhost:6379/0',
backend='redis://localhost:6379/0',
include=['tasks']
)
app.conf.update(
task_serializer='json',
accept_content=['json'],
result_serializer='json',
timezone='Asia/Shanghai',
enable_utc=True,
)
再新建 tasks.py:
from celery_app import app
import time
@app.task(bind=True, max_retries=3, default_retry_delay=10)
def send_email(self, to: str, subject: str, body: str):
try:
print(f"Sending email to {to} ...")
time.sleep(2) # 模拟网络请求
return {"to": to, "status": "sent"}
except Exception as exc:
raise self.retry(exc=exc, countdown=60)
@app.task
def generate_report(user_id: int, month: str):
print(f"Generating report for user {user_id}, month {month}")
time.sleep(5)
return f"/reports/{user_id}_{month}.pdf"
4.3 启动 Worker
celery -A celery_app worker -l info -P eventlet
Windows 用户建议加 -P eventlet 或 -P gevent,因为 Celery 默认的 prefork 在 Windows 下不支持。Linux/macOS 可以直接跑默认的 prefork,它基于进程,能利用多核 CPU。
4.4 调用任务
再开一个终端:
from tasks import send_email, generate_report
# 异步发送,立即返回 AsyncResult
result = send_email.delay("alice@example.com", "Hi", "Welcome!")
print(result.id) # 任务 ID
print(result.status) # PENDING / SUCCESS / FAILURE
# 生成报表
report = generate_report.delay(42, "2026-07")
print(report.get(timeout=10)) # 阻塞等待结果
五、任务调用高级技巧
delay() 是最简单的调用方式,实际生产中更常用 apply_async(),因为它支持更多参数:
send_email.apply_async(
args=["bob@example.com", "Order", "Your order is shipped"],
countdown=300, # 5 分钟后执行
eta=datetime(2026, 8, 5, 9, 0), # 指定时间点
expires=3600, # 1 小时内有效
queue='email', # 指定队列
priority=5, # 优先级(需 Broker 支持)
retry=True,
retry_policy={
'max_retries': 5,
'interval_start': 30,
'interval_step': 30,
'interval_max': 300,
}
)
5.1 签名(Signatures)与链式任务
有时候任务之间有依赖:先下载图片,再压缩,最后上传 CDN。Celery 用 Signature 可以把任务组合起来:
from celery import chain, group, chord
from tasks import download, compress, upload
# 顺序执行:download -> compress -> upload
pipeline = chain(download.s("url"), compress.s(), upload.s())
pipeline.delay()
# 并行执行 10 个下载
group(download.s(f"url_{i}") for i in range(10)).delay()
# 先并行处理,再汇总
chord(
(download.s(f"url_{i}") for i in range(10)),
compress.s()
).delay()
5.2 任务序列化与参数安全
默认 Celery 用 JSON 序列化消息。JSON 安全、可读、跨语言,但缺点是不支持 Python 的 datetime、Decimal、自定义对象。如果任务参数里有复杂对象,需要自定义序列化器,或者更简单:只传基础类型和 ID,让 Worker 自己查对象。
另一个常见坑是 pickle。虽然 pickle 能序列化任意 Python 对象,但它有安全风险,不要接收不可信来源的任务。生产环境建议显式限制 accept_content=['json'],关闭 pickle。
六、Worker 启动与监控
生产环境启动 Worker 不建议直接 celery worker,而是用进程管理器。下面是 systemd 单元文件示例:
[Unit]
Description=Celery Worker
After=network.target
[Service]
Type=forking
User=www
Group=www
EnvironmentFile=/etc/default/celeryd
WorkingDirectory=/var/www/myapp
ExecStart=/var/www/myapp/venv/bin/celery multi start worker1 \
-A celery_app --logfile=/var/log/celery/worker.log \
--pidfile=/var/run/celery/worker.pid
ExecStop=/var/www/myapp/venv/bin/celery multi stop worker1 \
--pidfile=/var/run/celery/worker.pid
ExecReload=/var/www/myapp/venv/bin/celery multi restart worker1 \
-A celery_app --logfile=/var/log/celery/worker.log \
--pidfile=/var/run/celery/worker.pid
[Install]
WantedBy=multi-user.target
监控方面,Celery 自带的 flower 非常够用:
pip install flower
celery -A celery_app flower --port=5555
Flower 提供 Web UI,可以实时查看任务状态、Worker 负载、失败重试,还能手动触发和撤销任务。配合 Nginx 反向代理和 Basic Auth,可以直接给团队用。
七、定时任务与 Celery Beat
Beat 是 Celery 的调度器。下面配置每分钟执行一次清理日志任务,每天上午九点发送日报:
from celery.schedules import crontab
app.conf.beat_schedule = {
'cleanup-every-minute': {
'task': 'tasks.cleanup_logs',
'schedule': 60.0,
'args': (),
},
'daily-report-at-9am': {
'task': 'tasks.send_daily_report',
'schedule': crontab(hour=9, minute=0),
'args': (),
},
}
app.conf.timezone = 'Asia/Shanghai'
启动 Beat:
celery -A celery_app beat -l info
注意:Beat 不要部署多份,否则同一任务会被重复调度。生产环境建议用 celery beat -s /var/run/celery/beat-schedule 指定持久化文件,并用 Supervisor 或 systemd 保证单实例。如果要求高可用,可以用 django-celery-beat 或 redbeat 把调度信息存到数据库或 Redis,支持多 Beat 抢锁。
八、路由、队列与优先级
不同任务对延迟和可靠性的要求不一样。邮件可以慢几秒,但支付回调必须立刻处理。Celery 支持把任务路由到不同队列:
app.conf.task_routes = {
'tasks.send_email': {'queue': 'email'},
'tasks.process_payment': {'queue': 'critical'},
'tasks.generate_report': {'queue': 'default'},
}
# 启动 Worker 只消费指定队列
celery -A celery_app worker -Q critical,email -l info
下面是常见队列设计对照表:
| 队列 | 任务类型 | Worker 数量 | 优先级策略 |
|---|---|---|---|
| critical | 支付回调、库存扣减 | 多实例、高并发 | 优先保证低延迟 |
| 注册验证、营销邮件 | 2-4 个 | 可容忍分钟级延迟 | |
| report | 报表生成、数据导出 | 1-2 个 | 错峰执行,避免打满 CPU |
| default | 普通业务任务 | 视流量而定 | 默认 FIFO |
九、重试、Acknowledgement 与幂等性
Celery 默认在任务开始执行时就确认消息(ACK)。如果 Worker 在执行中途崩溃,任务可能丢失。对可靠性要求高的任务,建议开启 task_acks_late=True,让 Worker 在任务完成后才确认:
app.conf.task_acks_late = True
app.conf.worker_prefetch_multiplier = 1
重试虽然好,但一定要保证任务幂等。也就是说,同一条任务执行两次,结果应该和执行一次一样。例如发送邮件前先在数据库里查一下"是否已发送",或者给任务参数加一个唯一 dedupe key。
@app.task(bind=True, max_retries=3)
def charge_order(self, order_id: str):
order = Order.get(order_id)
if order.status == "paid":
return {"id": order_id, "status": "already_paid"}
try:
payment_gateway.charge(order)
order.mark_as_paid()
return {"id": order_id, "status": "paid"}
except PaymentError as exc:
raise self.retry(exc=exc, countdown=30)
十、与 Web 框架集成
Celery 最常见的用法是和 Flask、Django、FastAPI 集成。以 FastAPI 为例,触发任务非常简单:
from fastapi import FastAPI
from tasks import generate_report
app = FastAPI()
@app.post("/reports")
def create_report(user_id: int, month: str):
task = generate_report.delay(user_id, month)
return {"task_id": task.id, "status": "submitted"}
@app.get("/reports/{task_id}")
def get_report_status(task_id: str):
result = generate_report.AsyncResult(task_id)
return {
"task_id": task_id,
"status": result.status,
"result": result.result if result.ready() else None
}
关键点:Web 进程只负责投递任务,不要在里面 .get() 阻塞等待结果。如果用户需要轮询状态,单独给一个接口查 Backend。
十一、测试 Celery 任务
单元测试 Celery 任务时,通常不希望真的连 Redis。可以用 task_always_eager=True 让任务同步执行:
import pytest
from celery_app import app
@pytest.fixture(scope='session')
def celery_config():
return {
'broker_url': 'memory://',
'result_backend': 'cache+memory://',
'task_always_eager': True,
}
def test_send_email(celery_app):
from tasks import send_email
result = send_email.delay("test@example.com", "Hi", "Body")
assert result.get()["status"] == "sent"
集成测试可以再开一个真实的 Redis,用 pytest fixture 启动临时 Worker,验证端到端流程。
十二、生产配置速查
下面列出 Celery 里最常用的一批配置项,部署时可以直接参考:
| 配置项 | 推荐值 | 说明 |
|---|---|---|
| broker_url | redis://host:6379/0 | 消息中间件连接地址 |
| result_backend | redis://host:6379/1 | 结果后端,建议与 Broker 分库 |
| task_serializer | json | 任务序列化方式,生产环境避免 pickle |
| task_acks_late | True | 任务完成后才确认消息,防丢失 |
| worker_prefetch_multiplier | 1 | 每个 Worker 每次只取一条任务 |
| worker_max_tasks_per_child | 1000 | 子进程处理多少任务后重启,防内存泄漏 |
| result_expires | 3600 | 结果保留 1 小时后自动清理 |
| task_time_limit | 3600 | 单个任务硬超时时间(秒) |
| task_soft_time_limit | 3000 | 软超时,任务可捕获异常做收尾 |
| timezone | Asia/Shanghai | Beat 调度时区 |
十三、常见陷阱与最佳实践
| 问题 | 原因 | 解决方案 |
|---|---|---|
| 任务丢失 | Worker 崩溃前未 ACK 或中间件配置不当 | task_acks_late=True;RabbitMQ 用持久化队列 |
| 任务重复执行 | 超时重试或 Beat 多实例 | 任务幂等;Beat 只跑单实例 |
| 内存泄漏 | Worker 长时间运行积累对象 | 设置 worker_max_tasks_per_child=1000 |
| 结果后端爆库 | 默认 backend 永久保存结果 | 设置 result_expires=3600 或不用 backend |
| 任务参数过大 | 把图片、大 JSON 直接塞进任务参数 | 传对象 ID 或文件路径,让 Worker 自己读取 |
| 时区混乱 | 未统一 timezone | enable_utc=True,显式设置 timezone |
| 启动失败 | Windows 默认 prefork 不支持 | 加 -P eventlet 或 -P gevent |
| 循环导入 | Celery 应用与任务互相 import | 把 app 和 task 分开放,用 include 注册 |
再补充几条生产经验:
- 不要用数据库做 Broker:SQLAlchemy 结果后端可以,但 Broker 请用 Redis 或 RabbitMQ,否则性能堪忧。
- 任务不要返回大对象:结果后端一般存 Redis,返回大对象会占内存。
- 日志要结构化:结合任务 ID 打日志,排查问题时能串起来。
- Worker 要可观测:Flower + Prometheus exporter,监控队列深度、失败率、执行耗时。
- 部署升级要优雅:先停 Beat,等 Worker 消费完再重启,避免任务中断。
- 参数校验前置:任务函数入口做类型和范围校验,避免无效任务进队列。
十四、总结
Celery 不是银弹,但在 Python 后端里处理异步任务、定时调度、削峰填谷,它仍然是最稳妥的选择。核心要记住三点:第一,任务参数尽量小,大的数据让 Worker 自己拉;第二,任务必须幂等,重试才不会导致业务错乱;第三,根据业务重要性拆分队列,不要把邮件和支付回调塞到同一个 Worker 里。
如果你在本地已经跑通了上面的示例,下一步可以试试把它接到 Flask 或 FastAPI 里,用接口触发 Celery,再用 Flower 观察任务流转。把异步链路跑通之后,很多曾经让你头疼的慢接口问题,都会变得清爽很多。