如何使用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-templatelocation:填写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
相关产品推荐
相关产品推荐

