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

如何使用Cloud Composer实现GCS文件到达时触发Dataflow作业

实现GCS文件触发Dataflow作业的Cloud Composer DAG配置方案

前置准备

  • 你的Cloud Composer环境已经正常运行,且绑定的服务账号拥有GCS存储对象读取、Dataflow作业创建、临时存储路径写入的权限
  • 待监听的GCS桶路径、Dataflow作业的模板路径(可以是预定义模板或者自行打包的自定义模板)、Dataflow运行所需的参数(比如输入输出路径、临时路径、区域等)已经确认无误

DAG核心实现逻辑

你只需要用到两个GCP provider自带的Airflow运算符即可完成需求:

1. GCS文件监听功能:使用GCSObjectsWithPrefixExistenceSensor

这个传感器专门用来监听指定GCS路径下是否有符合前缀的新文件抵达,核心配置参数如下:

  • bucket:填写要监听的GCS桶名,不需要带gs://前缀
  • prefix:填写要监听的路径前缀,比如要监听gs://my-bucket/input/下的所有文件,就填input/
  • poke_interval:每次轮询检查的间隔时间,单位为秒,POC阶段可以设为30,生产环境根据延迟要求调整
  • mode:推荐设为reschedule,避免长时间占用worker槽位

2. Dataflow作业触发功能:使用DataflowTemplatedJobStartOperator

如果使用自定义模板或者谷歌官方预定义Dataflow模板,用这个运算符即可直接触发作业,核心配置参数如下:

  • template:填写Dataflow模板的GCS全路径,比如gs://my-bucket/templates/my-etl-template
  • location:填写Dataflow作业运行的区域,和Composer区域保持一致可以减少跨区延迟
  • parameters:以字典格式传入Dataflow作业需要的运行参数,比如输入路径、输出路径等
  • wait_until_finished:如果需要DAG等待作业运行完成再结束就设为True,不需要的话设为False,作业触发后DAG直接标记成功

完整DAG示例代码

from airflow import DAG
from airflow.providers.google.cloud.sensors.gcs import GCSObjectsWithPrefixExistenceSensor
from airflow.providers.google.cloud.operators.dataflow import DataflowTemplatedJobStartOperator
from datetime import datetime, timedelta

# 配置DAG默认参数
default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

with DAG(
    'gcs_trigger_dataflow_poc',
    default_args=default_args,
    description='POC: GCS文件抵达自动触发Dataflow作业',
    schedule_interval=timedelta(minutes=10), # 调度周期可根据需求调整
    catchup=False,
    tags=['gcp', 'poc', 'etl'],
) as dag:

    # 任务1:监听GCS指定路径的新文件
    gcs_file_sensor = GCSObjectsWithPrefixExistenceSensor(
        task_id='gcs_listen_new_file',
        bucket='your-monitor-bucket-name', # 替换为实际监听桶名
        prefix='input/', # 替换为实际要监听的路径前缀
        poke_interval=30,
        mode='reschedule',
    )

    # 任务2:触发Dataflow模板作业
    run_dataflow_job = DataflowTemplatedJobStartOperator(
        task_id='trigger_dataflow_etl',
        template='gs://your-template-bucket/path/to/your-template', # 替换为实际模板路径
        location='us-central1', # 替换为实际作业运行区域
        parameters={
            'inputFile': 'gs://your-monitor-bucket-name/input/*', # 替换为实际输入路径参数
            'outputPath': 'gs://your-output-bucket/output/', # 替换为实际输出路径参数
            'tempLocation': 'gs://your-temp-bucket/tmp/', # 替换为实际临时文件路径
        },
        wait_until_finished=True,
    )

    # 定义任务依赖:监听到文件后再触发Dataflow作业
    gcs_file_sensor >> run_dataflow_job

注意事项

  • 传感器监听到文件后不会自动删除或者标记已处理,POC验证通过后如果需要避免重复处理相同文件,可以在DAG中新增一个GCS文件移动或者删除的任务,放到Dataflow作业运行成功之后执行
  • 如果使用Dataflow Flex模板运行作业,把运算符替换为DataflowFlexTemplateJobOperator即可,参数配置逻辑基本一致
  • 所有用到的GCS路径都要保证Composer的服务账号有对应的读写权限,否则会出现权限报错

内容的提问来源于stack exchange,提问作者Shriyut Jha

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 22:21:01