如何在同一Locust文件中实现REST与WebSocket的顺序性能测试?
在Locust中实现顺序执行REST POST与WebSocket调用
我需要编写一个性能测试,满足以下要求:
- 按顺序执行两个调用
- 第一个是REST POST请求
- 第二个是WebSocket调用
目前我写的代码里只有REST调用被执行,请问怎么在同一个Locust文件里实现这个需求?是否可行?
原代码如下:
import time import json import logging import re import gevent import websocket from locust import User, FastHttpUser, SequentialTaskSet, task, between headers = [ "auth: headers", ] class MySocketUserClass(User): def connect(self, host: str, header=[], **kwargs): self.ws = websocket.create_connection(host, header=header, **kwargs) gevent.spawn(self.receive_loop) def send(self, body, context={}): self.environment.events.request.fire( request_type="WSS", name='test', response_time=None, response_length=len(body), exception=None, context={**self.context(), **context}, ) logging.debug(f"WSS: {body}") self.ws.send(json.dumps(body)) def on_message(self, message): # override this method in your subclass for custom handling print(message) self.environment.events.request.fire( request_type="WSR", name='test', response_time=0, response_length=len(message), exception=None, context=self.context(), ) def receive_loop(self): while True: message = self.ws.recv() logging.debug(f"WSR: {message}") self.on_message(message) def sleep_with_heartbeat(self, seconds): while seconds >= 0: gevent.sleep(min(15, seconds)) seconds -= 15 self.send({}) class WebSocketCall(MySocketUserClass): @task def my_task(self): self.connect("wss://example.com", header=headers) # example of subscribe self.send({'example': 'payload'}) # wait for additional pushes, while occasionally sending heartbeats, like a real client would self.sleep_with_heartbeat(10) def on_message(self, message): # TO BE IMPLEMENTED pass class RestCall(FastHttpUser): default_headers = { "accept": "application/json", "accept-encoding": "gzip, deflate, br", "accept-language": "en-US,en;q=0.9,it;q=0.8", "content-type": "application/json", "sec-ch-ua": '"Chromium";v="110", "Not A(Brand";v="24", "Microsoft Edge";v="110"', "sec-ch-ua-mobile": "?0", "sec-ch-ua-platform": '"Windows"', "sec-fetch-dest": "empty", "sec-fetch-mode": "cors", "sec-fetch-site": "same-origin", "user-agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/110.0.0.0 Safari/537.36 Edg/110.0.1587.50", "origin": "https://example.com", "referer": "https://example.com", } def call_rest(self): self.iam = 'rest' with self.rest( "POST", "https://example.com", json={ "example": "payload" }, ) as resp: # CHECK RESPONSE print(resp) @task def sent1(self): self.call_rest() class UserBehaviour(SequentialTaskSet): tasks = [RestCall, WebSocketCall]
问题原因
原代码的核心错误是把RestCall和WebSocketCall都定义成了User子类,而SequentialTaskSet的tasks属性只能接受TaskSet子类或者task装饰的方法,不能直接放User类。这导致只有第一个User类的逻辑被执行,第二个被忽略。
可行方案:整合到同一个User的SequentialTaskSet中
Locust的虚拟用户(User类)应该代表一个完整的测试流程执行者,所以我们可以把REST和WebSocket的逻辑整合到同一个User类的SequentialTaskSet里,按顺序执行。
修改后的代码如下:
import time import json import logging import gevent import websocket from locust import FastHttpUser, SequentialTaskSet, task, between # 通用请求头 REST_HEADERS = { "accept": "application/json", "accept-encoding": "gzip, deflate, br", "accept-language": "en-US,en;q=0.9,it;q=0.8", "content-type": "application/json", "sec-ch-ua": '"Chromium";v="110", "Not A(Brand";v="24", "Microsoft Edge";v="110"', "sec-ch-ua-mobile": "?0", "sec-ch-ua-platform": '"Windows"', "sec-fetch-dest": "empty", "sec-fetch-mode": "cors", "sec-fetch-site": "same-origin", "user-agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/110.0.0.0 Safari/537.36 Edg/110.0.1587.50", "origin": "https://example.com", "referer": "https://example.com", } WS_HEADERS = [ "auth: headers", ] # WebSocket功能Mixin,方便复用 class WebSocketMixin: def connect_ws(self, host: str, headers=[]): self.ws = websocket.create_connection(host, header=headers) # 启动消息接收协程 gevent.spawn(self._ws_receive_loop) def send_ws(self, body): payload = json.dumps(body) # 记录发送请求事件 self.environment.events.request.fire( request_type="WSS", name="ws_subscribe", response_time=None, response_length=len(payload), exception=None, context=self.context() ) logging.debug(f"WSS发送: {payload}") self.ws.send(payload) def _ws_receive_loop(self): while True: try: message = self.ws.recv() # 记录接收请求事件 self.environment.events.request.fire( request_type="WSR", name="ws_receive", response_time=0, response_length=len(message), exception=None, context=self.context() ) logging.debug(f"WSR接收: {message}") self.handle_ws_message(message) except websocket.WebSocketConnectionClosedException: logging.warning("WebSocket连接已关闭") break def handle_ws_message(self, message): # 自定义消息处理逻辑,可按需重写 print(f"收到WebSocket消息: {message}") def ws_sleep_with_heartbeat(self, seconds): while seconds > 0: sleep_time = min(15, seconds) gevent.sleep(sleep_time) seconds -= sleep_time # 发送心跳 self.send_ws({}) # 顺序执行的任务集合 class UserTaskSequence(SequentialTaskSet): @task def execute_rest_post(self): # 执行REST POST请求 with self.client.post( "https://example.com", json={"example": "payload"}, headers=REST_HEADERS, name="rest_post_request" ) as resp: print(f"REST响应状态: {resp.status_code}") # 可添加响应断言逻辑 # resp.success() 或自定义判断 @task def execute_websocket_flow(self): # 建立WebSocket连接 self.connect_ws("wss://example.com", headers=WS_HEADERS) # 发送订阅消息 self.send_ws({"example": "payload"}) # 保持连接并发送心跳,模拟真实客户端行为 self.ws_sleep_with_heartbeat(10) # 关闭连接(可选) self.ws.close() # 虚拟用户类,整合HTTP和WebSocket功能 class PerformanceTestUser(FastHttpUser, WebSocketMixin): tasks = [UserTaskSequence] wait_time = between(1, 3) # 任务执行间隔时间
关键修改说明
- 整合功能:让虚拟用户类同时继承
FastHttpUser(提供HTTP客户端)和WebSocketMixin(提供WebSocket功能),这样一个用户可以完成整个流程。 - SequentialTaskSet:把REST和WebSocket调用分别放在
SequentialTaskSet的task方法中,确保按顺序执行。 - 事件记录:保留了WebSocket发送/接收的请求事件记录,这样Locust的统计面板能展示WebSocket的相关指标。
- 代码复用:用Mixin封装WebSocket逻辑,方便后续扩展或复用。
这样修改后,每个虚拟用户都会先执行REST POST请求,再执行WebSocket流程,完全符合需求。
内容的提问来源于stack exchange,提问作者Lorenzo Garuti
相关产品推荐
相关产品推荐

