文章目录
- 为什么需要任务编排?
- 什么是 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 当成永久工作流引擎
- 本篇总结
- 下一篇
- 参考资料
为什么需要任务编排?
真实业务很少只有一个孤立任务。以“生成并发送月度报表”为例:
- 查询多个业务模块的数据;
- 并行计算各模块指标;
- 汇总所有结果;
- 生成 Excel 或 PDF;
- 上传对象存储;
- 发送通知。
最直接但错误的写法,是让一个 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))执行过程:
- 十个
add任务并行执行; - Celery 等待全部完成;
- 将结果组成列表;
- 调用
total(results); - 最终返回汇总结果。
适合场景:
- 分片计算后汇总;
- 多来源数据抓取后生成报告;
- 多张图片处理后打包;
- 多个检测任务完成后做统一判断。
使用 Chord 的三个重要限制
必须有可用的 Result Backend
Chord 需要知道每个并行任务何时完成并读取结果。没有 Result Backend 无法正常汇总。
RPC Result Backend 不支持 Chord。选择 Backend 时要确认当前 Celery 版本的具体支持情况。
参与任务不能忽略结果
如果全局配置:
task_ignore_result=TrueChord 中的任务要显式保留结果:
@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 适合任务编排,但对持续数天、需要人工审批、复杂补偿、版本化状态机的流程,应评估专门工作流引擎。
本篇总结
- Signature 是 Celery 工作流的基本构件;
.s()接收上游结果,.si()忽略上游参数;- Chain 表达顺序依赖;
- Group 表达相互独立的并行任务;
- Chord 表达并行完成后的统一汇总;
- Chord 需要兼容的 Result Backend,且任务不能忽略结果;
- Map、Starmap 和 Chunks 用于控制批量任务的表达和粒度;
- 不要让 Celery 任务阻塞等待其他任务;
- 复杂工作流需要独立的业务状态、幂等性和可观测性设计。
下一篇
系列最后一篇将把 Celery 带入生产环境:Beat 定时任务、队列路由、长短任务隔离、Worker 并发、Flower 监控、安全配置、上线检查和 GitHub 源码阅读路线。
参考资料
- Canvas: Designing Work-flows
- Calling Tasks
- celery/celery:canvas.py
- celery/celery:result.py