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

基于文件变更触发Airflow DAG时watchdog未触发目标DAG问题

问题描述

我正在编写一套数据管道,需求为监听指定目录下JSON文件的修改事件,事件触发后运行另一个DAG。我尝试通过TriggerDagRunOperator实现跨DAG触发,同时使用watchdog库完成文件监控,但编辑被监控的文件时,目标hello_world_dag并未按预期运行。

当前实现包含两个DAG:

  • helper_dag:负责启动文件监控、触发目标DAG
  • hello_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个核心错误,导致监控逻辑完全不生效:

  1. Airflow Operator不是可直接调用的触发函数:on_modified回调里仅引用了trigger_other_dag这个Operator实例,既没有执行逻辑,也没有调用Airflow的DagRun触发接口——Operator是DAG定义阶段的任务模板,不是运行时可直接触发DAG的可执行对象。
  2. Watchdog启动逻辑写在__main__判断块中:Airflow加载DAG文件时以模块导入方式执行,不会运行if __name__ == "__main__"下的代码,Observer监控线程从未启动。
  3. 阻塞逻辑缺失:就算Observer成功启动,没有主线程阻塞逻辑的情况下,DAG跑完@once调度的任务后进程会直接退出,监控线程会被立即销毁。

修复方案

不要在DAG定义层直接运行Watchdog常驻进程,正确实现逻辑如下:

  1. 将文件监控逻辑封装为PythonOperator的执行函数,在函数内启动Observer、阻塞等待事件
  2. 放弃在回调中直接引用TriggerDagRunOperator实例,改用Airflow内置的trigger_dag工具方法,在监控回调中直接发起DAG触发请求
  3. 增加防抖逻辑,避免编辑器保存文件时触发多次修改事件,导致DAG重复运行
  4. 给监控任务配置专用资源池,避免长期占用通用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 11:06:31