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

如何在Locust中控制多用户类创建顺序或并行执行任务

两种实现方案

方案一:分阶段启动用户类(先Sync后Normal)

通过全局集合存储同步后的用户ID,结合Locust事件钩子控制用户启动顺序,保证NormalUser获取的ID是SyncUser的子集。

代码实现

from locust import HttpUser, task, events, run_single_user
import time
import random
from locust.runners import MasterRunner

# 全局存储SyncUser生成的userId
synced_user_ids = set()

class SyncUser(HttpUser):
    host = '<hostname>'
    wait_time = 0

    def on_start(self):
        self.userId = getsyncPullUser()
        synced_user_ids.add(self.userId)
        self.headers = {"Content-Type": "application/json"}

    @task
    def sync_task(self):
        params = {"timeout": 10}
        self.client.get("/sync/api", headers=self.headers, params=params, name="GET_sync/api")

class NormalUser(HttpUser):
    host = '<hostname>'
    wait_time = 0.5

    def on_start(self):
        # 从同步后的ID集合中选取,替代原getuser()
        if synced_user_ids:
            self.userId = random.choice(list(synced_user_ids))
        else:
            # 兜底逻辑,避免集合为空
            self.userId = getuser()
        self.headers = {"Content-Type": "application/json"}

    @task
    def get_messages(self):
        params = {"status":"ACCEPTED", "source":"postman"}
        self.client.get("/<api>", headers=self.headers, params=params, name="GET /<api>")

# 控制启动顺序的事件钩子
@events.test_start.add_listener
def on_test_start(environment, **kwargs):
    if isinstance(environment.runner, MasterRunner):
        # 分布式主节点:先启动SyncUser,延迟后启动NormalUser
        environment.runner.start_user_classes([SyncUser])
        time.sleep(30)  # 等待同步完成,时间按需调整
        environment.runner.start_user_classes([NormalUser])
    else:
        # 单机模式:先跑SyncUser实例,再跑NormalUser
        run_single_user(SyncUser)
        time.sleep(30)
        run_single_user(NormalUser)

运行命令

locust -f your_script.py --headless -u 10 -r 2

方案二:单用户类并行执行多任务

在同一个用户类中,用后台线程持续运行sync任务,主线程执行normal任务,确保sync逻辑先启动或并行运行。

代码实现

from locust import HttpUser, task
import threading
import time

class CombinedUser(HttpUser):
    host = '<hostname>'
    stop_event = None

    def on_start(self):
        self.userId = getsyncPullUser()
        self.headers = {"Content-Type": "application/json"}
        
        # 启动sync任务后台线程
        self.stop_event = threading.Event()
        sync_thread = threading.Thread(target=self.sync_loop, daemon=True)
        sync_thread.start()
        
        # 等待sync先执行一轮,确保ID已同步
        time.sleep(5)

    def sync_loop(self):
        params = {"timeout": 10}
        while not self.stop_event.is_set():
            self.client.get("/sync/api", headers=self.headers, params=params, name="GET_sync/api")
            time.sleep(1)  # 控制sync任务执行间隔

    @task
    def get_messages(self):
        params = {"status":"ACCEPTED", "source":"postman"}
        self.client.get("/<api>", headers=self.headers, params=params, name="GET /<api>")

    def on_stop(self):
        if self.stop_event:
            self.stop_event.set()

说明

  • 方案一适合大规模用户分角色执行任务的场景,严格保证ID子集关系
  • 方案二适合单用户同时承载两种任务的场景,资源占用更集中

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 04:25:06