第20章:Celery Canvas 进阶——chunks、嵌套、回调与 Stamping

发布时间:2026/9/4 22:40:37
第20章:Celery Canvas 进阶——chunks、嵌套、回调与 Stamping 0. 上一章思考题参考答案思考题 1不一样。chain(a|b)的消息数量更多chain 的每一步会被拆成独立消息逐步投递还有「链式回调」的隐藏消息且中间结果会被写进 Backend供下一步取用而手动a.delay()后b.delay()只有两条消息、结果不自动传递。所以 chain 的代价是「多几次消息 Backend 写入」换的是可观察、可重放。思考题 2header 中一张海报失败后其余任务继续完成计数器永远到不了 N失败的没计数→ body 不被触发 → 直到result_chord_join_timeoutcelery/app/defaults.py默认 3 秒轮询间隔的 join 超时后chord 进入失败路径ChordError。感知环节 join 超时轮询celery.chord_unlock 解锁任务第 36 章会读它的源码。生产上应配置 errback让失败在「计数层」就显形而不是等超时。1. 项目背景促销团队要发百万张优惠券小周准备用第 19 章的group一把梭——group(send_coupon.s(uid) for uid in million_uids)。大师看了一眼代码直接叫停「你这一个 group 就是100 万条消息同时压进 BrokerRedis 直接内存告警RabbitMQ 的交换机都要吐了更别说百万个 AsyncResult 的结果键Backend 先原地爆炸。」小周这才意识到group 是把「一批任务」当对象但消息量不会因为它是 group 就变少。百万条消息就是百万条消息Broker 要存、Worker 要消费、Backend 要记账。第二个问题接踵而至运营说「优惠券发错了批次要把 7 月这批撤销」——小周翻遍代码无法定位「哪些任务属于 7 月这批」因为任务消息上没有任何批次标识。撤销只能靠「任务名 时间窗」猜误伤了一堆别的任务。两个新需求 ① 百万级任务不能一条消息一个任务 → chunks 分片1000 人一片 ② 批次可治理任务要能「打标签」→ Stampingstamp 批次号 ③ 复杂图chain 里嵌 group结果怎么展开 → 嵌套 canvas ④ 部分失败一张失败不能拖垮整批 → errbacks / 部分失败策略本章目标用chunks把百万发放拆成 1000 人一片用stamp给画布节点打批次标支持按 stamp 撤销顺带把嵌套 canvas 与 errback 部分失败策略补齐。2. 项目设计场景大师在百万级任务上线前把四个坑逐一摆出来。小胖一百万条消息怎么了Redis 内存不是挺大的吗一百万也就几 GB忍忍就过去了而且一条消息一个任务多直观小白小胖你这个「忍忍」就是大促崩盘的剧本。百万条消息一次性压进 BrokerRedis 内存暴涨、单队列堆积几十分钟、Backend 百万个结果键第 8 章刚治过celery-task-meta-*。我想问chunks 到底是「消息变少了」还是「只是包装变了」1000 人一片到底发出去多少条消息大师chunks 是把 N 个任务「打包」成一条消息——chunks(1000 人)发出去的不是 1000 条消息而是1 条「包含 1000 个子任务的分片消息」。Worker 收到分片后在进程内循环执行这 1000 个子任务顺序或按配置。所以百万用户 1000 片 1000 条消息Broker 压力直接除以 1000。代价① 一片内的任务没有跨进程并行一个子进程串行跑 1000 个单片的吞吐受单个进程限制——1000 人 × 0.1s 100 秒/片靠多片并行摊平② 一片内一旦崩溃整片未执行的子任务要重投粒度变粗幂等重要性升级。xmap/xstarmap是 chunks 的兄弟把「一个参数列表映射到同一任务」的语法糖。技术映射chunks 把 1000 张菜票钉成一沓后厨一次拿一沓逐个做group 1000 张菜票散着递进窗口传菜员一次拿一张。小白那 Stamping 呢我查了examples/stamping/stamp 是给消息打标签——它和「任务参数里带个 batch_id」有什么区别大师本质区别在作用面与用途。参数里的 batch_id 是「业务字段」任务执行时才知道stamp 是消息元数据header 层面作用在「调度/撤销/监控」这些框架层revoke可以按 stamp 精确撤销一批、事件流可以按 stamp 过滤统计、日志可以按 stamp 检索。一句话参数是给任务函数看的stamp 是给框架和运维看的。给画布打 stamp 的姿势canvas.stamp(promo_batch, batch_id)——它会递归给画布里的所有节点都盖上章第 20 章源码在celery/canvas.py的 stamp 方法。小胖还有那个「chain 里嵌 group」我上次链里放了个 group结果下游收到的参数奇奇怪怪的是咋回事大师这是嵌套 canvas 的经典困惑group 在 chain 里默认是「扁平化」的——chain 里的 group 结果会被展开成多个参数调用下游f(g(1),g(2))变成 f 收到列表而 chord 在 chain 里则保持「一个结果」。规则一句话chain 里嵌 group 广播式展开chain 里嵌 chord 汇总式传递。写之前想清楚下游要「一个个」还是「一坨」。技术映射嵌套画布 流水线的「并线」——group 像「多车并线后分开进不同匝道」chord 像「多车并线后汇入主路走同一出口」。匝道和出口不同接法就不同。3. 项目实战3.1 环境准备沿用环境Redis Broker Backend。参考examples/stamping/与examples/resultgraph/。3.2 分步实现步骤 1用 chunks 发放百万优惠券1000 人一片目标消息量从 100 万降到 1000压测可观测。# coupon_tasks.pyimporttimefromceleryimportCelery appCelery(coupon,brokerredis://localhost:6379/0,backendredis://localhost:6379/1)app.task(namepromo.send_coupon,bindTrue)defsend_coupon(self,uid:int)-str:发一张优惠券模拟 0.05s。time.sleep(0.05)returnfcoupon-{uid}# 一百万用户的发放计划uidslist(range(1,1_000_001))# 不推荐group 一把梭 → 100 万条消息# app.group(send_coupon.s(u) for u in uids).apply_async()# 推荐chunks 1000 人一片 → 1000 条消息chunkedsend_coupon.chunks(zip(uids),1000)# chunks(参数列表, 每片大小)resultchunked.apply_async()print(片数:,len(result),即消息数)运行结果文字描述chunked.apply_async()投递的消息数 1000每片一条内含 1000 个子任务len(result)返回 1000 个 AsyncResult 的列表Worker 日志按片执行子任务每片 1000×0.05s≈50s 串行1000 片并行摊平。对比 group 方案的 100 万条消息Broker 压力三个数量级下降。步骤 2给画布打 stamp批次标签目标给发放任务盖上「批次章」为「按批次撤销」做准备。# stamp_demo.pyfromcoupon_tasksimportapp,send_coupon canvasapp.group(send_coupon.s(uid)foruidin[1,2,3])canvas.stamp(promo_batch,batch-2026-07)# 递归盖满所有节点rescanvas.apply_async()print(已打标批次: batch-2026-07任务数:,len(res))检查 stamp 是否生效用--logleveldebug看 Worker 日志消息头里出现stamps: {promo_batch: batch-2026-07}或抓一条消息对照第 35 章消息协议。步骤 3按 stamp 撤销某批次目标撤销「7 月批次」而不影响其他批次——这就是 stamp 的框架级价值。# revoke_by_stamp.pyfromceleryimportsignaturefromcelery.resultimportAsyncResultfromcoupon_tasksimportapp# 方式一按 stamp 撤销需要记录 stamp 对应的任务 ID 列表或事件流过滤revoked[]forrinresult:# 步骤 2 返回的结果列表ifr.statePENDING:r.revoke()# 未开始的任务撤销revoked.append(r.id)print(已撤销任务数:,len(revoked))# 方式二生产推荐事件流/监控层按 stamp 过滤后批量 revoke第 25 章运行结果文字描述处于 PENDING 的任务被撤销Worker 消费时发现撤销集合跳过执行并标记 REVOKED第 10 章已开始执行的不受影响——「批次治理」从猜时间窗升级为按标签精确打击。步骤 4嵌套 canvas——chain 里嵌 group 的展开规则目标用实验验证「chain 里嵌 group 结果展开」避免参数错位。# nest_demo.pyfromcoupon_tasksimportapp,send_coupon# chain 里嵌 group下游收到「展开后的多个参数」还是「一个列表」capp.chain(app.group(send_coupon.s(1),send_coupon.s(2)),summarize.s(),# 观察 summarize 收到什么)app.task(namepromo.summarize)defsummarize(*args,**kwargs):print(收到参数:,args,kwargs)returnsummarized运行结果文字描述summarize打印收到参数: (coupon-1, coupon-2)——group 在 chain 里被展开成两个位置参数若想收到列表用 chord第 19 章步骤 3或immutable组合。结论嵌套前先问「下游要一坨还是一个个」。步骤 5部分失败策略——errback 与 body 失败兜底目标一张失败不拖垮整批失败可感知、可重放。# partial_fail_demo.pyfromcoupon_tasksimportapp,send_couponapp.task(namepromo.on_partial_fail,bindTrue)defon_partial_fail(self,request,exc,traceback):errback收到失败回调把失败的子任务捞出来。print(f失败任务:{request}, 异常:{exc})returnloggedgapp.group(send_coupon.s(u)foruin[1,2,3])g.link_error(on_partial_fail.s())# 组级 errbackresg.apply_async()运行结果文字描述3 个任务全成功时无 errback 触发把其中一个参数改成触发异常如 uid0 触发 ValueErrorerrback 收到失败的子任务引用与异常可据此重放——「部分失败」不再黑盒。chord 的 body 失败同样走link_error第 19 章思考题 2 的落地。步骤 6chunks 与 xmap/xstarmap 的快速对照目标分清三个「批量语法糖」的适用姿势避免写错参数形态。# map_demo.pyfromcoupon_tasksimportapp,send_coupon# chunks参数列表分片返回「每片一个 AsyncResult」c1send_coupon.chunks(zip([1,2,3,4,5,6]),2)# 3 片# xmap单参数映射chunks 的别名语义更清晰c2send_coupon.xmap([1,2,3,4,5,6],chunk_size2)# 等价 3 片# xstarmap参数对映射每个子任务多个参数c3send_coupon.xstarmap([(1,),(2,)],chunk_size1)# 2 片print(c1 片数:,len(c1),| c2 片数:,len(c2),| c3 片数:,len(c3))运行结果文字描述三个写法都返回「片数 参数列表长度 / chunk_size」的结果列表xmap 适合单参数批量、xstarmap 适合多参数批量——写批量任务前先确认任务签名是(uid)还是(uid, channel)选错参数形态会静默错位。3.3 可能遇到的坑及解决方法坑现象解决group 一把梭打爆 Broker百万条消息一次性压入换 chunks消息数 ÷ 片大小chunks 一片崩溃整片重投片内任务未执行的部分重投片内子任务也要幂等片大小按「重投成本」权衡stamp 后 revoke 不生效撤销的是任务 ID 而非 stamp 本身stamp 是元数据撤销仍需映射到任务 ID用事件流按 stamp 过滤第 25 章chain 嵌 group 参数错位下游收到展开参数而非列表想收列表用 chord想广播用 chaingrouperrback 不触发link_error挂错层级组级挂 group、链级挂 chain参数契约确认errback 签名3.4 完整代码清单与测试验证清单coupon_tasks.pystamp_demo.pyrevoke_by_stamp.pynest_demo.pypartial_fail_demo.py。官方参考examples/stamping/revoke_example.py按 stamp 撤销的完整示例。测试验证# tests/test_canvas_advanced.pyfromcoupon_tasksimportapp,send_coupon app.conf.task_always_eagerTruedeftest_chunks_reduce_message_count():chunkedsend_coupon.chunks(zip(range(1,100_001)),1000)assertlen(chunked)100# 10 万用户 → 100 片消息数deftest_stamp_applies_to_nodes():canvasapp.group(send_coupon.s(1),send_coupon.s(2))canvas.stamp(promo_batch,b1)fornodeincanvas:assertnode.options.get(stamps,{}).get(promo_batch)b1\orpromo_batchinstr(node)# 打标递归生效deftest_nested_group_flattens():# chain 中 group 结果展开为位置参数契约验证captured{}fromcoupon_tasksimportappas_app_app.task(namepromo.capture)defcapture(*args,**kwargs):captured[args]argsreturnok_app.chain(_app.group(send_coupon.s(1)),capture.s()).apply()assertlen(captured[args])1python-mpytest tests/test_canvas_advanced.py-v# 3 passed4. 项目总结4.1 优点 缺点维度chunks分片批量group 一把梭消息量÷片大小1000 人/片1 人 1 条Broker 压力低打爆风险并行粒度片级并行片内串行任务级全并行重投粒度整片重投粗单条重投细适用百万级幂等批量百级以内精细控制4.2 适用场景适用① 百万级群发优惠券/通知/对账② 需要批次治理按 stamp 撤销/统计的场景③ 复杂流程需要嵌套编排chain 嵌 group/chord④ 需要「部分失败可感知可重放」的业务⑤ 批量任务需要可观测片级进度的运营场景。不适用① 片内任务必须独立并发且量级在百级group 更合适② 任务不可幂等chunks 整片重投会放大风险③ 需要严格实时进度片内串行执行进度粒度是片④ 参数形态复杂的异构批量先统一参数契约再用 xstarmap。4.3 注意事项chunks 的片大小是「吞吐 vs 重投粒度」的权衡越大 Broker 越轻、单片越慢、重投越粗。stamp 是框架级元数据业务字段放参数、治理字段放 stamp别混。嵌套 canvas 先画图再写码group 展开 / chord 汇总写之前想清楚下游签名。errback 本身也是任务它有队列、重试、失败策略别当「免费回调」。4.4 常见踩坑经验3 个生产故障故障百万优惠券发放把 Redis 打爆触发内存告警。根因group 一把梭 100 万条消息。对策chunks 1000 人/片1000 条。教训消息量是 Broker 的生命线分片是批量任务的必修课。故障撤销「7 月批次」误伤了 8 月任务。根因用「时间窗」猜批次无标签。对策stamp 打批次章 按 stamp 过滤撤销。教训没有元数据的批量任务 无法治理的批。故障chain 嵌 group 后下游疯狂报参数错。根因group 结果展开规则没吃透下游签名对不上。对策契约测试3.4 的 flatten 用例。教训嵌套画布的隐性契约靠测试钉死。4.5 思考题chunks(1000 人/片)里如果第 500 个用户的任务抛业务异常片内剩余 500 个会怎样如果 Worker 进程崩溃呢提示异常 vs 进程死亡的重投语义第 18 章stamp 打在 group 上会递归到子任务如果给 chain 打 stamp撤销「整个 chain 未开始的节点」时行为和你预期一致吗提示第 36 章源码读 canvas.stamp答案见第 21 章开头的「上一章思考题参考答案」。延伸阅读与资源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 实战修炼与源码剖析