源码基线Celery 5.6.2本仓库celery/目录本章定位从「看 Flower 绿点」到「有 SLA 的可观测性」0. 上一章思考题参考答案思考题 1确认型撤销 「撤销 回执 兜底」三件套① 下发 revoke 后用事件流本章监听该任务是否出现task-revoked事件回执② 未出现回执的在超时后重试 revoke③ 最终兜底——任务「复活」被重新执行时靠任务幂等键第 11 章让重复执行无副作用。教训分布式系统没有「一次命令万无一失」只有「命令 回执 幂等」三层防御。思考题 2autoscale适合常态波动无人值守、自动跟随积压pool_grow/shrink适合可预知的大促提前拉满、结束收回。混用风险自动伸缩与手动伸缩争夺子进程数量控制权——手动 grow 到 8 后autoscale 认为「空闲」又缩回 2人为操作被覆盖反之亦然。实践大促窗口关闭 autoscale或改用纯手动日常开 autoscale控制权要单一归属。1. 项目背景第 14 章搭了 Flower第 24 章学了远程控制但值班同学最想要的还是「出事前被提醒」队列积压超过 1 万条时有人告警Worker 心跳丢失时有人告警短信成功率掉到 90% 时有人告警。现在这些全靠「人肉盯 Flower 截图」——大促 0 点截图截图再截图眼皮都不敢眨结果还是漏了「report 队列在 00:05 开始暴涨」的信号。另一个痛点是「历史回放」昨天 14:00 的那次故障想要「当时所有失败任务、每个任务的耗时分布」——Flower 是实时视图翻不到昨天日志系统里能翻但要一条条拼。结论实时视图Flower不是可观测性可观测性是「指标可存、可查、可告警」。可观测性的三个层次本章从 1 到 3 走完 ① 事件流Celery 每一秒都在发事件task-sent/received/started/succeeded/failed、worker-heartbeat ② 指标把事件流聚合成数字成功率、P99、积压深度、重试率 ③ 告警指标过线就提醒积压 1 万、心跳丢失、成功率 90%本章目标读懂事件总线celery/events/消费事件写入 Prometheus配置两条核心告警——「队列积压 1 万」「Worker 心跳丢失」——从「看绿点」升级到「有 SLA」。2. 项目设计场景大促复盘会小周把「0 点漏报」的截图投上屏幕。小胖Flower 不是能看队列深度吗昨天它不是一直在吗怎么还说「漏报」我昨晚上眼皮打架确实没盯住——但那是我的问题大师不是你的问题是设计的问题——把「可靠性」押在人的眼皮上本身就是设计缺陷。Flower 是「看」不是「盯」它展示实时状态但没有阈值、没有历史、没有主动提醒。真正的盯梢是告警系统数值过线 → 通知人。所以问题的答案是把 Celery 的事件流变成指标再让指标系统替我们盯梢。小白事件流是啥跟第 24 章的 pidbox 是一回事吗我理解任务执行完的状态写在 Backend事件是不是就是 Backend 状态的一份拷贝大师不是拷贝是独立的流。Worker 在任务生命周期各节点主动发送事件celery/events/dispatcher.py的EventDispatchertask-sent投递、task-received领取、task-started、task-succeeded、task-failed、task-retried、task-revokedWorker 自己还发worker-heartbeat心跳。事件不是「写完 Backend 再抄一份」而是平行发出的广播流——所以task_receive之类的事件即使 Backend 不写状态也能观察到第 10 章说过的「事件与状态两个体系」。事件消费端是celery/events/receiver.py的EventReceivercelery/events/state.py的内存状态模型Flower 就是基于它。一句话Backend 是「账本」事件是「流水——账本记结果流水记全过程。技术映射Backend 银行对账单结果事件流 每一笔交易的流水日志全过程Flower 实时看流水屏Prometheus 把流水聚合成「每秒交易笔数、平均处理时长」的仪表盘 超线报警器。小白那指标怎么定义「成功率」「P99 执行时间」这些从事件里怎么算出来Prometheus 怎么对接大师事件流是「原始数据」指标是「聚合结果」中间需要一个转换器消费事件 → 更新计数器/直方图 → 暴露给 Prometheus。两个落地路线① 现成导出器如celery-exporter社区项目订阅事件直接吐出celery_task_succeeded_total等指标省事但定制弱② 自写消费者EventReceiver接收事件用 Prometheus client 库维护指标celery_events_total{state...}、celery_task_duration_seconds直方图、celery_queue_depth从事件流之外用队列命令采样。核心指标五件套本章实战目标积压深度、消费速率、成功率、P99 执行时间、重试率——对应 SLA 的「有没有积压、跑得完吗、成功吗、多快、重了几次」。小胖那「心跳丢失」告警怎么做心跳事件不是一直发吗怎么判断「丢了」大师心跳是worker-heartbeat事件celery/worker/heartbeat.py每隔一段时间默认 2 秒发一次。「心跳丢失」 超过 N 秒没收到某节点的任何事件转换器里维护「各节点最后心跳时间」worker_heartbeat_timestamp - now 阈值如 60s就置指标celery_worker_online 0Prometheus 据此告警。注意阈值比心跳间隔大一个数量级心跳 2s阈值 60s避免网络抖动误报。积压深度告警同理队列深度 10000 就告警——两条告警就是本章的交付物SLA 从这里开始。技术映射心跳 员工的打卡心跳丢失 打卡断了——但「断了 30 秒」可能是打卡机抽风断 5 分钟才是人没了。告警阈值就是「容忍度」的设定。3. 项目实战3.1 环境准备沿用环境Redis Broker Backend。新增依赖pipinstallprometheus-client3.2 分步实现步骤 1确认事件在发——celery events兜底验证目标先确认事件流是通的再谈消费。# 终端 AWorker 开事件发送生产默认配置需显式开启事件开关celery-Aorder_tasks worker--loglevelinfo--events--poolsolo# 终端 B盯事件流celery-Aorder_tasks events--dump运行结果文字描述终端 B 持续滚动事件——投一个短信任务后依次出现task-sent→task-received→task-started→task-succeededWorker 的心跳事件每 2 秒一条。事件源确认连通celery/bin/events.py的 dumper 是现成的 debug 工具。步骤 2写事件消费者聚合核心指标目标把事件流变成 Prometheus 指标自写转换器可控可扩展。# metrics_exporter.pyimporttimefromcollectionsimportdefaultdictfromprometheus_clientimportCounter,Histogram,Gauge,start_http_serverfromcelery.eventsimportEventReceiverfromcelery.events.stateimportStatefromorder_tasksimportapp# —— 指标定义 ——EVENTSCounter(celery_events_total,事件总数,[state])SUCCESSCounter(celery_task_succeeded_total,任务成功数,[task])FAILEDCounter(celery_task_failed_total,任务失败数,[task])DURATIONHistogram(celery_task_duration_seconds,任务执行耗时,[task],buckets(0.1,0.5,1,2,5,10,30))QUEUE_DEPTHGauge(celery_queue_depth,队列积压深度,[queue])LAST_HEARTBEATGauge(celery_worker_last_heartbeat_seconds,节点最后心跳时间戳,[hostname])WORKER_ONLINEGauge(celery_worker_online,节点是否在线,[hostname])stateState()last_heartbeatdefaultdict(float)HEARTBEAT_TIMEOUT60.0defon_event(event):etypeevent[type]EVENTS.labels(etype).inc()ifetypetask-succeeded:SUCCESS.labels(event[uuid]).inc()DURATION.labels(task).observe(event.get(runtime,0))elifetypetask-failed:FAILED.labels(event[uuid]).inc()elifetypeworker-heartbeat:hostevent[hostname]last_heartbeat[host]time.time()LAST_HEARTBEAT.labels(host).set(time.time())defsample_queues():采样队列深度积压指标Broker 层采样事件流没有队列深度。importredis rredis.Redis()forqin(celery,order,sms,report):try:QUEUE_DEPTH.labels(q).set(r.llen(q))exceptException:passdefheartbeat_sweep():心跳丢失扫描超过阈值没心跳的节点置离线。nowtime.time()forhost,tsinlast_heartbeat.items():WORKER_ONLINE.labels(host).set(1if(now-ts)HEARTBEAT_TIMEOUTelse0)if__name____main__:start_http_server(9100)# Prometheus 抓取端口withapp.connection_for_read()asconn:recvEventReceiver(conn,handlers{*:on_event},appapp)# 主循环收事件 周期性采样队列与心跳importthreading threading.Timer(10.0,lambda:(sample_queues(),heartbeat_sweep())).start()recv.capture(limitNone,timeoutNone)步骤 3起导出器 Prometheus 拉取验证目标验证指标真的被 Prometheus 采集到。# 终端 A启动导出器监听 9100python metrics_exporter.py# 终端 B灌一批任务产生事件foriin123;docelery-Aorder_tasks call orders.send_order_sms--args[$i];done# 终端 C本机验证指标无 Prometheus 时直接 curlcurlhttp://localhost:9100/metrics|Select-Stringcelery_运行结果文字描述curl输出包含celery_events_total{statetask-succeeded} 3.0、celery_queue_depth{queuecelery} 0.0、celery_worker_last_heartbeat_seconds{hostnameceleryDESKTOP} 1.7e9等指标——事件流 → 指标 的管道打通。步骤 4Prometheus Grafana 接线目标指标落地到监控平台配两条告警。# prometheus.yml简化scrape_configs:-job_name:celerystatic_configs:-targets:[localhost:9100]# 告警规则两条核心告警groups:-name:celeryrules:-alert:QueueBacklogHighexpr:celery_queue_depth{queuecelery}10000for:5mlabels:{severity:critical}annotations:summary:celery 队列积压超过 1 万当前 {{ $value }}-alert:WorkerHeartbeatLostexpr:celery_worker_online 0for:2mlabels:{severity:critical}annotations:summary:Worker 心跳丢失节点可能假死运行结果文字描述Grafana 面板出现三张图队列积压、成功率、P99celery_queue_depth 10000持续 5 分钟触发QueueBacklogHigh告警kill 一个 Worker 后 60 秒celery_worker_online变 02 分钟后触发WorkerHeartbeatLost——盯梢交给了系统人只管处理。3.3 可能遇到的坑及解决方法坑现象解决Flower/导出器收不到事件Worker 没开事件发送启动加--events或配置task_send_sent_event等事件开关心跳告警频繁误报阈值太接近心跳间隔阈值 ≥ 心跳间隔 × 10心跳 2s → 阈值 60s队列深度指标缺失事件流里没有队列深度队列深度是 Broker 层采样redis LLEN / 管理台 API不是事件事件消费者重启丢事件事件是广播流无持久化指标由「消费时刻」累计断点期间用拉取式补齐结果键/Backend 兜底高并发下事件风暴每秒几千事件消费者处理不过来消费者单独部署 批量聚合只订阅需要的类型3.4 完整代码清单与测试验证清单metrics_exporter.py事件消费 指标prometheus.yml 告警规则。SLA 指标五件套沉淀 Wiki指标事件来源告警示例积压深度Broker 采样 1 万 5 分钟消费速率task-received 计数速率为 0 且积压 0成功率succeeded/failed 计数 90% 10 分钟P99 执行时间runtime 直方图 5s 10 分钟重试率task-retried 计数 30% 10 分钟Worker 心跳worker-heartbeat丢失 60s测试验证# tests/test_events.pyimportjsonfrommetrics_exporterimportEVENTS,SUCCESS,on_eventdeftest_task_succeeded_event_metrics():beforeSUCCESS.labels(t)._value.get()on_event({type:task-succeeded,uuid:t,runtime:0.3})afterSUCCESS.labels(t)._value.get()assertafterbefore1deftest_heartbeat_event_recorded():importtimefrommetrics_exporterimportlast_heartbeat on_event({type:worker-heartbeat,hostname:celeryw1})assertceleryw1inlast_heartbeatdeftest_event_types_registered():fortin(task-sent,task-received,task-started,task-succeeded,task-failed,task-retried,task-revoked,worker-heartbeat):EVENTS.labels(t).inc()# 类型可写assertEVENTS._metricsisnotNonepython-mpytest tests/test_events.py-v# 3 passed4. 项目总结4.1 优点 缺点维度事件流 Prometheus本章纯 Flower 人肉盯主动告警✅ 阈值触发通知❌ 靠人看历史回放✅ 指标有历史曲线❌ 实时视图定制指标✅ 自写转换器❌ 固定视图运维成本多一个消费者 Prometheus零缺点事件流无持久化断点丢指标——4.2 适用场景适用① 生产 SLA 监控成功率/P99/积压② 大促容量预警积压/心跳③ 自定义业务指标按任务/按批次聚合④ 事件驱动的审计与血缘第 23 章配合⑤ 运维大盘与值班告警的「数据底座」。不适用① 需要精确计数不丢一条的指标事件流是尽力而为用 Backend 结果键计数兜底② 单机学习环境Flower 够用③ 需要完整事件历史的场景事件无持久化落库用 snapshotcelery/events/snapshot.py或自建存储。4.3 注意事项事件发送要显式开启Worker--events或对应配置不开 监控盲区。事件流「尽力而为」消费者崩溃期间的事件丢失无法追回——关键指标加 Backend 兜底采样。队列深度是 Broker 层数据RabbitMQ 用管理台 API/rabbitmqctlRedis 用 LLEN不是事件。告警阈值要带「for」持续窗口5m/2m避免瞬时抖动误报心跳阈值 间隔 × 10。事件消费者与 Worker 要「反亲和」部署消费者吃事件流与 Worker 抢 Broker 连接会互相拖累消费端崩溃不影响 Worker 生产事件事件是单向广播。时间线对齐事件里的local_received消费者本机时钟与timestamp发送方时钟跨机器可能不一致——跨主机对比事件时序前先做时钟对齐NTP 校验否则 P99 统计会被时钟差污染。告警阈值校准节奏每次大促后复盘告警的「漏报/误报」按实测数据校准阈值——阈值是活文档不是一次定死的常量。4.4 常见踩坑经验3 个生产故障故障大促 0 点积压 5 万无人报警。根因Worker 没开--events导出器收不到任何事件。对策事件开关写进启动模板 导出器自身加「零事件告警」。教训监控系统的盲区往往在「监控本身没起」。故障心跳告警每 10 分钟响一次值班脱敏。根因阈值 5s ≈ 心跳间隔 2s网络抖动即触发。对策阈值调到 60s。教训告警阈值调不好等于没有告警狼来了效应。故障成功率指标虚高。根因导出器消费的是事件流消费者重启期间失败的「事件」丢了只统计到成功的。对策用 Backend 结果键兜底对账成功率双源对比。教训指标的可信度要用第二数据源校验。4.5 思考题事件流是「尽力而为」的广播。设计「成功率 SLA 告警」时如何避免消费者重启造成「虚高/虚低」提示结果键兜底、双源对账celery events --dump显示事件里有local_received字段它和事件本身的timestamp有什么区别提示时钟源与到达时刻答案见第 26 章开头的「上一章思考题参考答案」。第 26 章 Signals 将用「不侵入任务代码」的方式补齐事件流覆盖不到的场景。### 4.6 推广计划提示延伸阅读与资源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 实战修炼与源码剖析
