将CSV文件上传至BigQuery分区表(从文件名生成分区键)
实现按文件名日期分区的BigQuery上传修改方案
以下是针对你的需求的具体代码修改步骤,分两种常用方案供你选择:
方案一:直接加载到指定日期分区(无需修改CSV内容)
这种方案不需要在CSV中新增字段,直接将文件数据写入对应日期的分区,表按天进行时间分区。
修改后的完整代码
import re from datetime import datetime from google.cloud import bigquery client = bigquery.Client.from_service_account_json(CREDENTIALS_LOCATION) def upload_from_gcs_to_bq(project_id, dataset_id, gsutil_uri, table_name,gcs_blob): # 从文件名提取日期 match = re.search(r'\d{4}-\d{2}-\d{2}', gcs_blob) if not match: raise ValueError(f"文件名{gcs_blob}中未找到有效日期格式(需符合YYYY-MM-DD)") partition_date_str = match.group() # 转换为BigQuery分区要求的YYYYMMDD格式 partition_suffix = datetime.strptime(partition_date_str, '%Y-%m-%d').strftime('%Y%m%d') # 构建带分区后缀的表ID base_table_id = f"{project_id}.{dataset_id}.{table_name}" partitioned_table_id = f"{base_table_id}${partition_suffix}" uri = f"{gsutil_uri}/{gcs_blob}.csv" job_config = bigquery.LoadJobConfig( schema=[ bigquery.SchemaField("filename", "STRING"), bigquery.SchemaField("sales_category", "STRING"), # 保留你原有的其他字段 ... ], skip_leading_rows=1, # 配置表为按天分区的时间分区表 time_partitioning=bigquery.TimePartitioning( type_=bigquery.TimePartitioningType.DAY, expiration_ms=7776000000, # 可选:90天后自动删除分区数据 ), # 写入策略:追加到指定分区(如果分区已存在则追加,不存在则创建) write_disposition=bigquery.WriteDisposition.WRITE_APPEND, ) # 执行加载任务到指定分区 load_job = client.load_table_from_uri( uri, partitioned_table_id, job_config=job_config ) load_job.result() # 等待任务完成 table = client.get_table(base_table_id) print(f"数据已成功加载到分区表 {base_table_id} 的 {partition_date_str} 分区,当前表总行数:{table.num_rows}") def main(): upload_from_gcs_to_bq(project_id, dataset_id, gsutil_uri, table_name,gcs_blob) if __name__ == '__main__': main()
方案说明
- 日期提取:从传入的
gcs_blob(文件名)中提取日期,并转换为BigQuery分区要求的YYYYMMDD格式,作为表名后缀(通过$分隔)。 - 分区配置:通过
time_partitioning将表设置为按天分区的时间分区表,首次加载时会自动创建该分区表。 - 写入策略:使用
WRITE_APPEND确保同一分区的数据不会被覆盖,而是追加写入。
方案二:基于表内date字段的分区(推荐用于需保留日期字段的场景)
如果需要在表中保留日期字段,且以该字段作为分区键,可以采用此方案,加载时自动从文件名提取日期并填充到新增的date字段。
修改后的完整代码
import re from datetime import datetime from google.cloud import bigquery client = bigquery.Client.from_service_account_json(CREDENTIALS_LOCATION) def upload_from_gcs_to_bq(project_id, dataset_id, gsutil_uri, table_name,gcs_blob): # 从文件名提取日期 match = re.search(r'\d{4}-\d{2}-\d{2}', gcs_blob) if not match: raise ValueError(f"文件名{gcs_blob}中未找到有效日期格式(需符合YYYY-MM-DD)") partition_date_str = match.group() base_table_id = f"{project_id}.{dataset_id}.{table_name}" uri = f"{gsutil_uri}/{gcs_blob}.csv" # 定义外部表配置,用于读取GCS中的CSV external_config = bigquery.ExternalConfig("CSV") external_config.source_uris = [uri] external_config.schema = [ bigquery.SchemaField("filename", "STRING"), bigquery.SchemaField("sales_category", "STRING"), # 保留你原有的其他字段 ... ] external_config.skip_leading_rows = 1 # 构建查询,新增date字段并填充从文件名提取的日期 sql = f""" SELECT *, DATE('{partition_date_str}') AS date FROM EXTERNAL_QUERY('{project_id}.{dataset_id}.temp_external_config', '') """ job_config = bigquery.QueryJobConfig( destination=base_table_id, # 配置按date字段进行天分区 time_partitioning=bigquery.TimePartitioning( type_=bigquery.TimePartitioningType.DAY, field="date", # 指定分区键为date字段 expiration_ms=7776000000, # 可选:90天后自动删除分区数据 ), # 写入策略:追加数据 write_disposition=bigquery.WriteDisposition.WRITE_APPEND, # 关联外部表配置 external_configs=[external_config], ) # 执行查询并写入分区表 query_job = client.query(sql, job_config=job_config) query_job.result() # 等待任务完成 table = client.get_table(base_table_id) print(f"数据已成功加载到分区表 {base_table_id},当前表总行数:{table.num_rows}") def main(): upload_from_gcs_to_bq(project_id, dataset_id, gsutil_uri, table_name,gcs_blob) if __name__ == '__main__': main()
方案说明
- 新增date字段:通过查询语句自动为每条数据添加
date字段,值为从文件名提取的日期。 - 字段分区:将
date字段设置为分区键,BigQuery会自动按该字段的日期值分配数据到对应分区。 - 外部表读取:通过
ExternalConfig直接读取GCS中的CSV,无需临时存储,简化流程。
内容的提问来源于stack exchange,提问作者Sana
相关产品推荐
相关产品推荐

