如何使用Airflow operator检测Google Cloud Storage的文件夹是否存在
实现方案
完全有现成的官方Operator可以直接用,不需要自定义开发。首先确保你已经安装了Airflow的Google Cloud Provider依赖:pip install apache-airflow-providers-google
核心Operator选择
你要实现的GCS指定路径存在性检查,直接用GCSPrefixSensor即可,这个Operator的作用就是检测GCS存储桶中指定前缀(也就是你说的路径)下是否存在至少1个对象。
要实现「路径不存在就直接结束DAG」的需求,只需要把参数soft_fail设置为True:当检查不到对应路径存在时,该检查任务会被标记为「已跳过」状态,所有依赖该任务的下游任务都会自动跳过,整个DAG就会直接结束。
核心参数说明
bucket:目标GCS存储桶的名称prefix:你要检查的GCS路径前缀,注意末尾可以加/避免匹配到相似前缀的其他路径gcp_conn_id:你在Airflow中配置的GCP连接ID,默认值为google_cloud_defaultpoke_interval:两次检查之间的间隔时间,单位为秒,按需配置即可timeout:检查任务的总超时时间,单位为秒,超过这个时长还没检测到路径存在就终止检查
示例DAG代码
from airflow import DAG from airflow.providers.google.cloud.sensors.gcs import GCSPrefixSensor from datetime import datetime # 你自己的下游任务,这里仅做示例 from your.operator.package import YourDownstreamOperator with DAG( dag_id="gcs_path_check_demo", start_date=datetime(2024, 1, 1), schedule_interval=None, catchup=False ) as dag: # GCS路径检查任务 check_gcs_path = GCSPrefixSensor( task_id="check_target_gcs_path_exist", bucket="your-bucket-name", prefix="your/target/path/", gcp_conn_id="your_custom_gcp_conn", soft_fail=True, poke_interval=60, timeout=300 ) # 路径存在后才会执行的下游任务 run_after_path_exist = YourDownstreamOperator( task_id="run_after_path_exist" ) # 配置任务依赖 check_gcs_path >> run_after_path_exist
补充场景
如果你要检查的是单个具体的GCS文件而非路径前缀,可以换用GCSObjectExistenceSensor,参数配置逻辑完全一致,只需要把prefix参数替换为object,填写对应文件的完整GCS路径即可。
内容的提问来源于stack exchange,提问作者Prithwiraj Samanta
相关产品推荐
相关产品推荐

