RuntimeError报错:多客户端连接时无法调用recv的问题
问题分析
错误RuntimeError: cannot call recv while another coroutine is already waiting for the next message的核心原因是:多个WebsocketConsumer实例共享同一个websocket连接,并且每个实例都试图独立从该连接接收消息。
ConnectionManager是单例实现,所有Consumer拿到的都是同一个traccar_connection对象。当第二个客户端连接时,第二个Consumer的handle_connection协程会再次调用async for traccar_data in traccar_connection,而此时第一个Consumer的协程已经在等待该连接的消息,websockets库不允许多个协程同时在同一个连接上执行接收操作,因此抛出错误。
另外代码还有两个潜在问题:
- 每次收到消息都调用
asyncio.ensure_future(self.process_message_queue()),会创建大量重复协程,导致消息被重复发送 requests.Session().post是同步阻塞调用,在异步环境中会阻塞事件循环
解决方案
重构代码,让单例连接统一负责接收消息,然后通过channel layer广播给所有在线Consumer,避免多个协程争抢同一个连接的接收权。
1. 修改ConnectionManager,统一管理消息接收与广播
import asyncio import aiohttp from urllib.parse import urlencode import websockets from channels.layers import get_channel_layer class ConnectionManager: _instance = None connection = None connection_lock = asyncio.Lock() is_running = False def __new__(cls, *args, **kwargs): if not cls._instance: cls._instance = super().__new__(cls, *args, **kwargs) return cls._instance async def get_connection(self): async with self.connection_lock: if self.connection is None or not self.connection.open: self.connection = await self.create_connection() # 仅启动一次消息接收协程 if not self.is_running: asyncio.create_task(self.listen_and_broadcast()) self.is_running = True return self.connection async def create_connection(self): email = "email" password = "Password" base_url = "url" async with aiohttp.ClientSession() as session: params = {"email": email, "password": password} headers = { "content-type": "application/x-www-form-urlencoded", "accept": "application/json", } async with session.post(f"{base_url}/api/session", data=params, headers=headers) as resp: cookies = resp.cookies token = cookies["JSESSIONID"].value return await websockets.connect( "ws://url", ping_interval=None, extra_headers=(("Cookie", f"JSESSIONID={token}"),), ) async def listen_and_broadcast(self): channel_layer = get_channel_layer() try: while True: if not self.connection or not self.connection.open: await self.get_connection() traccar_data = await self.connection.recv() modified_data = self.process_data(traccar_data) # 通过channel layer广播到组 await channel_layer.group_send( "group", { " "type": "send_to_client", "data": modified_data } ) except websockets.exceptions.ConnectionClosed: self.is_running = False self.connection = None # 尝试重连 await asyncio.sleep(5) await self.get_connection() @staticmethod def process_data(traccar_data): return f"{traccar_data} - Processed by Django"
2. 简化WebsocketConsumer,仅负责接收组消息并发送给客户端
from channels.generic.websocket import AsyncWebsocketConsumer class WebsocketConsumer(AsyncWebsocketConsumer): async def connect(self): await self.accept() await self.channel_layer.group_add("group", self.channel_name) async def disconnect(self, close_code): await self.channel_layer.group_discard("group", self.channel_name) # 处理组广播的消息 async def send_to_client(self, event): data = event["data"] await self.send(str(data)) async def receive(self, text_data): # 不需要处理客户端发送的消息可保留空实现 pass
关键改进点
- 单例
ConnectionManager仅启动一个协程负责接收消息,避免多个协程争抢连接 - 使用
channel_layer.group_send统一广播消息,所有在线Consumer都能收到 - 将同步的
requests替换为异步的aiohttp,避免阻塞异步事件循环 - 增加连接断开后的自动重连逻辑,提升稳定性
内容的提问来源于stack exchange,提问作者Vahid
相关产品推荐
相关产品推荐

