结合threading与queue运行Snowflake CSV查询的断连问题求助
问题核心原因
你当前的实现每执行一条查询都会单独创建1个新线程、建立1个新的Snowflake连接,Snowflake对单账号的并发连接数有默认上限(通常为几百量级),并发超过400后会被服务端主动拒绝连接,同时创建4000个线程本身也会极大消耗操作系统资源,完全没有必要。
优化后可直接运行的代码
import threading from csv import DictReader from concurrent.futures import ThreadPoolExecutor, as_completed import snowflake.connector as sf # 线程局部存储,每个线程持有独立的连接,避免多线程共用连接冲突 thread_local = threading.local() def get_sf_connection( sfPswd = 'psw', sfUser = 'user', sfAccount = 'account', sfRole = 'role', sfWarehouse = 'COMPUTE_WH', sfDatabase = 'db' ): # 每个线程只创建一次连接,复用执行所有分配到的查询 if not hasattr(thread_local, "connection"): if sfPswd == '': import getpass sfPswd = getpass.getpass('Password:') try: conn = sf.connect( user=sfUser, password=sfPswd, account=sfAccount ) # 一次性初始化会话参数,不用每次查询都执行 with conn.cursor() as cur: cur.execute(f'USE ROLE {sfRole}') cur.execute(f'USE WAREHOUSE {sfWarehouse}') cur.execute(f'USE DATABASE {sfDatabase}') cur.execute("ALTER SESSION SET QUERY_TAG = 'threadtest'") thread_local.connection = conn print('线程连接建立成功') except Exception as e: print(f'连接失败,请检查凭据:{str(e)}') raise return thread_local.connection def execute_query(query_item): query_id, sql_query = query_item try: conn = get_sf_connection() with conn.cursor() as cur: cur.execute(sql_query) print(f'查询{query_id}执行完成') return True, query_id except Exception as e: print(f'查询{query_id}执行失败,错误信息:{str(e)},查询语句:{sql_query[:100]}...') return False, query_id if __name__ == '__main__': # 1. 读取CSV中的查询语句 queryStatements = [] with open('result.csv', encoding='utf-8') as csvfile: dict_reader = DictReader(csvfile) for idx, member in enumerate(dict_reader): queryStatements.append((idx, f"{member['QUERYTXT']};")) # 2. 配置最大并发数,可根据你Snowflake账户的并发上限调整,建议不要超过200,避免触发连接限制 MAX_WORKERS = 32 success_cnt = 0 fail_cnt = 0 # 3. 线程池执行所有查询 with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: futures = [executor.submit(execute_query, item) for item in queryStatements] for future in as_completed(futures): success, qid = future.result() if success: success_cnt += 1 else: fail_cnt += 1 # 4. 关闭所有线程持有的连接 print(f'所有查询执行完毕,成功{success_cnt}条,失败{fail_cnt}条')
核心修改点说明
- 用
ThreadPoolExecutor管理线程生命周期,自动调度4000条查询,不用手动创建管理4000个线程 - 新增线程局部存储,每个线程只创建一次Snowflake连接,该线程下的所有查询都复用这一个连接,大幅减少连接创建开销,同时将并发连接数控制在你设置的
MAX_WORKERS数值,只要低于Snowflake的并发上限就不会出现断连问题 - 会话初始化操作(设置角色、仓库、查询标签等)只在连接创建时执行一次,不用每条查询重复执行,提升执行效率
- 新增异常捕获逻辑,查询失败时会打印错误信息和对应的查询内容,方便定位问题
- 所有查询执行完成后统一统计执行结果,清晰看到成功失败条数
内容的提问来源于stack exchange,提问作者wtfpython
相关产品推荐
相关产品推荐

