公司动态
第6章:Celery任务调用方式与参数契约
0. 上一章思考题参考答案思考题 1不会串。Task实例self虽被共享但self.request是执行上下文的线程/协程局部绑定——每次任务执行时trace 都会为当前执行线程创建新的Request并绑定上去celery/app/trace.py。你在__call__里改self.request.headers改的是「当前这次执行」的 Request其他并行执行中的同名任务各持各的 Request互不可见。思考题 2shared_task定义时只登记进全局注册表并不绑定任何 App真正绑定时机是某个 App 执行finalize()时把全局注册表里的任务并入「当时的 current_app」。所以多 App 同进程时shared_task 会被第一个 finalize 的 App收走后续创建的 App 拿不到它。代价shared_task 默认假设「一进程一主 App」多 App 场景要么各自显式app.task要么手动控制 finalize 顺序。1. 项目背景订单服务要上线了leader 要求把「下单成功后的异步动作」收敛成一个统一出口OrderService.place(order_id)团队 4 个人却写出了 4 种风格阿强传整个 ORM 对象send_order_sms.delay(order_obj)阿珍传字典{oid: 100, uid: 7}阿伟传位置参数send_order_sms.delay(100, 7)小周想传参数还带了个eta——结果eta被当成任务函数的关键字参数吞进了**kwargs短信任务直接TypeError。上个月还出过一次「生产大事故」有人在任务里传了 ORM 对象测试环境用了 pickle 序列化勉强能跑把整个 session 都拖进了消息体一上生产切成 JSON 直接Object of type Order is not JSON serializable——几百条消息发送失败任务堆积接口开始降级。复盘结论一句话任务参数不是想传什么就传什么调用方式也不是只有 delay 一种。跨进程的调用必须像 API 一样有契约。现状同一任务 4 种调用姿势 send_order_sms.delay(order_obj) # 传 ORM 对象 → 序列化炸 send_order_sms.delay({oid:100}) # 传字典 → Worker 侧取错字段 send_order_sms.delay(100, 7) # 位置参数 → 加参数即灾难 send_order_sms.delay(100, eta...) # eta 被当 kwargs → TypeError ↓ 治理 OrderService.place(order_id) # 唯一出口 参数校验 过期兜底2. 项目设计场景代码评审会上大师把四种调用姿势贴到屏幕上。小胖发任务不就一个delay吗我看官方文档第一页就在用你们搞出apply_async、send_task、signature一堆名词跟奶茶店的「标准杯、大杯、超大杯、桶装」一样纯属营销。小白我先问个技术问题delay(order_id)和apply_async(args[order_id])到底是不是一回事还有send_task(orders.send_order_sms, args[100])连任务对象都不用 import它和delay的边界在哪signature又是什么大师这四个是一张「调用方式谱系」按抽象程度排调用方式需要 import 任务对象能带的选项定位delay(*args, **kwargs)需要不能带只有参数日常最简单调用apply_async(args, kwargs, **options)需要全量countdown/queue/expires/link…需要控制投递行为时signature(name, args...)/s()不需要冻结参数可编排Canvas 编排的积木send_task(name, args, kwargs)不需要全量按任务名字符串投递记住两句话delay是apply_async的语法糖apply_async的完整签名才是「真身」celery/app/task.py:547send_task按字符串投递本地不需要任务定义适合跨服务解耦celery/app/base.py:846。signaturecelery/canvas.py:234则是把「任务参数选项」冻结成一个对象可以存、可以传、可以拼进 chain/group——它是第 19 章 Canvas 的基础积木。技术映射delay 到店点单报了菜名就走apply_async 点单备注忌口辣度、时间、包间send_task 电话叫外卖报地址就行不用进店signature 打包好的固定套餐券随时能用、还能拼团。小白那我追问选项的优先级任务装饰器里写过app.task(queuesms)全局配置也有task_default_queue调用时又传queueorder——最后进哪个队列还有expires到底是谁在过期谁来丢掉过期的任务大师优先级三条调用时选项 任务级选项 全局配置。queueorder赢。expires要注意一个反直觉的点过期判断发生在 Worker 收到消息之后——消息照样投递Worker 拿到一看「哦过期了」才丢。所以 Broker 已经堆积了几万条时expires救不了你该堵还是堵根治要靠第 21 章的队列治理。countdown是「现在N 秒后执行」eta是「具体时刻执行」两者互斥countdown的语义是「不早于」第 3 章思考题讲过。小胖挠头那回调呢我看有人写send_order_sms.apply_async(args[100], linksend_push.delay)——发完短信自动发个 push 推送这不开挂了吗大师对link是「成功后自动投递下一个任务」前一个任务的返回值会作为第一个参数传给回调link_error是「失败后投递告警任务」。不过别高兴太早回调任务也是独立投递的消息它不会继承原任务的事务、队列和路由出问题要单独排查。而且 link 的值应该传signaturesend_push.s(...)传delay的返回值已经是 AsyncResult就错了——这是新手最常犯的。技术映射link 电影放完自动接着放下一部link_error 停电了自动发道歉短信。但「自动」≠「免费」回调任务的队列、幂等都要自己管。3. 项目实战3.1 环境准备沿用第 3 章环境源码安装 Redis order_tasks.py。新增无外部依赖本章全部用标准库 已有代码。3.2 分步实现步骤 1建立参数契约——只传order_id拒绝一切对象目标定义OrderService.place()作为唯一投递出口参数先过校验。# order_service.pyimportjsonfromorder_tasksimportsend_order_sms,send_push,alert_opsdef_check_json_serializable(value,path$):递归校验参数可 JSON 序列化——把「运行时 TypeError」提前到「投递前报错」。try:json.dumps(value)exceptTypeErrorase:raiseValueError(f任务参数不可序列化{path}:{e})classOrderService:staticmethoddefplace(order_id:int,mobile:str)-str:下单成功后的统一异步出口。契约只传标量不传对象。ifnotisinstance(order_id,int)ororder_id0:raiseValueError(order_id 必须为正整数)_check_json_serializable({order_id:order_id,mobile:mobile})# 统一通过 apply_async 带选项过期 5 分钟兜底 失败告警回调rsend_order_sms.apply_async(args[order_id],expires300,# 5 分钟内没执行就作废秒linksend_push.s(mobilemobile),# 成功后发 pushsignature 对象link_erroralert_ops.s(order_idorder_id),# 失败进告警任务)returnr.id关键点expires300单位是秒link传的是send_push.s(...)签名对象而不是send_push.delay(...)的返回值。参数契约表随代码一并交付评审/测试/运维三方共用任务参数类型必填默认值过期回调orders.send_order_smsorder_idint是—300s成功→send_push失败→alert_opsorders.close_orderorder_idint是—600s失败→alert_opsorders.gen_invoiceorder_idint是—1800s—这张表就是「跨进程契约」的实体化第 28 章会把它升级成可执行契约测试直接读表断言。步骤 2补齐 Worker 侧防御——任务函数里二次校验目标契约不能只靠调用方自觉Worker 侧同样兜底。# order_tasks.py 中更新 send_order_smsapp.task(nameorders.send_order_sms,bindTrue)defsend_order_sms(self,order_id:int)-bool:ifnotisinstance(order_id,int)ororder_id0:raiseValueError(f非法 order_id:{order_id!r})# 参数错误不进重试...# 业务逻辑步骤 3演示「调用时选项覆盖任务选项」目标用真实运行验证优先级规则。# call_options.pyfromorder_tasksimportsend_order_sms# 任务定义在 orders 队列调用时覆盖到 sms 队列第 9 章正式拆队列r1send_order_sms.delay(1)r2send_order_sms.apply_async(args[2],queuesms,countdown5)r3send_order_sms.apply_async(args[3],expires2)# 2 秒内不执行即作废print(r1:,r1.id,| r2:,r2.id,(5s 后执行),| r3:,r3.id,(2s 过期))步骤 4用send_task演示「不 import 任务也能投递」目标体验按名字投递的解耦能力跨语言系统/管理后台常用。# send_task_demo.pyfromorder_tasksimportapp# 不 import 任何任务函数直接按名字投递app.send_task(orders.send_order_sms,args[999],countdown10)print(已按名字投递 orders.send_order_sms)运行结果文字描述Worker 日志约 10 秒后出现Task orders.send_order_sms[...] received说明消息已正常入队并执行。步骤 5验证 link / link_error 回调链目标亲手确认「成功回调收返回值、失败回调收告警」。# link_demo.pyfromorder_tasksimportsend_order_sms,alert_ops# 正常链路成功 → send_push 收到前任务的返回值 True 作为第一个参数send_order_sms.apply_async(args[101],linkalert_ops.s(reasonsms-ok))# 异常链路任务抛 ConnectionError → link_error 触发 alert_opssend_order_sms.apply_async(args[10001],link_erroralert_ops.s(reasonsms-fail))运行结果文字描述Worker 日志中先出现orders.send_order_sms[...] succeeded in ...随后自动出现orders.alert_ops[...] received——注意回调任务的日志里第一个参数是上游任务的返回值异常链路同理先 FAILURE 再 alert_ops。两个任务之间没有调用方参与全由 Worker 侧自动衔接——这就是 link 的价值把「等结果再发下一个」的轮询代码全部消灭。步骤 6契约演进演练——加参数与改语义的正确姿势目标把第 2 节说的契约演进原则变成可执行的动作。# 演进 1加参数 —— 新参数带默认值老调用方无感app.task(nameorders.send_order_sms,bindTrue)defsend_order_sms(self,order_id:int,channel:straliyun)-bool:...# 老调用方只传 order_id 依然可用# 演进 2改语义 —— 不碰老任务名新开 v2灰度切换app.task(nameorders.send_order_sms_v2,bindTrue)defsend_order_sms_v2(self,order_id:int,batch:listNone)-bool:...# 支持批量发送语义完全不同# 演进 3投递侧灰度 —— 读配置决定走 v1 还是 v2importos TASK_NAME(orders.send_order_sms_v2ifos.environ.get(SMS_V2)1elseorders.send_order_sms)app.send_task(TASK_NAME,args[102])运行结果文字描述SMS_V21环境变量切换后投递的任务名变为 v2Worker 侧执行新的批量逻辑——加参数靠默认值改语义靠新任务名 灰度开关契约演进全程不破坏存量调用方。3.3 可能遇到的坑及解决方法坑现象解决delay(eta...)报unexpected keyword argumentdelay 不支持选项eta 被当任务参数换apply_async(eta...)linksend_push.delay不生效传的是 AsyncResult 而非 signaturelinksend_push.s(...)expires300后任务仍执行Broker 堆积消息早就投递了expires 只在 Worker 侧生效治本靠队列治理第 21 章send_task拼错任务名消息成功入队但永远没人执行NotRegistered任务名常量化TASKS.SEND_SMS orders.send_order_sms传对象在测试环境正常、生产报错测试用 pickle、生产用 JSON统一序列化上线前跑_check_json_serializable自检3.4 完整代码清单与测试验证清单order_service.py唯一出口 校验、order_tasks.py任务 Worker 侧兜底、call_options.py、send_task_demo.py。测试验证# tests/test_contract.pyimportpytestfromorder_serviceimportOrderServicefromorder_tasksimportapp,send_order_sms app.conf.task_always_eagerTrue# 同步执行验证契约逻辑deftest_place_rejects_bad_order_id():withpytest.raises(ValueError):OrderService.place(-1)withpytest.raises(ValueError):OrderService.place(abc)deftest_place_rejects_unserializable_param():withpytest.raises(ValueError):OrderService.place(100,mobileobject())deftest_call_option_overrides_task_option():# 任务定义在默认队列调用时指定 queue 选项应胜出sigsend_order_sms.apply_async(args[1],queuesms)assertsig.options[queue]smsdeftest_delay_is_sugar():# delay 内部就是 apply_async(args, kwargs)rsend_order_sms.delay(1)assertr.task_idisnotNonepython-mpytest tests/test_contract.py-v# 4 passed4. 项目总结4.1 优点 缺点维度统一出口 显式契约OrderService随处裸 delay参数安全投递前校验非法参数当场报错序列化失败到运行时才炸可审计所有投递点收敛一处改选项只改一处四处改四处漏一致性expires/link 统一默认各写各的缺点 1多一层封装简单场景略啰嗦直接缺点 2出口成为单点改动影响面大——4.2 适用场景适用① 多团队共享任务契约必先于实现② 需要统一超时/告警/路由的订单域③ 跨语言系统投递send_task 任务名常量④ 管理后台按名字触发任务⑤ 任务参数需要灰度演进v1/v2 共存期。不适用① 脚本级一次性投递裸 delay 够用② 任务与调用方同进程且不会演进的内部工具③ 尚未定型、契约一天三改的探索期任务先放开 delay稳定后再收口。4.3 注意事项delay不能带任何选项带选项一律apply_async。expires/countdown单位是秒eta传datetime注意时区第 12 章。link 回调默认接收上一个任务的返回值作第一个参数签名要留好位置。参数契约三原则只传 ID 不传对象新参数带默认值语义变化换任务名v2。任务名与参数契约是「发布出去的 API」一旦有跨团队调用方就要像维护 HTTP API 一样维护它版本、文档、废弃窗口。4.4 常见踩坑经验3 个生产故障故障大促时任务参数序列化大面积失败。根因有人传 ORM 对象测试环境 pickle 掩盖了问题。对策统一 JSON 投递前校验本章落地。教训契约校验要发生在投递时刻不是消费时刻。故障delay(eta...)静默丢参数任务在错误时间执行。根因eta 进了**kwargs成为业务参数业务里没校验就忽略了。对策代码评审禁用delay(选项)写法。教训语法糖掩盖不了语义差异。故障回调任务把失败告警发成了短信轰炸。根因link_error的告警任务参数写错错误分支反而触发正常流程。对策回调任务单独队列、单独幂等键第 9 章落地。教训回调与主任务要有隔离。4.5 思考题send_task(orders.send_order_sms, args[1])与send_order_sms.delay(1)的消息在 Broker 里有什么不同为什么说 send_task 更「解耦」也更「危险」任务参数契约要随版本演进如短信任务要加channel参数怎么加才不破坏现有调用方提示默认值、位置参数数量、任务名策略答案见第 7 章开头的「上一章思考题参考答案」。延伸阅读与资源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 实战修炼与源码剖析