如何用相同代码在Airflow中创建并运行多标签同逻辑DAG
在Airflow中复用代码创建多配置DAG的实用方案
方法1:循环遍历配置列表生成DAG
这是最直接的实现方式——定义一个包含所有DAG差异化配置的列表,遍历每个配置生成独立的DAG对象,核心任务逻辑完全复用。
示例代码:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime # 定义所有DAG的差异化配置:dag_id、标签、路由、调度规则等 dag_configs = [ { "dag_id": "user_data_pipeline_us", "tags": ["us", "user_data"], "route_config": {"source": "s3://us-user-data", "dest": "redshift-us-cluster"}, "schedule_interval": "@daily" }, { "dag_id": "user_data_pipeline_eu", "tags": ["eu", "user_data"], "route_config": {"source": "s3://eu-user-data", "dest": "redshift-eu-cluster"}, "schedule_interval": "@daily" }, { "dag_id": "user_data_pipeline_ap", "tags": ["ap", "user_data"], "route_config": {"source": "s3://ap-user-data", "dest": "redshift-ap-cluster"}, "schedule_interval": "@weekly" } ] # 复用的核心任务逻辑 def process_user_data(**context): route = context["params"]["route_config"] print(f"Processing data from {route['source']} to {route['dest']}") # 这里编写实际的数据处理逻辑 # 遍历配置生成DAG for config in dag_configs: with DAG( dag_id=config["dag_id"], start_date=datetime(2024, 1, 1), tags=config["tags"], schedule_interval=config["schedule_interval"], params={"route_config": config["route_config"]}, catchup=False ) as dag: process_task = PythonOperator( task_id="process_user_data", python_callable=process_user_data, provide_context=True ) # 若有其他通用任务,在此统一定义即可 process_task # 必须将DAG对象存入全局变量,Airflow调度器才能识别 globals()[config["dag_id"]] = dag
方法2:用类封装DAG模板
如果需要更灵活的扩展能力,可以创建DAG基类封装通用逻辑,再通过不同配置实例化出独立DAG。适合核心逻辑稳定、但部分场景需要少量扩展的需求。
示例代码:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime class BaseDataPipelineDAG: def __init__(self, dag_id, tags, route_config, schedule_interval): self.dag_id = dag_id self.tags = tags self.route_config = route_config self.schedule_interval = schedule_interval self.dag = self._create_dag() # 复用的核心处理逻辑 def _process_data(self, **context): route = context["params"]["route_config"] print(f"Processing data for {self.dag_id}: {route['source']} -> {route['dest']}") # 核心业务逻辑 # 通用DAG构建逻辑 def _create_dag(self): with DAG( dag_id=self.dag_id, start_date=datetime(2024, 1, 1), tags=self.tags, schedule_interval=self.schedule_interval, params={"route_config": self.route_config}, catchup=False ) as dag: process_task = PythonOperator( task_id="process_data", python_callable=self._process_data, provide_context=True ) # 可添加其他通用任务节点 process_task return dag # 实例化不同配置的DAG us_dag = BaseDataPipelineDAG( dag_id="user_data_pipeline_us", tags=["us", "user_data"], route_config={"source": "s3://us-user-data", "dest": "redshift-us-cluster"}, schedule_interval="@daily" ) globals()[us_dag.dag_id] = us_dag.dag eu_dag = BaseDataPipelineDAG( dag_id="user_data_pipeline_eu", tags=["eu", "user_data"], route_config={"source": "s3://eu-user-data", "dest": "redshift-eu-cluster"}, schedule_interval="@daily" ) globals()[eu_dag.dag_id] = eu_dag.dag
方法3:用Airflow Variables/外部配置文件驱动
如果需要动态修改配置而不改动代码,可以将DAG配置存储在Airflow Variables或外部YAML/JSON文件中,读取配置后生成DAG。适合配置频繁变更的场景。
Airflow Variables实现示例
- 在Airflow UI的「Variables」中创建名为
multi_dag_configs的变量,值为JSON格式的配置列表:
[ {"dag_id": "user_data_pipeline_us", "tags": ["us", "user_data"], "route_config": {"source": "s3://us-user-data", "dest": "redshift-us-cluster"}, "schedule_interval": "@daily"}, {"dag_id": "user_data_pipeline_eu", "tags": ["eu", "user_data"], "route_config": {"source": "s3://eu-user-data", "dest": "redshift-eu-cluster"}, "schedule_interval": "@daily"} ]
- 代码读取变量生成DAG:
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.models import Variable from datetime import datetime import json # 读取Airflow Variables中的配置 dag_configs = json.loads(Variable.get("multi_dag_configs")) def process_user_data(**context): route = context["params"]["route_config"] print(f"Processing data from {route['source']} to {route['dest']}") for config in dag_configs: with DAG( dag_id=config["dag_id"], start_date=datetime(2024, 1, 1), tags=config["tags"], schedule_interval=config["schedule_interval"], params={"route_config": config["route_config"]}, catchup=False ) as dag: process_task = PythonOperator( task_id="process_user_data", python_callable=process_user_data, provide_context=True ) process_task globals()[config["dag_id"]] = dag
关键注意事项
- 必须将生成的DAG对象存入全局变量(
globals()[dag_id] = dag),否则Airflow调度器无法识别。 - 每个DAG的
dag_id必须唯一,不能重复。 - 使用外部配置文件时,需确保Airflow Worker/Scheduler有权限读取该文件,且路径在允许范围内。
内容的提问来源于stack exchange,提问作者Shivangi Singh
相关产品推荐
相关产品推荐

