Celery系列-04-Canvas任务编排

📅 2026/8/25 4:53:34
Celery系列-04-Canvas任务编排
文章目录为什么需要任务编排什么是 Signature.s() 与 .si() 有什么区别.s()允许接收前一步结果.si()不可变 SignatureChain顺序执行Group并行执行Group 回调为什么容易误用Chord并行完成后汇总使用 Chord 的三个重要限制必须有可用的 Result Backend参与任务不能忽略结果同步有成本Chord 中有任务失败会怎样将 Chain、Group 和 Chord 组合Map 与 StarmapMapStarmapChunks控制任务粒度工作流中的数据应该怎么传工作流的可观测性工作流设计的常见错误在任务里 .get() 等待子任务让任务之间传递巨大结果把数百万条小工作拆成数百万个任务假设 Group 的普通回调只执行一次忽略重复执行把 Canvas 当成永久工作流引擎本篇总结下一篇参考资料为什么需要任务编排真实业务很少只有一个孤立任务。以“生成并发送月度报表”为例查询多个业务模块的数据并行计算各模块指标汇总所有结果生成 Excel 或 PDF上传对象存储发送通知。最直接但错误的写法是让一个 Celery 任务发送其他任务并调用.get()等待app.taskdefbuild_report():resultcalculate.delay()dataresult.get()# 不推荐returnrender(data)这会占着一个 Worker 执行槽等待另一个 Worker容易降低吞吐甚至在并发资源不足时形成死锁。Celery Canvas 提供声明式工作流把“先做什么、哪些并行、何时汇总”表达成任务图由 Celery 调度而不是让任务阻塞等待。什么是 SignatureSignature 是“一次任务调用”的可序列化描述包含任务名称位置参数和关键字参数倒计时、队列、过期时间等执行选项回调、错误回调等关系。sigadd.s(2,3)此时并未执行任务。可以查看或修改后再发送sig.apply_async(countdown10)也可以使用完整写法sigadd.signature(args(2,3),options{queue:math},)Signature 是 Canvas 的基础。Chain、Group 和 Chord 都由 Signature 组合而成。.s()与.si()有什么区别.s()允许接收前一步结果add.s(10)在 Chain 中如果上一步返回5它会被调用为add(5,10)前一步结果会作为额外位置参数放在已有参数之前。.si()不可变 Signaturenotify.si(report completed)不可变 Signature 不接收前一步传来的参数只使用创建时给定的参数。等价的完整写法notify.signature(args(report completed,),immutableTrue,)当任务只是“流程完成后发通知”并不关心前一步返回值时使用.si()更安全。Chain顺序执行定义任务fromceleryimportCelery appCelery(demo)app.taskdefadd(x:int,y:int)-int:returnxyapp.taskdefmultiply(x:int,y:int)-int:returnx*y构建 Chainfromceleryimportchain workflowchain(add.s(2,3),# 5multiply.s(10),# multiply(5, 10) 50add.s(1),# add(50, 1) 51)resultworkflow.apply_async()print(result.get(timeout10))也可以使用管道操作符workflowadd.s(2,3)|multiply.s(10)|add.s(1)resultworkflow.delay()适合场景数据提取 → 转换 → 加载生成文件 → 上传 → 通知创建订单 → 风控检查 → 分配履约多阶段模型推理。Chain 的失败通常会阻止后续普通步骤继续执行。需要显式设计失败记录、补偿或错误回调。Group并行执行fromceleryimportgroup jobgroup(add.s(1,1),add.s(2,2),add.s(3,3),)resultjob.apply_async()print(result.get(timeout10))输出[2, 4, 6]Group 返回GroupResult常用方法result.ready()result.successful()result.failed()result.completed_count()result.get(timeout10)result.revoke()动态并行jobgroup(add.s(i,i)foriinrange(100))job.apply_async()Group 适合彼此独立的任务。并行度仍受 Worker 数量、并发数、队列和下游容量限制。一次发送一百万个极小任务消息开销可能比计算本身更高。Group 回调为什么容易误用Group 不是一个真正执行汇总逻辑的普通任务。把回调直接link到 Group可能不会得到预期的“所有任务完成后只执行一次”错误回调也可能因多个子任务失败而被调用多次。需要“全部完成后汇总”时应使用 Chord而不是依赖 Group 的普通链接行为。Chord并行完成后汇总Chord Group Callback。fromceleryimportchordapp.taskdeftotal(values:list[int])-int:returnsum(values)resultchord((add.s(i,i)foriinrange(10)),total.s(),).apply_async()print(result.get(timeout20))执行过程十个add任务并行执行Celery 等待全部完成将结果组成列表调用total(results)最终返回汇总结果。适合场景分片计算后汇总多来源数据抓取后生成报告多张图片处理后打包多个检测任务完成后做统一判断。使用 Chord 的三个重要限制必须有可用的 Result BackendChord 需要知道每个并行任务何时完成并读取结果。没有 Result Backend 无法正常汇总。RPC Result Backend 不支持 Chord。选择 Backend 时要确认当前 Celery 版本的具体支持情况。参与任务不能忽略结果如果全局配置task_ignore_resultTrueChord 中的任务要显式保留结果app.task(ignore_resultFalse)defcalculate(partition_id:int)-int:...同步有成本Chord 需要追踪一组任务是否全部完成。不同 Backend 的实现不同但都存在状态写入和同步成本。不要把几微秒的小计算拆成成千上万个 Chord 子任务。Chord 中有任务失败会怎样如果 Header 中某个任务失败Chord 的最终结果会进入失败状态通常表现为ChordError。其他已经发送的并行任务并不会因此自动取消。可以为最终回调配置错误处理app.taskdefon_workflow_error(request,exc,traceback):log_failure(request.id,str(exc))workflowchord([calculate.s(i)foriinrange(10)],total.s().on_error(on_workflow_error.s()),)workflow.apply_async()错误回调自身应幂等。复杂 Group 或 Chord 中的失败传播细节要结合所用 Celery 版本测试而不是只依赖直觉。将 Chain、Group 和 Chord 组合报表流程fromceleryimportchain,chord workflowchain(prepare_report.s(report_id),chord([calculate_sales.s(),calculate_inventory.s(),calculate_refunds.s(),],merge_report.s(),),upload_report.s(),send_notification.s(user_id),)workflow.apply_async()需要仔细检查参数传递。Chain 会把前一步结果传入下一步Group 中的可变 Signature 也会接收上游结果。如果某一步不需要上游结果workflowprepare.s()|cleanup.si(resource_id)选择.s()还是.si()本质上是在设计数据流。Map 与 StarmapMapnormalize.map([ A , B , C ])类似[normalize(item)foriteminitems]但 Map 通常只发送一个任务消息由一个任务顺序处理参数列表它与 Group 的多个并行任务不同。Starmap函数接受多个位置参数时add.starmap([(1,2),(3,4),(5,6),])类似[add(*args)forargsinitems]Chunks控制任务粒度如果有一百万条记录为每条记录创建一个消息可能产生巨大调度开销。Chunks 可以把数据分块itemszip(range(1000),range(1000))jobadd.chunks(items,50)job.apply_async()这会将一千组参数分成每块五十组减少消息数量。任务粒度需要权衡太小消息、序列化和调度开销过大太大并行度不足、单次失败重做成本高合适单任务耗时足以覆盖调度成本同时便于重试和扩展。工作流中的数据应该怎么传不要在消息中传输ORM 对象打开的文件句柄数据库连接超大二进制文件无法稳定 JSON 序列化的自定义对象。建议传输数据库主键对象存储地址版本号小型 JSON 数据业务幂等键和追踪 ID。例如不要返回整份百兆报表给下一任务而应上传到对象存储并返回{object_key:reports/2026/08/report-123.xlsx,checksum:...,version:1,}工作流的可观测性复杂工作流不仅要记录单个任务 ID还应记录root_id工作流根任务parent_id父任务业务流程 ID输入数据版本每一步状态和耗时最终产物地址失败步骤与补偿状态。可以在任务内访问app.task(bindTrue)defstep(self,payload):logger.info(task_id%s root_id%s parent_id%s,self.request.id,self.request.root_id,self.request.parent_id,)仅依赖AsyncResult不足以构建长期业务审计。关键流程应在业务数据库中保存可理解的状态。工作流设计的常见错误在任务里.get()等待子任务这会浪费 Worker 并发槽。用 Chain、Group 或 Chord 表达依赖。让任务之间传递巨大结果会增加 Broker 或 Backend 压力。改为传引用。把数百万条小工作拆成数百万个任务调度成本可能压垮系统。使用 Chunks 或合理批量。假设 Group 的普通回调只执行一次需要全量汇总时使用 Chord并验证错误行为。忽略重复执行工作流中的每一个有副作用步骤都应幂等。某一步重试时上游成功并不意味着下游不会重复。把 Canvas 当成永久工作流引擎Celery Canvas 适合任务编排但对持续数天、需要人工审批、复杂补偿、版本化状态机的流程应评估专门工作流引擎。本篇总结Signature 是 Celery 工作流的基本构件.s()接收上游结果.si()忽略上游参数Chain 表达顺序依赖Group 表达相互独立的并行任务Chord 表达并行完成后的统一汇总Chord 需要兼容的 Result Backend且任务不能忽略结果Map、Starmap 和 Chunks 用于控制批量任务的表达和粒度不要让 Celery 任务阻塞等待其他任务复杂工作流需要独立的业务状态、幂等性和可观测性设计。下一篇系列最后一篇将把 Celery 带入生产环境Beat 定时任务、队列路由、长短任务隔离、Worker 并发、Flower 监控、安全配置、上线检查和 GitHub 源码阅读路线。参考资料Canvas: Designing Work-flowsCalling Taskscelery/celerycanvas.pycelery/celeryresult.py