大数跨境

Celery 从入门到原理拆解:用两个真实案例讲透定时任务与异步任务

Celery 从入门到原理拆解:用两个真实案例讲透定时任务与异步任务 静远AI出海
2026-07-16
4
导读:后端两大刚需——定时任务和异步任务,Celery 一个全搞定。本文从订单超时关单、实时分析报告两个真实案例出发,层层拆解 beat/worker/broker 三角架构、crontab 调度原理、持久

后端开发有两大刚需:到点自动执行的定时任务,耗时操作不阻塞用户的异步任务。Celery 一个全搞定。

做后端开发,不管什么业务,几乎都会遇到这两类需求:

第一类:定时任务。 电商系统里,用户下单后 30 分钟未支付要自动关单;会员系统里,每天检查有没有会员到期;数据系统里,每天凌晨拉取最新行情数据。这些任务的共同点是——到点了自动跑,不需要人触发

第二类:异步任务。 用户点击"生成个性化分析报告",后台要跑好几秒甚至几分钟的计算(拉数据、算指标、调 AI 模型)。你不能让用户盯着转圈圈等几分钟,更不能让这个请求占着 Web 服务的线程。正确做法是把任务丢到后台,立即返回,用户稍后来查结果

Python 生态里,解决这两类需求的标准答案就是 Celery

这篇文章我用两个真实业务案例,带你从"怎么用"到"底层原理"层层拆解 Celery:

  • 案例一(定时任务):订单超时自动关单——讲 crontab 调度、beat 调度器、持久化机制
  • 案例二(异步任务):实时个性化分析报告——讲任务投递、worker 执行、进度上报、async 桥接

读完你应该能:独立设计任何定时/异步任务的实现方案,遇到任务不跑/狂跑/漏跑都能自己定位,彻底搞懂 Celery 的运行机制。


一、先认识 Celery:三个角色的三角架构

讲具体案例前,先建立心智模型。很多 Celery 的问题,根因都出在一个误解上:以为调度器是"执行"任务的

不是。Celery 的架构是三个独立角色,通过消息队列解耦:

  
  
  
┌──────────┐   ① 投递消息     ┌─────────┐   ② 拉取消息     ┌────────┐
│   beat   │ ───────────────> │  broker │ <─────────────── │ worker │
│ (定时调度)│   到队列          │ (Redis) │                  │(执行者)│
└──────────┘                  └─────────┘                  └────────┘
     │                                                          │
     │ ③ 更新本地                                              ④ 执行任务函数
     │   celerybeat-schedule.db                                记录结果到 broker
角色
职责
类比
beat
定时检查调度规则,到点了把"执行任务 X"的消息塞进队列
闹钟:到点提醒你,但它不替你干活
broker(Redis)
消息中转站,存待执行的任务消息
快递柜:beat 投递,worker 取件
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')
 的值
含义
minute
'*'
(默认!)
匹配每一分钟
hour
'*/1'
等价于 '*',匹配每个小时
AND 组合
每分钟触发

你以为只设了 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() 对每个任务做四件事:

  1. 算时间:根据上次执行时间 + crontab 规则,算出下一次该执行的时刻。
  2. 比较:如果"下一次时刻"≤ 当前时间,判定"该触发了"。
  3. 投递:向 Redis 队列发一条消息:
  
  
  
{
  "task": "membership.close_expired_orders",
  "id": "uuid-xxx",
  "args": [], "kwargs": []
}

注意:消息里只有任务名和参数,没有函数代码。worker 必须自己根据名字找到函数。

  1. 更新锚点:把"上次执行时间"更新为当前时间,写回本地文件。

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 秒。用户不想干等,希望:

  1. 点击后立即返回"正在生成"
  2. 前端轮询进度(0% → 50% → 100%)
  3. 完成后展示报告

这类用户手动触发、耗时较长、需要进度反馈的任务,就是异步任务的典型场景。

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"0if result.info else 0,
        "message": result.info.get("message"""if isinstance(result.info, dictelse "",
    })

前端用 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 时(不是任务被调用时)做两件事:

  1. 把函数包装成 celery.Task 实例(带 rundelayapply_asyncretry 等方法)。
  2. 把 (name, Task实例) 注册进全局 registry。

关键推论:worker 进程必须能 import 到被装饰的函数,否则注册不会发生。如果 worker 收到消息但 registry 里查不到对应任务名,会报 NotRegistered 错误。这就是为什么所有 @shared_task 应该集中管理——确保它们都会被 worker 加载。

4.2 调用契约:bind 与函数签名

bind 值
worker 调用方式
适用场景
False
(默认)
fn(*args, **kwargs)
简单任务,不需要 task 实例
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
ORM 对象(SQLAlchemy model 实例)
listdict
(值也需可序列化)
datetime
(用 ISO 字符串传)
tuple
(序列化为 list)
自定义类实例、函数、lambda

案例二里,任务参数都是 intstrlist[str],返回值是 {"message": "分析完成"} 这种纯 dict——全部 JSON-safe。

反例:如果 service 层直接 return portfolio(ORM 对象),worker 序列化结果时直接崩溃。

4.4 约束速查表

约束
说明
注册
必须被 @shared_task 装饰,且被 worker import
命名 name=
 全局唯一
bind
需要用 task 实例时必须 bind=True,首参为 self
参数
必须可 JSON 序列化
返回值
必须可 JSON 序列化
幂等性
beat 可能重投导致多次执行,建议幂等设计
超时
受 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')
 未锁 minute
改成 crontab(minute=0, hour='*/N')
定时任务不跑
beat 没启动 / 时刻没到 / worker 全挂
查 beat 进程、worker 日志、broker 队列
定时任务漏跑
beat 长时间宕机
查 last_run_at,必要时手动触发
worker 报 NotRegistered
任务函数没被 worker import
检查 include 配置,确保 tasks 模块被加载
异步任务结果序列化失败
返回了 ORM 对象
service 层转成 dict 再返回
自定义装饰器任务参数崩溃
用了 functools.wraps
改用手动复制 __name__/__doc__
慢任务卡住整个 worker
solo 模式串行执行
生产用 prefork 并发
async 任务报"事件循环已关闭"
用了 asyncio.run()
改用进程级持久事件循环
beat 重启报 shelve 损坏
上次非正常退出
删 .db 文件重启

七、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引路人!

【声明】内容源于网络
0
0
静远AI出海
AI SaaS应用开发脚手架项目Fast SaaS的作者。致力于分享各类AI工具和大型模型的使用经验,探索高效的AI应用开发方法,专注AI应用的出海与变现。
内容 390
粉丝 0
静远AI出海 AI SaaS应用开发脚手架项目Fast SaaS的作者。致力于分享各类AI工具和大型模型的使用经验,探索高效的AI应用开发方法,专注AI应用的出海与变现。
总阅读372
粉丝0
内容390