You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在asyncio中使用BigQuery客户端引发OpenSSL/SSL.py挂起问题

Python BigQuery 并发读取挂起问题

环境信息

  • Python版本:Python 3.8.10
  • Google BigQuery包版本:2.30.1

测试代码

from google.oauth2 import service_account
from google.cloud import bigquery

import logging
import asyncio
import faulthandler
import traceback
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor

logger = logging.getLogger('test')
logger.setLevel(logging.ERROR)

def get_big_query_client(source_details):
   print("get client")
   credentials_dict = source_details['auth_config']
   credentials = service_account.Credentials.from_service_account_info(credentials_dict)
   project_id = credentials_dict["project_id"]
   client = bigquery.Client(credentials=credentials, project=project_id)
   return client

async def check_and_run_dd_for_tenants(tenant_ids, source_detail, table_name):
   await asyncio.gather(*[flatten_column_details_of_dataset(source_detail, table_name) for tenant_id in tenant_ids])

async def flatten_column_details_of_dataset(source_details, table_name):

   curr_column_data = []
  
   client = get_big_query_client(source_details)
   schema_name = 'test_schema'
   query_table_name = f"{schema_name}.{table_name}"
   details = {'query_table_name':query_table_name}
   table_reference = details.get("query_table_name")

   try:
       table_details = client.get_table(table_reference)
       print("table_details",table_details)

       for _schema_data in table_details.schema:
           data_type = _schema_data.field_type

   except Exception as exp:
       logger.error(exp)
       logger.error(traceback.format_exc())
       pass
   finally:
       client.close()
   return curr_column_data

def main():
   logger = logging.getLogger("Run DD")
   logger.setLevel(logging.INFO)
   faulthandler.enable(all_threads=True)
   faulthandler.dump_traceback_later(180, repeat=True)

   source_details = {
       "auth_config": {
           "auth_provider_x509_cert_url": "*************",
           "auth_uri": "https://accounts.google.com/o/oauth2/auth",
           "client_email": "**************",
           "client_id": "****************",
           "client_x509_cert_url": "**********************",
           "private_key": "*****************",
           "private_key_id": "***************",
           "project_id": "**********************",
           "token_uri": "https://oauth2.googleapis.com/token",
           "type": "service_account"
       }
   }
   table_names = ['test1']
   tenant_id = []
   for i in range(1,30):
       tenant_id.append(i)
   for table in table_names:
       res = asyncio.run(check_and_run_dd_for_tenants(tenant_id, source_details, table))
       print(res)

if __name__ == "__main__":
   main()

问题现象

使用asyncio实现并发读取BigQuery表时,代码偶尔出现无响应挂起。通过faulthandler捕获的调用栈显示,挂起点为table_details = client.get_table(table_reference),最终阻塞在OpenSSL/SSL.py的recv_into函数中。注:涉及的目标表均为小表。

已尝试的无效解决方案

  • 引入线程锁控制并发
  • 为get_table方法添加timeout参数
  • 将Socket设置为非阻塞模式并配置超时
  • 将读取逻辑封装为独立函数,通过loop.run_in_executor调用

该问题已在BigQuery官方代码仓库提交相关反馈,现寻求能有效解决该挂起问题的方案。


内容的提问来源于stack exchange,提问作者Nandha

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.09 13:33:10