0. 上一章思考题参考答案
思考题 1:撤销集合是时序敏感的——「消费前必须知道哪些不能动」(Mingle 在 Tasks 前),因为消费动作会立刻执行任务;而 Worker 状态(谁活着)是最终一致即可——晚几秒知道邻居挂了不影响正确性(Gossip 在 Tasks 后,边跑边传)。顺序设计的原则:先同步「会阻止动作的记忆」,后传播「不影响动作的状态」。
思考题 2:TasksStep(celery/worker/consumer/tasks.py)在声明消费者后,通过task_consumer.qos(prefetch_count=N)设置 QoS——N = 并发 ×worker_prefetch_multiplier(第 18 章公式),在消费者声明时生效,由 kombu 底层转成 AMQP 的 basic.qos 帧(Redis transport 转成可见性窗口语义)。所以「预取数」是 Broker 侧的消费约束,不是 Worker 侧的逻辑。
1. 项目背景
高级篇第三站,也是最硬核的一站:任务执行引擎。第 5 章我们用过self.request、第 11 章用过self.retry()、第 26 章挂过task_prerun信号——但「任务函数被调用前后,内核到底做了什么」一直没有全景。小周的困惑来自一次「慢任务排查」:短信任务日志显示succeeded in 2.003s,但任务函数体里明明只 sleep 了 1 秒——另外 1 秒去哪了?
大师提示他去看celery/app/trace.py和celery/worker/request.py——任务执行不是一个「函数调用」,而是「一条流水线」:
消息 → Request(反序列化 + 上下文)→ strategy(投递策略)→ build_tracer(生成执行器) → trace_task(prerun → run → 成功/失败/重试 → postrun → 写 Backend)那 1 秒的「消失时间」就藏在流水线里:反序列化(消息体 JSON 解析)、Request 构造、tracer 的各环节钩子、Backend 写入——它们都是「任务耗时」的一部分,但任务函数体感知不到。本章目标:读懂这条流水线,并给trace_task增加一个慢任务火焰日志(执行 > 3s 打印当前栈),用真实短信任务验证——既解答「时间去哪了」,又拿到一个生产可用的排查工具。
一句话先记住本章结论:「发任务」在消息层完成,「执行任务」在 trace 层完成——第 34 章是中间那段「从消息到函数体」的全部秘密。
2. 项目设计
场景:小周把「2 秒 vs 1 秒」的差异摆在桌面,大师开始拆流水线。
小胖:2 秒对 1 秒,不就差 1 秒嘛,sleep 多了呗,有啥好查的?你们搞源码的,净纠结这些细节!
小白:小胖你这个「多了呗」就是「拍脑袋排障」——第 15 章批过无数次了。我看了下celery/worker/request.py,有一个Request类,还有celery/worker/strategy.py的default函数。我想先问:消息进来后,Request 是怎么被构造的?strategy 又管什么?
大师:顺序是:消息到达 → Tasks Step 把消息交给 strategy(投递策略)→ strategy 决定「何时执行、怎么执行」→ 交给 Request 构造执行上下文 → 进入 trace。Request(celery/worker/request.py)的职责是把消息反序列化成任务执行所需的上下文:task_id、args/kwargs、headers、retries、delivery_info(队列/交换机)——它就是第 5 章self.request背后的那个对象。strategy(celery/worker/strategy.py)是调度决策层:处理countdown/eta(放 Timer,第 21 章)、rate_limit(限速节奏)、acks_late的确认时机——「什么时候执行」由 strategy 决定,「执行成什么样」由 trace 决定。
技术映射:strategy = 前台接待(决定「这单什么时候进后厨」:延时单放 Timer 冷藏、限速单排队);Request = 点菜单(把顾客的话整理成后厨能懂的规格:菜名、忌口、桌号);trace = 后厨炒菜流水线(备菜→炒→出锅→上账)。
小白:那build_tracer和trace_task呢?我看了celery/app/trace.py——build_tracer(:344)看起来是「生成执行器」的工厂,trace_task(:689)是「入口」。它们的生命周期到底是啥?
大师:build_tracer是工厂:根据任务对象生成一个「带完整钩子的执行函数」——它把信号(第 26 章)、重试逻辑(第 11 章)、Backend 写入(第 8 章)、异常捕获全部编织进一个闭包;trace_task是入口:拿到消息后调用这个 tracer。完整的执行链(源码对照):
trace_task(:689) └─ build_tracer 生成 tracer(:344) ├─ task_prerun 信号(第 26 章) ← 第 5 章 self.request 就绪 ├─ 调用 task.run(*args, **kwargs) ← 任务函数体(1 秒 sleep) ├─ 成功 → task_postrun → 写 Backend(SUCCESS) + 结果 ├─ 失败 → task_failure 信号 → 判定「重试 or 终态」 │ ├─ 可重试 → task.retry() 抛 Retry → 状态 RETRY → 重投 │ └─ 终态 → 写 Backend(FAILURE) + traceback └─ 返回 (state, retval, ...)那 1 秒的「消失时间」= 信号钩子 + 反序列化 + Backend 写入 + 结果序列化的总开销——任务函数体只是流水线的一段,不是全部。
小胖:那异常呢?任务里抛了个ValueError,谁会捕获?捕获之后 ExceptionInfo 又是啥?听着像「异常的档案袋」?
大师:好比喻。trace 的异常处理链:task.run抛异常 → tracer 捕获 → 包装成ExceptionInfo(celery/app/trace.py的 ExceptionInfo 类:持有异常对象 + traceback 文本 + 可序列化)→ 分派:若匹配任务的重试条件(第 11 章 autoretry/手动 retry)→ 走 Retry 分支(任务状态 RETRY,消息重投);否则 →task_failure信号(第 26 章死信入口)→ 写 Backend FAILURE(get()时 propagate 回调用方,第 8 章)。ExceptionInfo 的价值:异常可以被序列化跨进程传递——Worker 进程的异常,能原样出现在调用方get()的堆栈里,靠的就是它。
技术映射:ExceptionInfo = 事故的「档案袋」——现场照片(traceback)+ 当事人(异常对象)装袋归档;袋子能跨部门(进程)寄送,调用方拆袋就能看到事故全貌。
3. 项目实战
3.1 环境准备
沿用环境(Redis Broker + Backend)。本章通过给 trace 增加慢任务火焰日志(源码级实验,可编辑安装),验证执行链路的每一段耗时。
3.2 分步实现
步骤 1:给trace_task增加慢任务火焰日志
目标:任务执行 > 3 秒时,打印当前栈与各阶段耗时。
# slow_trace.py —— 通过信号实现(不改框架源码,更安全的生产方案)# 但本章目标是「读懂 trace」,先做一个源码级插入:importtimefromceleryimportsignals# 方案 A(生产推荐):信号采集各阶段耗时(第 26 章)_stage_t0={}@signals.task_prerun.connectdeft0(sender,task_id,task,**kw):_stage_t0[task_id]=time.time()@signals.task_postrun.connectdeft1(sender,task_id,task,retval,state,**kw):t_start=_stage_t0.pop(task_id,None)ift_startand(time.time()-t_start)>3:print(f"[SLOW]{task.name}{task_id}耗时{time.time()-t_start:.2f}s")importtraceback traceback.print_stack()# 打印调用栈:慢在哪一段celery-Aorder_tasks worker--loglevel=info--pool=solo-Qsms# 投递一个 4 秒的慢任务(sleep 4)celery-Aorder_tasks call orders.slow_task--args='[4]'--queue=sms运行结果(文字描述):任务完成后打印[SLOW] orders.slow_task ... 耗时 4.02s+ 调用栈——慢任务被自动标记(栈里能看到执行路径:trace_task → tracer → task.run → sleep)。
步骤 2:拆解「任务耗时」的各阶段(回答 2s vs 1s)
目标:量化流水线各段开销,看清「消失的时间」。
# stage_timing.py —— 各阶段计时(信号方案)importtimefromceleryimportsignals _ctx={}@signals.task_prerun.connectdefpre(sender,task_id,task,**kw):_ctx[task_id]={"t0":time.time()}@signals.task_postrun.connectdefpost(sender,task_id,task,retval,state,**kw):c=_ctx.pop(task_id,{})c["postrun"]=time.time()# 这里只能看到 prerun→postrun 的框架开销;函数体内耗时由任务自己统计print(f"[STAGE]{task.name}: 框架开销 ≈{c['postrun']-c['t0']:.3f}s")运行结果(文字描述):短任务(函数体 1s)的框架开销约 5~30ms(信号+Backend 写入),长任务 2s vs 1s 的差异主要在「排队等待 + 预取 + 反序列化」,不在流水线钩子——结合inspect scheduled/active(第 15 章)就能定位「时间去哪了」。
步骤 3:验证异常处理链——Retry 与 ExceptionInfo
目标:读trace_task的分支,验证重试与终态的判定路径。
# retry_flow_demo.pyfromorder_tasksimportapp@app.task(name='orders.demo_retry',bind=True,max_retries=2,autoretry_for=(ValueError,))defdemo_retry(self,flag:str)->str:ifflag=="bad":raiseValueError("业务错误(应被自动重试)")return"ok"celery-Aorder_tasks worker--loglevel=debug--pool=solo-Qorder celery-Aorder_tasks call orders.demo_retry--args='["bad"]'# 观察 Worker 日志与结果celery-Aorder_tasks result<task_id>运行结果(文字描述,节选):
[DEBUG] Task orders.demo_retry[..] retry: Retry in 0s # 异常 → Retry 分支 [DEBUG] Task orders.demo_retry[..] succeeded in ... # 重试后成功 # 若把 max_retries=0 重跑: [ERROR] Task orders.demo_retry[..] raised ValueError # 终态分支 [INFO] Task orders.demo_retry[..] FAILURE # 写 Backend FAILURE对照 trace 源码:
retry: Retry in Ns对应build_tracer里捕获Retry异常的路径;raised ... FAILURE对应终态写 Backend 路径——源码里的每个分支,都在日志里有对应痕迹。
步骤 4:慢任务火焰日志的生产化(收敛为工具)
目标:把实验代码收敛成「生产可用的慢任务告警器」。
# slow_task_monitor.py —— 生产版(信号方案,第 26 章规范:只记账不打栈)importtimefromprometheus_clientimportCounterfromceleryimportsignals SLOW=Counter('celery_slow_task_total','慢任务计数(>3s)',['task'])_ctx={}@signals.task_prerun.connectdef_t0(sender,task_id,task,**kw):_ctx[task_id]=time.time()@signals.task_postrun.connectdef_t1(sender,task_id,task,retval,state,**kw):t0=_ctx.pop(task_id,None)ift0and(time.time()-t0)>3:SLOW.labels(task.name).inc()# 只计数,慢任务定位交给日志/追踪print(f"[SLOW]{task.name}耗时{time.time()-t0:.2f}s(task_id={task_id})")运行结果(文字描述):与第 25 章 Prometheus 体系对接——celery_slow_task_total指标进大盘,慢任务数突增 = 告警;「慢任务火焰日志」从实验工具升级为生产监控组件(排查时再用 traceback.print_stack 的调试版)。
步骤 5:子任务上下文栈——任务里发任务的「栈」
目标:理解「任务 A 里调用任务 B」时,trace 的上下文如何嵌套(第 3 章「子任务」的源码层)。
# context_demo.pyfromorder_tasksimportapp@app.task(name='orders.parent',bind=True)defparent(self,child_count:int)->str:print(f"[parent] 当前栈深:{len(self.request.children)ifhasattr(self.request,'children')else'n/a'}")foriinrange(child_count):child.delay(i)# 任务里发任务(第 3/6 章)return"parent-done"@app.task(name='orders.child',bind=True)defchild(self,idx:int)->str:returnf"child-{idx}"celery-Aorder_tasks worker--loglevel=info--pool=solo-Qorder celery-Aorder_tasks call orders.parent--args='[3]'运行结果(文字描述):parent 执行后,child任务被投递 3 次;self.request里能通过children等字段看到子任务关联——trace 的执行上下文是「栈式」的:每个任务执行时压入当前 Request,执行完弹出(第 5 章self.request线程局部的源码机理)。生产里「子任务结果汇总回父任务」就是靠这层栈 + 第 23 章结果树实现的。
3.3 可能遇到的坑及解决方法
| 坑 | 现象 | 解决 |
|---|---|---|
| 慢任务不触发日志 | 阈值 > 实际耗时 | 阈值按 P99 校准(第 30 章);先观察再定值 |
| 信号里 print 刷屏 | 每个任务都打 | 只打超阈值;或用指标计数(步骤 4) |
| Retry 分支看不到 | 日志级别不够 | --loglevel=debug看 retry/raised 痕迹 |
| ExceptionInfo 序列化失败 | 异常对象不可 pickle | 业务异常继承标准 Exception;traceback 文本兜底 |
| 框架开销误判为慢 | prerun→postrun 差 30ms 以为是框架问题 | 区分「队列等待 + 预取」与「执行流水线」两段 |
3.4 完整代码清单与测试验证
清单:slow_task_monitor.py(生产版)、stage_timing.py(阶段计时)、retry_flow_demo.py(异常链验证)。trace 生命周期速查(沉淀 Wiki,对照celery/app/trace.py):
| 阶段 | 源码位置 | 触发点 |
|---|---|---|
| 反序列化 | Request(worker/request.py) | 消息到达 |
| 投递决策 | strategy(worker/strategy.py) | ETA/限速/ack 时机 |
| 执行器生成 | build_tracer(trace.py:344) | 首次执行任务 |
| 执行 | trace_task(trace.py:689) | prerun→run→postrun |
| 重试分支 | Retry 异常捕获 | 匹配重试条件 |
| 终态写回 | Backend store_result | SUCCESS/FAILURE |
测试验证:
# tests/test_trace.pyfromslow_task_monitorimport_ctxfromorder_tasksimportapp,demo_retry app.conf.task_always_eager=Truedeftest_retry_task_defined():fromorder_tasksimportdemo_retryassertdemo_retry.max_retries==2assertValueErrorindemo_retry.autoretry_fordeftest_exception_goes_failure_when_no_retry():@app.task(name='orders.no_retry')defno_retry():raiseValueError("boom")r=no_retry.apply()assertr.failed()deftest_slow_monitor_ctx_clean_after_postrun():# 正常执行后 _ctx 应清空(无泄漏)fromorder_tasksimportsend_order_sms send_order_sms.apply(args=[1])assertlen(_ctx)==0python-mpytest tests/test_trace.py-v# 3 passed4. 项目总结
4.1 优点 & 缺点
| 维度 | 信号方案(不侵入) | 改 trace 源码方案 |
|---|---|---|
| 安全性 | 零侵入,生产可用 | 改框架,需还原 |
| 覆盖 | 钩子级(信号点) | 全链路(任意插入) |
| 维护 | 独立文件 | 与版本耦合 |
| 调试深度 | 看得到钩子 | 看得到全部 |
| 推荐 | 生产监控 | 学习/实验 |
4.2 适用场景
- 适用:① 慢任务定位(火焰日志/告警);② 理解「任务耗时」构成(排障、容量计算第 30 章);③ 异常链路的可观测(重试/终态判定);④ 生产级慢任务监控(指标 + 告警);⑤ 子任务上下文与结果树的源码理解(第 23 章血缘的延伸)。
- 不适用:① 业务任务逻辑(trace 是框架层,业务别碰);② 需要「修改异常处理语义」的场景(改 trace 风险大,优先用信号或自定义 Task 基类)。
4.3 注意事项
- 任务函数体 ≠ 任务耗时:排队、预取、反序列化、信号、Backend 写入都在耗时里——容量公式(第 30 章)用「端到端」而非函数体时长。
- 改源码调试后必须还原(git diff 检查),生产不带实验代码。
- ExceptionInfo 的序列化依赖异常对象可 pickle:自定义异常要兼容。
- 慢任务阈值按 P99 校准,别拍脑袋定 3s。
- trace 是「只读理解」的禁区:业务侧通过信号(第 26 章)介入,别在任务代码里改 trace 行为。
4.4 常见踩坑经验(3 个生产故障)
- 故障:报表任务「执行 1 秒、总耗时 10 分钟」。根因:任务在队列里排队 9 分 59 秒(预取被长任务占满,第 18 章)。对策:队列隔离 + 预取调 1。教训:耗时排查先分「排队段」与「执行段」(第 15 章三板斧)。
- 故障:慢任务告警天天响,值班脱敏。根因:阈值 3s 低于 P99。对策:用历史数据定 P99,阈值 = P99 × 1.5。教训:告警阈值脱离数据 = 狼来了。
- 故障:get() 抛出的异常和 Worker 里的对不上。根因:自定义异常类不可序列化,ExceptionInfo 降级为通用错误。对策:异常类保持简单可 pickle。教训:跨进程的异常也是「数据」,要可序列化。
- 故障:任务里发子任务,子任务全堆在默认队列没人消费。根因:父任务在 order 队列,子任务没指定队列(第 6 章调用选项没传)。对策:子任务显式带 queue 或路由按任务名覆盖。教训:「任务里发任务」的队列归属要显式声明,别指望继承父队列。
4.5 思考题
build_tracer是「按任务生成执行器」的工厂——同一个任务被两个 Worker 同时执行,tracer 会生成几次?(提示:缓存与线程局部)trace_task里「写 Backend」失败(Redis 挂了)会发生什么?任务会判失败吗?还是结果丢了任务照样成功?
答案见第 35 章开头的「上一章思考题参考答案」。执行引擎看完了,下一站是「消息在网上长什么样」——第 35 章报文与序列化。
延伸阅读与资源
Dify 从入门到进阶:LLM 应用平台实战修炼
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 实战修炼与源码剖析