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

Python Asyncio报错:Task异常未捕获与队列绑定不同事件循环

Python Asyncio 两个报错的解决办法

遇到的报错

  • Task exception was never retrieved
  • RuntimeError: <Queue ...> is bound to a different event loop

问题根源

  1. 事件循环不匹配:队列queue在主模块初始化时创建,外层又套了无限while循环反复调用asyncio.run(main1())。而asyncio.run()每次执行都会新建独立事件循环,第一次循环时队列绑定了第一个事件循环,第二次调用时用新循环,旧队列与新循环不兼容,触发循环绑定错误。
  2. 任务异常未捕获: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())

关键修改点

  1. 队列随循环重建:把队列创建放到内部while循环中,每次重连生成新队列,确保队列与当前事件循环绑定,解决循环不匹配问题。
  2. 任务内部异常捕获:给producer和consumer_algoritm添加try-except,捕获内部异常并打印,避免Task exception was never retrieved错误。
  3. 用gather管理任务:替换单独await任务的方式,asyncio.gather可同时等待多个任务,return_exceptions=True设置后,单个任务出错不会直接终止程序,方便后续重连处理。
  4. 添加task_done():消费者处理完队列元素后调用该方法,保证队列内部计数正确,避免内存泄漏。

内容的提问来源于stack exchange,提问作者Кирилл Маликов

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 03:23:18