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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 06:50:29