BigQuery Python客户端:如何判断带目标表的SELECT查询是否返回数据
问题:BigQuery Python客户端判断带目标表的SELECT查询实际返回行数
在使用BigQuery Python客户端时,需要判断SELECT查询在将结果插入目标表前是否返回数据,但原代码中results.total_rows返回的是目标表的总行数,而非查询实际返回的行数。原代码如下:
from google.cloud import bigquery import time client = bigquery.Client() query = """ SELECT col1, col2 FROM project_id.dataset_id.source_table WHERE col2 is not null """ job_config = bigquery.QueryJobConfig(destination="project_id.dataset_id.destination_table") query_job = client.query(query, job_config=job_config) while query_job.state != "DONE": time.sleep(1) query_job.reload() results = query_job.result() total_rows_returned = results.total_rows if total_rows_returned == 0: print("Query did not return any rows.") else: print(f"Query returned {total_rows_returned} rows. Inserting into the destination table.")
解决方案
有两种可靠的方式获取SELECT查询实际返回的行数:
方法1:通过查询任务的统计信息获取
当查询任务完成后,可以从query_job.statistics中提取原查询的返回行数,这个数值不受目标表原有数据影响。修改后的代码如下:
from google.cloud import bigquery client = bigquery.Client() query = """ SELECT col1, col2 FROM project_id.dataset_id.source_table WHERE col2 is not null """ job_config = bigquery.QueryJobConfig(destination="project_id.dataset_id.destination_table") query_job = client.query(query, job_config=job_config) query_job.result() # 自动等待任务完成,替代手动轮询 # 从统计信息中获取查询实际返回的行数 total_rows_returned = int(query_job.statistics['query']['totalRows']) if total_rows_returned == 0: print("Query did not return any rows.") else: print(f"Query returned {total_rows_returned} rows. Inserted into the destination table.")
说明:query_job.statistics['query']['totalRows']存储的是SELECT查询本身返回的结果行数,即使写入目标表,这个值也不会包含目标表原有的数据。
方法2:先执行dry_run预估行数(避免不必要的写入)
如果想在实际写入前就判断查询是否有返回数据,可以开启dry_run模式,提前获取预估的行数,再决定是否执行正式写入:
from google.cloud import bigquery client = bigquery.Client() query = """ SELECT col1, col2 FROM project_id.dataset_id.source_table WHERE col2 is not null """ # 先执行dry_run,仅预估行数不实际执行 dry_run_config = bigquery.QueryJobConfig(dry_run=True, use_query_cache=False) dry_run_job = client.query(query, job_config=dry_run_config) estimated_rows = dry_run_job.total_rows if estimated_rows == 0: print("Query will not return any rows. Skip inserting.") else: print(f"Query will return {estimated_rows} rows. Start inserting into destination table.") # 执行正式写入操作 job_config = bigquery.QueryJobConfig(destination="project_id.dataset_id.destination_table") query_job = client.query(query, job_config=job_config) query_job.result() print("Insert completed.")
说明:dry_run不会实际执行查询或写入,只会返回预估的行数,适合提前判断是否需要执行后续写入操作,节省计算资源。
内容的提问来源于stack exchange,提问作者fancyeffective
相关产品推荐
相关产品推荐

