如何触发DAG感知GCS指定路径的文件新增或更新?
监听GCS指定路径的文件变化触发Airflow DAG
方法一:使用GCSObjectsWithPrefixExistenceSensor(轮询方式)
GCSObjectUpdateSensor仅支持监听单个对象,若要监听整个路径,可替换为GCSObjectsWithPrefixExistenceSensor——该传感器能监听指定前缀(对应GCS路径)下的所有文件新增或变化。以下是修改后的DAG代码:
from airflow import DAG from airflow.providers.google.cloud.sensors.gcs import GCSObjectsWithPrefixExistenceSensor from airflow.operators.bash import BashOperator from datetime import datetime, timedelta DEFAULT_ARGS = { 'owner': 'airflow', 'depends_on_past': False, } with DAG( dag_id="test_sensor", default_args=DEFAULT_ARGS, dagrun_timeout=timedelta(hours=96), start_date=datetime(2023, 8, 16), schedule_interval=timedelta(minutes=5), # 轮询间隔,按需调整 catchup=False, ) as dag: gcs_prefix_sensor = GCSObjectsWithPrefixExistenceSensor( task_id='gcs_prefix_sensor', bucket="cs-us-ds", prefix='Stage/staging/', # 指定监听的路径前缀,末尾斜杠确保匹配该路径下的文件 timeout=300, # 超时时间,按需延长 mode='reschedule', # 释放worker资源,到下一次检查时间再调度,比默认poke模式更高效 soft_fail=True, poke_interval=60, # 每60秒检查一次,按需调整 dag=dag) def sleep_task(task_id): task = BashOperator( task_id=task_id, bash_command="sleep 20", ) return task wait_phase1 = sleep_task("wait_phase1") gcs_prefix_sensor >> wait_phase1
核心参数说明:
prefix:对应GCS的目标路径,末尾保留斜杠可避免匹配路径前缀相似的其他文件(如Stage/staging_file.txt)mode='reschedule':避免传感器持续占用worker资源,适合长周期监听场景poke_interval:设置检查间隔,平衡实时性与API调用频率
方法二:事件驱动触发(Cloud Function + Airflow API)
若想避免轮询、实现实时触发,可采用GCS事件触发Cloud Function,再通过Cloud Function调用Airflow API启动DAG:
- 创建Cloud Function,触发条件设为GCS桶
cs-us-ds中Stage/staging/路径下的文件创建/更新事件 - 在Cloud Function中编写代码调用Airflow API,示例代码如下:
import requests import os # 从环境变量读取配置 AIRFLOW_API_URL = os.environ.get('AIRFLOW_API_URL') AIRFLOW_DAG_ID = 'test_sensor' AIRFLOW_API_KEY = os.environ.get('AIRFLOW_API_KEY') def trigger_airflow_dag(event, context): # 过滤目标路径外的文件事件 file_path = event['name'] if not file_path.startswith('Stage/staging/'): return '非目标路径,跳过触发' # 调用Airflow API启动DAG headers = { 'Authorization': f'Bearer {AIRFLOW_API_KEY}', 'Content-Type': 'application/json' } url = f"{AIRFLOW_API_URL}/api/v1/dags/{AIRFLOW_DAG_ID}/dagRuns" payload = {} response = requests.post(url, headers=headers, json=payload) if response.status_code in (200, 201): return f"DAG {AIRFLOW_DAG_ID} 触发成功" else: return f"DAG触发失败: {response.text}"
该方法优势:实时响应文件变化,无需轮询,资源利用率更高,适合时效性要求高的场景。
内容的提问来源于stack exchange,提问作者Tom J Muthirenthi
相关产品推荐
相关产品推荐

