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
相关产品推荐
相关产品推荐

