人工智能AI 应用AI Agent【免费下载链接】Tutorial-Codebase-KnowledgePocket 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 的概念讲起完整演示用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(fSignature: {add_sig}) print(fTask name: {add_sig.task}) print(fArguments: {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(countdown10)表示 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头部一个并行执行的groupbody主体一个回调签名在 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(fFetching data from {url}...) time.sleep(1) # 模拟抓取数据 data fContent from {url} - {random.randint(1, 100)} print(fFetched: {data}) return data app.task def process_part_a(data): print(fProcessing Part A for: {data}) time.sleep(2) result_a fKeywords for {data} print(Part A finished.) return result_a app.task def process_part_b(data): print(fProcessing Part B for: {data}) time.sleep(3) # 模拟稍长的处理 result_b fLanguage for {data} print(Part B finished.) return result_b app.task def combine_results(results): # results 是一个列表包含 process_part_a 和 process_part_b 的返回值 print(fCombining results: {results}) time.sleep(1) final_output fCombined: {results[0]} | {results[1]} print(fFinal 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 # - headergroup 并行运行 process_part_a 与 process_part_b # 两个任务都会收到 fetch_data 传来的 data。 # - bodycombine_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(fWorkflow definition:\n{workflow}) # 启动工作流 print(\nSending workflow to Celery...) result_promise workflow.apply_async() print(fWorkflow 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的chordCelery 会聪明地把data分发给 group 中的每一个任务combine_results.s()chord 的 body 签名初始同样不需要参数因为 chord 会自动把 header group 的结果列表传给它chain(...)把fetch_data与chord串联chord(group(...), ...)声明 group 必须全部完成后才会调用combine_resultsworkflow.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如backendredis://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。Worker 1 执行第一个任务Worker 取到add(2, 2)的消息以参数(2, 2)执行add结果为4若配置了 Result Backend将结果4存入后端Worker 注意到原始消息中的link选项指向add.s(4)。Worker 1 投递第二个任务Worker 取出第一个任务的结果4使用链接的签名add.s(4)把结果4前置拼接到链接签名的参数中得到实际执行的add.s(4, 4)链定义中自带的那个4保留任务结果4插入到它前面向 Broker 发送一条新的add(4, 4)消息。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本质上是字典的子类持有task、args、kwargs、options等字段Task实例上的.s()方法位于celery/app/task.py是创建Signature的快捷入口apply_async通过调用_merge合并参数与选项然后委托给self.type.apply_async任务的方法或app.send_tasklink、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, appself._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.pychain 执行 class _chain(Signature): # ... __init__, __or__ ... def apply_async(self, argsNone, kwargsNone, **options): # ... 处理 always_eager ... return self.run(args, kwargs, appself.app, **options) def run(self, argsNone, kwargsNone, appNone, **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.groupapply_async遍历其tasks逐个freeze为每个任务分配共同的group_id发送消息并把收集到的AsyncResult组装成GroupResult它使用vine库中的barrier机制追踪整组完成状态。chord 类位于celery/canvas.py其task名为celery.chordapply_async/run与结果后端协同backend.apply_chord。典型流程是先运行 headergroup并配置它在完成时通知后端后端在计数达到预期任务数后触发 body 任务。从源码结构可以推断Canvas 的设计把「编排逻辑」下沉到了任务消息的link/chain字段与后端协调机制中因此应用进程只需投递首个任务即可后续流程完全由 Worker 自主接力——这正是它能把复杂工作流从应用代码中剥离出来的根本原因。小结与下一步Canvas 把普通任务升级为可组合的工作流构件Signaturetask.s()捕获单次任务调用的完整计划而不立即执行原语chain|、group、chord把签名组合成不同的执行拓扑chain顺序执行前一个的输出成为下一个的输入group并行执行chord并行执行后用全部结果触发一个回调任务你可以像搭乐高一样嵌套组合这些原语建模复杂的业务逻辑对工作流原语调用.apply_async()时Celery 只投递第一个任务剩余流程逻辑通过任务选项或后端协调完成。通过 Canvas你可以把复杂的编排逻辑从应用代码迁移进 Celery 本身让任务更模块化、系统更健壮。完成工作流构建后下一章将介绍如何实时监控任务启动、完成与失败的状态——见 Chapter 9: Events对本系列的整体架构与章节索引可参阅 Celery 教程首页。赞分享人工智能AI 应用AI Agent【免费下载链接】Tutorial-Codebase-KnowledgePocket Flow: Codebase to Tutorial项目地址https://gitcode.com/gh_mirrors/tu/Tutorial-Codebase-Knowledge点击查看免费下载相关推荐Celery Canvas 工作流设计指南从 Signature 到 Group、Chain、Chord 与 Stamping 全解析Celery Canvas 工作流设计指南从 Signature 到 Group、Chain、Chord 与 Stamping 全解析 导读 本文基于 Cel任务调度后端消息队列Celery工作流与任务组合Chain、Group与ChordCelery工作流与任务组合Chain、Group与Chord 本文深入探讨了Celery中三种核心工作流组合模式Chain任务链、Group任务组任务调度后端消息队列Celery任务链终极指南Chain、Group和Chord的深度应用Celery任务链终极指南Chain、Group和Chord的深度应用 Celery是一个强大的Python分布式任务队列库专门用于处理后台任务调度和分布式任务调度后端消息队列创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
