使用Python调用BigQuery API 基于YAML配置从存储桶文件自动创建BQ表
实现思路
- 配置层设计:用yml文件存储所有非动态参数,避免硬编码,包括GCP项目信息、GCS桶规则、BQ数据集配置、字段映射规则等,后续修改规则无需改动核心代码。
- GCS文件遍历:调用GCS SDK拉取指定桶下的所有文件,过滤掉临时文件、文件夹以及非允许后缀的文件,提取文件的纯名称(去掉路径前缀和后缀)作为BQ的目标表名。
- 文件解析规则适配:根据文件后缀(csv、parquet、json等)自动匹配BQ加载的对应格式规则,支持自动推断schema,也支持从yml配置中读取固定的表字段规则。
- BQ表创建执行:调用BQ SDK构造加载任务,源路径关联GCS文件URI,目标表用提取到的文件名,配置写入模式、分区规则等属性后执行任务,完成后校验表是否创建成功。
代码样例(Python实现)
1. YML配置文件示例(config.yml)
gcp: project_id: "你的GCP项目ID" credential_path: "你的服务账号密钥文件路径" gcs: bucket_name: "目标存储桶名称" file_prefix: "" # 可选,只处理指定前缀下的文件 allowed_suffix: [".csv", ".parquet", ".json"] # 允许处理的文件后缀列表 bigquery: dataset_id: "目标BQ数据集ID" write_disposition: "WRITE_TRUNCATE" # 可选值:WRITE_APPEND/WRITE_EMPTY autodetect_schema: True # 是否自动推断表结构,设为False则需要配置下方自定义schema # 自定义表结构配置,key为文件名(不含后缀) table_schema: user_info: - name: "user_id" type: "STRING" mode: "REQUIRED" - name: "register_time" type: "TIMESTAMP" mode: "NULLABLE"
2. 核心代码
import os import yaml from google.cloud import storage, bigquery from google.oauth2 import service_account # 加载yml配置 def load_config(config_path: str = "config.yml") -> dict: with open(config_path, "r", encoding="utf-8") as f: return yaml.safe_load(f) # 初始化GCP客户端 def init_gcp_clients(config: dict): credentials = service_account.Credentials.from_service_account_file( config["gcp"]["credential_path"] ) storage_client = storage.Client(project=config["gcp"]["project_id"], credentials=credentials) bq_client = bigquery.Client(project=config["gcp"]["project_id"], credentials=credentials) return storage_client, bq_client # 获取符合要求的GCS文件列表 def get_target_files(storage_client, config: dict) -> list: bucket = storage_client.get_bucket(config["gcs"]["bucket_name"]) blobs = bucket.list_blobs(prefix=config["gcs"]["file_prefix"]) target_files = [] allowed_suffix = tuple(config["gcs"]["allowed_suffix"]) for blob in blobs: # 过滤文件夹和非允许后缀的文件 if not blob.name.endswith("/") and blob.name.endswith(allowed_suffix): target_files.append(blob) return target_files # 从GCS文件创建BQ表 def create_bq_table(bq_client, blob, config: dict): # 提取表名:去掉路径和后缀 file_full_name = os.path.basename(blob.name) table_name = os.path.splitext(file_full_name)[0] table_id = f"{config['gcp']['project_id']}.{config['bigquery']['dataset_id']}.{table_name}" # 配置加载任务 job_config = bigquery.LoadJobConfig() job_config.write_disposition = config["bigquery"]["write_disposition"] # 匹配文件格式 file_suffix = os.path.splitext(file_full_name)[1].lower() if file_suffix == ".csv": job_config.source_format = bigquery.SourceFormat.CSV job_config.skip_leading_rows = 1 # 跳过CSV表头行,可按需调整 elif file_suffix == ".parquet": job_config.source_format = bigquery.SourceFormat.PARQUET elif file_suffix == ".json": job_config.source_format = bigquery.SourceFormat.NEWLINE_DELIMITED_JSON # 配置表结构 if config["bigquery"]["autodetect_schema"]: job_config.autodetect = True else: schema = [] for field in config["bigquery"]["table_schema"].get(table_name, []): schema.append(bigquery.SchemaField(field["name"], field["type"], mode=field["mode"])) job_config.schema = schema # 执行加载任务 load_job = bq_client.load_table_from_uri( source_uris=f"gs://{config['gcs']['bucket_name']}/{blob.name}", destination=table_id, job_config=job_config ) load_job.result() # 等待任务执行完成 print(f"表 {table_id} 创建完成,共加载 {bq_client.get_table(table_id).num_rows} 行数据") if __name__ == "__main__": config = load_config() storage_client, bq_client = init_gcp_clients(config) target_files = get_target_files(storage_client, config) for blob in target_files: create_bq_table(bq_client, blob, config) print("全部表创建任务执行完毕")
3. 依赖安装
pip install google-cloud-storage google-cloud-bigquery pyyaml
注意事项
- 运行前需要确保使用的服务账号拥有GCS桶的读取权限,以及对应BQ数据集的表创建、数据写入权限。
- 如果处理的是压缩文件,需要在LoadJobConfig中添加对应压缩格式的配置,例如GZIP格式可添加
job_config.compression = bigquery.Compression.GZIP。 - 需要表分区、分桶等高级属性的话,可以在LoadJobConfig中额外添加对应参数配置。
内容的提问来源于stack exchange,提问作者SKP
相关产品推荐
相关产品推荐

