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

如何使用Locust对基于PubSub的API执行负载测试?

实现Locust测试PubSub架构API的方案

核心思路

每个Locust任务生成唯一关联ID,发布消息到入口Topic时携带该ID;随后通过订阅输出Topic,匹配对应ID的响应消息,记录从发布到收到响应的端到端耗时,以此完成异步流程的负载测试。

具体实现步骤

1. 初始化PubSub客户端

在Locust用户类中初始化发布者(入口Topic)和订阅者(输出Topic的订阅),同时维护待匹配请求的状态:

from locust import HttpUser, task, between
from google.cloud import pubsub_v1
import uuid
import time
import json

class PubSubLoadTestUser(HttpUser):
    wait_time = between(1, 3)

    def on_start(self):
        # 初始化入口Topic发布者
        self.publisher = pubsub_v1.PublisherClient()
        self.entry_topic_path = self.publisher.topic_path("你的项目ID", "入口Topic名称")
        
        # 初始化输出Topic订阅者
        self.subscriber = pubsub_v1.SubscriberClient()
        self.output_sub_path = self.subscriber.subscription_path("你的项目ID", "输出Topic订阅名称")
        
        # 存储待匹配的请求ID与发布时间戳
        self.pending_requests = {}
        # 标记订阅监听是否已启动
        self.subscription_running = False

    def on_stop(self):
        self.publisher.close()
        if hasattr(self, "subscription_future"):
            self.subscription_future.cancel()
        self.subscriber.close()

2. 异步监听匹配响应(推荐)

通过订阅者的回调函数自动匹配关联ID,一旦找到对应消息就记录耗时并完成请求:

def _match_response_message(self, message):
    try:
        msg_data = json.loads(message.data.decode("utf-8"))
        request_id = msg_data.get("correlation_id")
        
        if request_id in self.pending_requests:
            # 计算端到端耗时并上报Locust统计
            elapsed_ms = (time.time() - self.pending_requests[request_id]) * 1000
            self.environment.events.request.fire(
                request_type="pubsub",
                name="entry_topic_to_output_topic",
                response_time=elapsed_ms,
                response_length=len(message.data),
                exception=None,
            )
            del self.pending_requests[request_id]
            message.ack()  # 确认消息已处理
        else:
            message.nack()  # 不匹配的消息重新入队
    except Exception as e:
        message.nack()
        self.environment.events.request.fire(
            request_type="pubsub",
            name="message_parse_error",
            response_time=0,
            response_length=0,
            exception=e,
        )

@task
def end_to_end_flow(self):
    # 生成唯一关联ID
    correlation_id = str(uuid.uuid4())
    # 构造带关联ID的负载消息
    payload = json.dumps({
        "correlation_id": correlation_id,
        # 其他业务负载字段...
    }).encode("utf-8")
    
    # 发布消息到入口Topic
    publish_future = self.publisher.publish(self.entry_topic_path, payload)
    publish_future.result()  # 确保消息发布成功
    
    # 记录发布时间,加入待处理队列
    self.pending_requests[correlation_id] = time.time()
    
    # 启动订阅监听(仅第一次执行任务时启动)
    if not self.subscription_running:
        self.subscription_future = self.subscriber.subscribe(
            self.output_sub_path, callback=self._match_response_message
        )
        self.subscription_running = True
    
    # 设置超时等待,避免无限阻塞
    timeout = 30  # 匹配API 20秒计算时间,预留冗余
    start_wait = time.time()
    while correlation_id in self.pending_requests:
        if time.time() - start_wait > timeout:
            del self.pending_requests[correlation_id]
            self.environment.events.request.fire(
                request_type="pubsub",
                name="response_timeout",
                response_time=timeout*1000,
                response_length=0,
                exception=TimeoutError("未在超时时间内收到对应响应"),
            )
            break
        time.sleep(0.5)  # 短间隔轮询状态

3. 备选方案:轮询拉取消息

如果不想用异步回调,可定期拉取输出Topic的消息并匹配关联ID:

@task
def end_to_end_flow_polling(self):
    correlation_id = str(uuid.uuid4())
    payload = json.dumps({
        "correlation_id": correlation_id,
        # 其他业务负载字段...
    }).encode("utf-8")
    
    # 发布消息
    publish_future = self.publisher.publish(self.entry_topic_path, payload)
    publish_future.result()
    publish_time = time.time()
    
    timeout = 30
    success = False
    
    # 轮询拉取消息
    while time.time() - publish_time < timeout:
        pull_response = self.subscriber.pull(
            request={"subscription": self.output_sub_path, "max_messages": 10}
        )
        ack_ids = []
        
        for received_msg in pull_response.received_messages:
            msg_data = json.loads(received_msg.message.data.decode("utf-8"))
            if msg_data.get("correlation_id") == correlation_id:
                # 记录耗时
                elapsed_ms = (time.time() - publish_time)*1000
                self.environment.events.request.fire(
                    request_type="pubsub",
                    name="entry_topic_to_output_topic",
                    response_time=elapsed_ms,
                    response_length=len(received_msg.message.data),
                    exception=None,
                )
                ack_ids.append(received_msg.ack_id)
                success = True
                break
            else:
                ack_ids.append(received_msg.ack_id)
        
        # 确认已处理的消息
        if ack_ids:
            self.subscriber.acknowledge(
                request={"subscription": self.output_sub_path, "ack_ids": ack_ids}
            )
        
        if success:
            break
        time.sleep(1)
    
    if not success:
        self.environment.events.request.fire(
            request_type="pubsub",
            name="response_timeout",
            response_time=timeout*1000,
            response_length=0,
            exception=TimeoutError("超时未收到对应响应"),
        )

关键注意事项

  • 订阅隔离:如果是多用户并发测试,可考虑为每个Locust用户创建独立的输出Topic订阅,避免消息被其他用户误处理;或确保关联ID全局唯一,在回调中严格过滤。
  • 资源清理:测试结束时务必关闭PubSub客户端,避免连接泄漏。
  • 统计准确性:通过environment.events.request.fire手动触发Locust统计事件,这样能在仪表盘上看到端到端的请求耗时、成功率等核心指标。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 04:16:39