第10章:Celery 任务状态机与生命周期
2026/9/4 9:04:47 网站建设 项目流程

0. 上一章思考题参考答案

思考题 1:消息被发进无人订阅的sms队列后,会永久停留在队列里(受 Broker 的消息/队列过期策略约束),Worker 侧无感知,任务在调用方视角永远是 PENDING。它被「发现」的时机只有三个:① 有人检查 Broker 队列深度发现异常堆积;② 按任务 ID 手动追查;③ 队列积压告警触发。所以队列规划必须配「积压告警」,否则黑洞消息可以沉默几个月。

思考题 2:预取数 = 并发数 ×worker_prefetch_multiplier(默认 4)。报表 Worker 预取 1×4=4 条(重任务 4 条已够堵 32 秒),短信 Worker 预取 4×4=16 条。预取过大的危害:消息被取走但堆在执行进程内存里排队,其他 Worker 即使空闲也拿不到——多 Worker 场景下造成「看起来都忙、实际不均」,甚至饿死新加入的 Worker。第 18 章会给不同任务的推荐倍数。


1. 项目背景

运营平台的「任务中心」要上线了:订单任务要展示时间线(提交 → 领取 → 执行中 → 成功/失败),报表导出要支持「取消还没开始的导出」。小周翻出代码,发现现状惨不忍睹:任务状态全靠猜——日志里搜 task_id,看到「succeeded」就算成功,看不到就默认「还在跑」;撤销功能压根没有,运营只能等 8 分钟的导出跑完再删文件。

更离谱的是,上次一个导出任务跑了 30 分钟没反应,小周在 Flower 上点了「revoke」,任务照样跑完——事后才知道:revoke 只对「还没开始执行」的任务生效,任务一旦被 Worker 拾取开始执行,revoke 就管不着了。而「terminate」倒是能杀,但他不知道有这回事,更不知道杀进程可能把数据库连接池的写一半的事务给「咔嚓」了。

现状:状态管理 = 翻日志 任务中心需要: ① 时间线:PENDING → RECEIVED → STARTED → SUCCESS / FAILURE / RETRY ② 撤销:区分「还没开始」(revoke)与「执行中」(terminate / AbortableTask) ③ 一致性:状态展示与任务真实生命周期严格对齐

本章目标:读懂 Celery 的状态机与状态常量(celery/states.py),做出任务中心的时间线与撤销功能,并搞清楚 revoke / terminate / AbortableTask 三者的适用边界。


2. 项目设计

场景:任务中心评审会,小周把「revoke 不生效」的现场复述了一遍。

小胖:状态不就俩吗——成功、失败。中间那些 RECEIVED、STARTED 是给谁看的?我玩游戏的读条也没这么细啊。

小白:我查了celery/states.py,状态常量有PENDINGRECEIVEDSTARTEDSUCCESSFAILURERETRYREVOKEDIGNORED一长串。我的问题是:这些状态是谁写的?写在哪儿?而且默认情况下好像只能看到 PENDING 和 SUCCESS,中间态去哪了?

大师:两个关键点。第一,状态写在 Backend 里(第 8 章的结果存储),Worker 在任务生命周期各节点调用 Backend 写状态。第二,中间态不是没有,是默认不写——task_track_started默认 False(celery/app/defaults.py),因为每次状态写入都是一次网络往返,写 STARTED 意味着任务执行前先多一次 Backend 写入,高频任务会显著增加 Backend 负载。所以设计上默认「只写两端(最终态),中间靠事件流(Events,第 25 章)」——RECEIVEDSTARTED其实是事件状态,默认不出现在 Backend 状态里。任务中心要时间线,要么开task_track_started(接受写放大),要么消费事件流(第 25 章的正道)。

技术映射:状态写在 Backend = 快递单上的「签收章」(少而关键);事件流 = 快递轨迹(每个节点都记,但量大事多)。两个体系各司其职。

小胖:那撤销呢?我那天 revoke 点下去,任务还跑得欢,是不是 bug?

大师:不是 bug,是语义边界。Celery 的撤销分三个档位:revoke(revoke 命令/API)——对尚未被 Worker 拾取的任务生效:Worker 消费时发现任务在撤销集合里,直接标记 REVOKED 跳过执行;已经跑起来的任务不中断terminate——对正在执行的任务发 SIGTERM 强杀子进程(prefork 池才有进程可杀),强制中断,但代价是任务写了一半的状态没人收尾(数据库连接、文件句柄)。AbortableTaskcelery/contrib/abortable.py)——协作式撤销:任务内部周期性检查self.is_aborted(),自己优雅收尾退出。一句话:revoke 拦「没上场的」,terminate 杀「在场上的」,AbortableTask 是「自己举手退场」。你们导出任务要支持取消,最优雅的方案是 AbortableTask。

小白:那 revoke 为什么有时对「还没开始」的任务也失效?我见过任务一直 PENDING,revoke 之后状态还是 PENDING。

大师:好问题,两个原因:①revoke 靠广播——control.revoke通过 pidbox 向所有 Worker 广播撤销集合,Worker 重启、广播窗口、新 Worker 加入都会漏(Mingle 机制就是启动时同步撤销集合,第 33 章讲);②PENDING 是「未知」的别名states.py注释:Task state is unknown),revoke 后任务仍可能显示 PENDING 直到 Worker 真正消费并写下 REVOKED。所以生产上撤销的验收标准是:看最终态是否 REVOKED,而不是看它「马上变」。另外还有个补丁:REVOKED是终态(在READY_STATES里),它和 FAILURE 一样会进结果集合。

技术映射:revoke = 球场广播「把没上场的 7 号换掉」(已经上场的听不见);terminate = 裁判直接红牌罚下(粗暴有效);AbortableTask = 球员自己申请换人(体面退场)。

小胖:哦——那我导出任务就上 AbortableTask。不过话说回来,状态这么多,任务中心判断「任务完没完」要看哪些?

大师:记住三个集合(celery/states.py):READY_STATES = {SUCCESS, FAILURE, REVOKED}(出最终结果了);UNREADY_STATES = {PENDING, RECEIVED, STARTED, REJECTED, RETRY}EXCEPTION_STATES = {RETRY, FAILURE, REVOKED}。判断「完没完」就查state in READY_STATES,等价于AsyncResult.ready()。还有个高级玩具:state类实现了按优先级比较(precedence),state(PENDING) < state(SUCCESS)为 True——写「状态回退」校验(不允许 SUCCESS 之后又出现 STARTED)时特别有用。


3. 项目实战

3.1 环境准备

沿用环境(Redis Broker + Backend)。本章给「导出对账单」任务升级为可撤销 + 可展示时间线

3.2 分步实现

步骤 1:开启track_started,补齐中间态

目标:让时间线出现 STARTED 节点(先理解代价再开启)。

# celeryconfig.py 追加task_track_started=True# 记录 STARTED 状态:每次任务执行前多一次 Backend 写入

步骤 2:把导出任务升级为 AbortableTask

目标:支持「执行中优雅取消」,而不是傻等 8 分钟。

# order_tasks.pyfromcelery.contrib.abortableimportAbortableTask@app.task(name='orders.export_statement',bind=True,base=AbortableTask,track_started=True)defexport_statement(self,date_str:str)->str:"""导出对账单:分批处理,每批检查是否被请求取消。"""importtime url=f"https://fs.xx.com/statement/{date_str}.csv"forbatchinrange(20):# 模拟 20 批写入ifself.is_aborted():# 协作式撤销:自己举手退场print(f"[export]{date_str}已取消于第{batch}批")self.update_state(state='ABORTED')# 自定义状态:任务中心可见return''time.sleep(0.5)print(f"[export]{date_str}完成:{url}")returnurl

步骤 3:实现任务中心的撤销接口

目标:区分「未开始 → revoke」「执行中 → abort 请求」。

# cancel_view.pyfromcelery.resultimportAsyncResultfromcelery.contrib.abortableimportAbortableAsyncResultfromorder_tasksimportapp,export_statementdefcancel_statement(task_id:str)->dict:r=AsyncResult(task_id,app=app)ifr.statein('PENDING','RECEIVED','STARTED'):ifr.state=='STARTED':AbortableAsyncResult(task_id,app=app).abort()# 执行中:请求协作取消return{"action":"abort-requested","state":r.state}r.revoke()# 未开始:直接撤销return{"action":"revoked","state":r.state}return{"action":"none","state":r.state,"reason":"任务已终态"}

步骤 4:全流程验证时间线与撤销

celery-Aorder_tasks worker--loglevel=info--pool=solo
# 终端 B:celery-Aorder_tasks call orders.export_statement--args='["2026-08-23"]'# 得到 task_id 后立刻查看状态(未执行时)celery-Aorder_tasks result<task_id># 显示状态 PENDING# 执行中取消python-c"from cancel_view import cancel_statement; print(cancel_statement('<task_id>'))"

运行结果(文字描述):

时间线(执行正常时):PENDING → STARTED → SUCCESS 时间线(中途取消时):PENDING → STARTED → ABORTED(Worker 日志打印「已取消于第 N 批」) 未开始即取消时:PENDING → REVOKED(Worker 日志打印 discarded 跳过执行)

步骤 5:演示 terminate 的暴力与风险(可选,仅 prefork)

# prefork 池(Linux):执行中的任务被 SIGTERM 强杀celery-Aorder_tasks inspect terminate<task_id># 发送强杀指令# Worker 日志出现 Terminated 15,任务无正常收尾

运行结果(文字描述):任务进程被信号杀死,没有SUCCESS/FAILURE 落库(或短暂 TERMINATED 后超时);若任务持有数据库事务,需依赖连接超时回滚。这就是 terminate 的代价。

步骤 6:用状态优先级做「回退校验」

目标:任务中心的时序守卫——防止脏数据把状态倒着写。

# state_guard.pyfromceleryimportstatesdefshould_accept(current:str,incoming:str)->bool:"""只允许状态「前进」:优先级低的旧状态不允许覆盖高优先级终态。"""returnstates.state(incoming)>=states.state(current)print(should_accept('SUCCESS','SUCCESS'))# True:幂等重复print(should_accept('SUCCESS','STARTED'))# False:终态后不允许回退print(should_accept('PENDING','STARTED'))# True:正常推进print(should_accept('FAILURE','REVOKED'))# True:REVOKED 优先级更高,可覆盖

运行结果:True / False / True / True。任务中心写入时间线前先过这个守卫,「明明成功了又显示执行中」的脏数据从源头断掉——这就是celery/states.pyprecedencestate比较器(支持<>=)的真实用途。

3.3 可能遇到的坑及解决方法

现象解决
revoke 后任务照跑任务已被 Worker 拾取开始执行执行中取消用terminateAbortableTask.abort()
revoke 后状态还是 PENDINGPENDING 是「未知」的别名,等 Worker 消费才写 REVOKED验收看最终态;必要时消费端强制刷新
abort() 后任务不停任务内没有周期检查is_aborted()AbortableTask 是协作式的,检查点必须写进任务体
开了 track_started,Redis 变慢每个任务多一次状态写入只在需要时间线的任务上开track_started=True,别全局开
terminate 后数据库锁没释放进程被强杀,事务连接残留依赖 DB 连接超时;关键写操作别依赖 terminate

3.4 完整代码清单与测试验证

清单:celeryconfig.py(track_started)、order_tasks.py(AbortableTask 导出任务)、cancel_view.py(撤销接口)。状态字典表(任务中心开发参照):

状态含义谁写是否终态
PENDING未知/未消费默认假设
RECEIVEDWorker 已领取(事件)Events
STARTED开始执行(需 track_started)Worker
SUCCESS成功Worker
FAILURE失败Worker
RETRY等待重试Worker
REVOKED已撤销Worker(消费时发现撤销)

测试验证:

# tests/test_states.pyimportpytestfromceleryimportstatesfromorder_tasksimportapp,export_statement app.conf.task_always_eager=Truedeftest_ready_states_are_terminal():assertstates.SUCCESSinstates.READY_STATESassertstates.REVOKEDinstates.READY_STATESassertstates.PENDINGinstates.UNREADY_STATESdeftest_state_precedence_blocks_rollback():# SUCCESS 之后不应再出现 STARTED:优先级比较可做校验assertstates.state(states.STARTED)<states.state(states.SUCCESS)deftest_export_flow_success():r=export_statement.apply(args=['2026-08-23'])assertr.state==states.SUCCESSdeftest_export_is_abortable():fromcelery.contrib.abortableimportAbortableTaskassertisinstance(export_statement,AbortableTask)
python-mpytest tests/test_states.py-v# 4 passed

4. 项目总结

4.1 优点 & 缺点

维度状态机驱动(Backend 状态 + 撤销体系)日志驱动(翻日志猜状态)
准确性状态是任务生命周期的一等公民靠 grep 与运气
可编程任务中心/告警可直接消费状态只能靠人看
撤销能力revoke/terminate/abort 三档
代价Backend 写放大(track_started)
缺点 1撤销广播有窗口期,非 100% 可靠——

4.2 适用场景

  • 适用:① 任务中心/管理后台展示时间线;② 长任务需人工取消(导出、批处理);③ 状态驱动的业务流转(订单超时关单依赖任务终态);④ 需要「状态回退守卫」的多系统状态同步场景。
  • 不适用:① 高频短任务开 track_started(写放大不划算,走事件流);② 需要「强一致取消」的场景(协作式 abort 有检查点延迟,terminate 粗暴——都不适合金融级一致性要求);③ 状态只在单进程内部消费的简单脚本(直接看返回值更快)。

4.3 注意事项

  • task_track_started只开在需要时间线的任务上,不要全局开
  • revoke 的验收标准是「最终态 REVOKED」,不是「立即变化」。
  • AbortableTask 的检查点要写在任务体循环里,密度与响应时间成反比。
  • terminate 是最后手段:先确认任务幂等(第 11 章)再用,否则半截数据要人工修复。
  • 自定义状态(如 ABORTED)记得注册进任务中心的状态字典,并明确「是否终态」——未登记的未知状态会破坏下游流转判断。

4.4 常见踩坑经验(3 个生产故障)

  1. 故障:运营点「取消导出」,任务继续跑完 8 分钟。根因:点了 revoke,但任务已 STARTED,revoke 不中断执行。对策:导出任务改 AbortableTask,取消走 abort()。教训:撤销要看任务处于哪个生命阶段,一个按钮不通用
  2. 故障:任务中心显示「执行中」与实际严重脱节。根因:全局开了 track_started,Backend 写入积压,状态更新延迟数秒。对策:track_started 按任务开 + Backend 容量评估。教训:可观测性本身也有成本,按需开启
  3. 故障:terminate 一个批量任务后,数据库锁挂了一小时。根因:SIGTERM 强杀,事务连接未正常回收,DB 锁超时才释放。对策:写操作类任务用 AbortableTask 协作取消,禁用 terminate。教训:强杀省 8 分钟,善后花 1 小时

4.5 思考题

  1. revoke 之后,任务为什么可能「先显示 REVOKED 又变回 PENDING」?广播撤销与 Worker 重启的关系是什么?(提示:Mingle 机制,第 33 章)
  2. 一个任务在执行中被 abort(),任务内部update_state(state='ABORTED')与直接return有什么区别?哪个更利于任务中心展示?

答案见第 11 章开头的「上一章思考题参考答案」。第 11 章将进入基础篇的核心章节「失败重试、超时与幂等」,状态机的 FAILURE 与 RETRY 分支将在那里全面展开。

延伸阅读与资源

Java 工程师进阶:从 JVM 生产排障到OpenJDK原理
NumPy 从入门到生产落地:全链路实战指南(科学计算/向量化)
Redis 8 实战精讲:从 CRUD 到源码,构建高可用缓存系统
Redis 实战修炼与原理进阶
Python 3实战精进:从脚本到高并发订单引擎
python入门:Rquests从菜鸟脚本到企业级SDK的网络实战圣经
Milvus向量数据库实战修炼:从 0 到 1精通向量检索与生产落地
MongoDB 实战进阶与内核修炼
后端工程师的 AI 转型第一课:Ollama 与私有化大模型实战
10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战
大型语言模型(LLM) vLLM 高性能推理落地实战
Agent开发之LlamaIndex 实战修炼与源码进阶
大语言模型Transformers 实战修炼与源码剖析

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询