news 2026/9/23 14:44:47

Celery Canvas 工作流编排实战:Signature、Chain、Group 与 Chord 从入门到原理

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Celery Canvas 工作流编排实战:Signature、Chain、Group 与 Chord 从入门到原理
  • 人工智能
  • AI 应用
  • AI Agent

【免费下载链接】Tutorial-Codebase-Knowledge

Pocket Flow: Codebase to Tutorial

项目地址:https://gitcode.com/gh_mirrors/tu/Tutorial-Codebase-Knowledge
点击查看免费下载

导读

Celery 是 Python 生态中最流行的分布式任务队列框架之一,而Canvas(画布)是 Celery 面向复杂任务工作流提供的编排组件。它由Signature(签名)原语(Primitives)组成,可以让你把单个任务组合成「顺序执行、并行执行、并行后汇聚」的复杂流程,而无需在应用代码里手工管理依赖与结果传递。本文将围绕docs/Celery/08_canvas__signatures___primitives_.md的核心内容,从 Signature 的概念讲起,完整演示用chaingroupchord构建一个真实的文章处理流水线,并深入celery/canvas.py的源码结构,剖析这些原语在 Broker 与 Worker 之间是如何一步步被执行的。读完本文,你将能够独立设计、提交并排查自己的 Celery 工作流。

为什么需要 Canvas:任务编排的痛点

在 Chapter 3: Task 中我们学会了如何用@app.task定义任务,并用.delay()/.apply_async()把单个任务投递给 Broker。但真实业务几乎不会只有一个独立任务。考虑这样一个场景:用户上传一篇文章后,系统需要:

  1. 从 URL 抓取文章内容;
  2. 对文本做关键词提取;
  3. 对文本做语言检测;
  4. 等上述两步都完成后,把文章与元数据一并保存到数据库。

如果简单地逐个投递任务,你无法表达「步骤 2 与 3 可以并行、但都必须在步骤 1 之后」这种依赖关系,更无法保证「保存」在两者都成功之后才执行。若在应用代码里手工轮询AsyncResult、拼装参数,代码会迅速变得脆弱且难以维护。

Canvas 解决的正是任务之间的依赖与流程控制问题。它允许你把工作流的拓扑结构直接声明出来:哪个任务先跑、哪些可以并行、哪个任务必须等齐所有并行结果后再跑,然后把整个流程交给 Celery 去执行。官方文档用一个非常形象的比喻来描述 Canvas:就像不同形状的乐高积木——

  • 有些积木代表单个任务;
  • 有些积木把任务首尾相连(顺序执行);
  • 有些积木让任务并排堆放(并行执行);
  • 还有些积木可以构建「多个并行步骤必须全部完成,才能拼上下一块」的结构。

核心概念一:Signature——单个任务的"预约单"

什么是 Signature

一个Signature封装了调用某个任务所需的全部信息:

  • 任务的名称task);
  • 位置参数args);
  • 关键字参数kwargs);
  • 以及执行选项(如countdowneta、队列名queue等)。

它并不立即执行任务,只是保存了一份「如何执行」的计划。你可以把它想象成一张预填好的请求表单或一张菜谱卡片:拿着它可以随时下单,但下单动作本身是独立的。

创建 Signature 最快捷的方式是任务函数上的.s()快捷方法:

# tasks.py from celery_app import app # 假设 app 在 celery_app.py 中定义 @app.task def add(x, y): return x + y # 为 add(2, 3) 创建签名 add_sig = add.s(2, 3) # add_sig 此刻只保存了执行 add(2, 3) 的"计划" print(f"Signature: {add_sig}") print(f"Task name: {add_sig.task}") print(f"Arguments: {add_sig.args}") # 真正运行时,需要在该签名上调用 .delay() 或 .apply_async() # result_promise = add_sig.delay()

输出示例:

Signature: tasks.add(2, 3) Task name: tasks.add Arguments: (2, 3)

关于执行选项,签名同样支持传入。例如add.s(2, 3).apply_async(countdown=10)表示 10 秒后再执行;你还可以在创建签名时直接写入queuerouting_key等选项,让该签名固定投递到特定队列,这与 Chapter 3: Task 中apply_async的选项体系是一致的。

Signature 的三个重要性质

  1. 可序列化:Signature 本质是一个可被序列化的结构(在 Celery 源码中它继承自字典),因此可以随任务消息在网络上传输——这正是它能被嵌入link选项、跨 Worker 传递的前提。
  2. 部分应用(partial application):你可以在创建签名时不填满所有参数,留待链式执行时由前一个任务的结果自动补全。这在后面的chain中会频繁用到。
  3. 可克隆(clone):签名支持clone()生成副本,在prepare_steps等内部逻辑中,Celery 会不断对签名进行克隆与参数合并,避免污染原始定义。

核心概念二:工作流原语——连接积木的四种方式

Canvas 提供了若干**原语函数(Primitives)**用于把签名组合成工作流。其中最核心的是chaingroupchord,另外还有chunksxmapstarmap等补充原语。

chain:顺序执行

chain把多个签名按顺序串联:前一个任务的返回值会被作为第一个参数传给后一个任务

  • 类比:一条流水线,每个工位把产出交给下一个工位。
  • 语法:(sig1 | sig2 | sig3)chain(sig1, sig2, sig3)
from celery import chain workflow = chain(add.s(2, 2), add.s(4)) # add(2, 2) 的结果 4 会拼进 add(4, ...) # 等价写法(管道操作符): workflow2 = add.s(2, 2) | add.s(4)

group:并行执行

group把一组签名并发投递,返回一个GroupResult特殊结果对象用于追踪整组任务。

  • 类比:同时雇佣多个工人做彼此独立、相似的工作。
  • 语法:group(sig1, sig2, sig3)
from celery import group parallel = group(add.s(1, 1), add.s(2, 2), add.s(3, 3))

chord:并行执行后汇聚回调

chord由两部分组成:

  • header(头部):一个并行执行的group

  • body(主体):一个回调签名,在 header 中所有任务都成功完成后才执行,并接收 header 全部结果组成的列表作为参数。

  • 类比:一个研究团队分头完成项目不同部分,全部完成后由组长汇总所有发现撰写最终报告。

  • 语法:chord(group(header_sigs), body_sig)

from celery import chord workflow = chord( group(add.s(1, 1), add.s(2, 2)), # header:并行 some_callback.s() # body:等 header 全部完成后执行 )

补充原语

  • chunks:把一批参数分块执行(如add.chunks(zip(range(100), range(100)), 10)分成 10 个任务,每个任务处理 10 对参数);
  • xmap/starmap:对一个参数列表映射执行同一个任务(xmap传单参数、starmap传参数元组)。

chaingroupchord是构建工作流最基础的三个原语,掌握它们足以覆盖绝大多数编排需求。

实战:构建文章处理工作流

回到开头的文章处理场景,我们用 Canvas 一步步实现Fetch →(并行处理 A 与 B)→ Combine(汇聚保存)

第 1 步:定义基础任务

# tasks.py from celery_app import app import time import random @app.task def fetch_data(url): print(f"Fetching data from {url}...") time.sleep(1) # 模拟抓取数据 data = f"Content from {url} - {random.randint(1, 100)}" print(f"Fetched: {data}") return data @app.task def process_part_a(data): print(f"Processing Part A for: {data}") time.sleep(2) result_a = f"Keywords for '{data}'" print("Part A finished.") return result_a @app.task def process_part_b(data): print(f"Processing Part B for: {data}") time.sleep(3) # 模拟稍长的处理 result_b = f"Language for '{data}'" print("Part B finished.") return result_b @app.task def combine_results(results): # 'results' 是一个列表,包含 process_part_a 和 process_part_b 的返回值 print(f"Combining results: {results}") time.sleep(1) final_output = f"Combined: {results[0]} | {results[1]}" print(f"Final Output: {final_output}") return final_output

这里的fetch_data模拟抓取,process_part_a模拟关键词提取,process_part_b模拟语言检测,combine_results模拟最终入库。为了让并行效果可见,process_part_b故意比process_part_a多睡 1 秒。

第 2 步:用 Canvas 组装工作流

# run_workflow.py from celery import chain, group, chord from tasks import fetch_data, process_part_a, process_part_b, combine_results # 要处理的 URL article_url = "http://example.com/article1" # 创建工作流结构 # 1. fetch_data 抓取数据,结果传递给下一步。 # 2. 下一步是一个 chord: # - header:group 并行运行 process_part_a 与 process_part_b, # 两个任务都会收到 fetch_data 传来的 data。 # - body:combine_results 接收 group 全部结果的列表。 workflow = chain( fetch_data.s(article_url), # 第 1 步:抓取 chord( # 第 2 步:chord group(process_part_a.s(), process_part_b.s()), # header:并行处理 combine_results.s() # body:汇聚结果 ) ) print(f"Workflow definition:\n{workflow}") # 启动工作流 print("\nSending workflow to Celery...") result_promise = workflow.apply_async() print(f"Workflow sent! Final result ID: {result_promise.id}") print("Run a Celery worker to execute the tasks.") # 可选:等待最终结果 # final_result = result_promise.get() # print(f"\nWorkflow finished! Final result: {final_result}")

关键点解读

  • fetch_data.s(article_url):为第一步创建签名;
  • process_part_a.s()process_part_b.s():为并行任务创建签名。注意这里故意不传data参数——chain会自动把fetch_data的结果传给序列中的下一个任务;而下一个任务是包含groupchord,Celery 会聪明地把data分发给 group 中的每一个任务
  • combine_results.s():chord 的 body 签名,初始同样不需要参数,因为 chord 会自动把 header group 的结果列表传给它;
  • chain(...)fetch_datachord串联;
  • chord(group(...), ...)声明 group 必须全部完成后才会调用combine_results
  • workflow.apply_async():只把第一个任务fetch_data)投递给 Broker,工作流其余部分被编码进任务选项(如link或 chord 信息),Celery 据此在每一步完成后自动触发下一步。

运行前请确保有一个正在运行的 Worker。执行后,从 Worker 日志中可以看到依赖与并行度完全符合预期:fetch_data先执行,随后process_part_aprocess_part_b并发执行,最后在 A、B 均完成后combine_results执行。

需要说明的执行前提:上述示例假设app已按 Chapter 1: Celery App 与 Chapter 2: Configuration 配置好 Broker;若要使用result_promise.get()获取最终结果,还需要配置 Result Backend(如backend='redis://localhost:6379/1')。

内部原理:一个 chain 的完整执行旅程

以更简单的工作流my_chain = (add.s(2, 2) | add.s(4))为例,逐环节追踪:

  1. 工作流定义:创建my_chain时,Celery 构造一个chain对象,内部保存两个签名add.s(2, 2)add.s(4)
  2. 提交(my_chain.apply_async()
    • Celery 取出链中的第一个任务add.s(2, 2)
    • 准备将该任务消息发送到 Broker Connection (AMQP);
    • 关键一步:它会在消息中加入一个特殊选项,通常称为link(在较新的协议中使用chain字段),该选项包含链中下一个任务的签名add.s(4)
    • 携带linkadd(2, 2)消息被发送到 Broker。
  3. Worker 1 执行第一个任务
    • Worker 取到add(2, 2)的消息;
    • 以参数(2, 2)执行add,结果为4
    • 若配置了 Result Backend,将结果4存入后端;
    • Worker 注意到原始消息中的link选项指向add.s(4)
  4. Worker 1 投递第二个任务
    • Worker 取出第一个任务的结果4
    • 使用链接的签名add.s(4)
    • 把结果4前置拼接到链接签名的参数中,得到实际执行的add.s(4, 4)(链定义中自带的那个4保留,任务结果4插入到它前面);
    • 向 Broker 发送一条新的add(4, 4)消息。
  5. Worker 2 执行第二个任务
    • 另一个(或同一个)Worker 取到add(4, 4)
    • 执行得到8,存入后端;
    • 消息中没有更多link,链结束。

group的实现相对直接:把组内所有任务消息并发投递。chord则复杂得多:它需要 Worker 之间通过 Result Backend 协调,统计 header 中已完成的任务数,达到阈值后才投递 body 回调任务。

整个流程可以用时序图直观呈现:

值得注意的细节是apply_async()返回的AsyncResult的 ID 指向的是链中最后一个任务,因此你可以直接对它.get()拿到整条链的最终结果,而不必关心中间任务。

源码视角:celery/canvas.py 中的关键实现

Canvas 的签名与原语逻辑主要集中在 Celery 源码的celery/canvas.py中。以下内容用于理解其内部结构,可作为阅读源码的路线图。

Signature 类

  • 定义于celery/canvas.py,本质上是字典的子类,持有taskargskwargsoptions等字段;
  • Task实例上的.s()方法(位于celery/app/task.py)是创建Signature的快捷入口;
  • apply_async:通过调用_merge合并参数与选项,然后委托给self.type.apply_async(任务的方法)或app.send_task
  • linklink_error:向options字典追加回调签名(成功回调与错误回调);
  • __or__:重载管道操作符|,根据右操作数的类型构造对应的_chain对象。
# 简化自 celery/canvas.py class Signature(dict): # ... 其他方法如 __init__, clone, set, apply_async ... def link(self, callback): # 把回调签名追加到 options 的 'link' 列表中 return self.append_to_list_option('link', callback) def link_error(self, errback): # 把错误回调签名追加到 options 的 'link_error' 列表中 return self.append_to_list_option('link_error', errback) def __or__(self, other): # 使用管道 '|' 运算符时被调用 if isinstance(other, Signature): # task | task -> chain return _chain(self, other, app=self._app) # ... 其他 group、chain 等情况 ... return NotImplemented

_chain 类

  • 同样位于celery/canvas.py,继承自Signature,其task名被硬编码为'celery.chain',真正的任务签名存放在kwargs['tasks']中;
  • apply_async/run:包含投递第一个任务、并把链的其余部分嵌入选项的逻辑(协议 1 用link,协议 2 用chain消息属性);
  • prepare_steps:这个较复杂的方法会递归展开嵌套原语(如链中嵌套链、需要升级为 chord 的 group),并在各步骤之间建立连接关系。
# 简化自 celery/canvas.py(chain 执行) class _chain(Signature): # ... __init__, __or__ ... def apply_async(self, args=None, kwargs=None, **options): # ... 处理 always_eager ... return self.run(args, kwargs, app=self.app, **options) def run(self, args=None, kwargs=None, app=None, **options): # ... 初始化 ... tasks, results = self.prepare_steps(...) # 展开并冻结任务 if results: # 如果有任务需要运行 first_task = tasks.pop() # 取出第一个任务(列表是逆序的) remaining_chain = tasks if tasks else None # 决定用 link 还是消息字段传递链信息 use_link = self._use_link # ... 判定逻辑 ... if use_link: # 协议 1:把第一个任务链接到第二个任务 if remaining_chain: first_task.link(remaining_chain.pop()) # (后续链接由 Worker 处理) options_to_apply = options # 透传原始选项 else: # 协议 2:把剩余逆序链嵌入选项 options_to_apply = ChainMap({'chain': remaining_chain}, options) # 只投递第一个任务 result_from_apply = first_task.apply_async(**options_to_apply) # 返回原链中最后一个任务的 AsyncResult return results[0]

group 类

  • 位于celery/canvas.py,其task名为'celery.group'
  • apply_async:遍历其tasks,逐个freeze(为每个任务分配共同的group_id),发送消息,并把收集到的AsyncResult组装成GroupResult;它使用vine库中的barrier机制追踪整组完成状态。

chord 类

  • 位于celery/canvas.py,其task名为'celery.chord'
  • apply_async/run:与结果后端协同(backend.apply_chord)。典型流程是先运行 headergroup,并配置它在完成时通知后端;后端在计数达到预期任务数后触发 body 任务。

从源码结构可以推断,Canvas 的设计把「编排逻辑」下沉到了任务消息的link/chain字段与后端协调机制中,因此应用进程只需投递首个任务即可,后续流程完全由 Worker 自主接力——这正是它能把复杂工作流从应用代码中剥离出来的根本原因。

小结与下一步

Canvas 把普通任务升级为可组合的工作流构件:

  • Signature(task.s()捕获单次任务调用的完整计划而不立即执行;
  • 原语chain|)、groupchord把签名组合成不同的执行拓扑:
    • chain:顺序执行(前一个的输出成为下一个的输入);
    • group:并行执行;
    • chord:并行执行后,用全部结果触发一个回调任务;
  • 你可以像搭乐高一样嵌套组合这些原语,建模复杂的业务逻辑;
  • 对工作流原语调用.apply_async()时,Celery 只投递第一个任务,剩余流程逻辑通过任务选项或后端协调完成。

通过 Canvas,你可以把复杂的编排逻辑从应用代码迁移进 Celery 本身,让任务更模块化、系统更健壮。完成工作流构建后,下一章将介绍如何实时监控任务启动、完成与失败的状态——见 Chapter 9: Events;对本系列的整体架构与章节索引,可参阅 Celery 教程首页。

  • 人工智能
  • AI 应用
  • AI Agent

【免费下载链接】Tutorial-Codebase-Knowledge

Pocket Flow: Codebase to Tutorial

项目地址:https://gitcode.com/gh_mirrors/tu/Tutorial-Codebase-Knowledge
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

SSH密钥管理与Git权限问题解决方案

1. 问题现象与初步诊断每次看到终端里跳出"Permission denied (publickey)"的红色错误提示,作为开发者都会心头一紧。这个看似简单的权限问题,实际上可能涉及SSH密钥管理、远程仓库配置、系统权限设置等多个技术环节的故障。最近在团队协作中&…

作者头像 李华
网站建设 2026/9/23 14:42:00

晶闸管整流直流电动机调速系统:主电路、双闭环与参数整定全解析

简介:面向电气工程及自动化专业学生与电力电子技术初学者,这是一份晶闸管整流直流电动机调速系统的课程设计文档。内容以三相桥式全控整流电路为依托,完整讲解转速电流双闭环控制结构、主电路参数计算、基于TCA785集成触发芯片的移相触发原理…

作者头像 李华
网站建设 2026/9/23 14:41:15

选GEO服务商比的不是价格:先淘汰只会发稿的

企业第一次接触 GEO 服务商,通常会收到两份截然不同的报价单:一份按篇数计费,承诺一个月发多少篇稿;另一份按项目计费,先要资料、先做诊断,报价看起来贵不少。多数人会本能地倾向第一份——单价清晰、交付明…

作者头像 李华