基于Cloud Composer-Airflow实现GCS CSV触发Dataflow转BigQuery方案咨询
GCS CSV触发Airflow + Dataflow转BigQuery 完整实现示例
核心问题解决思路
1. DataflowTemplateOperator的template参数说明
你需要先把Python Dataflow作业打包成Dataflow模板上传到GCS,template参数就是这个模板文件在GCS的完整路径(格式:gs://<你的存储桶>/<模板存放路径>/<模板名>)。
生成模板的命令示例:
python csv_to_bq_dataflow.py \ --runner DataflowRunner \ --project your-gcp-project-id \ --staging_location gs://your-bucket/staging \ --temp_location gs://your-bucket/temp \ --template_location gs://your-bucket/dataflow-templates/csv-to-bq-transform-template
2. 动态配置outputTable实现多CSV对应多BQ表
核心逻辑是从触发的CSV文件名中提取目标表名:比如GCS新增customer_records.csv,就对应BigQuery表your-project:target_dataset.customer_records。
Airflow的GCSObjectSensor会返回触发的文件路径,我们通过task_instance.xcom_pull()获取该路径,解析出文件名(去掉.csv后缀),再拼接成完整的BQ表ID。
完整代码示例
1. Python Dataflow转换脚本(csv_to_bq_dataflow.py)
负责读取GCS CSV、字段重命名、写入BQ,支持通过参数传入输入路径和输出表:
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, GoogleCloudOptions, StandardOptions class RenameFields(beam.DoFn): def process(self, element): # 示例:原字段映射为新字段,按需修改 return [{ 'user_id': element['old_user_id'], 'user_name': element['old_user_name'], 'register_date': element['old_reg_date'] }] def run(): import argparse parser = argparse.ArgumentParser() parser.add_argument('--input', help='GCS CSV文件路径,例如gs://your-bucket/input/user_data.csv') parser.add_argument('--output_table', help='BigQuery表ID,例如your-project:dataset.user_data') args, pipeline_args = parser.parse_known_args() pipeline_options = PipelineOptions(pipeline_args) google_cloud_options = pipeline_options.view_as(GoogleCloudOptions) pipeline_options.view_as(StandardOptions).runner = 'DataflowRunner' with beam.Pipeline(options=pipeline_options) as p: (p | '读取CSV文件' >> beam.io.ReadFromText(args.input, skip_header_lines=1) | '解析CSV为字典' >> beam.Map(lambda line: dict(zip(['old_user_id','old_user_name','old_reg_date'], line.split(',')))) | '字段重命名转换' >> beam.ParDo(RenameFields()) | '写入BigQuery' >> beam.io.WriteToBigQuery( args.output_table, schema='user_id:STRING, user_name:STRING, register_date:DATE', # 替换为你的目标表Schema write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED ) ) if __name__ == '__main__': run()
2. Airflow DAG脚本(gcs_csv_trigger_dag.py)
通过传感器监听新增CSV,动态生成Dataflow任务参数:
from airflow import DAG from airflow.providers.google.cloud.sensors.gcs import GCSObjectSensor from airflow.providers.google.cloud.operators.dataflow import DataflowTemplateOperator from airflow.utils.dates import days_ago # 配置GCP基础信息 PROJECT_ID = 'your-gcp-project-id' GCS_BUCKET = 'your-bucket' INPUT_PREFIX = 'input/' # 监听GCS该路径下的新增文件 DATAFLOW_TEMPLATE_PATH = 'gs://your-bucket/dataflow-templates/csv-to-bq-transform-template' TARGET_DATASET = 'target_dataset' # BigQuery目标数据集 default_args = { 'owner': 'airflow', 'start_date': days_ago(1), 'retries': 1, } with DAG( 'gcs_csv_to_bq_dag', default_args=default_args, schedule_interval=None, # 由传感器触发,无需定时 catchup=False, ) as dag: # 1. 监听GCS新增CSV文件 gcs_csv_sensor = GCSObjectSensor( task_id='gcs_csv_sensor', bucket=GCS_BUCKET, object_prefix=INPUT_PREFIX, wildcard='*.csv', # 仅匹配CSV格式文件 mode='poke', poke_interval=30, # 每30秒检查一次 timeout=3600, # 超时时间1小时 ) # 2. 触发Dataflow转换任务 run_dataflow_transform = DataflowTemplateOperator( task_id='run_dataflow_csv_transform', template=DATAFLOW_TEMPLATE_PATH, project_id=PROJECT_ID, parameters={ # 动态获取触发的CSV文件路径 'input': "{{ task_instance.xcom_pull(task_ids='gcs_csv_sensor') }}", # 动态生成输出表ID:从文件名提取表名,拼接成完整BQ表ID 'output_table': "{{ '{}:{}.{}'.format('{}', '{}', task_instance.xcom_pull(task_ids='gcs_csv_sensor').split('/')[-1].replace('.csv', '')) }}".format(PROJECT_ID, TARGET_DATASET) }, location='us-central1', # 替换为你的Dataflow运行区域 ) gcs_csv_sensor >> run_dataflow_transform
注意事项
- 确保Airflow服务账号拥有GCS读取、Dataflow运行、BigQuery写入的权限
- 生成Dataflow模板时,若依赖第三方包,需通过
--requirements_file参数指定依赖文件 - BigQuery表Schema需与Dataflow转换后的字段匹配,生产环境不建议开启自动Schema推断
内容的提问来源于stack exchange,提问作者Nhu Dao
相关产品推荐
相关产品推荐

