异步任务:Celery / RQ / APScheduler
Web 请求应该"快进快出"。发邮件、生成报表、调用慢接口、跑机器学习推理——这些活一旦塞进请求里,用户就得盯着转圈。这一章讲清楚:什么时候该上任务队列、Celery 的三个部件怎么配合、任务写重了会死在哪(幂等与重试)、以及那个几乎人人都踩过一次的"任务比事务先跑"的坑。
什么时候该用任务队列
是什么:任务队列把"耗时的活"从请求线程里挪出来,丢进一个队列,由独立的 worker 进程慢慢消费。
为什么需要:因为一个 HTTP 请求的资源是有限的——一个 Gunicorn worker 线程被你占了 30 秒发邮件,这 30 秒里它就服务不了别人。默认 Gunicorn 只开几个 worker,几个慢请求就能把整个站点拖到超时。任务队列把"响应时间"和"任务耗时"解耦:接口 50ms 返回"已提交",邮件在后台慢慢发。
| 信号 | 说明 | 该不该上队列 |
|---|---|---|
| 耗时长 | 单次执行 > 1 秒(发邮件、生成 PDF、调第三方 API) | 该,最典型 |
| 需要重试 | 调外部服务可能失败,要退避重试 | 该,队列自带重试机制 |
| 需要定时 | 每天凌晨对账、每小时清缓存 | 该,用 beat / cron |
| 可以批量聚合 | 1000 条日志攒一起写一次 | 该,省资源 |
| 只是查一次数据库 | 10 毫秒的事 | 不该,队列的收发开销比任务本身还大 |
| 结果必须立刻返回 | 用户下单要立刻知道库存够不够 | 不该,这是同步逻辑 |
论上了队列,你就多了一个"会失败的地方"
① 语义变了:同步调用"要么成功要么抛异常",异步任务"提交成功 ≠ 执行成功"。接口返回 200 只代表"任务进队列了",用户看到的"邮件已发送"可能是假的。② 要多建设三样东西:失败重试、任务状态查询(前端轮询任务 id)、监控告警(任务堆积了没人知道就是灾难)。③ 所以判断标准不是"能不能异步",而是"这个活的失败用户能不能容忍、能不能稍后补偿"。发通知邮件可以异步(失败重发即可);扣款不能异步(必须同步确认)。
常见翻车现场:用户注册接口把"创建用户"也扔进队列了——结果用户马上登录却提示"账号不存在"(任务还在排队)。原则:写数据库的核心链路留在请求里,可延迟的副作用(通知、统计、索引、缩略图)才扔队列。另外,队列不是"无限缓冲带"——broker 里堆几万个任务时,先查是不是有任务在无限重试,而不是忙着加 worker。
Celery 架构:broker / worker / result backend 三件套
是什么:Celery 由三个角色组成。生产者(你的 Django 进程)把任务描述丢进 broker(消息中转站);worker 进程从 broker 取任务执行;执行结果(返回值、状态)写进 result backend(存结果的地方)。
为什么要分成三个:因为职责不同,得能各自独立扩缩容。请求量大就多加 Django 进程,任务积压就多加 worker,broker 和 backend 是基础设施不用动。混在一起的话你没法"只扩 worker"。
| 部件 | 作用 | 常用选型 | 关键要求 |
|---|---|---|---|
| broker | 存放待执行的任务消息 | Redis、RabbitMQ | 快、可靠、别丢消息 |
| worker | 真正执行任务的进程 | celery -A proj worker | 可水平扩容,要能优雅重启 |
| result backend | 存任务返回值与状态 | Redis、数据库、可直接不配 | 只在需要查结果时才配 |
| beat | 定时触发器(第 5 节详述) | celery -A proj beat | 全局只能有一个实例 |
| flower | 任务监控面板 | celery -A proj flower | 生产建议加鉴权 |
论RabbitMQ 还是 Redis?一句话决定
① RabbitMQ:专业的消息中间件,支持消息确认、持久化、复杂路由、优先级队列。队列不丢消息是它的核心卖点。金融、订单这类"丢一条任务就是事故"的场景选它。② Redis:你已经用 Redis 做缓存了,那就直接用它当 broker——少装一个组件,运维成本低,性能也够。但要知道代价:Redis 作为 broker 在极端情况下(内存打满被 OOM Killer 干掉、未开 AOF 持久化)可能丢消息,而且不支持 Celery 的部分高级特性(如优先级)。③ 结论:小项目 / 内部系统用 Redis;对消息可靠性有硬要求用 RabbitMQ;不确定就先用 Redis,接口抽象好了将来换成本不高。
很多人复制模板时顺手配了 result_backend = "redis://...",然后每个任务的返回值(包括 None 和一大坨日志字符串)都写进 Redis,还默认不设过期时间。几天后 Redis 内存告警。正确做法:① 只在确实要查任务结果时才配 backend,其余场景把它去掉;② 配上 result_expires = 3600 让结果一小时过期;③ 任务里别返回大对象(返回 id,别返回整个 DataFrame)。
Redis 作 broker:最小可跑配置
怎么用:四步——装依赖、配 Django、写任务、起 worker。
启动(开发环境开两个终端)
Celery 4 之后不再支持 Windows 的默认多进程池。本地开发用 -P solo(celery -A config worker -l info -P solo)或在 WSL/Docker 里跑,别在 Windows 上死磕。另一个高频问题:任务明明 delay() 了却一直不动——按顺序查:① worker 起了吗(看终端有没有 ready);② 任务所在的 tasks.py 在 INSTALLED_APPS 的 app 里吗(autodiscover_tasks 只扫已注册的 app);③ 队列名对不对(用了 -Q 就得两边一致);④ 改了代码没重启 worker(worker 不会自动重载,开发时可以配 watchmedo 或 --autoreload)。
任务幂等与重试策略:任务一定会重复执行
是什么:幂等指一个任务执行一次和执行三次,业务结果一样。为什么必然需要:因为队列的可靠性保证是"至少一次(at-least-once)",不是"恰好一次"。worker 执行到一半被 kill、网络抖动导致 ack 丢失、你手动重放——任何一条都会让任务跑第二遍。Celery 默认 acks_late=False(收到就 ack),web 上老教程常推荐改 acks_late=True 提高可靠性,代价就是"重投变多"。
| 异常类型 | 该不该重试 | 理由 |
|---|---|---|
| 网络超时 / 连接被拒 | 该,指数退避 | 临时故障,等一下可能就好了 |
| 对方返回 429 / 503 | 该,配合 Retry-After | 限流,退避后重试合理 |
| 数据校验失败(参数非法) | 不该 | 重试一百次也是同样的错 |
| KeyError / TypeError 等代码 bug | 不该 | 重试只是拖延暴露问题的时间 |
| 数据库死锁 | 该,短退避 | 重试通常能过 |
| 余额不足、库存不足 | 不该 | 业务结果确定,重试无意义 |
论重试的正确姿势:退避 + 抖动 + 上限
① 为什么不能立刻重试:对方如果已经过载,你立刻重试等于火上浇油;一万个任务同时立刻重试就是"重试风暴",能把依赖服务打得再也起不来。② 指数退避:1s、2s、4s、8s……给对方恢复时间。③ 抖动 jitter:不加抖动时所有任务的重试时刻是同步的,尖峰仍然存在;加随机量把尖峰摊平。④ 上限与死信:超过 max_retries 后任务进入失败状态,必须有人知道——所以第 7 节要配 Flower 告警和死信处理。无限重试是最坏的设计:队列只涨不消,最后整个系统因为它瘫痪。
autoretry_for=(Exception,) 是最常见的偷懒写法——它把 KeyError、TypeError 这类代码缺陷也吞进重试循环,日志里只看到"retrying",看不到真正的 traceback,排查时间翻几倍。把异常列表收窄到你真的想重试的那几类;重试时把 exc 打到日志里;给任务加 expires(比如 expires=3600)避免"过期任务"在几小时后还去发一封昨天的邮件。
beat 定时任务:谁在"到点了"
是什么:celery beat 是一个独立进程,按配置的时间表往 broker 里投递任务。worker 只管执行,触发是 beat 的活。
为什么不用系统 cron:cron 调用 Django 需要额外包一层 management command,任务状态、重试、并发控制都得自己写。beat 复用整套 Celery 生态,而且和 worker 解耦——任务跑到哪台机器都行。
beat 用本地文件 celerybeat-schedule 记录"上次跑过没"。你部署时不小心起了两个 beat 进程(或容器副本数设成 2),它们各有各的 schedule 文件,同一个任务就被投递两次——"每天给用户发一封日报"变成两封。规则:beat 单实例部署,用 systemd 的 Restart=always 或 K8s 的 replicas: 1 保证;--schedule 路径放到持久化目录,容器重建也不会重复触发。任务侧仍然要写幂等(第 4 节),这是双保险。
事务与投递:transaction.on_commit,这个坑人人踩一次
现象:任务偶尔报"订单不存在"(Order matching query does not exist),而且只在生产环境偶发,本地从来复现不了。打日志看时间戳——任务执行时间比订单创建时间还早。
原因:Django 在 ATOMIC_REQUESTS=True 或 @transaction.atomic 下,视图里的所有写操作都在一个事务里,直到视图返回才真正 COMMIT。而 send_order_mail.delay(order.id) 这一行执行时,事务还没提交。worker 是另一个进程、另一个数据库连接,它读不到未提交的数据——查不到订单,任务失败;如果配了 autoretry_for=(Exception,),它会重试几次,恰好事务提交了就成功,所以表现为"偶发失败"。本地用 SQLite 单事务可能直接成功,所以你复现不了。
论为什么要回调,而不是"手动 commit 后再投"
① 事务可能根本没提交:视图后面某一步抛异常,整个事务回滚——订单压根不存在,但你前面已经把"发确认邮件"的任务投出去了。用户会收到一封"订单已创建"却查不到订单的邮件。② on_commit 保证的是"要么都做,要么都不做":提交成功 → 投递;回滚 → 回调被丢弃不执行。③ 测试里的坑:TestCase 把每个用例包在事务里然后回滚,on_commit 回调永远不会执行,你会以为任务没被调用。改用 TransactionTestCase,或 with self.captureOnCommitCallbacks(execute=True):。④ 另一个变体:异步视图(async def)里不能用同步的 transaction.on_commit,要用 sync_to_async(transaction.on_commit)(...) 或 await aon_commit(...)(Django 4.2+)。
for o in orders: transaction.on_commit(lambda: send(o.id))——循环变量 o 是晚绑定的,回调真正执行时 o 已经是最后一个对象,你给所有人发的都是同一封邮件。修法:lambda oid=o.id: send(oid) 用默认参数把值绑死,或者 functools.partial。这个坑和第 1 章"可变默认参数"是同一类问题:Python 的闭包捕获的是变量,不是值。
RQ:轻量替代方案,什么时候别上 Celery
是什么:RQ(Redis Queue)是只依赖 Redis 的极简任务队列:一个 enqueue 投递、一个 rq worker 消费。没有 beat、没有复杂的路由、配置不到 10 行。
为什么需要知道它:Celery 很强但很重——一个项目只有"发邮件 + 生成报表"两个异步需求,配 Celery 的认知成本可能超过收益。RQ 的定位就是"Redis 之上的 200 行级任务队列"。
| 维度 | Celery | RQ | APScheduler |
|---|---|---|---|
| 定位 | 完整分布式任务平台 | Redis 上的轻量队列 | 进程内定时调度器 |
| 依赖 | Redis / RabbitMQ | 只依赖 Redis | 无(跑在进程内) |
| 定时任务 | beat,功能完整 | 无内置,需外部触发 | 强项,cron 表达式直观 |
| 重试 | 丰富(退避/抖动/自定义) | 基础 Retry | 需自己写 |
| 跨机器 | 原生支持 | 支持 | 不支持(进程内) |
| 适用 | 中大型、任务种类多 | 小项目、需求简单 | 单机定时(爬虫、清理) |
| 上手成本 | 高 | 低 | 最低 |
APScheduler 是"进程内调度器"——你在 Django 里 BackgroundScheduler().start(),每个 gunicorn worker 进程都会启动一份自己的调度器。开 4 个 worker,你的"每天 3 点清理"就执行 4 次。它适合:单进程的爬虫脚本、常驻的数据采集进程、本地工具。不适合:多 worker 的 Web 服务(用 Celery beat),或需要跨机器协调的场景。
可观测:Flower、失败告警与死信重放
为什么是必修课:没有监控的任务队列等于"把定时炸弹藏起来"。任务从什么时候开始堆积、哪个任务一直在失败、worker 是不是全挂了——这些必须能一眼看到。
| 要监控的指标 | 怎么看 | 告警阈值建议 |
|---|---|---|
| 队列深度 | Flower / Redis LLEN / rq info | 持续 > 1000 或持续增长 |
| 任务失败率 | Flower 的 Failed 面板、task_failure 次数 | 5 分钟内失败 > 10 条 |
| 任务耗时 | Flower 的 runtime 曲线 | P95 超过预期 3 倍 |
| worker 存活 | 部署系统探针 + Flower 在线列表 | 在线 worker 数 = 0 |
| 定时任务是否按时跑 | beat 日志 + 任务开始时间 | 延迟超过一个周期 |
论死信与手动重放:失败的任务去哪儿了
① 死信(dead letter):重试超过 max_retries 后任务标记为 FAILURE,不再自动执行。它不会消失,但也不会自己恢复。② 怎么捞出来:开启 task_track_started 并配 result backend,Flower 的 Failed 列表里能看到参数与异常,点进去可手动 requeue;或者写一个 management command 扫 celery-task-meta-* 键筛出失败任务。③ 重放的前提:任务必须幂等(第 4 节)——否则重放一次就多加一次积分。④ 更稳的做法:关键业务不要在任务里"只做一次就扔",把任务做成"扫描 + 处理"的收敛式:任务失败不重试,由"每分钟扫描异常单据"的定时任务兜底,这样系统永远能自愈,且天然幂等。
worker 是独立进程,日志不会出现在 Django 的日志文件里——你在任务里 print() 是白印的,在视图里配的 logging 也不会自动继承。要单独配:在 worker 启动参数里加 --logfile=/var/log/celery/worker.log -l info,或用 celery.signals.setup_logging 把 Django 的 LOGGING 配置注入进去。排查顺序建议:① Flower 看任务状态与异常;② worker 日志看 traceback;③ 打开 -l debug 看任务收发细节。三样都没有,就只剩下猜了。
① 判断标准:耗时长、要重试、要定时、可批量 → 上队列;核心写链路和不许延迟的逻辑留在请求里。
② 三件套:broker 存消息、worker 干活、result backend 存结果(不需要就别配,配了就设 result_expires)。小项目 Redis 当 broker 够用。
③ 幂等是前提:队列保证"至少一次",重复执行是常态。用唯一键或状态机把重复变成无害。
④ 重试要克制:只重试值得重试的异常,配指数退避 + 抖动 + 上限,别用 (Exception,) 吞掉代码 bug。
⑤ transaction.on_commit:请求里的任务投递必须挂到事务提交之后,否则 worker 查不到数据——这是最经典的"本地不复现、线上偶发"故障。
⑥ 可观测:Flower + task_failure 告警 + beat 单实例,缺一个都会在半夜两点叫醒你。
1.(概念题)Celery 的 broker 和 result backend 各是干什么的?什么情况下可以不配 result backend?
查看答案
答案:broker 存放待执行的任务消息(Redis / RabbitMQ),worker 从这里取任务;result backend 存放任务的返回值与状态。如果业务不需要查询任务结果(发邮件、清缓存这类"投出去就不用管"的任务),就不配 backend——可以省下 Redis 内存,避免结果堆积把内存撑爆。
2.(概念题)为什么说"任务一定会重复执行"?重复执行会导致什么问题?怎么解决?
查看答案
答案:因为队列的投递语义是 at-least-once:worker 执行中崩溃、ack 丢失、acks_late=True 下重投、手动重放,都会让任务跑第二遍。重复会导致重复加积分、重复发通知、重复扣款。解决靠幂等:用唯一键(unique 约束 + 捕获 IntegrityError)或状态机(filter(status="pending").update(status="running") 抢锁)。
3.(代码题)下面这段代码为什么会出现"订单不存在"的偶发报错?怎么修?order = Order.objects.create(...) → send_mail.delay(order.id)(视图上有 @transaction.atomic)
查看答案
答案:事务还没提交,worker 是另一个连接,查不到未提交的订单。修法:transaction.on_commit(lambda: send_mail.delay(order.id)),让投递挂在提交成功之后;回滚时回调自动丢弃。
4.(概念题)autoretry_for=(Exception,) 有什么问题?你会怎么配?
查看答案
答案:它把代码 bug(KeyError、TypeError)也纳入重试,日志里只见"retrying"看不到真实 traceback,问题被拖延;还会造成无意义的重试风暴。应该收窄成 (ConnectionError, Timeout) 这类临时故障,配 retry_backoff=True、retry_jitter=True、合理的 max_retries 和 expires。
5.(思考题)业务只有"用户注册后发一封欢迎邮件",你在 Celery 和 APScheduler 之间选哪个?如果把整个 Django 换成单机 Flask 小工具呢?
查看答案
答案:Django 多进程部署时必须选 Celery(或 RQ)——APScheduler 在进程内,几个 worker 就会重复执行,且没有一个"可重试、可观测"的执行体。如果是单机上跑的小工具、只有一个进程,APScheduler 或 threading.Timer 就够;甚至可以在响应后直接开线程发邮件(但要接受进程退出会丢任务的代价)。核心判断依据是:有几个进程、能否接受任务丢失。