如何用Terraform(GCP)结合Python通过Dataflow将CSV导入BigQuery
使用Terraform部署自定义Dataflow作业将CSV加载至BigQuery
核心资源说明
要完成这个任务,Terraform需要配置以下GCP资源:
- 专用服务账号(给Dataflow作业提供权限)
- GCS存储桶(用于存放CSV文件、Dataflow临时文件/日志及自定义Python脚本)
- BigQuery数据集和目标表(与CSV字段结构匹配)
- Dataflow作业(指定自定义Python脚本、参数及运行配置)
完整Terraform脚本示例
1. GCP Provider配置
terraform { required_providers { google = { source = "hashicorp/google" version = "~> 4.0" } } } provider "google" { project = "your-gcp-project-id" # 替换为你的GCP项目ID region = "us-central1" # 替换为你的目标区域 zone = "us-central1-a" }
2. 创建Dataflow专用服务账号及权限
Dataflow需要访问GCS、BigQuery的权限,创建专用账号并绑定必要角色:
# 创建服务账号 resource "google_service_account" "dataflow_sa" { account_id = "dataflow-csv-loader-sa" display_name = "Dataflow CSV to BigQuery Service Account" } # 绑定BigQuery数据编辑权限 resource "google_project_iam_member" "dataflow_sa_bq_access" { project = "your-gcp-project-id" role = "roles/bigquery.dataEditor" member = "serviceAccount:${google_service_account.dataflow_sa.email}" } # 绑定GCS对象管理权限 resource "google_project_iam_member" "dataflow_sa_gcs_access" { project = "your-gcp-project-id" role = "roles/storage.objectAdmin" member = "serviceAccount:${google_service_account.dataflow_sa.email}" } # 绑定Dataflow工作者角色 resource "google_project_iam_member" "dataflow_sa_worker" { project = "your-gcp-project-id" role = "roles/dataflow.worker" member = "serviceAccount:${google_service_account.dataflow_sa.email}" }
3. 创建GCS存储桶(若已有可跳过)
用于存储CSV文件、Python脚本及Dataflow运行所需的临时文件:
resource "google_storage_bucket" "csv_work_bucket" { name = "your-unique-bucket-name" # GCS桶名需全局唯一 location = "US" force_destroy = true # 测试环境用,生产环境请关闭 } # 上传本地CSV文件到GCS(若文件已在GCS可跳过) resource "google_storage_bucket_object" "source_csv" { name = "input/source_data.csv" bucket = google_storage_bucket.csv_work_bucket.name source = "./local-path/to/your/data.csv" # 本地CSV文件路径 } # 上传自定义Python脚本到GCS resource "google_storage_bucket_object" "dataflow_script" { name = "scripts/csv_to_bq.py" bucket = google_storage_bucket.csv_work_bucket.name source = "./local-path/to/your/csv_to_bq.py" # 本地Python脚本路径 }
4. 创建BigQuery数据集和目标表
提前定义与CSV匹配的表结构:
# BigQuery数据集 resource "google_bigquery_dataset" "target_dataset" { dataset_id = "csv_import_dataset" location = "US" delete_contents_on_destroy = true # 测试环境用,生产环境请关闭 } # BigQuery目标表(schema需与CSV字段完全匹配) resource "google_bigquery_table" "target_table" { dataset_id = google_bigquery_dataset.target_dataset.dataset_id table_id = "csv_import_table" schema = <<EOF [ { "name": "id", "type": "INTEGER", "mode": "REQUIRED" }, { "name": "username", "type": "STRING", "mode": "NULLABLE" }, { "name": "score", "type": "FLOAT", "mode": "NULLABLE" } ] EOF }
5. 部署Dataflow作业
指定自定义Python脚本路径、运行参数及服务账号:
resource "google_dataflow_job" "csv_to_bq_job" { name = "csv-to-bigquery-loader" python_file = "gs://${google_storage_bucket.csv_work_bucket.name}/scripts/csv_to_bq.py" service_account_email = google_service_account.dataflow_sa.email region = "us-central1" # 传递给Python脚本的参数 parameters = { input_path = "gs://${google_storage_bucket.csv_work_bucket.name}/input/source_data.csv" output_table = "${google_bigquery_table.target_table.project}.${google_bigquery_table.target_table.dataset_id}.${google_bigquery_table.target_table.table_id}" delimiter = "," } temp_gcs_location = "gs://${google_storage_bucket.csv_work_bucket.name}/temp" staging_location = "gs://${google_storage_bucket.csv_work_bucket.name}/staging" # 确保依赖资源创建完成后再启动作业 depends_on = [ google_project_iam_member.dataflow_sa_bq_access, google_project_iam_member.dataflow_sa_gcs_access, google_bigquery_table.target_table ] }
自定义Python脚本示例(需提前编写)
用Apache Beam实现CSV读取与BigQuery写入逻辑:
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, SetupOptions class ParseCSV(beam.DoFn): def process(self, element, delimiter=','): # 按分隔符拆分CSV行,需与BigQuery schema对应 id, username, score = element.split(delimiter) return [{'id': int(id), 'username': username.strip(), 'score': float(score)}] def run(): import argparse parser = argparse.ArgumentParser() parser.add_argument('--input_path', required=True) parser.add_argument('--output_table', required=True) parser.add_argument('--delimiter', default=',') known_args, pipeline_args = parser.parse_known_args() options = PipelineOptions(pipeline_args) options.view_as(SetupOptions).save_main_session = True with beam.Pipeline(options=options) as p: (p | '读取CSV文件' >> beam.io.ReadFromText(known_args.input_path, skip_header_lines=1) | '解析CSV行' >> beam.ParDo(ParseCSV(), delimiter=known_args.delimiter) | '写入BigQuery' >> beam.io.WriteToBigQuery( known_args.output_table, write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED )) if __name__ == '__main__': run()
部署步骤
- 将上述Terraform代码保存为
main.tf,替换所有占位符为实际值 - 运行
terraform init初始化Terraform环境 - 运行
terraform plan预览要创建的资源 - 运行
terraform apply确认并创建所有资源,启动Dataflow作业
内容的提问来源于stack exchange,提问作者TEJAS
相关产品推荐
相关产品推荐

