如何用Python通过BigQuery API高效实现条件增删改操作
基于BigQuery MERGE操作实现JSON数据的增删改逻辑
输入JSON示例
{ "ID": "ABC34567", "NUM": "45678", "ORG_ID": "161", "CONTACT_NUMBER": null, "FLAG": "N" }
核心实现思路
利用BigQuery的MERGE语句,将传入的JSON数据作为临时数据源,与目标表通过ID字段关联,根据FLAG的值执行对应操作:
- 当
FLAG = 'N'时:- 若
ID不存在,插入新行 - 若
ID已存在,更新NUM、ORG_ID、CONTACT_NUMBER字段
- 若
- 当
FLAG = 'Del'时:- 若
ID已存在,删除该行 - 若
ID不存在,不执行任何操作
- 若
Python代码实现
步骤1:初始化BigQuery客户端
确保已安装google-cloud-bigquery依赖,初始化客户端:
from google.cloud import bigquery client = bigquery.Client()
步骤2:定义目标表与传入的JSON数据
# 目标表格式:project.dataset.table_name target_table = "your-project.your-dataset.your-table" # 传入的JSON数据 input_json = { "ID": "ABC34567", "NUM": "45678", "ORG_ID": "161", "CONTACT_NUMBER": None, "FLAG": "N" }
步骤3:构建参数化MERGE查询
使用参数化查询传递JSON数据,避免SQL注入,同时适配NULL值:
MERGE `{target_table}` AS target USING ( SELECT @id AS ID, @num AS NUM, @org_id AS ORG_ID, @contact_number AS CONTACT_NUMBER, @flag AS FLAG ) AS source ON target.ID = source.ID -- 处理FLAG为Del的情况:匹配到则删除 WHEN MATCHED AND source.FLAG = 'Del' THEN DELETE -- 处理FLAG为N的情况:匹配到则更新字段 WHEN MATCHED AND source.FLAG = 'N' THEN UPDATE SET NUM = source.NUM, ORG_ID = source.ORG_ID, CONTACT_NUMBER = source.CONTACT_NUMBER -- 处理FLAG为N且ID不存在的情况:插入新行 WHEN NOT MATCHED AND source.FLAG = 'N' THEN INSERT (ID, NUM, ORG_ID, CONTACT_NUMBER) VALUES (source.ID, source.NUM, source.ORG_ID, source.CONTACT_NUMBER)
步骤4:执行参数化查询
在Python中传递参数并执行:
# 构建查询语句,替换目标表占位符 query = f""" MERGE `{target_table}` AS target USING ( SELECT @id AS ID, @num AS NUM, @org_id AS ORG_ID, @contact_number AS CONTACT_NUMBER, @flag AS FLAG ) AS source ON target.ID = source.ID WHEN MATCHED AND source.FLAG = 'Del' THEN DELETE WHEN MATCHED AND source.FLAG = 'N' THEN UPDATE SET NUM = source.NUM, ORG_ID = source.ORG_ID, CONTACT_NUMBER = source.CONTACT_NUMBER WHEN NOT MATCHED AND source.FLAG = 'N' THEN INSERT (ID, NUM, ORG_ID, CONTACT_NUMBER) VALUES (source.ID, source.NUM, source.ORG_ID, source.CONTACT_NUMBER) """ # 定义参数,对应查询中的占位符 job_config = bigquery.QueryJobConfig( query_parameters=[ bigquery.ScalarQueryParameter("id", "STRING", input_json["ID"]), bigquery.ScalarQueryParameter("num", "STRING", input_json["NUM"]), bigquery.ScalarQueryParameter("org_id", "STRING", input_json["ORG_ID"]), bigquery.ScalarQueryParameter("contact_number", "STRING", input_json["CONTACT_NUMBER"]), bigquery.ScalarQueryParameter("flag", "STRING", input_json["FLAG"]), ] ) # 执行查询 query_job = client.query(query, job_config=job_config) query_job.result() # 等待执行完成
关键说明
- 若
CONTACT_NUMBER字段类型不是STRING,需调整ScalarQueryParameter中的类型参数(如INT64、FLOAT64等) - 批量处理多条JSON数据时,可将数据源改为
UNION ALL的多组值,或使用临时表加载批量数据后执行MERGE - 确保BigQuery客户端拥有目标表的读写权限
内容的提问来源于stack exchange,提问作者ditil
相关产品推荐
相关产品推荐

