You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用相同代码在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实现示例

  1. 在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"}
]
  1. 代码读取变量生成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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 03:35:33