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

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    
---&gt; 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.&lt;locals&gt;.run_until_complete(self, future)
     90     f._log_destroy_pending = False
     91 while not f.done():
---&gt; 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.&lt;locals&gt;._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)
--&gt; 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:
--&gt; 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):
--&gt; 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 06:52:07