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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 09:22:30