You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何触发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:

  1. 创建Cloud Function,触发条件设为GCS桶cs-us-ds中Stage/staging/路径下的文件创建/更新事件
  2. 在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 03:33:13