Python asyncio aiopg连接特定主机数据库时OSError异常捕获问题
批量异步查询PostgreSQL时的WinError 10038异常处理
问题背景
有155组主机+数据库的组合,在Jupyter Notebook中编写异步程序批量执行查询,代码如下:
SQL = """ SELECT 1 """ import asyncio import aiopg import nest_asyncio nest_asyncio.apply() loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) async def runQuery(host, db, sql): print(str(host), str(db)) if str(host) == 'nan': return 'No host' try: conn = await aiopg.connect(database=str(db), user='some_user', password='some_password', host=str(host), timeout = 2) print('xxxx') except (psycopg2.Error, asyncio.TimeoutError, psycopg2.OperationalError) as e: return 'Error creating connection exception' if not conn or conn.closed == 1: return 'Error creating connection' try: cur = await conn.cursor() if not cur: return 'No cursor' await cur.execute(sql) result = await cur.fetchall() await conn.close() return result except psycopg2.Error as e: await conn.close() return 'Error fetching' async def main(): tasks = [] for item in csvFile: tasks.extend([runQuery(item['host'], it, SQL) for it in item['dbs']]) results = [] i = 90 # to start with host+db number 90 d = 1 # how many to run at once while i <= len(tasks): r = await asyncio.gather(*tasks[i:(i+d)]) i = i+d print(i) results.extend(r) return results results = loop.run_until_complete(main()) print(results) loop.close()
程序大部分组合可正常运行,但第90、120等组合会触发错误,且该异常无法在runQuery函数的conn=await aiopg.connect...代码块中被捕获,需要忽略这些单连接异常而不终止整个进程。
触发的错误及回溯
错误信息:
OSError: [WinError 10038] An operation was attempted on something that is not a socket
完整回溯信息:
--------------------------------------------------------------------------- OSError Traceback (most recent call last) Cell In[9], line 70 68 #results = await asyncio.gather(*tasks) 69 return results ---> 70 results = loop.run_until_complete(main()) 71 print(results) 73 loop.close() File ~\AppData\Local\Programs\Python\Python312\Lib\site-packages\nest_asyncio.py:92, in _patch_loop.<locals>.run_until_complete(self, future) 90 f._log_destroy_pending = False 91 while not f.done(): ---> 92 self._run_once() 93 if self._stopping: 94 break File ~\AppData\Local\Programs\Python\Python312\Lib\site-packages\nest_asyncio.py:115, in _patch_loop.<locals>._run_once(self) 108 heappop(scheduled) 110 timeout = ( 111 0 if ready or self._stopping 112 else min(max( 113 scheduled[0]._when - self.time(), 0), 86400) if scheduled 114 else None) --> 115 event_list = self._selector.select(timeout) 116 self._process_events(event_list) 118 end_time = self.time() + self._clock_resolution File ~\AppData\Local\Programs\Python\Python312\Lib\selectors.py:323, in SelectSelector.select(self, timeout) 321 ready = [] 322 try: --> 323 r, w, _ = self._select(self._readers, self._writers, [], timeout) 324 except InterruptedError: 325 return ready File ~\AppData\Local\Programs\Python\Python312\Lib\selectors.py:314, in SelectSelector._select(self, r, w, _, timeout) 313 def _select(self, r, w, _, timeout=None): --> 314 r, w, x = select.select(r, w, w, timeout) 315 return r, w + x, [] OSError: [WinError 10038] An operation was attempted on something that is not a socket
解决方案
这个错误源于Windows系统下SelectSelector的socket资源处理问题,结合nest_asyncio的补丁可能引发资源泄漏,以下是针对性的修复方案:
1. 扩展异常捕获范围
在runQuery函数中加入OSError捕获,确保底层socket错误能被处理:
async def runQuery(host, db, sql): print(str(host), str(db)) if str(host) == 'nan': return 'No host' conn = None try: conn = await aiopg.connect(database=str(db), user='some_user', password='some_password', host=str(host), timeout=2) cur = await conn.cursor() await cur.execute(sql) result = await cur.fetchall() return result except (psycopg2.Error, asyncio.TimeoutError, psycopg2.OperationalError, OSError) as e: return f"Connection/Query error for {host}:{db} - {str(e)}" finally: # 确保连接被关闭,避免资源泄漏 if conn and not conn.closed: await conn.close()
2. 使用asyncio.gather的异常安全模式
通过return_exceptions=True参数,让单个任务的异常作为结果返回,而非终止整个异步进程:
async def main(): tasks = [] # 控制并发数,避免过多连接耗尽系统资源 semaphore = asyncio.Semaphore(10) # 封装带并发限制的查询函数 async def bounded_query(host, db, sql): async with semaphore: return await runQuery(host, db, sql) # 生成所有任务 for item in csvFile: tasks.extend([bounded_query(item['host'], it, SQL) for it in item['dbs']]) # 执行所有任务,单个异常不会中断整体 results = await asyncio.gather(*tasks, return_exceptions=True) # 处理结果:区分正常返回和异常 processed_results = [] for idx, res in enumerate(results): if isinstance(res, Exception): processed_results.append(f"Task {idx} failed: {str(res)}") else: processed_results.append(res) return processed_results
3. 优化事件循环初始化
简化事件循环的设置,避免重复创建可能引发的资源问题:
import asyncio import aiopg import nest_asyncio nest_asyncio.apply() # 使用默认事件循环即可 loop = asyncio.get_event_loop() results = loop.run_until_complete(main()) print(results) loop.close()
关键说明
Semaphore控制并发数:避免同时创建数百个数据库连接,减轻系统压力,也是降低WinError 10038触发概率的核心措施return_exceptions=True:确保单个任务的异常不会导致整个批量查询中断,异常会被作为结果返回,方便后续处理finally块:强制关闭数据库连接,防止socket资源泄漏,从根源上减少无效socket操作的可能
内容的提问来源于stack exchange,提问作者Lidbey
相关产品推荐
相关产品推荐

