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

如何在同一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)  # 任务执行间隔时间

关键修改说明

  1. 整合功能:让虚拟用户类同时继承FastHttpUser(提供HTTP客户端)和WebSocketMixin(提供WebSocket功能),这样一个用户可以完成整个流程。
  2. SequentialTaskSet:把REST和WebSocket调用分别放在SequentialTaskSet的task方法中,确保按顺序执行。
  3. 事件记录:保留了WebSocket发送/接收的请求事件记录,这样Locust的统计面板能展示WebSocket的相关指标。
  4. 代码复用:用Mixin封装WebSocket逻辑,方便后续扩展或复用。

这样修改后,每个虚拟用户都会先执行REST POST请求,再执行WebSocket流程,完全符合需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 22:12:33