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

复用未关闭的asyncpg连接池时触发“another operation is in progress”错误的原因咨询

复用未关闭的asyncpg连接池时触发“another operation is in progress”错误的原因咨询

我遇到了一个关于asyncpg连接池复用的问题,想请教大家原因和解决方法:

问题背景

我原本的异步计算代码可以通过多次调用asyncio.run()正常执行,但在加入asyncpg连接池后,第一次调用正常,第二次复用同一个连接池时就触发了错误。

环境信息

  • Windows 10 22 H2
  • Python 3.10.13
  • 依赖包:
    async-timeout==5.0.1
    asyncpg==0.30.0
    

正常工作的基础代码(无数据库操作)

这段代码两次调用asyncio.run()都能正常输出结果:

import asyncio
import random


async def d(c: int):
    await asyncio.sleep(random.random())
    return c**2


async def main(input: list[any]):
    result = await asyncio.gather(*(d(i) for i in input))
    return result


if __name__ == "__main__":
    print(f"Start 1 asyncio.run()")
    result = asyncio.run(main([1, 2, 3]))
    print(f"{result = }")
    print(f"Start 2 asyncio.run()")
    result = asyncio.run(main(result))
    print(f"{result = }")

执行结果:

Start 1 asyncio.run()
result = [1, 4, 9]
Start 2 asyncio.run()
result = [1, 16, 81]

加入asyncpg连接池后的出错代码

当我引入asyncpg连接池后,第一次调用正常,但第二次复用连接池就报错:

import asyncio
import asyncpg


async def d(c: str, pool: asyncpg.Pool):
    async with pool.acquire() as con:
        return await con.fetch("SELECT $1", c)


async def main(input: list[str], pool: asyncpg.Pool = None):
    if not pool:
        pool = await asyncpg.create_pool(dsn="Connection arguments")
    result = await asyncio.gather(*(d(i, pool) for i in input))
    print(f"{result = }")
    return pool


if __name__ == "__main__":
    print(f"Start 1 asyncio.run()")
    pool = asyncio.run(main(["1", "2", "3"]))
    print(f"{pool.is_closing() = }")
    print(f"{pool.get_size() = }")
    print(f"{pool.get_idle_size() = }")
    print(f"Start 2 asyncio.run()")
    asyncio.run(main(["2", "3", "4"], pool))

错误信息

Start 1 asyncio.run()
result = [[<Record ?column?='1'>], [<Record ?column?='2'>], [<Record ?column?='3'>]]
pool.is_closing() = False
pool.get_size() = 10
pool.get_idle_size() = 10
Start 2 asyncio.run()
Traceback (most recent call last):
  File "c:\python310\lib\asyncio\base_events.py", line 753, in call_soon    
    self._check_closed()
  File "c:\python310\lib\asyncio\base_events.py", line 515, in _check_closed
    raise RuntimeError('Event loop is closed')
RuntimeError: Event loop is closed

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "d:\ProgProject\python\learn_asyncpg\111.py", line 7, in d
    return await con.fetch("SELECT $1", c)
  File "d:\ProgProject\python\learn_asyncpg\venv\lib\site-packages\asyncpg\connection.py", line 690, in fetch
    return await self._execute(
  File "d:\ProgProject\python\learn_asyncpg\venv\lib\site-packages\asyncpg\connection.py", line 1864, in _execute
    result, _ = await self.__execute(
  File "d:\ProgProject\python\learn_asyncpg\venv\lib\site-packages\asyncpg\connection.py", line 1961, in __execute
    result, stmt = await self._do_execute(
  File "d:\ProgProject\python\learn_asyncpg\venv\lib\site-packages\asyncpg\connection.py", line 2024, in _do_execute
    result = await executor(stmt, None)
  File "asyncpg\\protocol\\protocol.pyx", line 206, in bind_execute
  File "asyncpg\\protocol\\protocol.pyx", line 192, in asyncpg.protocol.protocol.BaseProtocol.bind_execute
  File "asyncpg\\protocol\\coreproto.pyx", line 1020, in asyncpg.protocol.protocol.CoreProtocol._bind_execute
  File "asyncpg\\protocol\\coreproto.pyx", line 1008, in asyncpg.protocol.protocol.CoreProtocol._send_bind_message
  File "asyncpg\\protocol\\protocol.pyx", line 967, in asyncpg.protocol.protocol.BaseProtocol._write
  File "c:\python310\lib\asyncio\proactor_events.py", line 365, in write
    self._loop_writing(data=bytes(data))
  File "c:\python310\lib\asyncio\proactor_events.py", line 401, in _loop_writing
    self._write_fut = self._loop._proactor.send(self._sock, data)
AttributeError: 'NoneType' object has no attribute 'send'

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "d:\ProgProject\python\learn_asyncpg\111.py", line 25, in <module>
    asyncio.run(main(["2", "3", "4"], pool))
  File "c:\python310\lib\asyncio\runners.py", line 44, in run
    return loop.run_until_complete(main)
  File "c:\python310\lib\asyncio\base_events.py", line 649, in run_until_complete
    return future.result()
  File "d:\ProgProject\python\learn_asyncpg\111.py", line 13, in main
    result = await asyncio.gather(*(d(i, pool) for i in input))
  File "d:\ProgProject\python\learn_asyncpg\111.py", line 6, in d
    async with pool.acquire() as con:
  File "d:\ProgProject\python\learn_asyncpg\venv\lib\site-packages\asyncpg\pool.py", line 228, in release
    raise ex
  File "d:\ProgProject\python\learn_asyncpg\venv\lib\site-packages\asyncpg\pool.py", line 218, in release
    await self._con.reset(timeout=budget)
  File "d:\ProgProject\python\learn_asyncpg\venv\lib\site-packages\asyncpg\connection.py", line 1562, in reset
    await self.execute(reset_query)
  File "d:\ProgProject\python\learn_asyncpg\venv\lib\site-packages\asyncpg\connection.py", line 349, in execute
    result = await self._protocol.query(query, timeout)
  File "asyncpg\\protocol\\protocol.pyx", line 360, in query
  File "asyncpg\\protocol\\protocol.pyx", line 745, in asyncpg.protocol.protocol.BaseProtocol._check_state
asyncpg.exceptions._base.InterfaceError: cannot perform operation: another operation is in progress

我的疑问

为什么复用一个未关闭且有空闲连接的连接池时,会触发another operation is in progress错误?这里的“另一个操作”具体指的是什么?


问题原因分析

这个错误的核心和asyncio.run()的工作机制强相关:

  • 每次调用asyncio.run()时,都会创建全新的事件循环,执行完毕后自动关闭该循环。
  • 第一次调用时创建的asyncpg Pool,其底层I/O操作完全绑定在第一次的事件循环上。
  • 第二次调用asyncio.run()生成新事件循环后,传入的Pool仍和已关闭的旧循环关联:
    1. 旧循环关闭导致连接的底层I/O操作直接失败(对应报错中的AttributeError: 'NoneType' object has no attribute 'send')。
    2. 当asyncpg尝试处理该错误、重置连接时,又和新循环中的连接操作冲突,最终触发another operation is in progress——这里的“另一个操作”指的是连接在旧循环中残留的失效操作,和新循环中的当前操作产生了资源竞争。

而之前的无数据库代码能正常运行,是因为它没有依赖任何和事件循环绑定的资源,纯异步计算在新循环中可以独立执行。

解决方案

方式1:单事件循环管理完整工作流(推荐)

只调用一次asyncio.run(),在内部管理Pool的完整生命周期,确保所有操作都在同一个事件循环中执行:

import asyncio
import asyncpg


async def d(c: str, pool: asyncpg.Pool):
    async with pool.acquire() as con:
        return await con.fetch("SELECT $1", c)


async def main(input: list[str], pool: asyncpg.Pool):
    result = await asyncio.gather(*(d(i, pool) for i in input))
    print(f"{result = }")
    return result


async def full_workflow():
    # 在同一个事件循环中创建Pool、执行所有任务
    async with asyncpg.create_pool(dsn="Connection arguments") as pool:
        print(f"执行第一次任务")
        await main(["1", "2", "3"], pool)
        print(f"执行第二次任务")
        await main(["2", "3", "4"], pool)


if __name__ == "__main__":
    asyncio.run(full_workflow())

方式2:手动复用事件循环(不推荐,仅特殊场景使用)

如果必须分多次执行,可以手动创建并复用同一个事件循环,避免asyncio.run()的循环重建:

import asyncio
import asyncpg


async def d(c: str, pool: asyncpg.Pool):
    async with pool.acquire() as con:
        return await con.fetch("SELECT $1", c)


async def main(input: list[str], pool: asyncpg.Pool = None):
    if not pool:
        pool = await asyncpg.create_pool(dsn="Connection arguments")
    result = await asyncio.gather(*(d(i, pool) for i in input))
    print(f"{result = }")
    return pool


if __name__ == "__main__":
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    try:
        print(f"Start 1")
        pool = loop.run_until_complete(main(["1", "2", "3"]))
        print(f"{pool.is_closing() = }")
        print(f"Start 2")
        loop.run_until_complete(main(["2", "3", "4"], pool))
    finally:
        # 手动关闭Pool和事件循环
        loop.run_until_complete(pool.close())
        loop.run_until_complete(pool.wait_closed())
        loop.close()

总结

永远不要在不同事件循环之间复用asyncpg的Pool或Connection,它们的底层I/O和创建时的事件循环强绑定。推荐使用第一种方案,通过单asyncio.run()管理完整工作流,能更安全地处理异步资源的生命周期。

备注:内容来源于stack exchange,提问作者Donchack

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 17:43:09