后端开发有两大刚需:到点自动执行的定时任务,耗时操作不阻塞用户的异步任务。Celery 一个全搞定。
做后端开发,不管什么业务,几乎都会遇到这两类需求:
第一类:定时任务。 电商系统里,用户下单后 30 分钟未支付要自动关单;会员系统里,每天检查有没有会员到期;数据系统里,每天凌晨拉取最新行情数据。这些任务的共同点是——到点了自动跑,不需要人触发。
第二类:异步任务。 用户点击"生成个性化分析报告",后台要跑好几秒甚至几分钟的计算(拉数据、算指标、调 AI 模型)。你不能让用户盯着转圈圈等几分钟,更不能让这个请求占着 Web 服务的线程。正确做法是把任务丢到后台,立即返回,用户稍后来查结果。
Python 生态里,解决这两类需求的标准答案就是 Celery。
这篇文章我用两个真实业务案例,带你从"怎么用"到"底层原理"层层拆解 Celery:
-
案例一(定时任务):订单超时自动关单——讲 crontab 调度、beat 调度器、持久化机制 -
案例二(异步任务):实时个性化分析报告——讲任务投递、worker 执行、进度上报、async 桥接
读完你应该能:独立设计任何定时/异步任务的实现方案,遇到任务不跑/狂跑/漏跑都能自己定位,彻底搞懂 Celery 的运行机制。
一、先认识 Celery:三个角色的三角架构
讲具体案例前,先建立心智模型。很多 Celery 的问题,根因都出在一个误解上:以为调度器是"执行"任务的。
不是。Celery 的架构是三个独立角色,通过消息队列解耦:
┌──────────┐ ① 投递消息 ┌─────────┐ ② 拉取消息 ┌────────┐
│ beat │ ───────────────> │ broker │ <─────────────── │ worker │
│ (定时调度)│ 到队列 │ (Redis) │ │(执行者)│
└──────────┘ └─────────┘ └────────┘
│ │
│ ③ 更新本地 ④ 执行任务函数
│ celerybeat-schedule.db 记录结果到 broker
|
|
|
|
|---|---|---|
| beat |
|
|
| broker(Redis) |
|
|
| worker |
|
|
记住一句话:beat 只投递不执行,worker 才是真正干活的。 这一条理解了,后面所有原理都能自己推演。
注意:beat 是可选的,只有定时任务才需要。异步任务不需要 beat——它由代码主动投递消息,worker 消费执行。两个案例会分别讲清楚。
二、案例一:订单超时自动关单(定时任务)
2.1 业务场景
用户下单后 30 分钟未支付,系统自动关闭订单,释放库存。
为什么不用延时队列(如 RabbitMQ TTL+死信队列)? 那是更精确的方案(每笔订单独立倒计时),但架构更重。对于"每 30 分钟扫一次过期订单"这种粗粒度场景,定时任务足够简单可靠。
2.2 三步实现
第一步:定义任务函数
# tasks.py
from celery import shared_task
from app.core.async_runner import run_async
from app.modules.membership.tasks import _close_expired_orders
@shared_task(name="membership.close_expired_orders")
def celery_close_expired_orders() -> None:
"""订单超时关闭(每30分钟执行)"""
logger.warning("[Celery] membership.close_expired_orders started")
try:
result = run_async(_close_expired_orders())
logger.info("[Celery] 关闭了 %s 笔超时订单", result)
except Exception as e:
logger.exception("[Celery] membership.close_expired_orders failed: %s", e)
raise
注意三个要点:
-
@shared_task(name="...")装饰器是把函数注册为 Celery 任务的关键(后面详解)。 -
函数名 name是全局唯一标识,后续 beat 配置要用。 -
业务逻辑 _close_expired_orders()是 async 的,用run_async()桥接(异步桥接后面案例二详解)。
第二步:配置 beat 调度规则
# celery_app.py
from celery.schedules import crontab
celery_app.conf.beat_schedule = {
# 订单超时关闭(每30分钟执行)
'membership.close_expired_orders': {
'task': 'membership.close_expired_orders',
'schedule': crontab(minute='*/30'),
},
}
第三步:启动 beat 和 worker
# 启动 beat(调度器)
celery -A app.tasks.celery_app beat --loglevel=info
# 启动 worker(执行者)
celery -A app.tasks.celery_app worker -P prefork -c 5 --loglevel=info
两个独立进程,beat 负责到点投递,worker 负责拉取执行。搞定。
2.3 原理拆解一:crontab 配置的默认值陷阱
crontab(minute='*/30') 意思是"每 30 分钟"。但很多人写"每小时一次"时会踩坑:
# ❌ 期望每小时跑一次,实际每分钟都在跑
'schedule': crontab(hour='*/1')
为什么?因为 crontab 的字段默认值都是 '*'(匹配所有值),多字段之间是 AND 关系:
|
|
crontab(hour='*/1')
|
|
|---|---|---|
|
|
'*'
|
|
|
|
'*/1' |
'*',匹配每个小时
|
|
|
|
每分钟触发 |
你以为只设了 hour,但 minute 被默认值"填满"了。正确写法:
# ✅ 每小时整点
'schedule': crontab(minute=0)
# ✅ 每 N 小时
'schedule': crontab(minute=0, hour='*/N')
核心规则:凡是涉及"小时"以上粒度的调度,必须显式指定 minute,否则 minute 默认 * 会退化成每分钟触发。
2.4 原理拆解二:beat 是怎么调度的
beat 进程内部是一个循环,每隔几秒执行一次:
# beat 核心循环(简化版)
while running:
interval = scheduler.tick() # 检查所有任务
time.sleep(interval)
每个 tick() 对每个任务做四件事:
-
算时间:根据上次执行时间 + crontab 规则,算出下一次该执行的时刻。 -
比较:如果"下一次时刻"≤ 当前时间,判定"该触发了"。 -
投递:向 Redis 队列发一条消息:
{
"task": "membership.close_expired_orders",
"id": "uuid-xxx",
"args": [], "kwargs": []
}
注意:消息里只有任务名和参数,没有函数代码。worker 必须自己根据名字找到函数。
-
更新锚点:把"上次执行时间"更新为当前时间,写回本地文件。
2.5 原理拆解三:PersistentScheduler 持久化到底在存什么
你可能注意到 beat 有个配置:
celery_app.conf.beat_schedule_filename = 'store/celerybeat-schedule.db'
先纠正一个常见误解:这个 .db 文件存的是"运行状态",不是"调度规则"。 调度规则在你代码的 beat_schedule 字典里,不在文件里。
那文件里存什么?每个任务存三个字段:
|
|
|
|---|---|
last_run_at |
|
next_run_at |
|
total_run_count |
|
为什么必须持久化? beat 判断"该不该跑"依赖一个时间锚点——上次是什么时候跑的。这个锚点必须跨进程重启存活。
场景:beat 在 14:35 崩溃,15:10 重启
-
没有持久化:beat 不知道历史,假设"从现在 15:10 算"。15:00 这个整点就丢了——这就是漏跑。 -
有持久化:读到 last_run_at=14:00,算出 15:00 该跑。发现 15:00 < 现在 15:10,于是补跑一次,再更新锚点为 15:00。
持久化的根本目的就四个字:不漏不重。
运维提醒:beat 被
kill -9非正常退出可能导致 .db 损坏(shelve 文件没正确关闭)。下次启动报错时,删掉 .db 文件重启即可(代价是丢失锚点,会补跑一次,通常无害)。改了 schedule 结构后也建议删 .db 重新开始,避免旧任务残留。
三、案例二:实时个性化分析报告(异步任务)
3.1 业务场景
用户在交易分析页面点击"生成分析报告",后台要拉取交易数据、计算收益率/回撤/夏普比率等指标,耗时 10-30 秒。用户不想干等,希望:
-
点击后立即返回"正在生成" -
前端轮询进度(0% → 50% → 100%) -
完成后展示报告
这类用户手动触发、耗时较长、需要进度反馈的任务,就是异步任务的典型场景。
3.2 实现:投递 + 执行 + 查询
第一步:定义带进度上报的任务函数
# tasks.py
@shared_task(name="analysis.portfolio_realtime", bind=True)
def celery_portfolio_realtime_analysis(
self, # ← bind=True 注入的 task 实例
portfolio_id: int,
user_id: int,
benchmarks: list[str] | None = None,
start_date: str | None = None,
end_date: str | None = None,
) -> dict:
logger.info("[Celery] analysis.portfolio_realtime started, portfolio_id=%s", portfolio_id)
try:
async def _run():
from app.core.database import get_db_session
from app.modules.portfolio.analysis.service import AnalysisService
async with get_db_session() as db:
service = AnalysisService(db)
result = await service.calculate_realtime_analysis(
portfolio_id, user_id=user_id,
benchmarks=benchmarks,
start_date=start_date, end_date=end_date,
)
return result
run_async(_run())
return {"message": "分析完成"}
except Exception as e:
logger.exception("[Celery] analysis.portfolio_realtime failed: %s", e)
raise
注意和案例一的关键区别——bind=True:
-
bind=True让 Celery 把 task 实例作为第一个参数(self)注入。 -
通过 self可以调用self.update_state()上报进度、self.request.id拿到任务 ID、self.retry()重试。
如果需要进度上报(比如长任务),就加一行:
# 在耗时计算的回调中上报进度
self.update_state(
state="PROGRESS",
meta={"progress": 50, "message": "正在计算夏普比率..."}
)
第二步:API 层投递任务
# router.py
from celery.result import AsyncResult
from app.tasks.tasks_register import celery_portfolio_realtime_analysis
@router.post("/portfolios/{portfolio_id}/analysis")
async def trigger_analysis(portfolio_id: int, db: AsyncSession = Depends(get_db)):
# 投递任务到队列,立即返回 task_id(不等执行完成)
task = celery_portfolio_realtime_analysis.delay(
portfolio_id, user_id=current_user_id
)
return success_response(data={"task_id": task.id})
delay() 是 apply_async() 的简写,它做的是:把任务参数序列化成消息,塞进 Redis 队列,立即返回一个 AsyncResult 对象(含 task_id)。不阻塞,不等结果。
第三步:前端轮询任务状态
# router.py
@router.get("/tasks/{task_id}/status")
async def get_task_status(task_id: str):
result = AsyncResult(task_id)
return success_response(data={
"status": result.status, # PENDING / STARTED / PROGRESS / SUCCESS / FAILURE
"progress": result.info.get("progress", 0) if result.info else 0,
"message": result.info.get("message", "") if isinstance(result.info, dict) else "",
})
前端用 setInterval 每 2 秒轮询这个接口,拿到进度更新进度条,status=SUCCESS 时停止轮询并加载结果。或者用SSE的方式更佳。
3.3 原理拆解四:worker 怎么找到并执行任务函数
worker 启动时做两件关键事:
事1:建立任务注册表(task registry)
worker import 所有被 @shared_task 装饰的模块,建立一张"名字→函数"的全局字典:
{
"analysis.portfolio_realtime": <Task 实例>,
"membership.close_expired_orders": <Task 实例>,
...
}
delay() 投递的消息里的 task 字段,就是查这张表的 key。
事2:拉取消息并执行
worker 从 Redis 用阻塞方式(BRPOP)取消息 → 解析任务名 → 查 registry 找到函数 → 调用执行 → 结果写回 Redis。
完整时序:
T+0:00 API 调用 task.delay(portfolio_id) → 消息进 Redis 队列
T+0:00 API 立即返回 {task_id: "uuid"}
T+0:00 worker BRPOP 取到消息
T+0:00 worker 查 registry → 找到 celery_portfolio_realtime_analysis
T+0:00 worker 调用函数,self.request.id = "uuid"
T+0:05 ├─ service 层计算中,self.update_state(PROGRESS, 50%)
T+0:15 ├─ service 层计算中,self.update_state(PROGRESS, 80%)
T+0:20 └─ 计算完成
T+0:20 worker 写回 {status: SUCCESS, result: {...}}
T+0:22 前端轮询拿到 SUCCESS,停止轮询,展示报告
3.4 原理拆解五:async/await 怎么在同步 worker 里跑
案例二的业务逻辑全是 async(async SQLAlchemy、aioredis),但 Celery worker 是同步执行环境——它调用任务函数时是普通的同步调用。如果你的函数是 async def,worker 拿到的是一个协程对象而不是结果。
桥接方案:进程级持久事件循环
# async_runner.py
import asyncio
_event_loop = None
def get_loop():
"""获取进程级持久事件循环(懒初始化)。"""
global _event_loop
if _event_loop is None or _event_loop.is_closed():
_event_loop = asyncio.new_event_loop()
asyncio.set_event_loop(_event_loop)
return _event_loop
def run_async(coro):
"""在持久事件循环中运行协程。"""
loop = get_loop()
return loop.run_until_complete(coro)
任务函数用同步 def,内部用 run_async() 桥接:
@shared_task(name="analysis.portfolio_realtime", bind=True)
def celery_portfolio_realtime_analysis(self, portfolio_id, user_id, ...):
# ↑ 同步 def ↑ self 由 bind=True 注入
async def _run():
... # 你的 async 业务逻辑
run_async(_run()) # ← 桥接:在持久事件循环里跑协程
return {"message": "分析完成"}
为什么不用 asyncio.run()? 因为它每次创建并销毁一个新事件循环。但 SQLAlchemy 的 async engine、Redis 连接池都绑定到创建它们的事件循环——循环关了连接池跟着销毁,下次任务又得重建,既慢又容易报"跨循环访问连接"的错误。
get_loop() 保证整个 worker 进程生命周期内复用同一个循环,连接池只建一次。
prefork 并发模式的影响:生产环境
-c 5是 5 个独立子进程,进程间内存隔离,所以每个 worker 子进程各自有独立的事件循环和连接池。5 个子进程 = 5 套 DB/Redis 连接池,规划数据库连接池上限时要算进去这个乘数。
四、一个函数要满足什么约束才能成为 Celery 任务
两个案例都用到了 @shared_task,这里把任务定义的约束系统讲清楚。
4.1 硬约束:必须注册到 task registry
唯一的硬性要求:函数必须被 @shared_task(或 @app.task)装饰,登记到全局 task registry。
装饰器在模块被 import 时(不是任务被调用时)做两件事:
-
把函数包装成 celery.Task实例(带run、delay、apply_async、retry等方法)。 -
把 (name, Task实例)注册进全局 registry。
关键推论:worker 进程必须能 import 到被装饰的函数,否则注册不会发生。如果 worker 收到消息但 registry 里查不到对应任务名,会报 NotRegistered 错误。这就是为什么所有 @shared_task 应该集中管理——确保它们都会被 worker 加载。
4.2 调用契约:bind 与函数签名
|
|
|
|
|---|---|---|
False
|
fn(*args, **kwargs) |
|
True |
fn(task_instance, *args, **kwargs) |
self.update_state()、self.retry()、self.request.id
|
案例一(定时任务)不需要 bind,因为不用上报进度。案例二(异步任务)必须 bind=True,因为要用 self.update_state() 上报进度、用 self.request.id 关联任务。
4.3 序列化约束:参数和返回值必须 JSON-safe
这是最容易踩的坑。delay() 传的参数会序列化成消息,函数返回值会写回 result backend。默认 JSON 序列化器只认基本类型:
|
|
|
|---|---|
intfloatstrboolNone |
|
listdict
|
datetime
|
tuple
|
|
案例二里,任务参数都是 int、str、list[str],返回值是 {"message": "分析完成"} 这种纯 dict——全部 JSON-safe。
反例:如果 service 层直接 return portfolio(ORM 对象),worker 序列化结果时直接崩溃。
4.4 约束速查表
|
|
|
|---|---|
| 注册 |
@shared_task 装饰,且被 worker import
|
| 命名 | name=
|
| bind |
bind=True,首参为 self
|
| 参数 |
|
| 返回值 |
|
| 幂等性 |
|
| 超时 |
task_time_limit 约束,超时强杀
|
五、进阶:自定义装饰器与任务监控
实际项目中,你会发现每个 beat 定时任务都在重复写同样的模板代码:logger 开头、try/except、记录执行结果。写多了自然会想抽一个装饰器统一处理。
5.1 需求:统一的任务监控装饰器
期望效果:
@monitored_beat_task(name="stock.indicator_update")
def celery_stock_indicator_update():
# 只写业务逻辑,不用管 logger/try-except/执行记录
return indicator_update_service.update_daily(date.today())
装饰器自动注入:执行日志、幂等锁(防并发)、执行记录持久化、失败原因采集。
5.2 实现:装饰器委托模式
核心认知:自定义装饰器自己不产生 registry 条目,产生条目的是它内部调用的 celery 原生 shared_task。
from celery import shared_task
def monitored_beat_task(name, **celery_kwargs):
def decorator(fn):
# ① 把原函数包成带监控逻辑的函数
def bound_entry(task_self, *args, **kwargs):
return _execute_with_monitoring(fn=fn, task_name=name, ...)
bound_entry.__name__ = fn.__name__
bound_entry.__doc__ = fn.__doc__
# ② 核心:委托给 celery 原生 shared_task
return shared_task(name=name, bind=True, **celery_kwargs)(bound_entry)
return decorator
它做的全部工作,是把"你的原始函数"包成"带监控逻辑的函数"(bound_entry),然后把这个包装函数交给 shared_task。
为什么这个模式能成立?因为 shared_task 不挑食——它不关心传给它的是原始函数还是包装函数,只要拿到一个可调用对象,就登记成任务。对 Celery 而言,bound_entry 就是一个普通函数。
装饰器的本质模式就是:预处理 + 委托给底层机制。
你的代码
└─ @monitored_beat_task (预处理:包监控逻辑)
└─ @shared_task(bind=True) (底层:产生 registry 条目)
└─ bound_entry (实际被登记的 callable)
└─ _execute_with_monitoring
└─ fn() 原始业务函数
5.3 一个深坑:functools.wraps 导致参数错配
写自定义装饰器时,你可能会想用 @functools.wraps(fn) 把原函数的元信息复制过来——千万别这么做。
functools.wraps 会给包装函数打上 __wrapped__ 属性。Celery 内部用 inspect.signature() 分析函数签名时,会沿 __wrapped__ 链回溯,看到的是最里层的原函数签名(可能是零参),而不是包装函数的真实签名(带 task_self)。
Celery 据此认为"这个任务不需要 bind 注入 self",于是 worker 执行时参数对不上 → 运行时崩溃。
正确做法:只手动复制 __name__ 和 __doc__,不用 functools.wraps:
# ✅ 安全
bound_entry.__name__ = fn.__name__
bound_entry.__doc__ = fn.__doc__
六、踩坑速查表
把全文的问题汇总成一张表,遇到问题先查这里:
|
|
|
|
|---|---|---|
|
|
crontab(hour='*/N')
|
crontab(minute=0, hour='*/N')
|
|
|
|
|
|
|
|
last_run_at,必要时手动触发
|
NotRegistered
|
|
|
|
|
|
|
|
|
functools.wraps
|
__name__/__doc__
|
|
|
|
|
|
|
asyncio.run()
|
|
|
|
|
|
七、crontab 配置决策树
最后送一张决策树,以后写定时配置不用再猜:
你要表达什么频率?
│
├─ 每 N 分钟 → crontab(minute='*/N')
│ (只设 minute,安全)
│
├─ 每小时整点 → crontab(minute=0)
│ (最简写法)
│
├─ 每 N 小时 → crontab(minute=0, hour='*/N')
│ ⚠️ 必须显式锁 minute=0
│
├─ 每天 X 点 → crontab(hour=X, minute=M)
│
├─ 每周 X 的 Y 点 → crontab(hour=Y, minute=M, day_of_week='mon')
│
└─ 每月 X 号 → crontab(hour=Y, minute=M, day_of_month=X)
总结:Celery 的设计哲学
回头看,Celery 的设计哲学就一句话:调度与执行解耦,通过消息队列连接。
-
beat 是个可靠的闹钟(定时任务才需要) -
worker 是个勤奋的执行者 -
broker 是它们之间的快递柜
两个案例分别覆盖了 Celery 的两大用途:定时任务靠 beat 按规则投递,异步任务靠代码主动投递。无论哪种,worker 的执行机制都一样——查 registry、调函数、写结果。
理解了这个三角,持久化、并发池、async 桥接、装饰器委托,都是围绕这个三角的工程实现细节。
希望这篇拆解能帮你少走点弯路。如果觉得有用,转发给你身边还在和异步任务较劲的同事。
我是静远,一人公司的 CTO,致力于成为你的AI引路人!

