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

Locust多Worker场景下高效执行昂贵进程初始化的正确方式

Locust多Worker场景下大型数据集加载优化方案

针对多Worker(--processes=N)场景下重复下载、Worker阻塞心跳的问题,推荐以下几种落地解决方案:

方案1:主进程统一处理,Worker共享内存获取数据

  • 仅让主/本地Runner在init事件中完成下载+加载,通过进程间共享内存把数据同步给所有Worker,避免重复下载和磁盘IO竞争
  • 核心是利用multiprocessing.Manager创建可跨进程访问的共享数据结构,Worker无需独自处理耗时操作,也不会阻塞心跳
from locust import events
from multiprocessing import Manager
import requests

shared_dataset = None

@events.init.add_listener
def init_dataset(environment, **kwargs):
    global shared_dataset
    # 仅主进程执行下载和加载
    if environment.runner and environment.runner.local:
        print("主进程启动数据集下载...")
        # 下载大型数据集
        resp = requests.get("https://your-dataset-url")
        with open("raw_dataset.txt", "wb") as f:
            f.write(resp.content)
        # 加载并预处理数据
        with open("raw_dataset.txt", "r") as f:
            dataset = [line.strip() for line in f if line.strip()]
        # 初始化共享内存容器
        manager = Manager()
        shared_dataset = manager.list(dataset)
        # 把共享数据传递给Worker
        environment.runner.shared_dataset = shared_dataset

@events.test_start.add_listener
def setup_worker_data(environment, **kwargs):
    global shared_dataset
    # Worker直接从主进程获取共享数据
    if not environment.runner.local:
        shared_dataset = environment.runner.shared_dataset

方案2:Worker异步加载,避免阻塞心跳线程

  • 在test_start事件中启动独立线程执行CPU密集的加载操作,主线程继续处理Locust的心跳逻辑,防止Worker被标记为缺失
  • 可添加等待逻辑,确保数据加载完成后再启动用户任务
import threading
from locust import events
import time

dataset = None

def load_dataset_task():
    global dataset
    print("Worker开始加载数据集...")
    # 模拟CPU密集型加载操作
    with open("raw_dataset.txt", "r") as f:
        dataset = [line.strip() for line in f if line.strip()]

@events.test_start.add_listener
def init_worker_data(environment, **kwargs):
    if not environment.runner.local:
        # 启动后台线程加载数据
        load_thread = threading.Thread(target=load_dataset_task, daemon=True)
        load_thread.start()
        # 等待数据加载完成,避免任务启动时数据未就绪
        while dataset is None:
            time.sleep(0.1)

方案3:预处理数据集,降低Worker加载开销

  • 主进程下载后将数据集预处理为序列化格式(如pickle),大幅减少Worker加载时的CPU运算量,缩短阻塞时间
  • 序列化后的文件加载速度远快于原始文本解析,从根源上降低心跳阻塞风险
import pickle
from locust import events
import requests

@events.init.add_listener
def preprocess_dataset(environment, **kwargs):
    if environment.runner and environment.runner.local:
        print("主进程下载并预处理数据集...")
        resp = requests.get("https://your-dataset-url")
        raw_data = resp.text.splitlines()
        dataset = [line.strip() for line in raw_data if line.strip()]
        # 序列化保存预处理后的数据
        with open("dataset.pkl", "wb") as f:
            pickle.dump(dataset, f)

@events.test_start.add_listener
def load_worker_dataset(environment, **kwargs):
    if not environment.runner.local:
        global dataset
        # 快速加载序列化文件,CPU开销极低
        with open("dataset.pkl", "rb") as f:
            dataset = pickle.load(f)

额外注意事项

  • 确保数据集文件存放在所有Worker都能访问的路径(如本地共享目录)
  • 共享内存方案适合中小规模数据集,超大型数据建议用预处理+异步加载组合
  • 所有方案需在User任务中校验数据集是否就绪,避免空指针错误

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 05:25:05