Python asyncio事件循环冲突报错求助:两种启动方式均失败
解决Python asyncio事件循环冲突的RuntimeError问题
错误原因分析
你遇到的两种错误本质都是同一个线程中同时存在多个运行中的事件循环:
- 用
loop.run_until_complete报错:当前环境已经有一个启动的事件循环,无法重复启动新循环。 - 用
asyncio.run报错:asyncio.run会自动创建新循环,但如果当前已有运行中的循环,它会直接抛出错误(这是Python的设计限制,避免嵌套循环导致的混乱)。
针对性解决方案
情况1:在Jupyter/IPython环境中运行
这类环境默认已经启动了事件循环,不需要手动调用启动方法,直接在单元格中执行:
conn, c = connect_to_database() await execute_main(conn, c)
情况2:普通Python脚本环境
确保整个脚本只有一个循环启动入口,改写启动逻辑:
import asyncio import aiohttp from database_utils import connect_to_database from pytdx.hq import TdxHq_API async def execute_main(conn, c): api = TdxHq_API() api.connect('119.147.212.81', 7709) try: c.execute("SELECT stockCode, shsz FROM stocks") rows = c.fetchall() column_names = [description[0] for description in c.description] results = [dict(zip(column_names, row)) for row in rows] async with aiohttp.ClientSession() as session: tasks = [] for result in results: stock_code = result['stockCode'].strip() market = result['shsz'] tasks.append(fetch_quote(session, api, market, stock_code, conn, c)) await asyncio.gather(*tasks) except Exception as e: print(f"An error occurred: {e}") api.disconnect() print("Execution completed.") async def main(): # 将数据库连接逻辑也放到异步入口中 conn, c = connect_to_database() await execute_main(conn, c) if __name__ == "__main__": # 用asyncio.run作为唯一的循环启动入口 asyncio.run(main())
情况3:当前已有运行中的循环(比如依赖库启动的)
如果无法避免现有循环,直接将协程添加到当前循环中执行:
conn, c = connect_to_database() # 获取当前运行的循环 loop = asyncio.get_running_loop() # 提交协程任务到现有循环 task = loop.create_task(execute_main(conn, c)) # 如果需要等待任务完成,直接await(仅支持顶层await的环境) await task
额外优化:解决同步操作阻塞循环的问题
你的代码中api.connect、数据库c.execute都是同步操作,会阻塞异步事件循环,降低效率。可以用asyncio.to_thread包装同步调用:
async def execute_main(conn, c): api = TdxHq_API() # 包装同步连接操作 await asyncio.to_thread(api.connect, '119.147.212.81', 7709) try: # 包装同步数据库查询 rows = await asyncio.to_thread(lambda: (c.execute("SELECT stockCode, shsz FROM stocks"), c.fetchall())[1]) column_names = [description[0] for description in c.description] results = [dict(zip(column_names, row)) for row in rows] async with aiohttp.ClientSession() as session: tasks = [] for result in results: stock_code = result['stockCode'].strip() market = result['shsz'] tasks.append(fetch_quote(session, api, market, stock_code, conn, c)) await asyncio.gather(*tasks) except Exception as e: print(f"An error occurred: {e}") # 包装同步断开操作 await asyncio.to_thread(api.disconnect) print("Execution completed.")
内容的提问来源于stack exchange,提问作者Steven
相关产品推荐
相关产品推荐

