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

如何在Airflow中实现任意文件上传至DAG目录时触发任务?

方案推荐

Airflow原生没有直接支持监测目录下任意新文件上传的Operator,以下是几种可行的解决思路:

1. 自定义Sensor监测目录新增文件

基于Airflow的PythonSensor或自定义Sensor实现目录变化监测逻辑:

  • 核心思路:记录目录上次扫描的文件列表(可存在Airflow变量或本地文件中),每次扫描对比当前文件列表,发现新增文件就标记任务为成功,触发后续流程。
  • 示例代码:
from airflow.sensors.python import PythonSensor
from airflow.models import Variable
import os

def check_new_files(directory):
    # 从Airflow变量获取上次记录的文件列表
    last_files = Variable.get("last_scanned_files", default_var=[])
    current_files = os.listdir(directory)
    # 找出新增文件
    new_files = [f for f in current_files if f not in last_files]
    if new_files:
        # 更新变量为当前文件列表
        Variable.set("last_scanned_files", current_files)
        return True
    return False

# 在DAG中使用该Sensor
file_monitor_sensor = PythonSensor(
    task_id="monitor_directory_new_files",
    python_callable=check_new_files,
    op_kwargs={"directory": "/path/to/your/dag/directory"},
    poke_interval=30,  # 每30秒扫描一次
    mode="reschedule"
)

2. 外部脚本+Airflow触发API

利用文件系统的事件监听工具(Linux下的inotify-tools,Windows的FileSystemWatcher)编写脚本,监听目录的文件创建事件,一旦检测到新文件,直接调用Airflow的CLI或REST API触发DAG:

  • 示例bash脚本(基于inotify-tools):
#!/bin/bash
MONITOR_DIR="/path/to/your/dag/directory"
DAG_ID="your_target_dag_id"

inotifywait -m -e create --format '%f' "$MONITOR_DIR" | while read FILE
do
    echo "New file detected: $FILE"
    # 调用Airflow CLI触发DAG
    airflow dags trigger "$DAG_ID"
done

将该脚本作为后台服务运行,即可实时响应新文件上传事件。

3. 结合Airflow Dataset功能(Airflow 2.4+)

利用Airflow的Dataset特性实现事件驱动触发:

  • 步骤1:定义一个Dataset,指向监测的目录(或逻辑标识)
  • 步骤2:编写外部脚本,监听目录新文件,当有文件上传时,通过Airflow API标记该Dataset为更新
  • 步骤3:将目标DAG设置为依赖该Dataset,当Dataset更新时自动触发DAG运行

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 17:39:58