如何创建可被继承的Airflow样板DAG(Boilerplate DAG)
Airflow 2.5.1 可继承样板DAG的正确实现
原代码的问题点
create_dag方法未被调用,导致基类和子类的任务都无法添加到DAG实例中- 子类中多个任务使用重复的
task_id(run_after_loop),会触发Airflow任务ID冲突报错 - 基类的
default_args虽在__init__中传递,但因任务未初始化,无法实际生效
修复后的代码实现
from datetime import datetime from airflow import DAG from airflow.operators.empty import EmptyOperator from airflow.operators.bash import BashOperator class BaseDag(DAG): # 定义基类默认参数 default_args = { 'owner': 'airflow', 'start_date': datetime(2023, 6, 18), 'schedule_interval': None, } def __init__(self, dag_id): super().__init__( dag_id=dag_id, default_args=self.default_args, catchup=False ) # 初始化时调用create_dag,确保任务被创建并关联到DAG self.create_dag() def create_dag(self): with self: # 创建基类通用任务 start = EmptyOperator(task_id='start') end = EmptyOperator(task_id='end') start >> end class MyDag(BaseDag): def create_dag(self): # 先调用基类方法生成通用任务 super().create_dag() with self: # 子类自定义任务,确保task_id唯一 task1 = BashOperator( task_id='task_1', bash_command='echo 1', ) task2 = BashOperator( task_id='task_2', bash_command='echo 2', ) task3 = BashOperator( task_id='task_3', bash_command='echo 3', ) # 获取基类任务,设置子类任务的依赖关系 start = self.get_task('start') end = self.get_task('end') start >> task1 >> task2 >> task3 >> end # 实例化子类,生成最终DAG first_dag = MyDag(dag_id='first_dag')
关键修改说明
- 在基类
__init__中调用self.create_dag(),确保任务逻辑被执行,任务正确关联到DAG实例 - 子类中修改任务ID为唯一值,避免冲突
- 子类
create_dag方法中重新使用with self:上下文,确保自定义任务正确绑定到当前DAG - 基类的
default_args会通过父类构造函数传递,所有子类自动继承这些默认配置
内容的提问来源于stack exchange,提问作者Jag Singh
相关产品推荐
相关产品推荐

