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

如何在不阻塞主事件循环的前提下等待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中,避免单个客户端发送失败导致整个广播任务崩溃。

之前尝试方法的问题分析

  1. await asyncio.wait(client.recv()):asyncio.wait需要传入future或任务的列表,直接传协程会报错,且无法正确等待消息。
  2. asyncio.wait([client.recv()]):虽然能运行,但会导致每次循环只能处理一个客户端的一条消息,且持有锁时等待消息,仍然会阻塞广播。
  3. asyncio.wait([await client.recv()]):语法错误,await client.recv()返回的是消息内容,不是future,无法传入asyncio.wait。
  4. asyncio.run_coroutine_threadsafe:没有解决全局循环持有锁等待消息的问题,首个客户端的消息接收仍然会阻塞锁,导致广播延迟。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 05:55:04