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

如何动态创建Airflow任务并遍历CLI Variables中的DAG ID

动态监控Airflow中所有指定DAG完成状态的方案

核心思路

直接从Airflow Variables中读取需要监控的DAG ID列表,动态生成对应的ExternalTaskSensor任务,待所有传感器检测到对应DAG成功后,触发通知邮件任务。

具体实现步骤

  1. 存储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格式更规范易维护。
  2. 编写监控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
    
  3. 维护说明

    • 后续新增需要监控的DAG时,仅需在Airflow Variables中更新monitored_dag_ids的列表,无需修改监控DAG的代码。
    • 可根据业务场景调整传感器的poke_interval(检查间隔)、timeout(超时时间)等参数。

关键细节

  • external_task_id=None时,传感器会校验整个DAG的成功状态,即DAG下所有任务都执行成功后,传感器才会标记完成。
  • execution_date_fn确保传感器监控的是与当前监控DAG同日期的执行实例,避免跨日期匹配错误。
  • mode='reschedule'相比poke模式更节省Airflow worker资源,适合长时间等待的场景。

内容的提问来源于stack exchange,提问作者Khilesh Chauhan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 06:36:27