Locust无法记录返回流的gRPC接口的成功请求
我完全懂你现在的头疼——用Locust压测gRPC流接口时,失败请求能正常在UI里显示,成功请求却彻底没踪影,代码看起来没报错,服务器日志也说响应正常,偏偏Locust就像没收到成功信号一样,甚至加了等待时间还会触发超时,换Robot Framework又能正常跑,确实让人卡壳。
先帮你捋捋最可能的核心问题,再给你针对性的解决方案:
核心问题:流请求的循环阻塞了任务执行
你在GrpcClient.load_data里用for response in self.stub.LoadData(request):遍历响应流,这个循环会一直等待直到服务端主动关闭流。如果遇到以下两种情况,这个循环就会无限卡住,根本走不到后面上报成功的events.request.fire步骤:
- 你的
LoadData是持续推送的订阅型接口(比如服务端不会主动关闭流,一直推数据); - 服务端虽然是一次性流接口,但返回完响应后没有正确关闭流。
这也就能解释为什么你看不到成功记录——Locust的任务线程一直卡着循环,根本没机会执行上报逻辑。
针对性解决方案
方案1:先验证流是否能正常结束(一次性流接口场景)
如果你的LoadData是「返回指定数量响应后就结束」的一次性接口,先单独写个小脚本验证流会不会正常退出:
# 本地测试脚本,单独运行看看流是否能正常结束 import grpc import abcapi_pb2 import abcapi_pb2_grpc def test_stream_completion(): channel = grpc.insecure_channel("localhost:50051") stub = abcapi_pb2_grpc.abcapiIfStub(channel) # 用和Locust里一样的请求参数 request = abcapi_pb2.C_Request(timestamp=123, series=0, limit=2) print("开始调用流接口...") response_count = 0 try: for resp in stub.LoadData(request): response_count += 1 print(f"收到第{response_count}条响应") print(f"流正常结束!共收到{response_count}条响应") except Exception as e: print(f"调用过程出错:{str(e)}") if __name__ == "__main__": test_stream_completion()
如果这个脚本里循环不会结束,那问题出在服务端——需要让服务端返回完指定limit的响应后,主动关闭流连接。
方案2:给流遍历加退出条件(持续流接口场景)
如果LoadData是「一直推送数据的订阅型接口」,那你不能无限等待,得主动加退出逻辑,比如接收指定数量的响应后跳出,或者给请求加超时:
示例1:接收指定数量响应后主动退出
修改GrpcClient.load_data里的流遍历部分:
try: request = abcapi_pb2.C_Request( timestamp=int(request_data["timestamp"]), series=int(request_data["series"]), limit=int(request_data["limit"]) ) start_time = time.time() response_count = 0 # 达到请求里的limit数量就主动跳出循环 for response in self.stub.LoadData(request): response_count += 1 if response_count >= request_data["limit"]: break elapsed = time.time() - start_time return True, elapsed, response_count, "Success"
示例2:给整个流请求加超时限制
try: request = abcapi_pb2.C_Request( timestamp=int(request_data["timestamp"]), series=int(request_data["series"]), limit=int(request_data["limit"]) ) start_time = time.time() response_count = 0 # 给流请求加5秒超时,超时后自动抛出DEADLINE_EXCEEDED错误 for response in self.stub.LoadData(request, timeout=5): response_count += 1 elapsed = time.time() - start_time return True, elapsed, response_count, "Success" except grpc.RpcError as e: if e.code() == grpc.StatusCode.DEADLINE_EXCEEDED: # 超时前已经收到部分响应,可以根据业务需求算部分成功 elapsed = time.time() - start_time return True, elapsed, response_count, f"超时前收到{response_count}条响应" else: error_message = f"RPC Error: {e.details()}" print(error_message) return False, 0, 0, error_message
方案3:优化Locust的上报逻辑(可选)
你在task里的时间计算有点冗余,可以直接用load_data返回的elapsed转成毫秒,不用重新计算一次:
@task def load_data(self): request_data = { "timestamp": 123, "series": 0, "limit": 2 } name = "LoadData" success, elapsed, response_count, message = self.client.load_data(request_data) # 直接用返回的elapsed转成毫秒,更准确 total_time = int(elapsed * 1000) if success: events.request.fire( request_type="grpc", name=name, response_time=total_time, response_length=response_count, exception=None, context={ "stream_count": response_count, "stream_time": elapsed } ) print(f"Received {response_count} stream responses in {elapsed:.2f}s") else: events.request.fire( request_type="grpc", name=name, response_time=total_time, response_length=0, exception=Exception(message), context={} )
额外排查小技巧
- 启动Locust时加上
--loglevel DEBUG,看看有没有被过滤的日志,说不定能找到为什么没上报的线索; - 确认你的
C_RequestUser类没有被标记为abstract(看你代码里是对的,没问题)。
先从验证流是否能正常结束入手,这个应该就是解决问题的关键。
备注:内容来源于stack exchange,提问作者Sudeep Bidwai

