如何在Locust中采集分析API自定义请求级指标并实现多维度分析?
Locust自定义指标分析与数据导出方案
一、利用Request Context实现自定义指标的切片分析
你提到的context参数确实是上报自定义指标的核心方式,具体实现步骤:
- 在请求回调中解析API响应里的自定义指标(比如缓存命中状态、中间检索值),将这些值存入请求的
context字典 - Locust会自动把
context数据关联到对应请求记录,可通过两种方式实现切片分析:- Locust UI原生过滤(限版本≥2.15):在图表页过滤框输入
context.cache_hit=true,就能只查看缓存命中的请求延迟;也可同时对比cache_hit=true和cache_hit=false两组的延迟差异 - 自定义维度聚合:在Master进程中监听
request_success/request_failure事件,按自定义指标维度聚合统计(比如按缓存命中/未命中分组计算平均延迟、请求数),避免Worker进程做聚合消耗资源
- Locust UI原生过滤(限版本≥2.15):在图表页过滤框输入
示例代码片段:
from locust import HttpUser, task, events class APITestUser(HttpUser): @task def test_api(self): response = self.client.get("/your-api-endpoint") # 解析响应中的自定义指标 cache_hit = response.json().get("cache_hit", False) retrieval_value = response.json().get("retrieval_value", 0) # 将指标存入context self.client.get( "/your-api-endpoint", context={"cache_hit": cache_hit, "retrieval_value": retrieval_value} ) # 在Master进程中做维度聚合(避免Worker消耗资源) @events.request_success.add_listener def on_request_success(request_type, name, response_time, response_length, context, exception, **kwargs): if context: # 按cache_hit维度拆分统计 cache_key = f"cache_hit_{context['cache_hit']}" events.request_success.fire( request_type=request_type, name=f"{name}_{cache_key}", response_time=response_time, response_length=response_length )
二、非延迟指标的分布统计
要查看自定义指标的分布(比如缓存命中占比、检索值区间分布),可通过两种方式实现:
- 实时统计展示:在Master进程中维护计数器和直方图结构,监听
request_success事件更新统计,通过events.report_to_master和events.slave_report实现Worker到Master的数据同步,最后在UI自定义页面或终端打印结果 - 离线分析前置:收集数据时直接记录每个请求的自定义指标值,后续导出后用Pandas、Matplotlib等工具分析分布
三、低影响导出Parquet结构化数据
要导出Parquet且不影响Worker吞吐量,核心是异步+批量写入,避免每个请求同步写磁盘:
- Worker端异步缓存:在Worker进程中用线程安全队列缓存请求记录(包含自定义指标、延迟、时间戳等),设置批量阈值(比如每1000条或每5分钟触发一次写入)
- Master端统一写入:通过Locust的Master-Worker通信机制,Worker定期将批量数据发送到Master,由Master统一写入Parquet文件(Master进程通常不处理请求,资源充足)
- 用PyArrow高效写入:PyArrow的Parquet写入支持批量操作,性能优于单条写入,能减少IO开销
示例代码框架:
import pyarrow as pa import pyarrow.parquet as pq from locust import events from queue import Queue import threading import time # Worker端缓存队列,设置合理容量避免内存溢出 data_queue = Queue(maxsize=10000) batch_size = 1000 def worker_data_writer(): while True: batch = [] # 批量收集队列中的数据 while len(batch) < batch_size and not data_queue.empty(): batch.append(data_queue.get()) if batch: # 转换为PyArrow Table table = pa.Table.from_pylist(batch) # 按小时分区写入,避免单文件过大 pq.write_to_dataset( table, root_path="locust_test_data", partition_cols=["timestamp_hour"], append=True ) time.sleep(1) # 启动Worker端异步写入线程(守护线程随进程退出) threading.Thread(target=worker_data_writer, daemon=True).start() @events.request_success.add_listener def collect_request_data(request_type, name, response_time, response_length, context, exception, **kwargs): # 构造包含自定义指标的完整记录 record = { "timestamp": time.time(), "timestamp_hour": time.strftime("%Y%m%d%H", time.localtime()), "request_type": request_type, "name": name, "response_time": response_time, "response_length": response_length, **(context or {}) } # 队列未满时加入,满时可选择阻塞或丢弃旧数据 if not data_queue.full(): data_queue.put(record)
注意事项:
- 分布式测试时,需通过
events.slave_report将Worker队列的批量数据发送到Master,Master监听events.master_report接收后统一写入,避免多Worker写同一文件冲突 - 根据Worker内存情况调整队列容量和批量大小,平衡内存占用与写入频率
内容的提问来源于stack exchange,提问作者dingus
相关产品推荐
相关产品推荐

