拓冰建站拓冰建站
首页 / 资讯中心 / 正文

第5章:Celery Task 定义——绑定、命名、基类与请求上下文

0. 上一章思考题参考答案思考题 1懒加载让「创建 App」变得极轻① 任务模块一被 import 就初始化全部组件会引入循环依赖config_from_object需要模块路径而模块 import 又依赖 app② 测试无法在创建后注入配置第 4 章 3.2 的 pytest 全靠「先建后配」③ 用不到的组件如 Backend白白初始化启动变慢。所以 Celery 只在第一次访问配置/组件时才真正装配。思考题 2current_app线程局部默认指向最近创建的 Celery 实例。多 App 同进程时后创建者会「赢」。危险场景全局信号回调、定时任务或第三方库用current_app取配置却拿到了另一个 App 的实例——配置错乱难以定位。规避多 App 场景显式持有实例引用不用current_app。1. 项目背景订单短信任务上线后小周接到新需求订单域要新增「超时关单」「发票生成」「库存预占」三个异步任务。小周复制粘贴了三次app.task于是出现了三份几乎一样的代码每个任务里都手写「日志打印 task_id」「异常捕获翻译」「重试次数拼接」——改一次重试策略要改三个文件。更麻烦的是leader 要求所有订单任务自动携带链路 IDtrace_id方便跨系统追踪用户「下单 → 短信 → 发票」要能一条链查到底。小周还遇到两个诡异问题任务函数里想读「当前任务是谁、重试第几次了」发现函数里根本没有这些信息——只能靠闭包 hack。一个任务里发起了另一个任务订单任务里发短信短信任务的 logger 打出来的task_id是短信任务自己的还是订单任务传的他完全说不清。这些问题都指向同一个认知升级Task 不只是「被装饰的函数」它是一个对象有选项、有上下文、可以被继承。把「重复的横切逻辑」收敛进自定义基类就是本章的实战目标。现状三份复制粘贴 send_order_sms手写日志重试拼接 close_order 手写日志重试拼接漏了 trace_id 一行 gen_invoice 手写日志重试拼接重试策略和另两个不一致 ↓ 抽象后 class OrderTask(Task): # 统一trace_id 注入、重试策略、日志格式 ...2. 项目设计场景小周在午休时把「三份复制粘贴」的代码铺开给大师看。小胖复制粘贴咋了三个任务又不会同时改改一个是一个我打王者都没你这么纠结。小白小胖别打岔。我看了下文档任务函数有个bindTrue参数绑定了之后函数第一个参数是self能拿到self.request——里面有任务 ID、重试次数。但我不理解self到底是什么是那个Task实例吗那同一个任务被 100 个 Worker 并行执行self会串吗大师问得漂亮这正是新手最容易懵的点。bindTrue后函数的第一个参数self指向任务类的实例如OrderTask但注意Worker 每次执行任务时是「按消息创建新的 Request 上下文」。self实例本身可以共享它承载的是定义而self.request是线程/进程局部的——每次执行都会把当前 Request 塞进去。所以「100 个任务并行self.request.id各是各的」因为每个执行线程的 Request 都不同。技术映射Task 类 食堂的「菜谱册」大家都翻同一本self.request 正在炒的这张点菜单每个灶台各拿各的不串台。小周那任务级选项是啥意思我看到name、serializer、ignore_result、acks_late、track_started、max_retries都能写进装饰器它们和app.conf.task_*全局配置啥关系大师任务级选项覆盖全局配置只对本任务生效。比如全局task_acks_lateFalse但「扣库存」任务可以单独app.task(acks_lateTrue)。优先级是任务级 调用时选项 全局配置。其中name是任务名契约第 3 章讲过、serializer是编码方式第 12 章、ignore_result是不写结果省 Redis第 8 章、acks_late是晚确认第 18 章、track_started是记录 STARTED 状态第 10 章。这些不用现在全懂但要知道装饰器括号里的东西是「任务属性」不是摆设。小白最后一个问题shared_task和app.task有什么区别为什么 Django 项目里推荐前者大师app.task是把任务注册到指定 app上shared_taskcelery/shared_task.py是先记在全局注册表里等 app 创建后再补绑——这样任务模块可以先被 import不依赖任何具体 app 实例。Django 场景下任务模块经常在 app 初始化之前就被 import比如 models 相互引用所以shared_task是安全选择源码在celery/contrib/django/task.py第 27 章展开。两者最终都是「任务注册进 app.tasks」区别只在注册时机与耦合方式。3. 项目实战3.1 环境准备沿用前几章环境源码安装 Redis。新建项目骨架order_service/ ├── celeryconfig.py # 第 4 章配置 ├── order_tasks.py # 本章主角OrderTask 基类 订单任务 └── call_demo.py # 运行验证脚本3.2 分步实现步骤 1定义OrderTask基类——统一 trace_id 注入、日志、重试策略目标把「横切逻辑」收敛到基类任务函数只剩业务。# order_tasks.pyimportlogging,uuidfromceleryimportCeleryfromcelery.app.taskimportTask appCelery(order_tasks)app.config_from_object(celeryconfig)loggerlogging.getLogger(orders)classOrderTask(Task):订单域任务基类自动注入 trace_id统一日志格式与重试策略。# 任务级选项写在这里 所有子任务继承可被具体任务覆盖acks_lateTrue# 晚确认防止执行中崩溃丢任务第 18 章详解max_retries3# 最多重试 3 次第 11 章详解ignore_resultFalsedef__call__(self,*args,**kwargs):Worker 真正执行任务的入口prerun 后调用 run 之前。 在这里统一注入 trace_id比在业务函数里逐行写要可靠得多。trace_idself.request.headers.get(trace_id)ifself.request.headers \elseNonetrace_idtrace_idorstr(uuid.uuid4())self.request.headers{**(self.request.headersor{}),trace_id:trace_id}logger.info(task%s id%s trace%s args%s,self.name,self.request.id,trace_id,args)returnsuper().__call__(*args,**kwargs)defon_failure(self,exc,task_id,args,kwargs,einfo):统一失败日志任何订单任务抛异常都会走到这里。logger.error(task%s id%s failed: %s,self.name,task_id,exc)returnsuper().on_failure(exc,task_id,args,kwargs,einfo)步骤 2用基类定义三个订单任务业务函数保持纯净目标任务函数里只写业务横切逻辑全在基类。# order_tasks.py追加app.task(baseOrderTask,bindTrue,nameorders.send_order_sms)defsend_order_sms(self,order_id:int)-bool:短信任务绑定后可通过 self 访问请求上下文。retriesself.request.retries# 当前是第几次重试logger.info(order%s retries%s trace%s,order_id,retries,self.request.headers.get(trace_id))# 模拟短信网关调用ifretries0andorder_id%70:# 故意让 1/7 的首次调用失败raiseConnectionError(短信网关超时)returnTrueapp.task(baseOrderTask,bindTrue,nameorders.close_order)defclose_order(self,order_id:int)-None:超时关单任务任务里再发起子任务展示 request 上下文传递。logger.info(close order%s trace%s,order_id,self.request.headers.get(trace_id))send_order_sms.delay(order_id)# 子任务调用app.task(baseOrderTask,nameorders.gen_invoice,ignore_resultTrue)defgen_invoice(order_id:int)-None:发票生成任务ignore_resultTrue 不写结果省 Redis。logger.info(gen invoice order%s,order_id)关键点①bindTrue后函数第一个参数必须是selfdelay(order_id)的实参会自动对齐到第二个参数② 任务级选项可以写在基类属性上被子任务继承也可以写在装饰器里单独覆盖——两者都能生效基类更利于统一治理。步骤 3验证 trace_id 传递与 self.request 内容目标用真实运行验证「self.request 是执行时上下文」的心智模型。# call_demo.pyfromorder_tasksimportsend_order_sms,close_order,gen_invoice,app r1send_order_sms.delay(21)# 首次必失败 → 触发重试r2close_order.delay(100)# 任务里再发子任务r3gen_invoice.delay(200)# ignore_result结果不写 Backendprint(r1:,r1.id,| r2:,r2.id,| r3:,r3.id)celery-Aorder_tasks worker--loglevelinfo--poolsolo python call_demo.py运行结果文字描述r1 首次执行失败日志出现 taskorders.send_order_sms id... failed: 短信网关超时 约 1 秒后默认退避重试再次执行成功日志中 retries1 r2 执行时打印 close order100随后 Worker 收到 orders.send_order_sms 的子任务消息 其 trace_id 与父任务相同——子任务通过消息头接力传递链路 ID r3 执行正常但 celery result r3.id 返回 Noneignore_result 生效。步骤 4用app.tasks验证注册与继承关系目标把「基类属性继承」固化成可断言的结论。# verify_tasks.pyfromorder_tasksimportapp,send_order_sms,close_order,gen_invoice tapp.tasks[orders.close_order]print(基类:,t.__class__.__mro__[1].__name__)# OrderTaskprint(继承 acks_late:,t.acks_late)# True继承自基类print(gen_invoice ignore_result:,app.tasks[orders.gen_invoice].ignore_result)print(send_order_sms bind 后签名:,send_order_sms.__name__)# 任务对象而非函数步骤 5老任务零改造治理——task_annotations按任务名批量注入选项目标历史任务没写基类通过配置批量设置任务选项改造成本为零。# celeryconfig.py 中追加task_annotations{orders.*:{max_retries:5,acks_late:True},# 按任务名前缀批量生效orders.gen_invoice:{ignore_result:True},# 精确覆盖}说明task_annotations源码在celery/app/annotations.py在任务执行前按规则把属性注入任务实例支持orders.*前缀通配。它是「存量任务统一治理」的兜底手段优先级低于装饰器参数与基类属性——别用它覆盖已有显式配置否则会出现「明明配了却不生效」的困惑。步骤 6任务级选项速查表写任务前过一眼目标装饰器括号里能写什么心里有数。选项作用默认值详解章节name任务名契约模块.函数名第 3 章bind函数第一个参数变 self可访问请求上下文False本章serializer本任务消息编码全局 json第 12 章ignore_result不写结果省存储False第 8 章acks_late执行完再确认防丢但可能重复False第 18 章track_started记录 STARTED 状态False第 10 章max_retries最大重试次数3第 11 章rate_limit任务执行速率限制无第 21 章3.3 可能遇到的坑及解决方法坑现象解决bindTrue后调用报参数错delay(order_id)报 missing selfbindTrue后self是第一个参数实参从第二个对齐检查调用处基类属性不生效在装饰器里app.task(baseOrderTask, max_retries5)与基类属性混用时优先级困惑规则装饰器参数 基类属性 全局配置同层后者覆盖on_failure不触发异常被业务函数内 try/except 吞了基类只兜底「未捕获异常」业务内捕获就自己处理self.request.headers为 None生产者没传 headers判空后再or {}第 3 步代码已处理子任务 trace_id 丢失子任务没继承父任务头子任务调用时显式apply_async(headers{trace_id: ...})或全局信号统一注入第 26 章3.4 完整代码清单与测试验证清单order_tasks.pyOrderTask 基类 3 个任务、call_demo.py、verify_tasks.py。生产规范补充订单任务一律用baseOrderTask代码评审见baseTask裸直接打回。测试验证# tests/test_tasks.pyfromunittestimportmockfromorder_tasksimportapp,send_order_sms,close_order app.conf.task_always_eagerTrue# 同步执行模式deftest_order_task_injects_trace_id():withmock.patch(order_tasks.logger)asm:send_order_sms.run(order_id1)# 基类 __call__ 一定打过日志且 trace_id 已注入 request.headersassertany(traceinc[0][0]forcinm.info.call_args_list)deftest_task_inherits_base_options():tapp.tasks[orders.close_order]assertt.acks_lateisTrue# 继承自 OrderTaskassertt.max_retries3deftest_bind_self_request_available():withmock.patch(order_tasks.logger)asm:send_order_sms.run(order_id1)# run() 直接执行函数体self.request 是同步上下文retries0assertany(retries0inc[0][0]forcinm.info.call_args_list)deftest_shared_task_macro():# shared_task 无需实例即可定义测试注册时机fromceleryimportshared_taskshared_task(nameorders.any)defany_task():return1assertorders.anyinapp.taskspython-mpytest tests/test_tasks.py-v# 4 passed4. 项目总结4.1 优点 缺点维度自定义 Task 基类收敛横切逻辑裸app.task 手写重复代码一致性重试/日志/链路全站统一三份代码三种行为可观测基类一处注入 trace_id全链路可查靠自觉漏一个断一截可维护改策略只动基类改一处漏两处学习成本需要理解 bind/call/MRO零成本缺点 1基类改出 bug 影响所有订单任务无传染面缺点 2调试时心智模型多一层简单直观4.2 适用场景适用① 同业务域多个任务共享日志/重试/链路需求订单域、支付域② 需要统一失败兜底告警、死信表的平台③ 需要强制注入 trace_id 的跨系统追踪体系。不适用① 只有一两个任务的微型项目过度设计② 任务行为差异极大、共性几乎没有的场景硬抽象反而僵硬。4.3 注意事项bindTrue改变函数签名代码评审时注意调用侧参数对齐。任务级选项继承链装饰器参数 基类属性 全局配置统一治理时优先用基类属性。ignore_resultTrue后AsyncResult.get()永远返回 None别在需要结果的流程里误开。task_annotations只做兜底治理新代码一律走基类两种机制同时用时先查优先级再动手改配置。4.4 常见踩坑经验3 个生产故障故障大促订单短信全部丢失排查发现任务抛异常但无任何日志。根因业务代码把异常吞掉且没写日志基类没有 on_failure。对策基类统一 on_failure 告警本章落地。教训失败路径必须有一条兜底日志/告警。故障链路追踪断链订单任务有 trace_id短信任务没有。根因子任务调用没带 headers。对策基类 全局信号统一注入第 26 章根治。教训链路 ID 必须在「发消息」这一刻注入而不是在任务里补。故障某任务重试策略被「悄悄改大」导致短信连发 10 次。根因有人在该任务的装饰器里覆盖了 max_retries10。对策治理规约——重试策略只放基类装饰器不允许覆盖。教训统一治理要连「覆盖口子」一起堵住。4.5 思考题bindTrue的self.request是线程局部的。那在OrderTask.__call__里修改self.request.headers会不会影响其他并行执行的同名任务提示区分「实例共享」与「Request 局部」shared_task与app.task最终都注册进app.tasks那「多 App 同进程」时shared_task会注册进哪个app这暴露了 shared_task 的什么代价答案见第 6 章开头的「上一章思考题参考答案」。延伸阅读与资源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 实战修炼与源码剖析
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门