在数据处理平台中,经常需要把业务端生成的文件上传到指定位置,随后自动启动清洗、校验或入库任务。Airflow 作为主流编排工具,提供了多种机制将文件到达事件转化为工作流执行。通过合理的 Sensor 设计与调度配置,可以避免人工干预和粗粒度定时轮询带来的资源浪费。

基于 FileSensor 的本地目录监听
Airflow 自带的 FileSensor 是最简单的文件触发方式。它会周期性检查某个文件路径是否存在,一旦命中便放行后续任务。这种方式适合单节点或共享挂载卷的场景,比如 NFS 目录。
使用 FileSensor 时需要注意,文件路径支持精确文件或目录。若监控目录,需结合 file_pattern 参数过滤。下面的例子展示了一个等待 CSV 文件出现的 DAG:
from airflow import DAG
from airflow.sensors.filesystem import FileSensor
from airflow.operators.dummy import DummyOperator
from datetime import datetime
with DAG(
dag_id='file_upload_trigger_demo',
start_date=datetime(2023, 1, 1),
schedule_interval=None,
catchup=False
) as dag:
wait_file = FileSensor(
task_id='wait_csv_file',
filepath='/shared/inbound/',
file_pattern='*.csv',
poke_interval=30,
timeout=3600,
soft_fail=False
)
process = DummyOperator(task_id='process_file')
wait_file >> process
上述代码中,poke_interval 控制轮询频率,timeout 定义最长等待时间。若超过一小时文件仍未到达,任务会失败并触发告警。这种机制的优点是零额外组件,缺点是在分布式 Worker 环境下,必须保证所有 Worker 都能访问同一文件系统,否则会出现某节点感知不到文件的问题。
另外,FileSensor 只判断存在性,不感知文件写入完成。如果文件很大,上传过程中就可能被检测到,导致后续任务读取到不完整数据。实践中常要求上传方写临时名再重命名,或配合文件大小稳定判断。
自定义 Sensor 监听对象存储上传
当文件存放在 S3、MinIO 或兼容 S3 的对象存储时,可以使用 S3KeySensor,或者编写自定义 Sensor 轮询 API。对于内网系统,如果存储提供事件回调,也可以由外部服务调用 Airflow 的 REST API 触发 DAG,实现真正的事件驱动。
下面示例展示一个自定义 Sensor,用于检查某个 HTTP 可访问目录中是否出现标记文件:
import requests
from airflow.sensors.base import BaseSensorOperator
from airflow.utils.decorators import apply_defaults
class HttpMarkerSensor(BaseSensorOperator):
@apply_defaults
def __init__(self, marker_url, *args, **kwargs):
super().__init__(*args, **kwargs)
self.marker_url = marker_url
def poke(self, context):
try:
resp = requests.head(self.marker_url, timeout=10)
return resp.status_code == 200
except requests.RequestException:
return False
# 在 DAG 中使用
from airflow import DAG
from datetime import datetime
with DAG('http_marker_trigger', start_date=datetime(2023, 1, 1), schedule_interval=None) as dag:
wait_marker = HttpMarkerSensor(
task_id='wait_marker',
marker_url='http://192.168.0.1:9000/inbound/done.txt',
poke_interval=60,
timeout=7200
)
自定义 Sensor 继承自 BaseSensorOperator,只需实现 poke 方法返回布尔值。这样可以将任意协议或业务接口接入 Airflow 触发逻辑。相比单纯文件监听,它更灵活,也能绕过共享文件系统限制。
需要强调的是,无论哪种 Sensor,都应设置合理的超时与重试,避免因为上游延迟造成 DAG 阻塞。同时,触发后的任务要具备幂等性,防止同一文件被重复处理。
通过 REST API 由上传服务直接触发
如果上传程序本身可控,最佳实践是上传完成后调用 Airflow REST API 触发指定 DAG,并传入文件路径作为参数。这样就彻底去掉轮询,实时性最高。
示例中使用 curl 触发,实际代码中可用 Python 的 requests 库携带认证信息调用。DAG 侧通过 dag_run.conf 接收参数:
curl -X POST
http://127.0.0.1:8080/api/v1/dags/upload_pipeline/dagRuns
-H 'Content-Type: application/json'
-d '{"conf": {"file_path": "/inbound/2023_data.csv"}}'
对应的 DAG 读取参数并传递给处理任务:
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
def handle_file(**kwargs):
path = kwargs['dag_run'].conf.get('file_path')
print('processing', path)
with DAG('upload_pipeline', start_date=datetime(2023, 1, 1), schedule_interval=None) as dag:
run = PythonOperator(
task_id='handle_file',
python_callable=handle_file,
provide_context=True
)
这种方案将触发职责交给上游,Airflow 只负责编排执行。它的延迟最低,也最容易做权限隔离。不过需要保证 API 安全,避免未授权触发。结合网关注册与令牌校验,可满足生产要求。
综合来看,小规模本地环境可用 FileSensor 快速落地;对象存储或跨网络场景适合自定义 Sensor 或事件回调。选择时需权衡实时性、运维成本与系统耦合度。
常见误区与处理建议
一个典型误区是认为 Sensor 能感知文件内容完整。实际上多数轮询方式只检查存在或前缀,大文件上传中途就可能被捕获。应通过约定命名规范或上传端写入完成标记来规避。
另一个问题是把 Sensor 挂在默认调度池,导致大量等待任务占用槽位。可以为 Sensor 单独配置资源池,或降低 poke_interval 以减轻数据库压力。正确配置后,文件上传触发任务能在 Airflow 中稳定高效运行。