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

如何在Airflow环境启动或重启时自动触发DAG Run

Airflow启动时自动触发DAG运行的实现方案

以下是两种生产环境常用的实现方式,可根据你的部署方式选择:


方案1:基于Airflow插件的启动钩子(兼容性最强)

该方式不依赖外部部署工具,完全通过Airflow原生能力实现,适合所有部署场景:

  • 操作步骤:
    1. 进入你的Airflow插件目录(默认路径为$AIRFLOW_HOME/plugins),新建Python文件如startup_dag_trigger.py
    2. 写入以下逻辑,仅在Scheduler主进程启动时触发一次目标DAG,同时做幂等校验避免重复运行:
    from airflow.plugins_manager import AirflowPlugin
    from airflow.models import DagRun
    from airflow.api.common.trigger_dag import trigger_dag
    from datetime import datetime, timedelta
    import os
    
    # 仅Scheduler进程启动时执行
    if "scheduler" in os.environ.get("AIRFLOW__CORE__EXECUTOR", "").lower() or os.environ.get("AIRFLOW_SCHEDULER_START", "0") == "1":
        TARGET_DAG_ID = "替换为你要触发的DAG ID"
        # 校验近10分钟内该DAG是否已有运行记录,避免重复触发
        recent_runs = DagRun.find(
            dag_id=TARGET_DAG_ID,
            execution_date=datetime.now() - timedelta(minutes=10)
        )
        if not recent_runs:
            trigger_dag(dag_id=TARGET_DAG_ID, execution_date=datetime.now())
    
    class StartupTriggerPlugin(AirflowPlugin):
        name = "startup_dag_trigger"
    
    1. 启动Scheduler时添加环境变量标记:AIRFLOW_SCHEDULER_START=1 airflow scheduler,每次Scheduler启动就会自动触发目标DAG

方案2:基于系统服务启动后置命令(适合systemd/Docker部署)

如果你的Airflow是通过systemd管理服务,或者用自定义Docker镜像启动,可直接在启动流程末尾加触发命令:

  • systemd部署场景:在airflow-scheduler.service的[Service]段落添加如下配置:
    # 等待30秒确保Airflow完全启动后再触发DAG
    ExecStartPost=/bin/sh -c 'sleep 30 && airflow dags trigger 替换为你的DAG ID'
    
  • Docker部署场景:在你的启动脚本最后,添加airflow dags trigger 替换为你的DAG ID命令即可,注意要放在Scheduler启动逻辑之后,加适当的等待时间避免请求报错

注意事项

  • 提前将目标DAG设置为激活(unpaused)状态,否则触发会失败
  • 测试类DAG建议配置catchup=False,避免服务重启时积压历史任务
  • 如果你的Airflow开启了认证,用API触发时需要带上对应的身份校验信息

内容的提问来源于stack exchange,提问作者A Buyer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 16:06:00