第29章:Celery 容器化、K8s 与 Helm 部署

第29章:Celery 容器化、K8s 与 Helm 部署 0. 上一章思考题参考答案思考题 1celery_worker默认会话级session scope——因为每个 Worker 是真实子进程函数级启动 每个用例都冷启动一个进程秒级成本 × 用例数CI 直接慢成蜗牛。会话级「起一次、跑全部」代价是「用例间共享 Worker 状态」需任务隔离独立 task_id 断言独立。权衡稳定性优先选函数级性能优先选会话级——生产 CI 默认会话级 污染防护。思考题 2eager 模式下「双投断言」测的是业务逻辑层的幂等去重表真的只生效一次真实 Worker 下测的是消息重投 幂等的联合语义acks_late 崩溃重投后去重表依然兜住。前者保「函数正确」后者保「分布式正确」——契约层写断言、集成层写场景两层缺一不可。1. 项目背景第 16 章的订单平台要上生产了K8s 平台的同事直接把 Web 的部署模板抄给了 Worker「Deployment、replicaCount 4、滚动发布走起。」上线第一周就出事滚动发布期间支付回调任务丢了——旧 Pod 被 SIGTERM 杀掉时Worker 进程直接退出正在执行的任务早确认连同消息一起没了新 Pod 刚起来还没连上 Broker中间窗口的投递也丢了。运维对照 Web 的认知排查了半天Web 无状态重启无所谓Worker 是有状态的消费进程重启窗口就是丢任务窗口。再看他们的编排又发现两个隐患Beat 也被 Deployment 双副本部署第 22 章的「双发」教科书现场一个容器里 worker、beat、flower 全开改配置要全部重启扩容只能整体扩小任务也被拉大。错误抄作业的三个症状 ① 滚动发布即丢任务无优雅退出SIGTERM 直接杀 ② Beat 双副本双发第 22 章教训复发 ③ 一容器多角色无法独立伸缩与发布本章目标用「分角色镜像 优雅退出 正确的副本策略」重写部署——3 组 Worker按队列 单 Beat Flower做一次滚动发布验证发布期间零丢支付回调。2. 项目设计场景支付回调丢失事故复盘三人对着 K8s 的 Deployment 模板。小胖Web 和 Worker 不都是「容器里跑个 Python 进程」吗咋就一个有状态一个没状态了我觉着抄模板没啥问题就是运气差了点小白区别在「进程自己记不记状态」Web 无状态——挂了重启请求重发就行Worker 有状态——它手里攥着「正在执行的任务」进程内存里的 Request 上下文 未确认的消息SIGTERM 直接杀 任务和消息一起蒸发第 18 章早确认的代价。所以第一个问题优雅退出到底是怎样的SIGTERM之后 Worker 会做什么大师看celery/apps/worker.py的退出逻辑收到SIGTERM后Worker 进入**「温暖关闭warm shutdown」① 停止接收新任务consumer.cancel② 让在途任务执行完**受worker_shutdown_timeout限制③ 全部收尾后进程退出。所以「优雅」「不接了 跑完 再走」。关键配套K8s 的terminationGracePeriodSeconds必须 worker_shutdown_timeout否则 K8s 到点发 SIGKILL优雅白优雅且acks_late 的队列在优雅退出时消息未确认 → Broker 重投第 18 章这是「发布期间不丢」的最后一层保险。技术映射优雅退出 店铺打烊——门口挂「停止接单」cancel让正在吃的客人吃完在途任务再关灯锁门退出K8s 的 gracePeriod 消防通道的「宽限时间」——太短警察SIGKILL就来砸门了。小白那「分角色镜像」呢为什么 worker、beat、flower 要拆开我理解 Beat 单实例第 22 章但 flower 不就是个监控吗不能放一起大师拆开的理由是生命周期不同、伸缩维度不同、故障域不同Worker按队列组扩容第 9 章sms 组 4 副本、report 组 1 副本Beat是有状态调度器第 22 章数据库 Scheduler 主备单副本 故障重启Flower是旁路监控挂了不影响业务但单独部署可以独立更新、独立监控。一个容器全开 扩 Worker 顺带扩 Beat双发风险、Flower 崩溃连累 Worker故障域耦合。正确姿势一个容器一个角色、一个角色一个 Deployment/StatefulSetBeat 也可用 Deployment 单副本 外部锁。小胖那滚动发布到底怎么搞才不丢任务我看 K8s 默认就是「先起新的、再杀旧的」不都这样吗大师滚动发布本身没错错在没给「新旧交接」留窗口。不丢任务的发布三件套① 优雅退出上文② 发布窗口内双跑——maxUnavailable: 0先起新、等 Ready、再杀旧新旧 Worker 同时消费消息被谁消费都行任务无状态第 18 章幂等兜底③ 发布时机——低峰 关键队列先放行。验证方法发布期间投递一批「支付回调」任务断言零丢失零重复Backend 状态 去重表对账第 28 章集成用例的部署版。本仓库helm-chart/的deployment.yaml是官方参考——注意它replicaCount与探针liveness/readiness 用inspect().stats()判活的写法我们在此基础上拆角色。技术映射发布 换班——新班组先来、与旧班组同堂双跑、旧班组把手里活干完再走优雅退出「先杀旧再起新」 旧班组下班了才通知新班组上班中间半小时没人接电话。3. 项目实战3.1 环境准备# 需要Docker kind/minikube 或真实集群 helmkind create cluster--namecelery-demo# 参考本仓库 docker/Dockerfile、helm-chart/官方部署资源3.2 分步实现步骤 1分角色镜像——一个容器一个角色目标worker / beat / flower 三个镜像入口互不耦合。# Dockerfile.worker参考 docker/Dockerfile 精简 FROM python:3.11-slim WORKDIR /app COPY requirements.txt . RUN pip install -r requirements.txt celery redis COPY . . # 入口Worker队列由启动参数决定 CMD [celery, -A, order_tasks, worker, -Q, order,sms,report, -c, 4, --loglevelinfo]# Dockerfile.beat FROM python:3.11-slim WORKDIR /app COPY requirements.txt . RUN pip install -r requirements.txt celery redis COPY . . # 入口Beat数据库 Scheduler 主备第 22 章 CMD [celery, -A, order_tasks, beat, -S, db_scheduler.DatabaseScheduler, --loglevelinfo]# Dockerfile.flower FROM mher/flower:2.0 CMD [flower, --brokerredis://redis:6379/0, --port5555]运行结果文字描述三个镜像各自构建docker build -f Dockerfile.worker -t order-worker:1.0 .等任一角色崩溃/升级不影响其他角色故障域隔离。步骤 2Helm 编排——3 组 Worker 单 Beat Flower目标按第 9 章队列规划拆 3 组 WorkerBeat 单副本Flower 旁路。# values.yaml基于本仓库 helm-chart/values.yaml 改造workerGroups:sms:# 短信组短任务水平扩展replicaCount:4command:[celery,-A,order_tasks,worker,-Q,sms,-c,4]resources:{requests:{cpu:500m,memory:512Mi}}order:# 订单组核心链路replicaCount:2command:[celery,-A,order_tasks,worker,-Q,order,-c,8]report:# 报表组重任务小副本replicaCount:1command:[celery,-A,order_tasks,worker,-Q,report,-c,1]beat:enabled:truereplicaCount:1# ★ 单副本第 22 章主备语义scheduler:db_scheduler.DatabaseSchedulerflower:enabled:truereplicaCount:1graceful:worker_shutdown_timeout:60# Worker 在途任务宽限秒terminationGracePeriodSeconds:90# ★ 必须 上面# deployment.yaml 模板节选基于 helm-chart/templates/deployment.yamlapiVersion:apps/v1kind:Deploymentmetadata:{name:{{.Release.Name}}-{{$group}}}spec:replicas:{{$wg.replicaCount}}strategy:type:RollingUpdaterollingUpdate:maxUnavailable:0# ★ 先起新、等 Ready、再杀旧发布不丢maxSurge:1template:spec:terminationGracePeriodSeconds:{{$.Values.graceful.terminationGracePeriodSeconds}}containers:-name:workerimage:{{$.Values.image.repository}}:{{$.Values.image.tag}}command:{{$wg.command}}env:-{name:CELERY_BROKER_URL,valueFrom:{configMapKeyRef:{name:celery,key:CELERY_BROKER_URL}}}livenessProbe:probe# 官方探针inspect().stats() 判活exec:command:[python,-c,from celery.task.control import inspect; inspect().stats()]运行结果文字描述helm install order order-chart后kubectl get pods看到sms-worker-xxx4 副本、order-worker-xxx2 副本、report-worker-xxx1 副本、beat-xxx1 副本、flower-xxx1 副本——一个角色一组 Pod扩容互不影响。步骤 3滚动发布演练——发布期间零丢支付回调目标用第 28 章的「发布用例」验证不丢消息。# 1) 发布前置任务后台持续投递 500 条支付回调任务模拟真实流量python flood_pay_callbacks.py# 2) 触发滚动发布镜像升级helm upgrade order order-chart--setimage.tag1.1# 3) 发布完成后对账500 条全部 SUCCESS、去重表无重复python verify_pay_callbacks.py运行结果文字描述发布期间旧 Pod 收到 SIGTERM → 停止接单 → 在途任务跑完60s 宽限→ 退出 新 Pod 先启动、Ready 后加入消费maxUnavailable: 0 对账结果500/500 SUCCESS去重表 500 行零丢失、零重复 若 worker_shutdown_timeout 配成 5s 而任务要 30s在途任务被 SIGKILL 打断 早确认队列直接丢任务——对账立刻出现缺口实验对照组。步骤 4验证 Beat 单副本与优雅退出参数目标确认「单主」与「宽限窗口」配置生效。kubectl get pods|Select-String beat# 确认只有 1 个 beatkubectl logs deploy/beat|Select-StringSending due# 无重复派发kubectl get deploy order-oyaml|Select-String terminationGracePeriodSeconds# 90运行结果文字描述Beat 只有一个 Pod双发不可能terminationGracePeriodSeconds: 90配置可见——三个关键参数单副本/宽限期/优雅退出全部落位。步骤 5回滚演练——发布失败如何安全回退目标发布有风险回滚要预演发布安全的另一半。# 发布 1.1 后对账发现失败率飙升 → 立即回滚helm rollback order0# 验证旧版本 Worker 以同样策略滚动接管同样 maxUnavailable: 0kubectl get pods-lapporder|Select-StringRunning# 回滚期间继续灌支付回调 → 对账零丢失运行结果文字描述helm rollback触发反向滚动——新 Pod 先起、就绪后旧 Pod 才被替换同样优雅退出发布/回滚两个方向都零丢任务回滚时间与发布相当约 1~2 分钟。教训「能回滚」和「回滚不丢」是两件事——回滚预案要用同样的宽限窗口与对账脚本验证。3.3 可能遇到的坑及解决方法坑现象解决发布丢任务SIGTERM 直接杀无优雅退出worker_shutdown_timeout terminationGracePeriodSeconds 配对Beat 双副本双发对账任务跑两遍beat replicaCount1 数据库 Scheduler第 22 章maxUnavailable1 丢窗口先杀旧再起新关键队列 maxUnavailable: 0探针误杀 Workerliveness 用 HTTP 探针用官方 inspect().stats() 命令探针无 HTTP 端口一容器多角色扩缩容互相牵连分角色镜像步骤 1回滚后队列里出现旧格式消息新旧镜像任务参数不兼容发布前契约评审第 28 章回滚配合消息 TTL 过期兜底3.4 完整代码清单与测试验证清单Dockerfile.worker/beat/flower、values.yaml、deployment.yaml模板 flood_pay_callbacks.py/verify_pay_callbacks.py发布对账脚本。发布安全核对表沉淀 Wiki检查项配置验证命令优雅退出worker_shutdown_timeout60发布日志看 warm shutdownK8s 宽限terminationGracePeriodSeconds90kubectl get deploy -o yaml零丢窗口maxUnavailable: 0对账脚本 500/500Beat 单主replicaCount1kubectl get pods分角色每角色独立 Deploymentkubectl get deploy测试验证# tests/test_deploy_config.py —— 部署配置的静态断言CI 可跑deftest_grace_period_gt_shutdown_timeout():importyaml valuesyaml.safe_load(open(values.yaml))assert(values[graceful][terminationGracePeriodSeconds]values[graceful][worker_shutdown_timeout])deftest_beat_is_single_replica():importyaml valuesyaml.safe_load(open(values.yaml))assertvalues[beat][replicaCount]1deftest_worker_groups_isolated():importyaml valuesyaml.safe_load(open(values.yaml))assertset(values[workerGroups]){sms,order,report}python-mpytest tests/test_deploy_config.py-v# 3 passed4. 项目总结4.1 优点 缺点维度分角色 优雅退出本章一容器全开 抄 Web 模板发布安全零丢任务宽限双跑丢任务直接杀伸缩粒度按队列组独立扩缩只能整体扩故障域角色互不牵连一崩全崩运维复杂度角色多、模板多简单成本镜像/编排文件多少4.2 适用场景适用① 生产订单/支付关键链路发布零丢是硬要求② 多队列多角色的大型任务平台③ 需要独立扩缩容与故障隔离的团队④ 与第 22 章数据库 Scheduler 配合的 Beat 高可用⑤ 需要「发布与回滚双向零丢」的成熟发布体系。不适用① 单机学习环境docker compose 足够第 16 章② 任务量极小且角色单一的简单服务③ 无 K8s 的团队先上 docker compose systemd 管理。4.3 注意事项terminationGracePeriodSeconds必须 worker_shutdown_timeout否则宽限形同虚设。长任务队列report的优雅退出窗口要按最长任务时长预留不是平均时长。发布对账脚本flood/verify是发布的标准动作没有对账的发布 没有发布。Beat 用 Deployment 单副本时重启窗口由数据库 Scheduler 的锁续约兜底第 22 章别用 StatefulSet 硬撑。回滚与发布要「同规格」回滚同样走优雅退出与 maxUnavailable: 0发布演练必须包含「失败回滚」方向步骤 5。4.4 常见踩坑经验3 个生产故障故障滚动发布丢 200 条支付回调。根因无优雅退出 maxUnavailable1。对策宽限窗口 maxUnavailable: 0 对账本章落地。教训Worker 的发布策略是「先交班后下班」不是「先下班后交班」。故障升级镜像后所有任务 NotRegistered。根因镜像里-A指向的模块路径变了。对策镜像内模块路径与启动命令对齐 发布前inspect registered预检。教训发布的第一个检查项是「任务注册表还在吗」。故障Beat 随 Worker 扩容变 4 副本对账跑 4 遍。根因一容器多角色 整体扩容。对策分角色 beat 单副本。教训扩 Worker 的冲动永远不要带节奏给 Beat。故障回滚后旧 Worker 消费不了新格式消息任务报错。根因镜像参数契约不一致回滚没校验兼容性。对策发布前契约评审第 28 章 回滚走同规格演练。教训回滚不是「回到旧版本」就完事消息格式的兼容窗口才是关键。4.5 思考题maxUnavailable: 0会让发布变慢先起后杀大促时想快发布怎么办提示牺牲窗口与容量的权衡按队列分级策略Worker 的readinessProbe为什么不能用 HTTP 探针inspect().stats()探针在 Worker「假死但进程活着」时能探出来吗提示第 15 章心跳与假死答案见第 30 章开头的「上一章思考题参考答案」。延伸阅读与资源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 实战修炼与源码剖析