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

基于websockets+urwid+asyncio的服务端/客户端退出异常排查

解决WebSocket+Urwid应用的未捕获异常与退出问题

让我们一步步拆解并解决你的问题,核心问题分为两类:服务端的未捕获Task异常和客户端的终端退出异常与连接管理问题。

问题根源分析

服务端侧

  1. 连接未正确移除:你的consumer_handler仅在捕获到异常时才从connected集合移除WebSocket,但当客户端正常关闭连接(按Q键)时,async for message in websocket会正常结束,不会进入except块,导致已关闭的连接始终留在connected集合中,后续广播任务持续尝试向其发送消息,抛出未捕获异常。
  2. 广播任务未处理单个连接异常:使用asyncio.wait时,单个WebSocket发送失败的异常会被封装在Task中,外层的try-except无法捕获这些Task内部的异常,从而出现「Task exception was never retrieved」警告。

客户端侧

  1. 事件循环逻辑错误:原代码中urwid_loop.start()是阻塞调用,导致consumer_handler直到Urwid退出后才启动,这与你描述的“客户端启动后服务端显示连接数1”矛盾,说明实际运行中可能存在逻辑错位,正确的做法是让Urwid循环与WebSocket客户端任务同时运行。
  2. 退出时未优雅关闭连接:按Q键或Ctrl+C时,未主动关闭WebSocket连接,导致服务端无法及时感知连接状态变化;同时Urwid未正确重置终端模式,引发终端异常。

服务端代码修复

1. 确保连接始终被移除

修改consumer_handler,将连接移除逻辑放在finally块中,无论正常退出还是异常退出,都从connected集合中清理无效连接:

async def consumer_handler(websocket, path):
    global connected
    connected.add(websocket)
    try:
        async for message in websocket:
            print(f"Received: {message}")
    except Exception as e:
        print(f"Consumer error: {str(e)}")
    finally:
        # 确保无论是否异常,都移除连接
        if websocket in connected:
            connected.remove(websocket)
            print("Unregistered websocket, current connections:", len(connected))

2. 广播任务处理单个连接异常

改用asyncio.gather并设置return_exceptions=True,将每个发送任务的异常作为结果返回,而非抛出,同时清理发送失败的连接:

async def broadcast_task():
    global connected
    while True:
        data = await handle_read_external_data()
        data_json = json.dumps(data)
        print("Current connections:", len(connected))
        
        if connected:
            # 复制集合避免遍历中集合变化
            connections_copy = list(connected)
            # 批量发送,返回异常而非抛出
            results = await asyncio.gather(
                *[ws.send(data_json) for ws in connections_copy],
                return_exceptions=True
            )
            # 检查结果,移除无效连接
            for ws, result in zip(connections_copy, results):
                if isinstance(result, Exception):
                    print(f"Send failed for connection: {str(result)}")
                    if ws in connected:
                        connected.remove(ws)
        
        await asyncio.sleep(1)

完整服务端代码:

import asyncio, websockets, random, json

connected = set()

##########################
# end-less loop routines
##########################
async def handle_read_external_data():
    await asyncio.sleep(1)
    data = {'value': (random.random() * 0.4) - 0.2 + 40}
    return data

async def broadcast_task():
    global connected
    while True:
        data = await handle_read_external_data()
        data_json = json.dumps(data)
        print("Current connections:", len(connected))
        
        if connected:
            connections_copy = list(connected)
            results = await asyncio.gather(
                *[ws.send(data_json) for ws in connections_copy],
                return_exceptions=True
            )
            for ws, result in zip(connections_copy, results):
                if isinstance(result, Exception):
                    print(f"Send failed for connection: {str(result)}")
                    if ws in connected:
                        connected.remove(ws)
        
        await asyncio.sleep(1)

##########################
# websockets
##########################
async def consumer_handler(websocket, path):
    global connected
    connected.add(websocket)
    try:
        async for message in websocket:
            print(f"Received: {message}")
    except Exception as e:
        print(f"Consumer error: {str(e)}")
    finally:
        if websocket in connected:
            connected.remove(websocket)
            print("Unregistered websocket, current connections:", len(connected))

##########################
# Asyncio event loop
##########################
def start_event_loop():
    loop = asyncio.get_event_loop()
    tasks = asyncio.gather(
        broadcast_task(),
        websockets.serve(consumer_handler, '0.0.0.0', 8765),
    )
    try:
        loop.run_until_complete(tasks)
    except KeyboardInterrupt:
        print("\nShutting down server...")
    finally:
        loop.close()

if __name__ == "__main__":
    start_event_loop()

客户端代码修复

1. 同时运行Urwid与WebSocket任务

使用asyncio.gather让Urwid循环和WebSocket客户端任务同时执行,避免阻塞;添加退出事件和信号处理,确保优雅关闭连接并恢复终端。

完整客户端代码:

import asyncio, websockets, json, urwid
from signal import SIGINT, SIGTERM

# 全局退出事件,用于通知所有任务退出
exit_event = asyncio.Event()
txt = urwid.Text(u"Hello World")

def show_or_exit(key):
    if key in ('q', 'Q'):
        exit_event.set()
        raise urwid.ExitMainLoop()
    return key

def handle_signal():
    """处理Ctrl+C等信号,触发退出逻辑"""
    exit_event.set()
    raise urwid.ExitMainLoop()

async def consumer_handler(websocket):
    """处理WebSocket消息接收"""
    while not exit_event.is_set():
        try:
            # 添加超时,避免阻塞在recv上无法响应退出事件
            message = await asyncio.wait_for(websocket.recv(), timeout=0.5)
            data = json.loads(message)
            txt.set_text(f"data: {data['value']:.4f}")
        except asyncio.TimeoutError:
            continue
        except websockets.exceptions.ConnectionClosed:
            print("WebSocket connection closed")
            break

async def websocket_client():
    """WebSocket客户端主逻辑"""
    try:
        async with websockets.connect('ws://localhost:8765') as websocket:
            await consumer_handler(websocket)
    except Exception as e:
        print(f"WebSocket client error: {str(e)}")
    finally:
        exit_event.set()

async def main():
    global txt
    fill = urwid.Filler(txt, 'top')
    
    loop = asyncio.get_running_loop()
    # 注册信号处理,捕获Ctrl+C
    loop.add_signal_handler(SIGINT, handle_signal)
    loop.add_signal_handler(SIGTERM, handle_signal)
    
    # 使用AsyncioEventLoop整合Urwid到asyncio循环
    evl = urwid.AsyncioEventLoop(loop=loop)
    urwid_loop = urwid.MainLoop(fill, event_loop=evl, unhandled_input=show_or_exit)
    
    # 将Urwid循环包装为async任务
    urwid_task = loop.create_task(urwid_loop.run())
    # 启动WebSocket客户端任务
    ws_task = loop.create_task(websocket_client())
    
    # 等待任一任务完成(退出事件触发或Urwid退出)
    await asyncio.wait([urwid_task, ws_task], return_when=asyncio.FIRST_COMPLETED)

if __name__ == "__main__":
    try:
        asyncio.run(main())
    except Exception:
        pass
    finally:
        # 强制恢复终端正常模式,避免Ctrl+C后终端异常
        import os
        os.system('stty sane')
        print("\nClient exited successfully")

验证效果

  1. 启动服务端后,控制台持续显示当前连接数;
  2. 启动客户端后,服务端显示连接数变为1,客户端界面实时接收数据;
  3. 按Q键退出客户端:服务端会打印「Unregistered websocket」,连接数变为0,无未捕获异常;
  4. 按Ctrl+C退出客户端:终端恢复正常,服务端正确清理连接。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:31:22