如何创建DAG中的Sensor检测GCS指定位置匹配通配符的对象
GCS通配符文件检测Sensor及对应DAG实现
前置依赖
首先安装Airflow Google Cloud官方provider包:
pip install apache-airflow-providers-google>=10.0.0
核心配置逻辑
- 检测频率:将Sensor的检测间隔设为3600秒,满足每1小时执行一次检测的要求
- 通配符匹配:直接在对象路径参数中传入glob格式的通配符表达式,底层会自动拉取GCS路径下的对象列表完成匹配
- 待调度状态配置:将Sensor运行模式设为
reschedule,未检测到匹配文件时会主动释放Worker资源,将任务标记为待调度状态,间隔1小时后自动重新触发检测,不会长期占用计算资源 - 下游触发逻辑:Sensor检测到匹配文件后会自动标记为成功状态,按照DAG依赖关系触发后续任务执行
完整可运行DAG代码
from datetime import datetime from airflow import DAG from airflow.providers.google.cloud.sensors.gcs import GCSObjectExistenceSensor from airflow.operators.python import PythonOperator # 按需修改配置项 GCS_TARGET_BUCKET = "替换为你的GCS桶名" GCS_FILE_MATCH_PATTERN = "业务路径/*/dt=*/data_*.parquet" # 支持*、**等glob通配符 CHECK_INTERVAL = 3600 # 检测间隔,单位秒,对应1小时 default_args = { "owner": "data-team", "start_date": datetime(2024, 1, 1), "retries": 0, } def file_process_logic(): """检测到匹配文件后执行的下游业务逻辑""" print("已检测到符合规则的GCS文件,启动后续处理流程") # 此处替换为实际业务逻辑,如文件加载、数据清洗、入库等 with DAG( dag_id="gcs_wildcard_file_monitor", default_args=default_args, schedule_interval=None, catchup=False, tags=["gcs", "monitor", "sensor"], ) as dag: # GCS通配符文件检测任务 gcs_file_check = GCSObjectExistenceSensor( task_id="check_matched_gcs_file", bucket=GCS_TARGET_BUCKET, object=GCS_FILE_MATCH_PATTERN, poke_interval=CHECK_INTERVAL, mode="reschedule", google_cloud_conn_id="google_cloud_default", # 提前在Airflow中配置GCS连接凭证 ) # 下游处理任务 file_process_task = PythonOperator( task_id="process_target_file", python_callable=file_process_logic, ) gcs_file_check >> file_process_task
配置注意事项
- 提前在Airflow连接管理中配置GCS服务账号,确保绑定的服务账号拥有目标GCS桶的
storage.objects.list权限,否则无法完成对象匹配检测 - 必须设置
mode="reschedule",如果使用默认的poke模式,任务会持续占用Worker槽位直到检测到文件,不适合1小时间隔的长周期检测场景 - 通配符匹配遵循标准glob规则:
*匹配单层路径下的任意字符,**匹配递归多层路径,可根据实际文件命名规则调整匹配表达式 - 若需要限制最长检测周期,可给Sensor传入
timeout参数(单位秒),超过设定时间未检测到文件则任务标记为失败,避免无限等待
内容的提问来源于stack exchange,提问作者Himanshu Sharma
相关产品推荐
相关产品推荐

