如何在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
相关产品推荐
相关产品推荐

