- 人工智能
- AI 应用
- AI Agent
【免费下载链接】Tutorial-Codebase-Knowledge
Pocket Flow: Codebase to Tutorial
导读
Celery 是 Python 生态中最流行的分布式任务队列框架之一,而Canvas(画布)是 Celery 面向复杂任务工作流提供的编排组件。它由Signature(签名)与原语(Primitives)组成,可以让你把单个任务组合成「顺序执行、并行执行、并行后汇聚」的复杂流程,而无需在应用代码里手工管理依赖与结果传递。本文将围绕docs/Celery/08_canvas__signatures___primitives_.md的核心内容,从 Signature 的概念讲起,完整演示用chain、group、chord构建一个真实的文章处理流水线,并深入celery/canvas.py的源码结构,剖析这些原语在 Broker 与 Worker 之间是如何一步步被执行的。读完本文,你将能够独立设计、提交并排查自己的 Celery 工作流。
为什么需要 Canvas:任务编排的痛点
在 Chapter 3: Task 中我们学会了如何用@app.task定义任务,并用.delay()/.apply_async()把单个任务投递给 Broker。但真实业务几乎不会只有一个独立任务。考虑这样一个场景:用户上传一篇文章后,系统需要:
- 从 URL 抓取文章内容;
- 对文本做关键词提取;
- 对文本做语言检测;
- 等上述两步都完成后,把文章与元数据一并保存到数据库。
如果简单地逐个投递任务,你无法表达「步骤 2 与 3 可以并行、但都必须在步骤 1 之后」这种依赖关系,更无法保证「保存」在两者都成功之后才执行。若在应用代码里手工轮询AsyncResult、拼装参数,代码会迅速变得脆弱且难以维护。
Canvas 解决的正是任务之间的依赖与流程控制问题。它允许你把工作流的拓扑结构直接声明出来:哪个任务先跑、哪些可以并行、哪个任务必须等齐所有并行结果后再跑,然后把整个流程交给 Celery 去执行。官方文档用一个非常形象的比喻来描述 Canvas:就像不同形状的乐高积木——
- 有些积木代表单个任务;
- 有些积木把任务首尾相连(顺序执行);
- 有些积木让任务并排堆放(并行执行);
- 还有些积木可以构建「多个并行步骤必须全部完成,才能拼上下一块」的结构。
核心概念一:Signature——单个任务的"预约单"
什么是 Signature
一个Signature封装了调用某个任务所需的全部信息:
- 任务的名称(
task); - 位置参数(
args); - 关键字参数(
kwargs); - 以及执行选项(如
countdown、eta、队列名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 秒后再执行;你还可以在创建签名时直接写入queue、routing_key等选项,让该签名固定投递到特定队列,这与 Chapter 3: Task 中apply_async的选项体系是一致的。
Signature 的三个重要性质
- 可序列化:Signature 本质是一个可被序列化的结构(在 Celery 源码中它继承自字典),因此可以随任务消息在网络上传输——这正是它能被嵌入
link选项、跨 Worker 传递的前提。 - 部分应用(partial application):你可以在创建签名时不填满所有参数,留待链式执行时由前一个任务的结果自动补全。这在后面的
chain中会频繁用到。 - 可克隆(clone):签名支持
clone()生成副本,在prepare_steps等内部逻辑中,Celery 会不断对签名进行克隆与参数合并,避免污染原始定义。
核心概念二:工作流原语——连接积木的四种方式
Canvas 提供了若干**原语函数(Primitives)**用于把签名组合成工作流。其中最核心的是chain、group、chord,另外还有chunks、xmap、starmap等补充原语。
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传参数元组)。
chain、group、chord是构建工作流最基础的三个原语,掌握它们足以覆盖绝大多数编排需求。
实战:构建文章处理工作流
回到开头的文章处理场景,我们用 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的结果传给序列中的下一个任务;而下一个任务是包含group的chord,Celery 会聪明地把data分发给 group 中的每一个任务;combine_results.s():chord 的 body 签名,初始同样不需要参数,因为 chord 会自动把 header group 的结果列表传给它;chain(...)把fetch_data与chord串联;chord(group(...), ...)声明 group 必须全部完成后才会调用combine_results;workflow.apply_async():只把第一个任务(fetch_data)投递给 Broker,工作流其余部分被编码进任务选项(如link或 chord 信息),Celery 据此在每一步完成后自动触发下一步。
运行前请确保有一个正在运行的 Worker。执行后,从 Worker 日志中可以看到依赖与并行度完全符合预期:fetch_data先执行,随后process_part_a与process_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))为例,逐环节追踪:
- 工作流定义:创建
my_chain时,Celery 构造一个chain对象,内部保存两个签名add.s(2, 2)与add.s(4)。 - 提交(
my_chain.apply_async()):- Celery 取出链中的第一个任务
add.s(2, 2); - 准备将该任务消息发送到 Broker Connection (AMQP);
- 关键一步:它会在消息中加入一个特殊选项,通常称为
link(在较新的协议中使用chain字段),该选项包含链中下一个任务的签名add.s(4); - 携带
link的add(2, 2)消息被发送到 Broker。
- Celery 取出链中的第一个任务
- Worker 1 执行第一个任务:
- Worker 取到
add(2, 2)的消息; - 以参数
(2, 2)执行add,结果为4; - 若配置了 Result Backend,将结果
4存入后端; - Worker 注意到原始消息中的
link选项指向add.s(4)。
- Worker 取到
- Worker 1 投递第二个任务:
- Worker 取出第一个任务的结果
4; - 使用链接的签名
add.s(4); - 把结果
4前置拼接到链接签名的参数中,得到实际执行的add.s(4, 4)(链定义中自带的那个4保留,任务结果4插入到它前面); - 向 Broker 发送一条新的
add(4, 4)消息。
- Worker 取出第一个任务的结果
- Worker 2 执行第二个任务:
- 另一个(或同一个)Worker 取到
add(4, 4); - 执行得到
8,存入后端; - 消息中没有更多
link,链结束。
- 另一个(或同一个)Worker 取到
group的实现相对直接:把组内所有任务消息并发投递。chord则复杂得多:它需要 Worker 之间通过 Result Backend 协调,统计 header 中已完成的任务数,达到阈值后才投递 body 回调任务。
整个流程可以用时序图直观呈现:
值得注意的细节是apply_async()返回的AsyncResult的 ID 指向的是链中最后一个任务,因此你可以直接对它.get()拿到整条链的最终结果,而不必关心中间任务。
源码视角:celery/canvas.py 中的关键实现
Canvas 的签名与原语逻辑主要集中在 Celery 源码的celery/canvas.py中。以下内容用于理解其内部结构,可作为阅读源码的路线图。
Signature 类
- 定义于
celery/canvas.py,本质上是字典的子类,持有task、args、kwargs、options等字段; Task实例上的.s()方法(位于celery/app/task.py)是创建Signature的快捷入口;apply_async:通过调用_merge合并参数与选项,然后委托给self.type.apply_async(任务的方法)或app.send_task;link、link_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(|)、group、chord把签名组合成不同的执行拓扑:chain:顺序执行(前一个的输出成为下一个的输入);group:并行执行;chord:并行执行后,用全部结果触发一个回调任务;
- 你可以像搭乐高一样嵌套组合这些原语,建模复杂的业务逻辑;
- 对工作流原语调用
.apply_async()时,Celery 只投递第一个任务,剩余流程逻辑通过任务选项或后端协调完成。
通过 Canvas,你可以把复杂的编排逻辑从应用代码迁移进 Celery 本身,让任务更模块化、系统更健壮。完成工作流构建后,下一章将介绍如何实时监控任务启动、完成与失败的状态——见 Chapter 9: Events;对本系列的整体架构与章节索引,可参阅 Celery 教程首页。
- 人工智能
- AI 应用
- AI Agent
【免费下载链接】Tutorial-Codebase-Knowledge
Pocket Flow: Codebase to Tutorial
相关推荐
Celery Canvas 工作流设计指南:从 Signature 到 Group、Chain、Chord 与 Stamping 全解析
Celery Canvas 工作流设计指南:从 Signature 到 Group、Chain、Chord 与 Stamping 全解析 导读 本文基于 Cel
任务调度后端消息队列Celery工作流与任务组合:Chain、Group与Chord
Celery工作流与任务组合:Chain、Group与Chord 本文深入探讨了Celery中三种核心工作流组合模式:Chain(任务链)、Group(任务组)
任务调度后端消息队列Celery任务链终极指南:Chain、Group和Chord的深度应用
Celery任务链终极指南:Chain、Group和Chord的深度应用 Celery是一个强大的Python分布式任务队列库,专门用于处理后台任务调度和分布式
任务调度后端消息队列
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考