Google BigQuery Update比Insert慢70倍,如何优化?
问题描述
使用Scrapy爬虫结合Google BigQuery存储数据,现有Insert和Update两个管道:
- Insert单条耗时0.05秒,效率正常
- Update单条耗时3.56秒,每分钟仅能更新20条,速度比Insert慢70倍
当前表规模约2万条,未来预计达50万条,需每日更新记录,求Update速度慢的原因及优化方案。
原Update代码:
# Define the update query query = f""" UPDATE `{self.dataset_id}.{self.table_id}` SET `Sold Status` = '{data['Sold Status']}', `Amount of Views` = '{data['Amount of Views']}', `Amount of Likes` = '{data['Amount of Likes']}', `Sold Date & Time` = '{data['Sold Date & Time']}' WHERE `Item number` = '{data['Item number']}' """ start_time = time.time() # Run the update query job = self.client.query(query) # Wait for the job to complete job.result() # Check if the query was successful if job.state == 'DONE': print('Update query executed successfully.') else: print('Update query failed.') end_time = time.time() execution_time = end_time - start_time logging.info(execution_time) return item
原Insert代码:
start_time = time.time() data = item slug = data['slug'] if slug in self.ids_seen: raise DropItem("Duplicate item found: {}".format(slug)) else: data.pop('slug', None) self.ids_seen.add(slug) table_ref = self.client.dataset(self.dataset_id).table(self.table_id) # Define the rows to be inserted rows = [ data ] # Insert rows into the table errors = self.client.insert_rows_json(table_ref, rows) if errors == []: print("Rows inserted successfully.") else: print("Encountered errors while inserting rows:", errors) end_time = time.time() execution_time = end_time - start_time logging.info(execution_time) return item
问题原因分析
- 单条更新的作业开销:BigQuery执行
client.query()会启动完整的查询作业,作业本身有固定的启动/初始化开销(约2-3秒),这直接导致单条Update耗时大部分浪费在作业启动上;而insert_rows_json是轻量批量写入接口,无此开销。 - 字符串拼接SQL的缺陷:
- 存在SQL注入风险,若数据含单引号等特殊字符,会直接导致SQL语法错误
- BigQuery无法缓存动态生成的查询计划,每次更新都要重新解析、优化查询,进一步增加耗时
- 未利用批量处理特性:Scrapy管道逐条处理Item,原Update代码也逐条发起请求,未攒批处理,重复触发多次作业启动开销。
优化方案
1. 核心优化:批量更新
将多条待更新数据攒成一批,通过一次MERGE操作完成,大幅减少作业启动次数。推荐临时表+MERGE的方式,这是BigQuery处理大规模更新的最优实践。
优化后的批量Update管道代码
from google.cloud import bigquery import time import logging from scrapy.exceptions import DropItem class BigQueryBatchUpdatePipeline: def __init__(self): self.client = bigquery.Client() self.dataset_id = "你的数据集ID" self.table_id = "你的表ID" self.update_batch = [] self.batch_size = 100 # 可测试调整,推荐100-500条/批 def process_item(self, item, spider): # 提取更新字段与主键,确保类型与表字段匹配 update_data = { "Item number": item["Item number"], "Sold Status": item["Sold Status"], "Amount of Views": str(item["Amount of Views"]), "Amount of Likes": str(item["Amount of Likes"]), "Sold Date & Time": item["Sold Date & Time"] } self.update_batch.append(update_data) # 达到批量阈值时执行更新 if len(self.update_batch) >= self.batch_size: self._execute_batch_update() return item def close_spider(self, spider): # 爬虫结束时处理剩余未批量的数据 if self.update_batch: self._execute_batch_update() def _execute_batch_update(self): start_time = time.time() temp_table_name = "_temp_update_batch" temp_table_ref = self.client.dataset(self.dataset_id).table(temp_table_name) try: # 1. 将批量数据写入临时表 errors = self.client.insert_rows_json(temp_table_ref, self.update_batch) if errors: logging.error(f"临时表写入失败: {errors}") return # 2. 执行MERGE语句更新主表 merge_query = f""" MERGE `{self.dataset_id}.{self.table_id}` main_table USING `{self.dataset_id}.{temp_table_name}` temp_table ON main_table.`Item number` = temp_table.`Item number` WHEN MATCHED THEN UPDATE SET `Sold Status` = temp_table.`Sold Status`, `Amount of Views` = temp_table.`Amount of Views`, `Amount of Likes` = temp_table.`Amount of Likes`, `Sold Date & Time` = temp_table.`Sold Date & Time` """ job = self.client.query(merge_query) job.result() # 等待作业完成 execution_time = time.time() - start_time logging.info(f"批量更新{len(self.update_batch)}条记录耗时: {execution_time:.2f}秒") finally: # 3. 删除临时表(BigQuery临时表会自动过期,也可手动清理) self.client.delete_table(temp_table_ref, not_found_ok=True) self.update_batch = []
2. 辅助优化措施
- 调整批量大小:测试不同的
batch_size(如100、200、500),找到平衡写入速度与作业开销的最优值 - 聚类表优化:若
Item number是频繁过滤字段,将主表设置为按Item number聚类,MERGE时会大幅减少扫描的数据量,提升更新效率 - 参数化查询:若不使用临时表,改用参数化构造USING子句,避免SQL注入并允许BigQuery缓存查询计划,示例代码如下:
def _execute_param_batch_update(self): start_time = time.time() sql_segments = [] params = [] for idx, data in enumerate(self.update_batch): param_suffix = f"_{idx}" sql_segments.append(f""" SELECT @item_num{param_suffix} AS `Item number`, @sold_status{param_suffix} AS `Sold Status`, @views{param_suffix} AS `Amount of Views`, @likes{param_suffix} AS `Amount of Likes`, @sold_time{param_suffix} AS `Sold Date & Time` """) # 添加参数 params.extend([ bigquery.ScalarQueryParameter(f"item_num{param_suffix}", "STRING", data["Item number"]), bigquery.ScalarQueryParameter(f"sold_status{param_suffix}", "STRING", data["Sold Status"]), bigquery.ScalarQueryParameter(f"views{param_suffix}", "STRING", data["Amount of Views"]), bigquery.ScalarQueryParameter(f"likes{param_suffix}", "STRING", data["Amount of Likes"]), bigquery.ScalarQueryParameter(f"sold_time{param_suffix}", "STRING", data["Sold Date & Time"]), ]) # 组装MERGE语句 merge_query = f""" MERGE `{self.dataset_id}.{self.table_id}` main_table USING ( {' UNION ALL '.join(sql_segments)} ) update_data ON main_table.`Item number` = update_data.`Item number` WHEN MATCHED THEN UPDATE SET `Sold Status` = update_data.`Sold Status`, `Amount of Views` = update_data.`Amount of Views`, `Amount of Likes` = update_data.`Amount of Likes`, `Sold Date & Time` = update_data.`Sold Date & Time` """ # 执行参数化查询 job_config = bigquery.QueryJobConfig(query_parameters=params) job = self.client.query(merge_query, job_config=job_config) job.result() execution_time = time.time() - start_time logging.info(f"参数化批量更新{len(self.update_batch)}条记录耗时: {execution_time:.2f}秒") self.update_batch = []
3. 额外建议
- 若每日更新记录占总表比例较高(如超过30%),可考虑先写入新数据到临时表,再替换主表,比批量更新更高效
- 确保BigQuery表字段类型与爬虫输出数据类型一致,避免类型转换带来的额外开销
内容的提问来源于stack exchange,提问作者JBJ
相关产品推荐
相关产品推荐

