
0. 上一章思考题参考答案思考题 1celery/app/管的是App 生命周期——配置、任务注册表、消息路由随进程常驻celery/worker/管的是消费运行时——谁在消费、怎么确认、请求上下文怎么包装生命周期随 Worker 启停celery/backends/管的是结果存储——状态与返回值落哪、何时过期。三者生命周期完全不同配置要稳定、消费要可控、结果要过期拆开才能独立演进与运维。思考题 2-e是 editable可编辑模式pip 在 site-packages 里写入一个指向仓库的路径钩子Windows 下是__editable__.celery...pth文件。Python 导入celery时按该路径直接找到仓库里的源码因此改仓库文件立即生效、无需重装。1. 项目背景上一章环境跑通后leader 给小周派了第一个正经需求下单成功后在 2 秒内给用户发一条订单确认短信但下单接口必须在 50ms 内返回。小周第一反应是加个线程池在接口里threading.Thread(targetsend_sms, args(order_id,)).start()。上线的第一个晚上就出事了服务重启时丢了几百条短信线程里的任务随进程蒸发短信网关限流时线程池被打爆主进程线程数飙到 800接口开始排队最终还是卡了用户。小周这才意识到「把代码放到别的线程跑」不等于「异步」。异步系统的第一性原理是任务要和调用方进程隔离并且任务本身要被持久化、可确认、可重试。而这正是「把函数升级成 Task」的意义。那问题来了app.task到底对一个函数做了什么为什么调task.delay(...)之后接口立刻返回而send_order_sms真的在 2 秒后执行了本章用一个小而完整的例子把「函数 → 分布式执行单元」的变身过程拆开看。下单接口50ms 返回 │ send_order_sms.delay(order_id) ▼ 消息 → RedisBroker→ Worker 拾取 → 包装成 Request → 执行函数体约 2s→ 结果后端2. 项目设计场景小周把自己的线程池方案讲给大师和小胖听。小胖线程池咋了我吃火锅的时候锅底和配菜不是同步上的——锅底先上涮肉慢慢来这不就是异步吗你那个Thread(...).start()我觉得挺好简单小白但是小周你也说了进程一重启线程里的任务全没了。我还在想另一个问题Worker 收到任务的时候怎么知道要执行哪个函数函数在 Web 进程里是一个内存对象到了 Worker 进程里同一个函数对象并不存在——那消息里到底传的是什么大师小白问到了核心。先回答小胖线程池里「任务」的生命周期等于进程的生命周期进程挂了任务就蒸发而 Celery 里任务被序列化成一条消息落到 Broker 里Worker 挂了消息还在换台机器照样执行。这就叫「任务的持久化」。再回答小白消息里传的从来不是函数而是函数的「名字」。send_order_sms在 Web 进程里被app.task装饰时Celery 做三件事把函数包装成一个celery.app.task.Task实例按「模块名 函数名」算出任务名如tasks.send_order_sms并注册进app.tasks注册表在 Task 实例上挂delay()、apply_async()等方法。Worker 进程启动时也会加载同样的模块、执行同样的注册——所以 Worker 拿到消息里的任务名到自己的注册表里一查就找到那个函数了。两边代码一样注册表就能对上暗号。技术映射任务名 菜单上的菜名跨进程的稳定标识任务函数体 后厨菜谱每个厨房都有一份消息 写有菜名的点菜单。小周那我调delay(order_id)的时候消息里是怎么装我的参数的order_id是 Python 对象Worker 那边拿到的是同一个对象吗大师消息体里放的是args和kwargs的序列化结果——默认 JSON。所以参数必须是 JSON 可序列化的数字、字符串、列表、字典都行ORM 对象、socket、文件句柄一概不行第 6 章讲参数契约。Worker 反序列化后重建参数再包进一个叫Request的上下文对象里里面带任务 ID、重试次数、投递信息等最后才调用你的函数。所以你在任务函数里能拿到self.request.id就是这个 Request。这一套「消息 → Request → 函数」的链路源码在celery/worker/request.py和celery/app/trace.py第 34 章精读。小胖那delay()和apply_async()有啥区别我百度看到俩都能发任务一个带个框一个不带大师笑delay(*args, **kwargs)就三行代码——return self.apply_async(args, kwargs)celery/app/task.py:534。delay 是 apply_async 的语法糖只管传参不管别的选项要加countdown、queue、expires这些控制项就得用apply_async(..., countdown10, queuesms)。另外还有个apply()——它是同步执行不经过 Broker直接在当前进程跑函数测试和调试用的别在生产用。三个方法的边界delay 传参、apply_async 传参选项、apply 同步执行。技术映射delay 只点菜不交代默认做法apply_async 点菜 交代口味忌口队列、延时apply 当着你的面现炒同步。3. 项目实战3.1 环境准备沿用第 2 章环境源码可编辑安装 Redis。新增依赖无HTTP 演示用 Python 标准库。dockercompose-fdocker/docker-compose.yml up-dredis# 已启动可跳过cdexamples/tutorial3.2 分步实现步骤 1定义第一个业务任务send_order_sms目标用app.task把「发短信函数」升级为「可跨进程执行的任务」。# examples/tutorial/order_tasks.pyimporttimefromceleryimportCelery appCelery(order_tasks,brokerredis://localhost:6379/0)app.task(nameorders.send_order_sms)# 显式命名防止模块改名后任务名漂移defsend_order_sms(order_id:int)-bool:模拟调用短信网关耗时约 2 秒返回是否发送成功。time.sleep(2)# 模拟三方网关 I/Oprint(f[SMS] 订单{order_id}的确认短信已发出)returnTrue关键点显式nameorders.send_order_sms——任务名是跨进程契约模块路径变了名字可能漂移显式命名是生产规范第 6、28 章扩展。命名建议遵循三段式约定业务域.动作如orders.send_order_sms、billing.gen_invoice跨团队共享的任务再加团队前缀如promo.orders.issue_coupon。任务名一旦对外发布就不要改第 4 章讲过「在途消息」会因改名而NotRegistered。步骤 1.5任务多模块时用autodiscover_tasks自动发现目标任务分散在多个业务模块时一键注册告别手工维护 import 清单。# order_tasks.py 中追加模块拆分场景app.autodiscover_tasks([order,billing],related_nametasks)# 等价于自动 import order.tasks、billing.tasks并注册其中的任务说明autodiscover_tasks源码在celery/app/base.py会按「包名.related_name」批量 import 任务模块业务模块约定统一叫tasks.py即可被自动发现。它是第 27 章 Django 集成的标配姿势注意它静默失败——路径写错不报错排查时先打印app.tasks.keys()核对。步骤 2用「下单接口」演示同步与异步的差别目标同一个接口对比内联调用阻塞 2 秒与delay()毫秒返回。# examples/tutorial/web_app.pyimportjson,timefromhttp.serverimportBaseHTTPRequestHandler,HTTPServerfromorder_tasksimportapp,send_order_smsclassHandler(BaseHTTPRequestHandler):defdo_POST(self):ifself.path/api/orders:order_idint(time.time()*1000)%1000000t0time.time()send_order_sms.delay(order_id)# 异步只发消息立即返回costround((time.time()-t0)*1000,1)self._respond(200,{order_id:order_id,async_cost_ms:cost})else:self._respond(404,{error:not found})def_respond(self,code,body):datajson.dumps(body).encode()self.send_response(code)self.send_header(Content-Type,application/json)self.send_header(Content-Length,str(len(data)))self.end_headers()self.wfile.write(data)if__name____main__:HTTPServer((127.0.0.1,8000),Handler).serve_forever()步骤 3启动 Worker用 curl 验证「接口快、任务慢」目标验证 Web 侧 50ms 内返回Worker 侧约 2 秒后完成。# 终端 A启动 WorkerWindows 加 --poolsolocelery-Aorder_tasks worker--loglevelinfo--poolsolo# 终端 B启动下单接口python web_app.py# 终端 C发一个下单请求curl-XPOST http://127.0.0.1:8000/api/orders运行结果文字描述curl 响应{order_id: 123456, async_cost_ms: 3.2} # 接口 3 毫秒返回远小于 50ms 终端 A Worker 日志 [INFO/MainProcess] Task orders.send_order_sms[...] received [SMS] 订单 123456 的确认短信已发出 # 2 秒后才打印 [INFO/MainProcess] Task orders.send_order_sms[...] succeeded in 2.003s步骤 4用AsyncResult查询任务状态目标给调用方一个「查单」的口子验证第 1 章「Backend 记账」的闭环。先给 App 配上结果后端。# order_tasks.py 中修改appCelery(order_tasks,brokerredis://localhost:6379/0,backendredis://localhost:6379/1)# 终端 Ccelery-Aorder_tasks call orders.send_order_sms--args[789]celery-Aorder_tasks result任务ID运行结果result 输出True # 任务执行成功返回值可查结果存在 Redis 1 号库30 分钟过期步骤 5用apply()与apply_async(countdown...)体验另两种调用目标对比三种调用方式的语义差异。# examples/tutorial/call_demo.pyfromorder_tasksimportsend_order_sms r1send_order_sms.apply(args(1,))# 同步执行阻塞 2 秒适合测试r2send_order_sms.apply_async(args(2,),countdown5)# 5 秒后执行可加 queue/expires 等选项print(同步结果:,r1.get(),| 异步任务ID:,r2.id)python call_demo.py# 预期先阻塞 2 秒打印 [SMS] 订单 1 ...再打印同步结果与异步任务 ID# Worker 日志约 5 秒后出现订单 2 的执行记录。同步与异步的适用边界小结① 需要立刻拿到结果才能继续 → 走同步REST 调用或apply()仅限测试调试② 结果几秒后才用 → 异步 AsyncResult查询③ 只发不管结果 → 异步 ignore_result。选择标准一句话调用方能不能等、结果要不要拿。拿不准时优先异步 查询因为它对调用方最无侵入。3.3 可能遇到的坑及解决方法坑现象解决NotRegistered: tasks.send_order_sms任务名对不上显式命名后必须完全一致Worker 重启加载最新模块任务函数里传了 ORM 对象序列化报TypeError: Object of type Order is not JSON serializable只传主键/ID对象由 Worker 进程自己查库第 6 章AsyncResult.get()一直 PENDING忘了配 backendCelery(..., backend...)或app.conf.result_backend redis://...Web 进程 import 任务模块时触发副作用接口启动变慢任务模块只定义任务不执行业务初始化import order_tasks放在接口入口处Windows 下 curl 没输出可能是启动顺序问题先 Worker 再 WebWorker 用--poolsoloautodiscover_tasks没找到任务也不报错模块路径或related_name写错静默失败用app.autodiscover_tasks([order], forceTrue)并打印app.tasks核对3.4 完整代码清单与测试验证清单order_tasks.py任务定义 配置、web_app.py下单接口、call_demo.py调用对比共 3 个文件。生产规范任务定义与调用方分属不同模块调用方只import不from ... import *。测试验证单元级用task_always_eager让任务在测试进程内同步执行# tests/test_order_sms.pyfromunittestimportmockfromorder_tasksimportapp,send_order_sms app.conf.task_always_eagerTrue# 测试模式不发 Broker直接同步执行deftest_send_order_sms_returns_true():assertsend_order_sms.run(order_id1)isTrue# run() 是任务函数体deftest_send_order_sms_registered_with_explicit_name():assertorders.send_order_smsinapp.tasksmock.patch(order_tasks.send_order_sms.delay)deftest_web_calls_delay_once(mock_delay):fromweb_appimportHandler# 触发模块加载mock_delay.assert_called_once_with(123456)# 调用方只走 delay不直接执行python-mpytest tests/test_order_sms.py-v# 3 passed4. 项目总结4.1 优点 缺点维度Celery Task消息驱动执行单元线程池 Thread().start()持久性消息落 Broker进程崩溃可恢复任务随进程蒸发扩展性任务可跨机器执行加 Worker 即扩容受限于单进程线程数可观测性有任务 ID、状态机、事件只能打日志触发方式消息驱动无共享内存共享内存、GIL 限制缺点 1参数必须可序列化参数随便传缺点 2有重复执行可能需幂等天然不重复缺点 3引入中间件运维成本零依赖4.2 适用场景适用① 通知类短信/邮件/推送任务② 需要跨进程解耦的业务动作③ 需要重试与状态追踪的写操作④ 定时任务配合 Beat第 13 章。不适用① 必须在请求内拿到结果才能继续的强同步流程② 任务参数含不可序列化对象需先改造为 ID 引用③ 极端低延迟微任务进程间通信开销大于收益。4.3 注意事项任务名一旦对外发布跨团队调用、线上消息在途改名要灰度兼容否则在途消息将NotRegistered。apply()同步执行绕过 Broker千万别在 Web 请求里用它「图省事」。结果后端不是免费的每个任务结果占 Redis 键记得result_expires默认 24 小时第 8 章调优。4.4 常见踩坑经验3 个生产故障故障升级任务模块后消息报NotRegistered。根因在途消息里的任务名还是旧名而新 Worker 只注册了新名。对策显式任务名 双版本窗口期同时注册。教训任务名是跨版本契约轻易不改。故障接口偶发 2 秒慢请求定位发现任务函数里有人直接requests.get调短信网关。根因把业务逻辑写进了「接口内联调用」分支。对策代码评审强制接口只允许 delay/apply_async。教训异步接口不允许同步调用远程服务。故障get()卡死超时排查发现把 ORM 对象传进了任务参数。根因参数被 pickle 时拖进整个 Session。对策只传主键Worker 内重新查询。教训任务参数必须可序列化、要最小化。4.5 思考题delay()与apply_async()的区别是什么为什么说「delay 只适合最简单的调用」提示看celery/app/task.py:534的实现如果两个 Worker 同时消费同一个队列任务 A 的countdown10会保证「不早于 10 秒执行」吗它是怎么被实现的提示任务消息是立即投递的延迟是在哪一侧实现的答案见第 4 章开头的「上一章思考题参考答案」。延伸阅读与资源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 实战修炼与源码剖析