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

Python 3.6下Asyncio WebSocket断连重连报错问题求助

问题分析与解决方案

错误原因解析

你遇到的RuntimeError: Event loop is closed和协程未被await的警告,根源在于这几个核心问题:

  • 事件循环关闭后无法复用:当检测到断连并设置break_task=True后,你调用了loop.stop()和loop.close(),但外层的while True循环会再次尝试用这个已经关闭的事件循环运行新任务——而关闭后的事件循环无法再创建或执行任务,直接触发错误。
  • 协程管理混乱:在break_task=True分支里将tasks = None,但之前的协程可能还处于未完成状态,加上事件循环被强制关闭,导致协程没有被正确await,触发警告。
  • 潜在的变量未赋值问题:在value_2函数中,跳出循环后直接执行data = json.loads(res),如果是因为ConnectionClosed异常跳出的,res根本没有被赋值,后续会触发NameError。

修复后的代码方案

我用asyncio.Event替代全局变量(比全局变量更安全、更符合asyncio的设计逻辑),同时避免关闭事件循环,让任务在断连后自动重连:

import asyncio
import json
import websockets
from copy import deepcopy

# 假设你的MESG和PAYLOAD是预先定义好的
MESG = {"type": "subscribe", "data": "value1"}
PAYLOAD = {"type": "auth", "token": "your_token"}

async def value_1(stop_event):
    while not stop_event.is_set():
        try:
            msg_dict = deepcopy(MESG)
            async with websockets.connect('wss://api.xxxxx') as ws:
                await ws.send(json.dumps(msg_dict))
                while not stop_event.is_set():
                    res = await ws.recv()
                    # 这里处理收到的消息,比如打印或存储
                    print(f"value_1 received: {res}")
        except websockets.exceptions.ConnectionClosed:
            print('CONNECTION CLOSED raised exception in value_1() -> TRYING TO RECONNECT !!!')
            await asyncio.sleep(1)  # 重连前等待1秒,避免频繁重试
        except Exception as e:
            print(f"Unexpected error in value_1: {e}")
            await asyncio.sleep(1)

async def value_2(stop_event):
    while not stop_event.is_set():
        try:
            auth_dict = deepcopy(PAYLOAD)
            async with websockets.connect('wss://api.xxxxx') as ws:
                await ws.send(json.dumps(auth_dict))
                while not stop_event.is_set():
                    await asyncio.sleep(0.05)
                    res = await ws.recv()
                    # 处理收到的消息,确保只有成功接收时才解析
                    data = json.loads(res)
                    print(f"value_2 received: {data}")
        except websockets.exceptions.ConnectionClosed:
            print('CONNECTION CLOSED raised exception in value_2() -> TRYING TO RECONNECT !!!')
            await asyncio.sleep(1)
        except json.JSONDecodeError:
            print("Failed to decode JSON in value_2")
        except Exception as e:
            print(f"Unexpected error in value_2: {e}")
            await asyncio.sleep(1)

async def main():
    stop_event = asyncio.Event()
    try:
        # 同时运行两个WebSocket任务
        await asyncio.gather(value_1(stop_event), value_2(stop_event))
    except KeyboardInterrupt:
        # 捕获Ctrl+C,优雅停止所有任务
        print("Stopping tasks...")
        stop_event.set()

if __name__ == "__main__":
    loop = asyncio.get_event_loop()
    try:
        loop.run_until_complete(main())
    finally:
        loop.close()

关键改进点

  • 用asyncio.Event替代全局变量:stop_event可以安全地在多个协程间传递停止信号,避免全局变量的竞态问题。
  • 每个任务自行处理重连:每个WebSocket任务内部都有循环,当连接断开时,等待一段时间后自动重新建立连接,不需要关闭整个事件循环。
  • 避免事件循环复用错误:整个程序只在最终退出时关闭事件循环,中间重连时复用同一个事件循环。
  • 完善错误处理:捕获了更多异常类型,避免程序意外崩溃,同时添加重连延迟,防止频繁请求给服务器造成压力。
  • 修复变量未赋值问题:将data = json.loads(res)放到连接正常的循环内部,只有成功收到消息时才解析,避免未赋值的错误。

内容的提问来源于stack exchange,提问作者zitoun

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 09:54:46