如何动态创建Airflow任务并遍历CLI Variables中的DAG ID
动态监控Airflow中所有指定DAG完成状态的方案
核心思路
直接从Airflow Variables中读取需要监控的DAG ID列表,动态生成对应的ExternalTaskSensor任务,待所有传感器检测到对应DAG成功后,触发通知邮件任务。
具体实现步骤
存储DAG ID列表到Airflow Variables
- 在Airflow UI的
Admin > Variables中创建变量(比如命名为monitored_dag_ids),值推荐用JSON数组格式(如["dag_sales_report","dag_user_analytics","dag_inventory_stat"]),也可以用逗号分隔的字符串(如dag_sales_report,dag_user_analytics),JSON格式更规范易维护。
- 在Airflow UI的
编写监控DAG代码
以下是可直接复用的示例代码:from airflow import DAG from airflow.sensors.external_task import ExternalTaskSensor from airflow.operators.email import EmailOperator from airflow.models import Variable from datetime import datetime, timedelta # 从Variables读取监控DAG列表 def get_monitored_dags(): # 读取并解析JSON格式的变量值为列表 return Variable.get("monitored_dag_ids", deserialize_json=True) default_args = { 'owner': 'data_team', 'depends_on_past': False, 'email_on_failure': True, 'email_on_retry': False, 'retries': 2, 'retry_delay': timedelta(minutes=10), 'start_date': datetime(2024, 1, 1), } with DAG( 'monitor_all_dags_completion', default_args=default_args, description='监控指定DAG全部完成后发送业务通知', schedule_interval='0 10 * * *', # 调度时间晚于被监控DAG的每日执行时间 catchup=False, tags=['monitoring', 'business'], ) as dag: # 动态生成所有DAG的传感器任务 sensor_tasks = [] for dag_id in get_monitored_dags(): dag_sensor = ExternalTaskSensor( task_id=f"wait_for_{dag_id}", external_dag_id=dag_id, external_task_id=None, # 监控整个DAG的成功状态,无需指定单个任务 execution_date_fn=lambda dt: dt, # 匹配当日的执行实例 mode='reschedule', # 定期检查,节省资源 timeout=3600 * 12, # 设置超时时间,比如12小时 poke_interval=300, # 每5分钟检查一次状态 ) sensor_tasks.append(dag_sensor) # 所有DAG完成后发送通知邮件 send_business_notice = EmailOperator( task_id='send_success_notification', to='business_analytics@yourcompany.com', subject='每日DAG全部执行完成通知', html_content='<p>今日所有业务相关DAG已成功执行完毕,可开始报表数据分析工作。</p>', ) # 设置依赖:所有传感器并行执行,全部完成后触发邮件任务 sensor_tasks >> send_business_notice维护说明
- 后续新增需要监控的DAG时,仅需在Airflow Variables中更新
monitored_dag_ids的列表,无需修改监控DAG的代码。 - 可根据业务场景调整传感器的
poke_interval(检查间隔)、timeout(超时时间)等参数。
- 后续新增需要监控的DAG时,仅需在Airflow Variables中更新
关键细节
external_task_id=None时,传感器会校验整个DAG的成功状态,即DAG下所有任务都执行成功后,传感器才会标记完成。execution_date_fn确保传感器监控的是与当前监控DAG同日期的执行实例,避免跨日期匹配错误。mode='reschedule'相比poke模式更节省Airflow worker资源,适合长时间等待的场景。
内容的提问来源于stack exchange,提问作者Khilesh Chauhan
相关产品推荐
相关产品推荐

