基于GCP的CSV文件自动化上传ETL方案需求(含历史与增量文件)
GCP CSV自动化上传至BigQuery方案(Dataproc/Dataflow)
前置准备
- 确认GCP服务器已安装并配置好GCP SDK(已完成认证,可正常执行
gcloud/gsutil命令) - 创建GCS存储桶(用于临时存储CSV文件):
gsutil mb gs://your-bucket-name - 在BigQuery创建目标数据集及表,提前定义好与CSV匹配的Schema(如字段名、数据类型)
- 确保操作账号拥有以下权限:GCS对象创建/读取、BigQuery数据编辑、Dataproc/Dataflow作业提交
方案一:基于Dataproc的自动化流程
1. 历史CSV批量上传
步骤1:创建Dataproc集群
通过GCP控制台或命令行创建单节点/多节点集群(测试用单节点足够):
gcloud dataproc clusters create csv-upload-cluster \ --region us-central1 \ --single-node \ --master-machine-type n1-standard-2 \ --service-account your-service-account@your-project.iam.gserviceaccount.com
步骤2:编写PySpark批量处理脚本
将以下代码保存为historical_upload.py,上传至GCS存储桶(如gs://your-bucket/scripts/):
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 替换为你的BigQuery表Schema csv_schema = StructType([ StructField("user_id", StringType(), nullable=True), StructField("order_amount", IntegerType(), nullable=True), StructField("order_date", StringType(), nullable=True) ]) spark = SparkSession.builder \ .appName("HistoricalCSVtoBQ") \ .config("spark.jars.packages", "com.google.cloud.spark:spark-bigquery-with-dependencies_2.12:0.35.1") \ .getOrCreate() # 读取本地所有CSV文件 df = spark.read.csv("C:/myDrive/*.csv", header=True, schema=csv_schema) # 可选:将原始CSV备份到GCS df.write.mode("overwrite").csv("gs://your-bucket/historical_backup/", header=True) # 写入BigQuery(覆盖现有表) df.write \ .format("bigquery") \ .option("table", "your-project.your-dataset.your-target-table") \ .option("temporaryGcsBucket", "gs://your-bucket/temp") \ .mode("overwrite") \ .save() spark.stop()
步骤3:提交Dataproc作业
gcloud dataproc jobs submit pyspark gs://your-bucket/scripts/historical_upload.py \ --cluster csv-upload-cluster \ --region us-central1
步骤4:清理资源(可选)
作业完成后删除集群节省成本:
gcloud dataproc clusters delete csv-upload-cluster --region us-central1
2. 每日增量CSV自动上传
步骤1:编写增量处理脚本
保存为incremental_upload.py并上传至GCS:
import os from datetime import datetime, timedelta from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType csv_schema = StructType([ StructField("user_id", StringType(), nullable=True), StructField("order_amount", IntegerType(), nullable=True), StructField("order_date", StringType(), nullable=True) ]) spark = SparkSession.builder \ .appName("IncrementalCSVtoBQ") \ .config("spark.jars.packages", "com.google.cloud.spark:spark-bigquery-with-dependencies_2.12:0.35.1") \ .getOrCreate() # 筛选最近24小时内新增/修改的CSV文件 cutoff_time = datetime.now() - timedelta(hours=24) target_files = [] for root, _, files in os.walk("C:/myDrive"): for file in files: if file.endswith(".csv"): file_path = os.path.join(root, file) mtime = datetime.fromtimestamp(os.path.getmtime(file_path)) if mtime >= cutoff_time: target_files.append(file_path) if not target_files: print("No new CSV files to process") spark.stop() exit() # 读取增量文件并写入BigQuery(追加模式) df = spark.read.csv(target_files, header=True, schema=csv_schema) df.write \ .format("bigquery") \ .option("table", "your-project.your-dataset.your-target-table") \ .option("temporaryGcsBucket", "gs://your-bucket/temp") \ .mode("append") \ .save() spark.stop()
步骤2:配置Cloud Scheduler定时触发
- 打开GCP Cloud Scheduler控制台,创建新任务
- 频率:输入
0 0 * * *(每天凌晨执行,可按需调整) - 目标类型:选择「Dataproc作业」
- 集群配置:选择「创建临时集群」(避免长期运维成本),指定区域、机器类型
- 作业配置:
- 作业类型:PySpark
- 主Python文件:
gs://your-bucket/scripts/incremental_upload.py - 服务账号:选择拥有Dataproc、GCS、BigQuery权限的账号
- 保存任务,每日将自动执行增量上传
方案二:基于Dataflow的自动化流程
1. 历史CSV批量上传
步骤1:编写Dataflow批处理脚本
保存为historical_dataflow.py:
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, GoogleCloudOptions, StandardOptions def parse_csv(line): # 按CSV分隔符拆分字段,替换为你的字段顺序 user_id, order_amount, order_date = line.split(',') return { "user_id": user_id, "order_amount": int(order_amount), "order_date": order_date } def run(): options = PipelineOptions() gcp_options = options.view_as(GoogleCloudOptions) gcp_options.project = "your-project" gcp_options.job_name = "historical-csv-to-bq" gcp_options.staging_location = "gs://your-bucket/staging" gcp_options.temp_location = "gs://your-bucket/temp" options.view_as(StandardOptions).runner = "DataflowRunner" options.view_as(StandardOptions).region = "us-central1" with beam.Pipeline(options=options) as p: (p | "ReadCSVFiles" >> beam.io.ReadFromText("C:/myDrive/*.csv", skip_header_lines=1) | "ParseCSV" >> beam.Map(parse_csv) | "WriteToBQ" >> beam.io.WriteToBigQuery( table="your-project.your-dataset.your-target-table", schema="user_id:STRING, order_amount:INTEGER, order_date:STRING", write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED ) ) if __name__ == "__main__": run()
步骤2:提交Dataflow作业
python historical_dataflow.py
2. 每日增量CSV自动上传
步骤1:编写增量Dataflow脚本
保存为incremental_dataflow.py:
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, GoogleCloudOptions, StandardOptions import os from datetime import datetime, timedelta def filter_recent_files(file_name): cutoff_time = datetime.now() - timedelta(hours=24) file_path = os.path.join("C:/myDrive", file_name) mtime = datetime.fromtimestamp(os.path.getmtime(file_path)) return mtime >= cutoff_time def parse_csv(line): user_id, order_amount, order_date = line.split(',') return { "user_id": user_id, "order_amount": int(order_amount), "order_date": order_date } def run(): options = PipelineOptions() gcp_options = options.view_as(GoogleCloudOptions) gcp_options.project = "your-project" gcp_options.job_name = "incremental-csv-to-bq" gcp_options.staging_location = "gs://your-bucket/staging" gcp_options.temp_location = "gs://your-bucket/temp" options.view_as(StandardOptions).runner = "DataflowRunner" options.view_as(StandardOptions).region = "us-central1" with beam.Pipeline(options=options) as p: (p | "ListCSVFiles" >> beam.Create([f for f in os.listdir("C:/myDrive") if f.endswith(".csv")]) | "FilterRecent" >> beam.Filter(filter_recent_files) | "ReadCSV" >> beam.io.ReadFromText(lambda f: os.path.join("C:/myDrive", f), skip_header_lines=1) | "ParseCSV" >> beam.Map(parse_csv) | "AppendToBQ" >> beam.io.WriteToBigQuery( table="your-project.your-dataset.your-target-table", schema="user_id:STRING, order_amount:INTEGER, order_date:STRING", write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED ) ) if __name__ == "__main__": run()
步骤2:配置Cloud Scheduler定时触发
- 打开Cloud Scheduler控制台,创建新任务
- 频率:输入
0 0 * * * - 目标类型:选择「HTTP」
- URL:填写Dataflow作业的触发URL(或使用
gcloud dataflow jobs run命令作为HTTP请求的内容) - 保存任务,每日自动执行增量上传
方案对比
- Dataproc:适合大规模数据处理,支持复杂数据转换,Spark性能优异;但需要管理集群(或使用临时集群),运维成本略高
- Dataflow:完全托管的Serverless服务,自动扩缩容,代码简洁,适合标准化ETL流程;对本地文件的依赖处理需注意权限配置
内容的提问来源于stack exchange,提问作者James Bond
相关产品推荐
相关产品推荐

