Python Asyncio报错:Task异常未捕获与队列绑定不同事件循环
Python Asyncio 两个报错的解决办法
遇到的报错
Task exception was never retrievedRuntimeError: <Queue ...> is bound to a different event loop
问题根源
- 事件循环不匹配:队列
queue在主模块初始化时创建,外层又套了无限while循环反复调用asyncio.run(main1())。而asyncio.run()每次执行都会新建独立事件循环,第一次循环时队列绑定了第一个事件循环,第二次调用时用新循环,旧队列与新循环不兼容,触发循环绑定错误。 - 任务异常未捕获:
producer和consumer_algoritm都是无限循环协程,内部抛出异常后没有捕获逻辑,导致出现Task exception was never retrieved错误。
修复方案及修改后的代码
把重连逻辑移到异步函数内部,每次重连创建新队列,同时给任务添加异常捕获:
import asyncio import json import websockets import datetime # 业务相关函数示例,根据实际情况替换 def USDT_balance(): return 100.0 def cancel_order(order_id, coin): pass def get_list_order(): return [], "dummy_order_id" def candleTelo(coin, timeframe): return "initial_ts" # 全局配置变量 timeframe = "1" coin = "SOL" volume = 1.0 ts = "initial_ts" async def TickersChannel_ws(timeframe_): """ Websocket 订阅价格 """ url = "wss://ws.okx.com:8443/ws/v5/business" async with websockets.connect(url) as ws: subs = { "op": "subscribe", "args": [ dict(channel="mark-price-candle" + timeframe_ + "m", instId="SOL-USDT-SWAP") ] } await ws.send(json.dumps(subs)) async for msg in ws: msg = json.loads(msg) if "event" not in msg: yield msg.get("data")[0][0] async def producer(que_): try: async for ticker in TickersChannel_ws(timeframe): await que_.put(ticker) await asyncio.sleep(0.1) except Exception as e: print(f"生产者出错: {e}") async def consumer_algoritm(que_, coin_, volume_): global ts try: while True: ts_ = await que_.get() if ts != ts_: ts = ts_ await asyncio.sleep(0.1) que_.task_done() # 标记队列任务完成,避免积压 except Exception as e: print(f"消费者出错: {e}") async def main(): print(f'启动时间 = {datetime.datetime.today().strftime("%Y-%m-%d-%H.%M.%S")}, 初始余额 = {USDT_balance()}') print("----------------------------------------") cancel_order(get_list_order()[1], coin) # 重连逻辑放在异步函数内部 while True: try: # 每次重连创建新队列,保证与当前事件循环绑定 queue = asyncio.Queue(1) ts = candleTelo(coin, timeframe) task1 = asyncio.create_task(producer(queue)) task2 = asyncio.create_task(consumer_algoritm(queue, coin, volume)) # 用gather统一管理任务,单个任务出错不终止整个程序 await asyncio.gather(task1, task2, return_exceptions=True) print("任务异常退出,准备重连...") await asyncio.sleep(5) except Exception as e: print(f"主循环出错: {e}") await asyncio.sleep(5) if __name__ == "__main__": asyncio.run(main())
关键修改点
- 队列随循环重建:把队列创建放到内部while循环中,每次重连生成新队列,确保队列与当前事件循环绑定,解决循环不匹配问题。
- 任务内部异常捕获:给
producer和consumer_algoritm添加try-except,捕获内部异常并打印,避免Task exception was never retrieved错误。 - 用gather管理任务:替换单独await任务的方式,
asyncio.gather可同时等待多个任务,return_exceptions=True设置后,单个任务出错不会直接终止程序,方便后续重连处理。 - 添加task_done():消费者处理完队列元素后调用该方法,保证队列内部计数正确,避免内存泄漏。
内容的提问来源于stack exchange,提问作者Кирилл Маликов
相关产品推荐
相关产品推荐

