如何从GCS无重复加载数据到BigQuery?求性能成本最优方案
问题
希望从Google Cloud Storage(Bucket)中获取CSV文件,加载到BigQuery表中并避免重复数据,要求实现性能与成本最优。当前代码如下:
def load_data_in_BQT(): job_config = bigquery.LoadJobConfig( schema=[ bigquery.SchemaField("id", "INTEGER"), bigquery.SchemaField("name", "STRING"), ], # The source format defaults to CSV, so the line below is optional. source_format=bigquery.SourceFormat.CSV, skip_leading_rows=1, autodetect=True, write_disposition=bigquery.WriteDisposition.WRITE_APPEND, # (Addition of the data (possibility of having duplications) # write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE, # (Formatting of the table and insertion of the new data (Loss of the old data)) ) uri = "gs://mybucket/myfolder/myfile.csv" load_job = self.client.load_table_from_uri( uri, self.table_ref["object"], job_config=job_config, )
原本思路是将CSV和BigQuery表数据都转为Pandas DataFrame后去重,再用WRITE_TRUNCATE重新插入,但海量数据下存在弊端,求更好方案。
最优解决方案
以下方案均基于BigQuery原生能力实现,无需拉取数据到本地,兼顾性能与成本:
方案1:临时表+MERGE语句(海量数据首选)
完全在BigQuery内部完成数据加载与去重,避免本地内存瓶颈,是处理大规模数据的最优方式。
步骤逻辑
- 将CSV加载到临时表(用
WRITE_TRUNCATE保证临时表每次都是最新的CSV数据) - 执行
MERGE语句,根据唯一键(比如id)合并临时表与目标表:不存在的插入,存在的可选择忽略或更新
代码示例
def load_data_without_duplicates(): # 1. 配置临时表加载任务 temp_table_ref = self.client.dataset("your_dataset").table("temp_csv_data") job_config = bigquery.LoadJobConfig( schema=[ bigquery.SchemaField("id", "INTEGER"), bigquery.SchemaField("name", "STRING"), ], source_format=bigquery.SourceFormat.CSV, skip_leading_rows=1, write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE, # 每次覆盖临时表 ) uri = "gs://mybucket/myfolder/myfile.csv" load_job = self.client.load_table_from_uri( uri, temp_table_ref, job_config=job_config ) load_job.result() # 等待加载完成 # 2. 执行MERGE语句去重合并 merge_query = f""" MERGE INTO `{self.table_ref["object"]}` AS target USING `{temp_table_ref.full_table_id}` AS source ON target.id = source.id WHEN NOT MATCHED THEN INSERT (id, name) VALUES (source.id, source.name) -- 如需更新已存在的记录,可添加:WHEN MATCHED THEN UPDATE SET name = source.name """ query_job = self.client.query(merge_query) query_job.result() # 等待合并完成 # 可选:删除临时表(不需要保留则执行) self.client.delete_table(temp_table_ref)
方案2:追加后批量去重(适合增量小数据)
如果CSV是增量数据且重复比例低,可先追加数据到目标表,再通过BigQuery的窗口函数批量去重,操作简单且成本较低。
代码示例
def append_and_deduplicate(): # 1. 先追加CSV数据到目标表 job_config = bigquery.LoadJobConfig( schema=[ bigquery.SchemaField("id", "INTEGER"), bigquery.SchemaField("name", "STRING"), ], source_format=bigquery.SourceFormat.CSV, skip_leading_rows=1, write_disposition=bigquery.WriteDisposition.WRITE_APPEND, ) uri = "gs://mybucket/myfolder/myfile.csv" load_job = self.client.load_table_from_uri( uri, self.table_ref["object"], job_config=job_config ) load_job.result() # 2. 执行去重查询,覆盖原表 deduplicate_query = f""" CREATE OR REPLACE TABLE `{self.table_ref["object"]}` AS SELECT id, name FROM ( SELECT *, ROW_NUMBER() OVER (PARTITION BY id ORDER BY _PARTITIONTIME DESC) AS rn FROM `{self.table_ref["object"]}` ) WHERE rn = 1 """ query_job = self.client.query(deduplicate_query) query_job.result()
注:如果是分区表,用_PARTITIONTIME保留最新记录;非分区表可替换为业务时间字段,或直接取任意一条重复记录。
方案3:外部表+物化视图(适合频繁更新的CSV)
如果CSV文件会定期更新,且需要实时查询去重后的数据,可创建外部表指向GCS的CSV,再通过物化视图自动维护去重数据。
步骤示例
- 创建外部表:
def create_external_table(): external_config = bigquery.ExternalConfig("CSV") external_config.source_uris = ["gs://mybucket/myfolder/myfile.csv"] external_config.skip_leading_rows = 1 external_config.schema = [ bigquery.SchemaField("id", "INTEGER"), bigquery.SchemaField("name", "STRING"), ] table = bigquery.Table(self.client.dataset("your_dataset").table("external_csv")) table.external_data_configuration = external_config self.client.create_table(table)
- 创建去重物化视图(通过SQL执行):
CREATE MATERIALIZED VIEW `your_project.your_dataset.deduplicated_view` AS SELECT id, name FROM ( SELECT *, ROW_NUMBER() OVER (PARTITION BY id ORDER BY _FILE_NAME DESC) AS rn FROM `your_project.your_dataset.external_csv` ) WHERE rn = 1
物化视图会自动刷新,保证数据与源CSV同步,无需手动执行加载操作。
内容的提问来源于stack exchange,提问作者RashidLadj_Winux
相关产品推荐
相关产品推荐

