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

Locust 2.15.1中RabbitMQ自定义客户端无法绑定用户的问题求助

解决Locust自定义Kombu客户端不被识别的问题

针对你遇到的Locust无法识别RabbitMQ消息发送事务、不显示用户运行的问题,核心原因是误用了专为HTTP场景设计的HttpUser类,且自定义客户端未向Locust统计系统上报请求事件。以下是修正后的实现方案:

关键修改点

  • 改用Locust基础User类替代HttpUser,无需HTTP相关的冗余逻辑
  • 在自定义客户端的消息发送方法中,通过events.request.fire上报请求状态,让Locust捕捉并统计事务
  • 在用户生命周期方法中初始化/销毁客户端,实现连接复用

修正后的代码

import os
from locust import User, task, events, run_single_user
from kombu import Connection, Exchange, Queue


class KombuClient:
    def __init__(self):
        self.connection = Connection(
            hostname=os.getenv("LOCUST_AMQP_CONFIG", "amqp://guest:guest@localhost:5672//")
        )
        self.channel = self.connection.channel()
        # 提前初始化Exchange和Queue,避免每次发送重复声明
        self.exchange = Exchange(
            name=os.getenv("LOCUST_AMQP_EXCHANGE", "test_exchange"),
            type="direct",
        )
        self.queue = Queue(
            name=os.getenv("LOCUST_AMQP_QUEUE", "test_queue"),
            exchange=self.exchange,
            routing_key=os.getenv("LOCUST_AMQP_KEY", "test_key"),
        )
        self.queue.maybe_bind(self.connection)
        self.queue.declare()
        self.producer = self.connection.Producer(
            exchange=self.exchange,
            routing_key=os.getenv("LOCUST_AMQP_KEY", "test_key"),
        )

    def send_message(self, message):
        request_meta = {
            "request_type": "AMQP",
            "name": "send_message",
            "response_time": 0,
            "response_length": len(message),
            "exception": None,
            "context": {},
        }
        try:
            self.producer.publish(message)
            print(f"INFO: Sent message: '{message}'")
        except Exception as e:
            request_meta["exception"] = e
        finally:
            # 上报请求事件给Locust,触发统计
            events.request.fire(**request_meta)

    def close_connection(self):
        self.connection.release()


class RabbitMQUser(User):
    min_wait = 5000
    max_wait = 9000

    def on_start(self):
        # 用户启动时创建客户端实例,复用连接
        self.client = KombuClient()

    def on_stop(self):
        # 用户停止时关闭连接
        self.client.close_connection()

    @task
    def send_hello_message(self):
        self.client.send_message("Hello World!")


if __name__ == "__main__":
    run_single_user(RabbitMQUser)

核心说明

  1. 基础User类适配:HttpUser是User的子类,自带HTTP客户端逻辑,自定义非HTTP协议的压测场景时,直接使用User类更简洁高效。
  2. 事务统计触发:Locust的UI和统计系统依赖events.request.fire传递的事件信息,只有上报请求状态,才能识别用户运行情况和事务数据。
  3. 连接复用优化:在on_start中初始化客户端,每个用户复用一个RabbitMQ连接,避免频繁创建销毁连接带来的性能损耗。
  4. 简化代码结构:Locust 2.x版本推荐直接在User类上使用@task装饰器,替代旧版的TaskSet,代码更直观。

修改后,Locust就能正确显示用户运行状态和RabbitMQ消息发送的事务统计数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 20:20:23