Locust WebSocket单用户测试正常但多用户并发执行报错‘socket already closed’问题咨询
解决Locust多用户并发WebSocket测试的"socket already closed"问题
我帮你排查了下,单用户跑正常但多用户并发就报连接关闭,核心是代码里踩了几个WebSocket和Locust协作的坑,咱们一个个来填:
1. 所有用户共用同一个user_id,服务器直接踢掉旧会话
你代码里的user_id是类级别的固定值,所有并发用户都用同一个ID发起请求。服务器端一般会对同一个用户的会话做唯一性限制——新会话建立后,旧的就被强制关闭了,这就是你看到"socket already closed"的头号原因。
修正方式:给每个测试用户生成唯一标识,比如用uuid模块,把user_id改成实例属性,在用户初始化时生成:
import uuid # 记得导入uuid模块 class WebsiteUser(WebsocketLocust): root="wss://xxxx.com/" # 把URI改成模板,后面用user_id动态填充 uri_start_conversation_template = f"{root}/xxx?bot_integer_id=xx" uri_send_audio_template = f"{root}messages/{{user_id}}/ws/bot/xxx?bot_integer_id=xx" @task def conversation(self): # ... 后续代码不变,只是引用self.uri_send_audio self.speech(get_url=self.uri_send_audio, user_id=self.user_id, message=send_audio_file) def on_start(self): # 每个用户生成唯一ID self.user_id = str(uuid.uuid4()) # 填充模板生成当前用户的专属连接地址 self.url = self.uri_start_conversation_template self.uri_send_audio = self.uri_send_audio_template.format(user_id=self.user_id) websocket.enableTrace(True) self.client.connect(self.url) init_msg = json.dumps({ "user_id": self.user_id, "message": "", "source": "web" }) self.client.send(init_msg)
2. Asyncio和Locust的Gevent协程打架,连接管理乱套
Locust是基于Gevent的协程框架,你却在代码里混用了asyncio的事件循环和websockets异步库——这俩协程模型不兼容,会导致连接状态混乱,时不时就出现连接被意外关闭的情况。
修正方式:统一用同步的WebSocket客户端(比如你已经导入的websocket-client)来处理第二个音频连接,放弃asyncio那套,和Locust的模型保持一致:
# 把原来的async speech方法改成同步实现 def speech(self, get_url: str, user_id: str, message): Flag = True ws = None try: print(Fore.YELLOW + f"[speech] 为用户 {user_id} 建立音频WS连接") ws = create_connection(get_url) print(Fore.GREEN + f"[audio] 用户 {user_id} 音频WS连接成功") while Flag: ws.send(message) packet = ws.recv() packet = json.loads(packet) if packet["reply"] != ["Complete"]: question = int(packet.get("reply")[0].split("|")[0]) message = audio_files(index_value=question) else: Flag = False print(Fore.BLUE + f"[audio] 用户 {user_id} 会话结束") except Exception as e: print(Fore.RED + f"用户 {user_id} 音频WS出错: {e}") finally: # 无论成功失败,都要关闭连接,避免资源泄漏 if ws: ws.close() print(Fore.CYAN + f"[audio] 用户 {user_id} 音频WS已关闭")
然后在conversation任务里直接调用这个同步方法就行,不用再折腾asyncio:
@task def conversation(self): while True: try: recv = json.loads(self.client.recv()) print(recv) question = int(recv.get("reply")[0].split("|")[0]) send_audio_file = audio_files(index_value=question) # 直接调用同步方法 self.speech(get_url=self.uri_send_audio, user_id=self.user_id, message=send_audio_file) except Exception as e: print(f"会话任务出错: {e}") break # 或者根据需求做重试逻辑
3. 连接没正确释放,积累无效连接导致服务器主动关闭
之前的代码如果遇到异常,可能会导致WebSocket连接没被关闭,这些挂着的无效连接积累多了,服务器会主动清理,进而影响正常用户的连接。
修正方式:给所有WebSocket操作加上try...finally块,确保连接一定会被关闭——比如上面的speech方法已经加了,你也可以给on_start里的主连接加上异常处理:
def on_start(self): self.user_id = str(uuid.uuid4()) self.url = self.uri_start_conversation_template self.uri_send_audio = self.uri_send_audio_template.format(user_id=self.user_id) enableTrace(True) try: self.client.connect(self.url) init_msg = json.dumps({ "user_id": self.user_id, "message": "", "source": "web" }) self.client.send(init_msg) print(f"用户 {self.user_id} 会话初始化完成") except Exception as e: print(f"用户 {self.user_id} 初始化会话失败: {e}")
整合后的完整代码
把上面的修改整合后,完整代码大概是这样:
import json import uuid from websocket import create_connection, enableTrace from locust import WebsocketLocust, task from colorama import Fore def audio_files(index_value): # 保留你的音频文件获取逻辑,返回二进制数据 return b"test_audio_data" class WebsiteUser(WebsocketLocust): root="wss://xxxx.com/" uri_start_conversation_template = f"{root}/xxx?bot_integer_id=xx" uri_send_audio_template = f"{root}messages/{{user_id}}/ws/bot/xxx?bot_integer_id=xx" min_wait = 1000 # 任务间隔最小1秒 max_wait = 3000 # 任务间隔最大3秒 @task def conversation(self): while True: try: recv = json.loads(self.client.recv()) print(recv) question = int(recv.get("reply")[0].split("|")[0]) send_audio_file = audio_files(index_value=question) self.speech(get_url=self.uri_send_audio, user_id=self.user_id, message=send_audio_file) except Exception as e: print(f"会话任务出错: {e}") break def on_start(self): self.user_id = str(uuid.uuid4()) self.url = self.uri_start_conversation_template self.uri_send_audio = self.uri_send_audio_template.format(user_id=self.user_id) enableTrace(True) try: self.client.connect(self.url) init_msg = json.dumps({ "user_id": self.user_id, "message": "", "source": "web" }) self.client.send(init_msg) print(f"用户 {self.user_id} 会话初始化完成") except Exception as e: print(f"用户 {self.user_id} 初始化会话失败: {e}") def speech(self, get_url: str, user_id: str, message): Flag = True ws = None try: print(Fore.YELLOW + f"[speech] 为用户 {user_id} 建立音频WS连接") ws = create_connection(get_url) print(Fore.GREEN + f"[audio] 用户 {user_id} 音频WS连接成功") while Flag: ws.send(message) packet = ws.recv() packet = json.loads(packet) if packet["reply"] != ["Complete"]: question = int(packet.get("reply")[0].split("|")[0]) message = audio_files(index_value=question) else: Flag = False print(Fore.BLUE + f"[audio] 用户 {user_id} 会话结束") except Exception as e: print(Fore.RED + f"用户 {user_id} 音频WS出错: {e}") finally: if ws: ws.close() print(Fore.CYAN + f"[audio] 用户 {user_id} 音频WS已关闭")
最后验证的几个关键点
- 每个用户都有唯一的
user_id,不会再出现会话冲突 - 统一用同步WebSocket客户端,和Locust的Gevent模型兼容,不会再出现协程打架的问题
- 所有连接都有明确的关闭逻辑,避免资源泄漏
- 加了异常处理,方便你调试并发时的各种问题
内容的提问来源于stack exchange,提问作者Boitumelo
相关产品推荐
相关产品推荐

