
0. 上一章思考题参考答案思考题 1双源对账 自我监控① 事件流指标succeeded/failed 计数与 Backend 结果键计数celery-task-meta-*按状态统计做差量对账——消费者重启期间事件丢失但结果键result_expires 内不丢恢复后按结果键修正事件流计数② 告警带for持续窗口过滤瞬时抖动③ 给导出器本身配「零事件」告警——监控系统挂了自己也必须能报警。三层下来「虚高/虚低」被限制在「消费者宕机窗口」内且可被对账发现。思考题 2timestamp是发送方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信号具体在什么时候发、能拿到什么数据大师先给分工基类适合「任务自身的行为」重试策略、默认队列、bindTrue的上下文信号适合**「与任务无关的横切关注点」**链路 ID、耗时、审计、告警——因为信号不依赖任务继承关系老任务不用改造就能生效。再看数据task_prerun拿到task_id、task任务对象、args/kwargstask_postrun额外拿到retval返回值和statetask_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 的数据库写入任务本身的执行时间就被拖成 500msbefore_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_idWeb 层可预置_ctxthreading.local()defset_trace_id(tid:str):Web 请求入口调用把 trace_id 塞进线程上下文。_ctx.trace_idtidsignals.before_task_publish.connectdefinject_trace_id(sender,headersNone,**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):tidgetattr(task.request,headers,{}).get(trace_id,-)print(f[trace] task{task.name}id{task_id}trace{tid})运行结果文字描述投递任意老任务Worker 日志自动出现[trace] task... traceuuid行同一个 trace_id 贯穿「生产者投递 → Worker 执行」全链路第 5 章手动实现的升级版且零侵入。步骤 2用task_prerun/task_postrun全站统计耗时目标所有任务自动记录耗时写 Prometheus第 25 章指标的直接补充。# observability_signals.py 追加importtimefromprometheus_clientimportHistogram TASK_DURATIONHistogram(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{taskorders.send_order_sms}...——40 个老任务自动全部接入耗时监控任务函数一行未改。步骤 3失败入死信表——信号里只记账落库交给死信任务目标task_failure信号捕获失败 → 投递「死信记录任务」异步落库信号不做重活。# observability_signals.py 追加 dlq_tasks.pyfromdlq_tasksimportrecord_dead_lettersignals.task_failure.connectdefcapture_failure(sender,task_id,task,args,kwargs,exc,einfo,**kw):失败入死信信号里只做『发任务』这一个轻动作。try:record_dead_letter.delay(task_idtask_id,task_nametask.name,args_reprstr(args)[:200],errorstr(exc)[:300])exceptException:# 信号必须异常隔离pass# dlq_tasks.py —— 死信任务真正落库的地方慢但不阻塞任何信号fromceleryimportCelery appCelery(dlq,brokerredis://localhost:6379/0)app.task(nameops.dead_letter,bindTrue)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})returnrecorded运行结果文字描述制造一个失败任务如raise ValueErrorWorker 日志立即出现[DLQ] ... failed: ...由死信任务执行异步落库task_failure信号本身只做了一次 delay 投递毫秒级任务执行耗时几乎不受影响。验证「重试不触发死信」带 autoretry 的任务第一次失败只走task_retry最终失败才出现[DLQ]。步骤 4Worker 生命周期信号——启动/退出通知目标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 为 Nonebefore_task_publish 的 headers 未传显式传 headers{}老调用方或判空兜底3.4 完整代码清单与测试验证清单observability_signals.pytrace/耗时/失败/生命周期四组信号、dlq_tasks.py死信任务。信号速查表沉淀 Wiki信号发送进程触发时机常用场景before_task_publish生产者消息投递前TraceID 注入、投递校验after_task_publish生产者消息投递后投递计数task_prerunWorker执行前耗时起点、线程上下文task_postrunWorker执行后成败都发耗时终点、结果审计task_retryWorker进入重试重试率指标task_failureWorker最终失败死信入表、告警task_revokedWorker被撤销撤销审计worker_ready / shutting_downWorker启动完成/退出前注册中心上下线heartbeat_sentWorker每次心跳心跳增强指标测试验证# tests/test_signals.pyfromunittestimportmockfromobservability_signalsimport_task_start,TASK_DURATIONdeftest_prerun_starts_timer():fromceleryimportsignalswithmock.patch(observability_signals.time.time,return_value100.0):signals.task_prerun.send(senderNone,task_idt1,taskmock.Mock(namex),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(senderNone,task_idt2,taskmock.Mock(nameorders.x),args[],kwargs{},excValueError(bad),einfoNone)m.assert_called_once()deftest_before_publish_injects_trace_id():fromceleryimportsignals headers{}signals.before_task_publish.send(senderorders.x,headersheaders)asserttrace_idinheaderspython-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 实战修炼与源码剖析