如何在Airflow DAG的循环中实现任务串行运行?
DAG循环生成串行任务实现方案
这个需求完全可以通过循环实现,无需手动逐个声明任务依赖。
你原有代码生成的任务默认并行,是因为仅在循环中实例化了任务对象,没有给任务设置上下游执行依赖,DAG调度器会默认并行执行所有无依赖的独立任务。
修改后的代码如下:
variables = ['first', 'second', 'third'] def run_dag_task(variable): task = dag_task(variable) return task prev_task = None for variable in variables: current_task = run_dag_task(variable) # 非首次循环时,给当前任务添加上游依赖 if prev_task: # Airflow 中使用 >> 定义依赖,其他DAG框架替换为对应依赖设置语法即可 prev_task >> current_task # 更新上一个任务为当前任务,供下一轮循环使用 prev_task = current_task
核心逻辑说明:
- 额外定义一个变量存储上一轮循环生成的任务实例,初始值设为空
- 第一次循环仅生成第一个任务,不需要设置依赖
- 从第二个任务开始,每个新生成的任务都绑定依赖,必须等前一个任务执行完成后才会启动
- 所有任务的串行依赖关系会随着循环自动完成绑定,哪怕后续你要新增更多变量到列表里,也不需要修改循环逻辑
内容的提问来源于stack exchange,提问作者Stijnvandenb
相关产品推荐
相关产品推荐

