复用未关闭的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仍和已关闭的旧循环关联:- 旧循环关闭导致连接的底层I/O操作直接失败(对应报错中的
AttributeError: 'NoneType' object has no attribute 'send')。 - 当asyncpg尝试处理该错误、重置连接时,又和新循环中的连接操作冲突,最终触发
another operation is in progress——这里的“另一个操作”指的是连接在旧循环中残留的失效操作,和新循环中的当前操作产生了资源竞争。
- 旧循环关闭导致连接的底层I/O操作直接失败(对应报错中的
而之前的无数据库代码能正常运行,是因为它没有依赖任何和事件循环绑定的资源,纯异步计算在新循环中可以独立执行。
解决方案
方式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
相关产品推荐
相关产品推荐

