0. 上一章思考题参考答案
思考题 1:双源对账 + 自我监控:① 事件流指标(succeeded/failed 计数)与 Backend 结果键计数(celery-task-meta-*按状态统计)做差量对账——消费者重启期间事件丢失,但结果键(result_expires 内)不丢,恢复后按结果键修正事件流计数;② 告警带for持续窗口过滤瞬时抖动;③ 给导出器本身配「零事件」告警——监控系统挂了自己,也必须能报警。三层下来,「虚高/虚低」被限制在「消费者宕机窗口」内且可被对账发现。
思考题 2:timestamp是发送方(Worker/生产者)时钟打的事件时间;local_received是接收方(消费者)本机时钟记录的到达时刻。跨机器时钟不一致时两者有偏差;计算端到端延迟(received - sent)前必须先做时钟对齐(NTP),否则会算出负延迟或虚高延迟——第 25 章注意事项已提醒,这里确认了它的来源。
1. 项目背景
第 5 章的OrderTask基类解决了「新任务统一注入 trace_id」的问题,但老代码不买账:仓库里有 40 多个历史任务,分散在 6 个模块,有的没走基类、有的自己写日志、有的失败连个痕迹都没有。leader 的要求很明确:「trace_id 全站覆盖、执行耗时全站统计、失败任务全站入死信表——不许改任务函数。」
小周数了数:给 40 个任务逐个加装饰器、改基类、改日志……两星期起步,还容易漏。大师提示他看一个东西:celery/signals.py——Celery 在任务生命周期的每个关键节点都会「喊一嗓子」,任何人可以注册监听,在不改任务代码的前提下插入逻辑:
任务执行时发生的事(信号发送点) before_task_publish ── 任务消息即将投递(生产者侧) after_task_publish ── 任务消息已投递 task_prerun ── Worker 开始执行任务前 task_postrun ── 任务执行完毕(成功/失败都触发) task_retry ── 任务进入重试 task_failure ── 任务最终失败 task_revoked ── 任务被撤销 worker_ready ── Worker 启动完成 worker_shutting_down ── Worker 即将退出 heartbeat_sent ── 心跳已发送本章目标:用信号把「trace_id 注入、耗时统计、失败入死信表」三件事做成横切逻辑——任务函数保持干净,新老任务一视同仁。
2. 项目设计
场景:小周把「改 40 个任务」的方案和信号方案摆上桌。
小胖:这不就是「观察者模式」嘛!我当年面试背过!但我不懂,这玩意和我在任务里自己写try/except有啥区别?最后不都是「执行前打日志、执行后打日志」吗?
小白:区别在归属:自己写 = 每个任务自带一份(改一处漏两处);信号 = 全站一份(改一处全站生效)。但我有个疑问:信号和自定义 Task 基类(第 5 章)都能做横切逻辑,它们的分工是什么?还有,task_prerun和task_postrun信号具体在什么时候发、能拿到什么数据?
大师:先给分工:基类适合「任务自身的行为」(重试策略、默认队列、bind=True的上下文),信号适合**「与任务无关的横切关注点」**(链路 ID、耗时、审计、告警)——因为信号不依赖任务继承关系,老任务不用改造就能生效。再看数据:task_prerun拿到task_id、task(任务对象)、args/kwargs;task_postrun额外拿到retval(返回值)和state;task_failure拿到exc(异常对象)、einfo(异常信息包装)、traceback。注意 task_postrun 无论成功失败都会发(state 区分),而 task_failure 只在最终失败时发(重试中的失败走 task_retry,不触发 task_failure——这是新手最容易踩的误判点)。
技术映射:信号 = 店里的广播喇叭——任务执行像「客人点菜」,喇叭喊「上菜了」「菜上了」「这道菜退了」;后厨(横切逻辑)只听喇叭做事,不用每道菜(任务)都贴一张流程卡。
小白:那before_task_publish呢?它和task_prerun有什么不同?我搜到的例子有人用前者做「投递前校验」,有人用后者做「执行前埋点」,晕了。
大师:发送点不同,所属进程不同:before_task_publish/after_task_publish在生产者进程发(发消息的那一刻);task_prerun/task_postrun在Worker 进程发(执行的那一刻)。所以:链路 ID 注入放before_task_publish(消息还没发,可以改 headers——这是 TraceID 传递的关键时机);耗时统计放task_prerun/task_postrun(在 Worker 侧配对计时)。生产者的信号能拿到 sender 和消息体(body/headers),这是「投递前拦截校验」(比如禁止在非生产队列投递)的唯一时机。
小胖:那我在信号里干点「重活」行不行?比如失败入死信表——往数据库插一行,这不就是一次普通写库吗,能有多重?
大师:这正是信号的性能红线。信号是同步执行的:task_prerun在 Worker 的执行线程里跑,task_failure同样——信号里做一次 500ms 的数据库写入,任务本身的执行时间就被拖成 +500ms;before_task_publish在生产者进程同步跑,信号慢 = 发任务慢 = 下单接口变慢(第 16 章 P99 直接破防)。所以信号里的重活三原则:① 尽量只做「记账」级操作(内存计数器、简单 Redis INCR);② 必须落库的动作放进信号发送的「另一个任务」(死信任务走队列,异步落库);③ 信号处理器本身要做异常隔离——信号里抛异常会向上传播,把任务执行/消息投递打断(用 try/except 包住,或raise后由框架捕获——推荐前者)。
技术映射:信号 = 走廊里的「值班记录本」——记一笔很快;但要是每记一笔都要打电话汇报总部(重活),整个走廊的人都得等你打完。
3. 项目实战
3.1 环境准备
沿用环境(Redis Broker + Backend)。本章不修改任何任务函数——全部通过信号接入。
3.2 分步实现
步骤 1:用before_task_publish统一注入 TraceID(生产者侧)
目标:所有任务(新老一律)在投递前自动带 trace_id,无需改任务代码。
# observability_signals.pyimportuuid,threadingfromceleryimportsignals# 线程局部:记录「当前请求上下文」的 trace_id(Web 层可预置)_ctx=threading.local()defset_trace_id(tid:str):"""Web 请求入口调用:把 trace_id 塞进线程上下文。"""_ctx.trace_id=tid@signals.before_task_publish.connectdefinject_trace_id(sender,headers=None,**kwargs):"""投递前:headers 里补 trace_id(链路 ID 的传递时机)。"""headers['trace_id']=getattr(_ctx,'trace_id',None)orstr(uuid.uuid4())# 任务侧读取(任务函数可通过 self.request.headers 拿到,第 5 章做法)# 但为了「老任务不改造」也能记录,在 task_prerun 统一打日志:@signals.task_prerun.connectdeflog_trace(sender,task_id,task,args,kwargs,**kw):tid=getattr(task.request,'headers',{}).get('trace_id','-')print(f"[trace] task={task.name}id={task_id}trace={tid}")运行结果(文字描述):投递任意老任务,Worker 日志自动出现[trace] task=... trace=<uuid>行;同一个 trace_id 贯穿「生产者投递 → Worker 执行」全链路(第 5 章手动实现的升级版,且零侵入)。
步骤 2:用task_prerun/task_postrun全站统计耗时
目标:所有任务自动记录耗时,写 Prometheus(第 25 章指标的直接补充)。
# observability_signals.py 追加importtimefromprometheus_clientimportHistogram TASK_DURATION=Histogram('celery_signal_task_duration_seconds','任务执行耗时(信号采集)',['task'])_task_start={}@signals.task_prerun.connectdefstart_timer(sender,task_id,task,**kw):_task_start[task_id]=time.time()@signals.task_postrun.connectdefstop_timer(sender,task_id,task,retval,state,**kw):t0=_task_start.pop(task_id,None)ift0:TASK_DURATION.labels(task.name).observe(time.time()-t0)运行结果(文字描述):curl localhost:9100/metrics出现celery_signal_task_duration_seconds_bucket{task="orders.send_order_sms"}...——40 个老任务自动全部接入耗时监控,任务函数一行未改。
步骤 3:失败入死信表——信号里只记账,落库交给死信任务
目标:task_failure信号捕获失败 → 投递「死信记录任务」异步落库(信号不做重活)。
# observability_signals.py 追加 + dlq_tasks.pyfromdlq_tasksimportrecord_dead_letter@signals.task_failure.connectdefcapture_failure(sender,task_id,task,args,kwargs,exc,einfo,**kw):"""失败入死信:信号里只做『发任务』这一个轻动作。"""try:record_dead_letter.delay(task_id=task_id,task_name=task.name,args_repr=str(args)[:200],error=str(exc)[:300])exceptException:# 信号必须异常隔离pass# dlq_tasks.py —— 死信任务:真正落库的地方(慢,但不阻塞任何信号)fromceleryimportCelery app=Celery('dlq',broker='redis://localhost:6379/0')@app.task(name='ops.dead_letter',bind=True)defrecord_dead_letter(self,task_id,task_name,args_repr,error):# 生产:INSERT INTO dead_letter(task_id, task_name, args, error, created_at)print(f"[DLQ]{task_name}{task_id}failed:{error}")return"recorded"运行结果(文字描述):制造一个失败任务(如raise ValueError),Worker 日志立即出现[DLQ] ... failed: ...(由死信任务执行,异步落库);task_failure信号本身只做了一次 delay 投递(毫秒级),任务执行耗时几乎不受影响。验证「重试不触发死信」:带 autoretry 的任务第一次失败只走task_retry,最终失败才出现[DLQ]。
步骤 4:Worker 生命周期信号——启动/退出通知
目标:worker_ready与worker_shutting_down用于「上线注册/下线摘流」。
# observability_signals.py 追加@signals.worker_ready.connectdefon_ready(sender,**kw):print(f"[worker]{sender.hostname}ready, 注册到注册中心")@signals.worker_shutting_down.connectdefon_shutdown(sender,**kw):print(f"[worker]{sender.hostname}shutting down, 摘流完成")运行结果(文字描述):Worker 启动完成打印 ready、SIGTERM 优雅退出时打印 shutting down——注册中心的上下线钩子(第 29 章优雅退出的前置知识)。
3.3 可能遇到的坑及解决方法
| 坑 | 现象 | 解决 |
|---|---|---|
| 信号里做重活 | 任务执行时间暴涨/接口变慢 | 只记账;落库走死信任务(步骤 3) |
| 信号里抛异常 | 任务执行被打断 | 信号处理器 try/except 隔离 |
| task_failure 不触发 | 任务在重试中 | 重试失败走 task_retry,最终失败才 task_failure |
| 信号重复注册 | 模块被 import 两次 | 信号模块只在一个入口 import(如 celeryconfig 或 app 模块) |
| headers 为 None | before_task_publish 的 headers 未传 | 显式传 headers={}(老调用方)或判空兜底 |
3.4 完整代码清单与测试验证
清单:observability_signals.py(trace/耗时/失败/生命周期四组信号)、dlq_tasks.py(死信任务)。信号速查表(沉淀 Wiki):
| 信号 | 发送进程 | 触发时机 | 常用场景 |
|---|---|---|---|
| before_task_publish | 生产者 | 消息投递前 | TraceID 注入、投递校验 |
| after_task_publish | 生产者 | 消息投递后 | 投递计数 |
| task_prerun | Worker | 执行前 | 耗时起点、线程上下文 |
| task_postrun | Worker | 执行后(成败都发) | 耗时终点、结果审计 |
| task_retry | Worker | 进入重试 | 重试率指标 |
| task_failure | Worker | 最终失败 | 死信入表、告警 |
| task_revoked | Worker | 被撤销 | 撤销审计 |
| worker_ready / shutting_down | Worker | 启动完成/退出前 | 注册中心上下线 |
| heartbeat_sent | Worker | 每次心跳 | 心跳增强指标 |
测试验证:
# tests/test_signals.pyfromunittestimportmockfromobservability_signalsimport_task_start,TASK_DURATIONdeftest_prerun_starts_timer():fromceleryimportsignalswithmock.patch('observability_signals.time.time',return_value=100.0):signals.task_prerun.send(sender=None,task_id='t1',task=mock.Mock(name='x'),args=[],kwargs={})assert_task_start.get('t1')==100.0deftest_failure_capture_publishes_dlq():fromceleryimportsignalsfromdlq_tasksimportrecord_dead_letterwithmock.patch.object(record_dead_letter,'delay')asm:signals.task_failure.send(sender=None,task_id='t2',task=mock.Mock(name='orders.x'),args=[],kwargs={},exc=ValueError('bad'),einfo=None)m.assert_called_once()deftest_before_publish_injects_trace_id():fromceleryimportsignals headers={}signals.before_task_publish.send(sender='orders.x',headers=headers)assert'trace_id'inheaderspython-mpytest tests/test_signals.py-v# 3 passed4. 项目总结
4.1 优点 & 缺点
| 维度 | 信号机制(横切) | 每个任务手写 |
|---|---|---|
| 覆盖面 | 新老任务一律生效 | 漏改一个就漏一处 |
| 侵入性 | 零(不改任务代码) | 每个任务都要动 |
| 维护 | 一处修改全站生效 | 处处改处处漏 |
| 风险 | 信号异常会连坐(需隔离) | 各任务独立 |
| 性能 | 同步执行需克制 | 同样有开销 |
4.2 适用场景
- 适用:① 全站链路 ID/耗时统计/失败治理(老代码多、改造难的场景尤其适合);② 与任务无关的横切关注点(审计、告警、注册中心);③ 生产者侧拦截校验(禁止误投队列);④ 生命周期钩子(上下线、心跳增强)。
- 不适用:① 任务自身的行为差异(重试策略、队列归属——用基类/装饰器);② 需要同步返回值的横切逻辑(信号无法「拦截并改写」任务参数,只能观测);③ 重活逻辑(必须异步化后才适合信号)。
4.3 注意事项
- 信号是同步的:处理器耗时直接叠加到任务/投递耗时上,重活一律异步化。
- 信号处理器异常隔离:包 try/except,别让横切逻辑打死业务。
- 信号与基类分工:任务行为用基类,横切关注用信号——别在信号里做「只有部分任务该做」的事。
- 信号模块注册一次:集中在 app 模块或 celeryconfig 里 import,防重复注册(重复注册 = 重复执行副作用)。
4.4 常见踩坑经验(3 个生产故障)
- 故障:全站任务执行时间突然 +800ms。根因:task_failure 信号里同步写 MySQL,赶上锁等待。对策:改投递死信任务异步落库(步骤 3)。教训:信号是「记账本」,不是「搬运工」。
- 故障:某个任务执行一半断了,信号把异常吞了没暴露。根因:信号处理器抛异常且没隔离,打断了任务执行。对策:处理器全部 try/except + 错误日志。教训:横切逻辑不能成为新的故障源。
- 故障:trace_id 一半任务有、一半没有。根因:before_task_publish 只在部分入口注册(模块 import 不全)。对策:信号注册集中化 + CI 断言注册数量。教训:零侵入的能力,也要有「注册即全站」的纪律。
4.5 思考题
task_postrun无论成败都会发送,而task_failure只在最终失败时发送——「重试中的失败」这两个信号各会触发几次?耗时统计放在 task_postrun 会不会把重试次数算进去?before_task_publish信号里能不能直接「拦截」消息不投递?(提示:信号能否改变消息内容/中止发布,看 signal 的返回值约定)
答案见第 27 章开头的「上一章思考题参考答案」。
延伸阅读与资源
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 实战修炼与源码剖析