如何使用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
相关产品推荐
相关产品推荐

