如何在Airflow环境启动或重启时自动触发DAG Run
Airflow启动时自动触发DAG运行的实现方案
以下是两种生产环境常用的实现方式,可根据你的部署方式选择:
方案1:基于Airflow插件的启动钩子(兼容性最强)
该方式不依赖外部部署工具,完全通过Airflow原生能力实现,适合所有部署场景:
- 操作步骤:
- 进入你的Airflow插件目录(默认路径为
$AIRFLOW_HOME/plugins),新建Python文件如startup_dag_trigger.py - 写入以下逻辑,仅在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"- 启动Scheduler时添加环境变量标记:
AIRFLOW_SCHEDULER_START=1 airflow scheduler,每次Scheduler启动就会自动触发目标DAG
- 进入你的Airflow插件目录(默认路径为
方案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
相关产品推荐
相关产品推荐

