基于文件变更触发Airflow DAG时watchdog未触发目标DAG问题
问题描述
我正在编写一套数据管道,需求为监听指定目录下JSON文件的修改事件,事件触发后运行另一个DAG。我尝试通过TriggerDagRunOperator实现跨DAG触发,同时使用watchdog库完成文件监控,但编辑被监控的文件时,目标hello_world_dag并未按预期运行。
当前实现包含两个DAG:
helper_dag:负责启动文件监控、触发目标DAGhello_world_dag:JSON文件变更时需要触发运行的主管道
现有代码
helper_dag 代码
from airflow import DAG from airflow.operators.dagrun_operator import TriggerDagRunOperator from watchdog.observers import Observer from watchdog.events import PatternMatchingEventHandler from datetime import datetime dag = DAG( dag_id="helper_dag", start_date=datetime(2022, 5, 28), catchup=False, schedule_interval="@once", ) trigger_other_dag = TriggerDagRunOperator( task_id="trigger_dagrun", trigger_dag_id="hello_world_dag", dag=dag, retries=0, ) class Handler(PatternMatchingEventHandler): def on_modified(self, event): trigger_other_dag if __name__ == "__main__": my_event_handler = Handler( patterns=["*.json"], ignore_patterns=None, ignore_directories=False, case_sensitive=True, ) my_observer = Observer() my_observer.schedule( event_handler=my_event_handler, path="/home/airflowsrv/airflow/", recursive=True ) my_observer.start()
hello_world_dag 代码
from datetime import datetime from airflow import DAG from airflow.operators.empty import EmptyOperator from airflow.operators.python_operator import PythonOperator def print_hello(): return str(datetime.now()) dag = DAG( "hello_world_dag", start_date=datetime(2022, 5, 28), catchup=False, schedule_interval=None, ) hello_operator = PythonOperator(task_id="hello_task", python_callable=print_hello, dag=dag) hello_operator
问题原因
代码存在3个核心错误,导致监控逻辑完全不生效:
- Airflow Operator不是可直接调用的触发函数:
on_modified回调里仅引用了trigger_other_dag这个Operator实例,既没有执行逻辑,也没有调用Airflow的DagRun触发接口——Operator是DAG定义阶段的任务模板,不是运行时可直接触发DAG的可执行对象。 - Watchdog启动逻辑写在
__main__判断块中:Airflow加载DAG文件时以模块导入方式执行,不会运行if __name__ == "__main__"下的代码,Observer监控线程从未启动。 - 阻塞逻辑缺失:就算Observer成功启动,没有主线程阻塞逻辑的情况下,DAG跑完
@once调度的任务后进程会直接退出,监控线程会被立即销毁。
修复方案
不要在DAG定义层直接运行Watchdog常驻进程,正确实现逻辑如下:
- 将文件监控逻辑封装为PythonOperator的执行函数,在函数内启动Observer、阻塞等待事件
- 放弃在回调中直接引用TriggerDagRunOperator实例,改用Airflow内置的
trigger_dag工具方法,在监控回调中直接发起DAG触发请求 - 增加防抖逻辑,避免编辑器保存文件时触发多次修改事件,导致DAG重复运行
- 给监控任务配置专用资源池,避免长期占用通用Worker槽阻塞其他任务
修复后的helper_dag代码
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.api.common.experimental.trigger_dag import trigger_dag from watchdog.observers import Observer from watchdog.events import PatternMatchingEventHandler from datetime import datetime, timedelta import time # 监控配置 MONITOR_PATH = "/home/airflowsrv/airflow/" TARGET_DAG_ID = "hello_world_dag" # 防抖间隔,单位:秒 DEBOUNCE_INTERVAL = 2 last_trigger_time = 0 class JsonFileHandler(PatternMatchingEventHandler): def on_modified(self, event): global last_trigger_time # 跳过目录事件 if event.is_directory: return # 防抖校验 now = time.time() if now - last_trigger_time < DEBOUNCE_INTERVAL: return last_trigger_time = now print(f"检测到JSON文件变更: {event.src_path}, 触发目标DAG") # 调用Airflow原生API触发DAG trigger_dag( dag_id=TARGET_DAG_ID, run_id=f"file_trigger_{datetime.now().strftime('%Y%m%d%H%M%S')}", conf={"modified_file": event.src_path}, execution_date=datetime.now(), replace_microseconds=False ) def monitor_dir_func(): event_handler = JsonFileHandler( patterns=["*.json"], ignore_directories=True, case_sensitive=True ) observer = Observer() observer.schedule(event_handler, MONITOR_PATH, recursive=True) observer.start() print(f"启动文件监控,监听路径: {MONITOR_PATH}") try: # 阻塞保持监控运行 while True: time.sleep(1) except KeyboardInterrupt: observer.stop() observer.join() with DAG( dag_id="helper_dag", start_date=datetime(2022, 5, 28), catchup=False, schedule_interval="@once", default_args={ "execution_timeout": timedelta(days=365), "retries": 0 } ) as dag: monitor_task = PythonOperator( task_id="monitor_json_files", python_callable=monitor_dir_func, pool="file_monitor_pool" # 建议提前创建仅分配1个槽的专用池给监控任务 )
补充说明
hello_world_dag代码无需修改,保持原有schedule_interval=None的配置即可正常被触发- 生产环境优先使用Airflow原生
FileSensor搭配TriggerDagRunOperator实现文件感知,稳定性高于自定义Watchdog实现,无需自行维护常驻线程 - 分布式部署场景下,需要将监控任务固定调度到挂载了目标监控目录的Worker节点上运行,同时确保Airflow运行用户对监控目录有读权限
内容的提问来源于stack exchange,提问作者enamya
相关产品推荐
相关产品推荐

