Cloud Composer新手求助:基于Airflow传感器实现GCS-CSV-Dataflow-BigQuery流程
Cloud Composer + Dataflow + BigQuery 自动处理GCS新增CSV文件实现方案
一、前置环境准备
- 启用Google Cloud相关服务:Cloud Composer、Dataflow、BigQuery、Cloud Storage API
- 创建两个GCS存储桶:一个用于存放源文件(如
your-input-bucket),一个用于Dataflow临时文件存储(如your-temp-bucket) - 在BigQuery中创建目标数据集(如
csv_processing)和对应表,表结构需匹配转换后的字段,示例SQL:CREATE TABLE csv_processing.target_table ( original_id STRING, user_name STRING, order_amount INT64 ) - 确保Cloud Composer的服务账号拥有以下权限:GCS存储桶读写权限、Dataflow作业提交权限、BigQuery表写入权限
二、Dataflow 转换代码实现
编写Python脚本csv_to_bq_transform.py,完成CSV读取、字段重命名、格式转换及BigQuery写入:
import argparse import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, GoogleCloudOptions from apache_beam.io.gcp.bigquery import WriteToBigQuery import csv class TransformCsvRow(beam.DoFn): def process(self, element): # 解析源CSV行(假设源字段为:id, name, amount) reader = csv.DictReader([element], fieldnames=['id', 'name', 'amount']) row = next(reader) # 字段重命名+格式转换 transformed_row = { 'original_id': row['id'], 'user_name': row['name'].strip(), 'order_amount': int(row['amount']) if row['amount'].isdigit() else None } yield transformed_row def run(): parser = argparse.ArgumentParser() parser.add_argument('--input_gcs_path', required=True, help='GCS源文件路径,如gs://your-input-bucket/test.csv') parser.add_argument('--output_bq_table', required=True, help='BigQuery目标表,如your-project-id:csv_processing.target_table') parser.add_argument('--temp_gcs_location', required=True, help='Dataflow临时文件路径,如gs://your-temp-bucket/temp/') known_args, pipeline_args = parser.parse_known_args() pipeline_options = PipelineOptions(pipeline_args) google_cloud_options = pipeline_options.view_as(GoogleCloudOptions) google_cloud_options.temp_location = known_args.temp_gcs_location with beam.Pipeline(options=pipeline_options) as p: ( p | '读取CSV文件' >> beam.io.ReadFromText(known_args.input_gcs_path, skip_header_lines=1) | '字段转换与重命名' >> beam.ParDo(TransformCsvRow()) | '写入BigQuery' >> WriteToBigQuery( known_args.output_bq_table, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED ) ) if __name__ == '__main__': run()
将该脚本上传到GCS的固定路径,比如gs://your-temp-bucket/dataflow-scripts/csv_to_bq_transform.py
三、Airflow DAG 编写(核心触发逻辑)
编写DAG文件gcs_csv_trigger_dag.py,通过GCSObjectSensor监听新增CSV文件,自动触发Dataflow作业:
from airflow import DAG from airflow.providers.google.cloud.sensors.gcs import GCSObjectSensor from airflow.providers.google.cloud.operators.dataflow import DataflowPythonOperator from datetime import datetime, timedelta 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_csv_auto_trigger', default_args=default_args, description='监听GCS新增CSV文件,自动触发Dataflow转换至BigQuery', schedule_interval=None, # 依赖传感器触发,不设置定时调度 catchup=False, tags=['gcs', 'dataflow', 'bigquery'], ) as dag: # 监听GCS中新增的CSV文件 gcs_csv_sensor = GCSObjectSensor( task_id='wait_for_new_csv', bucket='your-input-bucket', object_pattern='*.csv', # 匹配所有CSV后缀文件 mode='reschedule', # 文件未出现时释放worker资源,定时重试 poke_interval=30, # 每30秒检查一次 timeout=3600, # 超时时间1小时,可按需调整 ) # 触发Dataflow转换作业 execute_dataflow = DataflowPythonOperator( task_id='run_csv_transform', py_file='gs://your-temp-bucket/dataflow-scripts/csv_to_bq_transform.py', options={ 'input_gcs_path': 'gs://your-input-bucket/{{ task_instance.xcom_pull(task_ids="wait_for_new_csv") }}', 'output_bq_table': 'your-project-id:csv_processing.target_table', 'temp_gcs_location': 'gs://your-temp-bucket/temp/', 'project': 'your-project-id', 'region': 'us-central1', # 替换为你的GCP区域 }, dag=dag, ) # 设置任务依赖:传感器检测到文件后触发Dataflow作业 gcs_csv_sensor >> execute_dataflow
四、部署与测试
- 部署DAG:
- 通过GCS浏览器将DAG文件上传到Cloud Composer的DAG目录:
gs://your-composer-bucket/dags/ - 或使用gcloud命令:
gcloud composer environments storage dags import --environment your-composer-env --location your-region --source gcs_csv_trigger_dag.py
- 通过GCS浏览器将DAG文件上传到Cloud Composer的DAG目录:
- 验证DAG加载:登录Airflow UI,确认
gcs_csv_auto_triggerDAG已启用 - 测试流程:
- 上传一个符合格式的CSV文件到
your-input-bucket - 在Airflow UI中查看
wait_for_new_csv任务状态,成功后会自动触发run_csv_transform - 前往Dataflow控制台确认作业运行,完成后检查BigQuery目标表是否有数据写入
- 上传一个符合格式的CSV文件到
关键注意事项
- 重复文件处理:若需避免重复处理,可在Dataflow作业完成后将源文件移动至归档目录,或在传感器中通过文件名前缀/后缀过滤已处理文件
- CSV格式兼容:如果CSV包含引号、自定义分隔符,建议改用
beam.io.ReadFromCsv或csv.DictReader解析,避免手动拆分出错 - 资源优化:传感器使用
reschedule模式比poke模式更节省worker资源,适合长期监听场景
内容的提问来源于stack exchange,提问作者Nhu Dao
相关产品推荐
相关产品推荐

