news 2026/8/25 4:53:24

Celery系列-04-Canvas任务编排

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Celery系列-04-Canvas任务编排

文章目录

  • 为什么需要任务编排?
  • 什么是 Signature?
  • `.s()` 与 `.si()` 有什么区别?
    • `.s()`:允许接收前一步结果
    • `.si()`:不可变 Signature
  • Chain:顺序执行
  • Group:并行执行
  • Group 回调为什么容易误用?
  • Chord:并行完成后汇总
  • 使用 Chord 的三个重要限制
    • 必须有可用的 Result Backend
    • 参与任务不能忽略结果
    • 同步有成本
  • Chord 中有任务失败会怎样?
  • 将 Chain、Group 和 Chord 组合
  • Map 与 Starmap
    • Map
    • Starmap
  • Chunks:控制任务粒度
  • 工作流中的数据应该怎么传?
  • 工作流的可观测性
  • 工作流设计的常见错误
    • 在任务里 `.get()` 等待子任务
    • 让任务之间传递巨大结果
    • 把数百万条小工作拆成数百万个任务
    • 假设 Group 的普通回调只执行一次
    • 忽略重复执行
    • 把 Canvas 当成永久工作流引擎
  • 本篇总结
  • 下一篇
  • 参考资料

为什么需要任务编排?

真实业务很少只有一个孤立任务。以“生成并发送月度报表”为例:

  1. 查询多个业务模块的数据;
  2. 并行计算各模块指标;
  3. 汇总所有结果;
  4. 生成 Excel 或 PDF;
  5. 上传对象存储;
  6. 发送通知。

最直接但错误的写法,是让一个 Celery 任务发送其他任务并调用.get()等待:

@app.taskdefbuild_report():result=calculate.delay()data=result.get()# 不推荐returnrender(data)

这会占着一个 Worker 执行槽等待另一个 Worker,容易降低吞吐,甚至在并发资源不足时形成死锁。

Celery Canvas 提供声明式工作流:把“先做什么、哪些并行、何时汇总”表达成任务图,由 Celery 调度,而不是让任务阻塞等待。

什么是 Signature?

Signature 是“一次任务调用”的可序列化描述,包含:

  • 任务名称;
  • 位置参数和关键字参数;
  • 倒计时、队列、过期时间等执行选项;
  • 回调、错误回调等关系。
sig=add.s(2,3)

此时并未执行任务。可以查看或修改后再发送:

sig.apply_async(countdown=10)

也可以使用完整写法:

sig=add.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():不可变 Signature

notify.si("report completed")

不可变 Signature 不接收前一步传来的参数,只使用创建时给定的参数。等价的完整写法:

notify.signature(args=("report completed",),immutable=True,)

当任务只是“流程完成后发通知”,并不关心前一步返回值时,使用.si()更安全。

Chain:顺序执行

定义任务:

fromceleryimportCelery app=Celery("demo")@app.taskdefadd(x:int,y:int)->int:returnx+y@app.taskdefmultiply(x:int,y:int)->int:returnx*y

构建 Chain:

fromceleryimportchain workflow=chain(add.s(2,3),# 5multiply.s(10),# multiply(5, 10) = 50add.s(1),# add(50, 1) = 51)result=workflow.apply_async()print(result.get(timeout=10))

也可以使用管道操作符:

workflow=add.s(2,3)|multiply.s(10)|add.s(1)result=workflow.delay()

适合场景:

  • 数据提取 → 转换 → 加载;
  • 生成文件 → 上传 → 通知;
  • 创建订单 → 风控检查 → 分配履约;
  • 多阶段模型推理。

Chain 的失败通常会阻止后续普通步骤继续执行。需要显式设计失败记录、补偿或错误回调。

Group:并行执行

fromceleryimportgroup job=group(add.s(1,1),add.s(2,2),add.s(3,3),)result=job.apply_async()print(result.get(timeout=10))

输出:

[2, 4, 6]

Group 返回GroupResult,常用方法:

result.ready()result.successful()result.failed()result.completed_count()result.get(timeout=10)result.revoke()

动态并行:

job=group(add.s(i,i)foriinrange(100))job.apply_async()

Group 适合彼此独立的任务。并行度仍受 Worker 数量、并发数、队列和下游容量限制。一次发送一百万个极小任务,消息开销可能比计算本身更高。

Group 回调为什么容易误用?

Group 不是一个真正执行汇总逻辑的普通任务。把回调直接link到 Group,可能不会得到预期的“所有任务完成后只执行一次”,错误回调也可能因多个子任务失败而被调用多次。

需要“全部完成后汇总”时,应使用 Chord,而不是依赖 Group 的普通链接行为。

Chord:并行完成后汇总

Chord = Group + Callback。

fromceleryimportchord@app.taskdeftotal(values:list[int])->int:returnsum(values)result=chord((add.s(i,i)foriinrange(10)),total.s(),).apply_async()print(result.get(timeout=20))

执行过程:

  1. 十个add任务并行执行;
  2. Celery 等待全部完成;
  3. 将结果组成列表;
  4. 调用total(results)
  5. 最终返回汇总结果。

适合场景:

  • 分片计算后汇总;
  • 多来源数据抓取后生成报告;
  • 多张图片处理后打包;
  • 多个检测任务完成后做统一判断。

使用 Chord 的三个重要限制

必须有可用的 Result Backend

Chord 需要知道每个并行任务何时完成并读取结果。没有 Result Backend 无法正常汇总。

RPC Result Backend 不支持 Chord。选择 Backend 时要确认当前 Celery 版本的具体支持情况。

参与任务不能忽略结果

如果全局配置:

task_ignore_result=True

Chord 中的任务要显式保留结果:

@app.task(ignore_result=False)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))workflow=chord([calculate.s(i)foriinrange(10)],total.s().on_error(on_workflow_error.s()),)workflow.apply_async()

错误回调自身应幂等。复杂 Group 或 Chord 中的失败传播细节要结合所用 Celery 版本测试,而不是只依赖直觉。

将 Chain、Group 和 Chord 组合

报表流程:

fromceleryimportchain,chord workflow=chain(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 也会接收上游结果。

如果某一步不需要上游结果:

workflow=prepare.s()|cleanup.si(resource_id)

选择.s()还是.si(),本质上是在设计数据流。

Map 与 Starmap

Map

normalize.map([" A "," B "," C "])

类似:

[normalize(item)foriteminitems]

但 Map 通常只发送一个任务消息,由一个任务顺序处理参数列表;它与 Group 的多个并行任务不同。

Starmap

函数接受多个位置参数时:

add.starmap([(1,2),(3,4),(5,6),])

类似:

[add(*args)forargsinitems]

Chunks:控制任务粒度

如果有一百万条记录,为每条记录创建一个消息可能产生巨大调度开销。Chunks 可以把数据分块:

items=zip(range(1000),range(1000))job=add.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(bind=True)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 适合任务编排,但对持续数天、需要人工审批、复杂补偿、版本化状态机的流程,应评估专门工作流引擎。

本篇总结

  1. Signature 是 Celery 工作流的基本构件;
  2. .s()接收上游结果,.si()忽略上游参数;
  3. Chain 表达顺序依赖;
  4. Group 表达相互独立的并行任务;
  5. Chord 表达并行完成后的统一汇总;
  6. Chord 需要兼容的 Result Backend,且任务不能忽略结果;
  7. Map、Starmap 和 Chunks 用于控制批量任务的表达和粒度;
  8. 不要让 Celery 任务阻塞等待其他任务;
  9. 复杂工作流需要独立的业务状态、幂等性和可观测性设计。

下一篇

系列最后一篇将把 Celery 带入生产环境:Beat 定时任务、队列路由、长短任务隔离、Worker 并发、Flower 监控、安全配置、上线检查和 GitHub 源码阅读路线。

参考资料

  • Canvas: Designing Work-flows
  • Calling Tasks
  • celery/celery:canvas.py
  • celery/celery:result.py
版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/25 4:53:08

FigmaCN 中文插件安装教程:3 分钟完成 Figma 界面汉化

FigmaCN 中文插件安装教程:3 分钟完成 Figma 界面汉化 【免费下载链接】figmaCN 中文 Figma 插件,设计师人工翻译校验 项目地址: https://gitcode.com/gh_mirrors/fi/figmaCN 第一次打开 Figma 做设计时,很多设计师会停下来数一数界面…

作者头像 李华
网站建设 2026/8/25 4:52:12

GEO优化:全域流量服务的实际效果与行业应用分析

引言随着人工智能技术的发展,生成式引擎优化(GEO)逐渐成为企业数字化营销的新趋势。区别于传统的搜索引擎优化(SEO),GEO通过构建标准化、可校验、可迭代的数字档案,使企业在大模型的认知体系中获…

作者头像 李华
网站建设 2026/8/25 4:47:31

计算机考研复试机试准备指南与核心算法解析

1. 复试机试准备的核心要点 复试机试是计算机相关专业研究生选拔的重要环节,通常考察编程能力、算法基础和计算机基础知识。不同于初试的理论考核,机试更注重实际动手能力和问题解决能力。根据我的经验,有效的机试准备需要系统性地覆盖以下几…

作者头像 李华
网站建设 2026/8/25 4:44:03

AI Agent驱动燃气行业效率革命:从被动响应到主动智防

1. 从“跑断腿”到“动动脑”:燃气行业的效率困局与AI破局点干了十几年能源信息化,我见过太多燃气公司的运维现场:调度中心里电话此起彼伏,一线巡检员顶着风雨在管线旁记录数据,抢修队像救火队员一样四处奔波。整个行业…

作者头像 李华
网站建设 2026/8/25 4:43:00

GEO系统选型:为什么缺失闭环能力的工具,后期维护成本更高?

很多企业在探索品牌在生成式AI搜索中的可见性时,容易陷入一个选型悖论:初始采购门槛低的工具,为何在后续运营中反而显得愈发“昂贵”?大多数采购决策者在初期会将软件订阅费或买断价格作为核心评价维度。然而,在AI搜索…

作者头像 李华