Locust分布式模式报错:未知消息类型_distributor_request致用户停滞
问题场景
使用Locust分布式(1主多Worker)执行负载测试,借助locust-plugins库的Distributor工具为每个虚拟用户预分配唯一资源,配置流程如下:
- Master节点:测试启动时生成对应总用户数的资源列表,通过自定义消息广播给所有Worker
- Worker节点:接收资源列表后初始化Distributor实例供用户任务调用
- 用户任务:在
use_model任务中通过next(items_models_dist)获取资源
执行到model = next(items_models_dist)时,出现警告:Unknown message type received from worker worker.some.id (index 0): _distributor_request,且虚拟用户停滞。未显式注册该消息处理程序,原以为插件会自动处理。
版本信息:Locust 2.32.3,Python 3.12.3,locust-plugins 4.5.3;启动命令:locust -f test.py --processes -1
解决办法
1. 在Master节点注册Distributor消息处理程序
_distributor_request是locust-plugins内部用于Worker向Master请求资源分配的消息类型,必须在Master端注册对应处理函数才能被识别。在测试文件的全局区域添加:
from locust_plugins.distributor import register_distributor_messages # 注册Distributor所需的消息处理器,仅Master节点会生效 register_distributor_messages()
2. 改用Distributor内置的分布式同步机制
不要手动通过自定义消息传递资源列表,利用插件自带的同步逻辑更可靠:
- Master端在全局测试启动事件中初始化Distributor并传入资源列表
- Worker端无需手动接收消息,直接调用Distributor实例即可
示例代码:
from locust import HttpUser, task, events from locust_plugins.distributor import Distributor, register_distributor_messages # 全局注册消息处理器 register_distributor_messages() # 全局Distributor实例 items_models_dist = None @events.test_start.add_listener def on_test_start(environment, **kwargs): if environment.master: # Master节点生成对应总用户数的资源列表 total_users = environment.runner.target_user_count resource_list = [f"model_{i}" for i in range(total_users)] # 初始化Distributor,自动同步资源列表到所有Worker global items_models_dist items_models_dist = Distributor(resource_list, environment=environment) class TestUser(HttpUser): @task def use_model(self): global items_models_dist model = next(items_models_dist) # 执行业务逻辑,例如调用模型接口 self.client.get(f"/api/model/{model}")
3. 确认启动环境依赖
确保所有进程(Master和Worker)都正确安装了locust-plugins库,启动命令--processes -1会自动按CPU核心数启动Worker,无需额外手动配置Worker节点。
原因说明
你之前手动传递资源列表的方式绕开了locust-plugins的内置同步机制,Worker向Master请求资源时发送的_distributor_request消息未被Master识别(未注册处理程序),导致Master无法响应资源分配请求,最终造成虚拟用户停滞。
内容的提问来源于stack exchange,提问作者PloniStacker

