Airflow中如何跨DAG获取另一DAG默认参数的email配置
跨DAG获取目标DAG default_args中email字段的实现方案
直接通过Airflow内置的DagBag类加载目标DAG对象即可读取配置,不需要依赖目标DAG的运行实例,实现逻辑非常简单。
实现代码
你可以新建一个独立的查询DAG,代码如下,直接替换目标DAG的ID/文件路径即可使用:
from airflow import DAG from airflow.operators.python_operator import PythonOperator from airflow.models.dagbag import DagBag from datetime import datetime, timedelta def get_target_dag_email(**context): # 方式1:加载全量DAG后按ID取目标DAG dag_bag = DagBag() # 方式2:指定目标DAG文件路径,仅加载单个DAG,性能更好 # dag_bag = DagBag(dag_folder="/你的dags目录路径/main_dag.py") target_dag = dag_bag.get_dag(dag_id="main") # 替换为你的目标DAG ID if not target_dag: raise Exception("目标DAG不存在,请检查DAG ID和文件路径是否正确") # 直接读取default_args中的email字段 email_config = target_dag.default_args.get("email") print(f"读取到的目标DAG email配置:{email_config}") # 如果需要给同DAG其他任务传值,可推入XCom context["ti"].xcom_push(key="target_email", value=email_config) return email_config default_args = { 'owner': 'airflow', 'depends_on_past': False, 'retries': 1, 'retry_delay': timedelta(minutes=5) } with DAG( dag_id="fetch_dag_email", default_args=default_args, schedule_interval=None, start_date=datetime(2024, 1, 1), catchup=False ) as dag: fetch_email_task = PythonOperator( task_id="fetch_email_task", python_callable=get_target_dag_email, provide_context=True # Airflow 1.x版本需开启,2.x版本可删除 )
注意事项
- 确保两个DAG都放在Airflow配置的
dags_folder路径下,DagBag才能正常扫描加载 - Airflow 2.3+版本中
PythonOperator的导入路径改为from airflow.operators.python import PythonOperator,核心逻辑无需调整 - 读取到的email值和你在目标DAG中配置的格式完全一致,对应你提供的示例DAG,拿到的返回值为
['example@gmail.com'] - 不建议通过
DagRun对象取配置,DagRun仅对应DAG的某次具体运行实例,若DAG从未触发过运行会拿不到值,直接从DagBag加载DAG定义是最稳定的方式
内容的提问来源于stack exchange,提问作者Kabir Dureja
相关产品推荐
相关产品推荐

