Docker Airflow ETL管道:如何配置S3 Sensor仅检测新上传文件触发任务
Airflow 检测S3新增文件触发ETL管道方案
当然可以实现这个需求,Airflow提供了原生的S3相关Sensor,通过配置参数就能忽略已有文件,仅响应新增文件:
1. 使用S3KeySensor结合last_modified_dt参数
S3KeySensor是Airflow官方的S3文件检测Sensor,通过last_modified_dt参数可以指定只检测特定时间之后上传/修改的文件,完美实现忽略已有文件的效果。
示例代码:
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor from datetime import datetime from airflow import DAG with DAG( dag_id="s3_new_file_trigger_etl", schedule_interval=None, # 若要完全由S3事件触发,设为None;也可配合定时轮询 start_date=datetime(2024, 1, 1), catchup=False, # 禁用回溯执行,避免检测历史文件 tags=["s3", "etl"] ) as dag: detect_new_s3_file = S3KeySensor( task_id="detect_new_s3_file", bucket_name="your-target-bucket", bucket_key="data/incoming/*", # 支持通配符匹配目标文件路径 wildcard_match=True, # 只检测DAG执行时间之后的文件,确保忽略已有文件 last_modified_dt=lambda context: context["execution_date"], poke_interval=60, # 每60秒轮询一次 mode="reschedule", # 无文件时释放Worker资源,适合长期等待场景 timeout=3600, # 超时时间(可选) ) # 后续ETL任务示例 # etl_task = PythonOperator(task_id="run_etl", python_callable=your_etl_function) # detect_new_s3_file >> etl_task
参数核心说明:
last_modified_dt:接收可调用对象,返回时间戳,Sensor仅检测该时间之后修改的文件。这里用DAG的execution_date作为阈值,确保只响应DAG启动后的新增文件。catchup=False:必须设置,防止Airflow回溯执行历史DAG run,避免误检测旧文件。mode="reschedule":相比默认的poke模式,会在没有文件时释放Worker资源,更适合需要长期等待的场景。
2. 更高效的实时触发方案(结合S3事件通知+Lambda)
如果需要更低延迟的触发,推荐用S3事件通知配合Lambda调用Airflow API触发DAG:
- 配置S3存储桶的事件通知,当有新文件上传时,发送事件到AWS Lambda。
- Lambda函数调用Airflow的REST API(比如
/api/v1/dags/<dag_id>/dagRuns)触发目标DAG。 - DAG启动后,用
S3KeySensor验证文件存在(此时文件是刚上传的,自然跳过已有文件)。
这种方式无需Sensor持续轮询,资源消耗更低,触发更实时。
Docker部署注意事项
- 确保Airflow容器已配置AWS权限:可以通过环境变量注入AWS密钥,或者挂载
~/.aws/credentials文件;若在AWS ECS上部署,可直接绑定IAM角色。 - 安装
apache-airflow-providers-amazon包:Docker部署时需确保Airflow镜像已包含该Provider,否则需在requirements.txt中添加并重新构建镜像。
内容的提问来源于stack exchange,提问作者nikhil kumar
相关产品推荐
相关产品推荐

