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

基于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定时触发

  1. 打开GCP Cloud Scheduler控制台,创建新任务
  2. 频率:输入0 0 * * *(每天凌晨执行,可按需调整)
  3. 目标类型:选择「Dataproc作业」
  4. 集群配置:选择「创建临时集群」(避免长期运维成本),指定区域、机器类型
  5. 作业配置:
    • 作业类型:PySpark
    • 主Python文件:gs://your-bucket/scripts/incremental_upload.py
    • 服务账号:选择拥有Dataproc、GCS、BigQuery权限的账号
  6. 保存任务,每日将自动执行增量上传

方案二:基于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定时触发

  1. 打开Cloud Scheduler控制台,创建新任务
  2. 频率:输入0 0 * * *
  3. 目标类型:选择「HTTP」
  4. URL:填写Dataflow作业的触发URL(或使用gcloud dataflow jobs run命令作为HTTP请求的内容)
  5. 保存任务,每日自动执行增量上传

方案对比

  • Dataproc:适合大规模数据处理,支持复杂数据转换,Spark性能优异;但需要管理集群(或使用临时集群),运维成本略高
  • Dataflow:完全托管的Serverless服务,自动扩缩容,代码简洁,适合标准化ETL流程;对本地文件的依赖处理需注意权限配置

内容的提问来源于stack exchange,提问作者James Bond

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 22:50:22