如何使用Python将Cloud Storage数据自动批量导入BigQuery
Python自动化实现GCS到BigQuery的分批次加载流程
前置准备
- 安装依赖包:
pip install google-cloud-bigquery google-cloud-storage - 确保你使用的账号/服务账号拥有以下权限:
- Cloud Storage存储桶的对象读取权限
- BigQuery对应数据集的编辑权限、作业创建权限
- 提前整理好27个固定的CSV文件名列表,按处理顺序排列(第一个为需要覆盖全量的文件)
核心实现代码
from google.cloud import bigquery # -------------- 配置参数,按需修改 -------------- PROJECT_ID = "你的GCP项目ID" DATASET_ID = "你的BigQuery数据集ID" TABLE_ID = "目标表名" GCS_BUCKET_NAME = "你的Cloud Storage存储桶名" # 按处理顺序排列的27个固定文件名,第一个为需要覆盖全量的文件 FIXED_FILENAMES = [ "file1.csv", "file2.csv", # 依次补充剩下的25个文件名 "file27.csv" ] # 可选:如果CSV有表头,设置为1,没有则为0 SKIP_HEADER_ROWS = 1 # --------------------------------------------- # 初始化BigQuery客户端 bq_client = bigquery.Client(project=PROJECT_ID) table_full_path = f"{PROJECT_ID}.{DATASET_ID}.{TABLE_ID}" for idx, filename in enumerate(FIXED_FILENAMES): gcs_uri = f"gs://{GCS_BUCKET_NAME}/{filename}" # 配置加载任务参数 job_config = bigquery.LoadJobConfig( source_format=bigquery.SourceFormat.CSV, field_delimiter=";", # 指定分号为分隔符 skip_leading_rows=SKIP_HEADER_ROWS, # 第一个文件覆盖原有表,后续文件追加 write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE if idx == 0 else bigquery.WriteDisposition.WRITE_APPEND, # 可选:自动探测表结构,如果表结构固定也可以手动指定schema,规避自动识别的类型误差 autodetect=True ) # 触发加载任务 load_job = bq_client.load_table_from_uri( gcs_uri, table_full_path, job_config=job_config ) # 等待任务完成 load_job.result() print(f"文件{filename}加载完成,处理模式:{'覆盖全量' if idx ==0 else '追加写入'}") print("所有27个文件全部加载完成")
优化建议
- 针对超过150GB的总数据量,你可以在
LoadJobConfig中添加priority=bigquery.QueryPriority.BATCH参数,使用批量加载模式,不会占用交互式资源配额,成本更低 - 可以添加异常捕获逻辑,记录加载失败的文件名称和错误信息,避免单次失败导致整个流程中断
- 加载完成后可以增加数据校验步骤,比如查询目标表的总行数、核心字段的非空值占比,和预期值对比确认数据完整性
- 如需实现全自动化调度,可以将脚本部署到Cloud Run或Cloud Functions,搭配Cloud Scheduler设置每周固定时间触发执行
内容的提问来源于stack exchange,提问作者Marcelo Soares
相关产品推荐
相关产品推荐

