如何以Pythonic方式链式串联Airflow可复用装饰任务(非动态映射)
Airflow 可复用任务链式串联实现方法
你当前的代码中,start_1至start_3任务直接依赖start,导致并行执行。要实现链式串联的依赖关系,只需在循环中维护一个跟踪当前上游任务的变量,依次创建新任务并设置依赖即可,无需使用动态映射。
修改后的代码如下:
from datetime import DateTime from airflow.decorators import dag, task @task def add_task(x, y): print(f"Task args: x={x}, y={y}") return x + y @dag(start_date=DateTime(2022, 1, 1), schedule=None, catchup=False) def mydag(): # 初始化起始任务 current_task = add_task.override(task_id="start_0")(1, 2) # 循环创建链式任务 for x, y, task_id_str in zip([1,3,5],[2,4,6],["start_1", "start_2", "start_3"], strict=True): # 创建新任务 new_task = add_task.override(task_id=task_id_str)(x, y) # 设置当前任务为新任务的上游 current_task >> new_task # 更新当前任务为新任务,用于下一次循环的依赖设置 current_task = new_task mydag()
代码说明
- 用
current_task变量跟踪当前的上游任务,初始值为start_0 - 每次循环创建新任务后,通过
current_task >> new_task建立依赖关系 - 更新
current_task为刚创建的新任务,确保下一个任务依赖于当前任务,最终形成start_0 >> start_1 >> start_2 >> start_3的链式结构
内容的提问来源于stack exchange,提问作者noirjaunerouge
相关产品推荐
相关产品推荐

