如何在不阻塞主事件循环的前提下等待WebSocket客户端消息?
WebSocket广播被客户端消息接收阻塞的解决方法
我正在编写一个简单的WebSocket脚本,实现客户端注册与注销功能,并每隔5秒向客户端广播随机消息,但遇到一个问题:一旦有客户端连接,广播功能就会停止,事件循环卡在await client.recv()语句处。需要实现监听客户端消息但不阻塞主事件循环的效果。
原代码如下:
import asyncio, websockets, random, string import websockets.asyncio.server class web_socket(websockets.asyncio.server.ServerConnection): pass connections: set[web_socket] = set() connections_lock = asyncio.Lock() async def register(websocket: web_socket): async with connections_lock: connections.add(websocket) async def unregister(websocket: web_socket): async with connections_lock: connections.discard(websocket) async def cleanup_connections(): while True: async with connections_lock: closed_clients = [client for client in connections if client.close_code == 1000] for client in closed_clients: connections.discard(client) await asyncio.sleep(1) async def handler(websocket:web_socket): await register(websocket) try: await websocket.wait_closed() except Exception as e: print(f"Error: {e}") async def random_messages(): while True: message = '-'.join(random.choices(string.ascii_letters + string.digits, k=10)) async with connections_lock: send_tasks = [client.send(message) for client in connections if client.close_code != 1000] if send_tasks: await asyncio.gather(*send_tasks) print(f"Broadcasted message: {message}") await asyncio.sleep(5) async def respond_to_messages(): while True: async with connections_lock: if connections: for client in connections: message = await client.recv() if message: print(message) else: await asyncio.sleep(1) async def main(): async with websockets.serve(handler, "0.0.0.0", 8765): asyncio.create_task(random_messages()) asyncio.create_task(cleanup_connections()) asyncio.create_task(respond_to_messages()) await asyncio.Future() asyncio.run(main())
问题根源
问题出在respond_to_messages函数里:
- 持有
connections_lock时调用await client.recv(),这会导致锁被长时间占用(直到客户端发送消息),而广播任务random_messages需要获取这个锁来遍历连接发送消息,因此被阻塞。 - 遍历所有客户端并逐个等待消息,只要有一个客户端不发送消息,整个任务就会卡在该客户端的
recv()调用上,无法处理其他客户端或释放锁。
正确的实现方式
不需要单独写全局的respond_to_messages函数,而是在每个客户端连接的handler里启动一个独立的任务来处理该客户端的消息接收,这样每个客户端的消息处理是独立的,不会互相阻塞,也不会长时间持有连接锁。
修改后的代码:
import asyncio, websockets, random, string import websockets.asyncio.server class web_socket(websockets.asyncio.server.ServerConnection): pass connections: set[web_socket] = set() connections_lock = asyncio.Lock() async def register(websocket: web_socket): async with connections_lock: connections.add(websocket) async def unregister(websocket: web_socket): async with connections_lock: connections.discard(websocket) async def cleanup_connections(): while True: async with connections_lock: closed_clients = [client for client in connections if client.close_code == 1000] for client in closed_clients: connections.discard(client) await asyncio.sleep(1) # 处理单个客户端的消息接收 async def handle_client_messages(websocket: web_socket): try: while True: message = await websocket.recv() if message: print(f"Received from client: {message}") except websockets.exceptions.ConnectionClosed: # 连接关闭时无需额外处理,handler里会注销 pass except Exception as e: print(f"Message handling error: {e}") async def handler(websocket: web_socket): await register(websocket) # 启动独立任务处理该客户端的消息 message_task = asyncio.create_task(handle_client_messages(websocket)) try: await websocket.wait_closed() finally: # 确保消息处理任务被取消 message_task.cancel() await unregister(websocket) async def random_messages(): while True: message = '-'.join(random.choices(string.ascii_letters + string.digits, k=10)) async with connections_lock: # 复制连接集合,避免遍历过程中集合被修改 active_clients = list(connections) # 过滤出未关闭的客户端并发送消息 send_tasks = [client.send(message) for client in active_clients if client.close_code != 1000] if send_tasks: await asyncio.gather(*send_tasks, return_exceptions=True) print(f"Broadcasted message: {message}") await asyncio.sleep(5) async def main(): async with websockets.serve(handler, "0.0.0.0", 8765): asyncio.create_task(random_messages()) asyncio.create_task(cleanup_connections()) await asyncio.Future() asyncio.run(main())
关键修改说明
- 每个客户端独立处理消息:在
handler中为每个新连接创建handle_client_messages任务,专门处理该客户端的消息接收,避免全局循环阻塞。 - 避免锁长时间占用:广播任务中先获取锁复制连接集合,然后释放锁再处理发送,减少锁的持有时间,不会被消息接收阻塞。
- 任务清理:在
handler的finally块中取消消息处理任务并注销客户端,确保资源正确释放。 - 容错处理:使用
return_exceptions=True在asyncio.gather中,避免单个客户端发送失败导致整个广播任务崩溃。
之前尝试方法的问题分析
await asyncio.wait(client.recv()):asyncio.wait需要传入future或任务的列表,直接传协程会报错,且无法正确等待消息。asyncio.wait([client.recv()]):虽然能运行,但会导致每次循环只能处理一个客户端的一条消息,且持有锁时等待消息,仍然会阻塞广播。asyncio.wait([await client.recv()]):语法错误,await client.recv()返回的是消息内容,不是future,无法传入asyncio.wait。asyncio.run_coroutine_threadsafe:没有解决全局循环持有锁等待消息的问题,首个客户端的消息接收仍然会阻塞锁,导致广播延迟。
内容的提问来源于stack exchange,提问作者Himanshu
相关产品推荐
相关产品推荐

