在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
相关产品推荐
相关产品推荐

