使用Nvidia Triton Server计算吞吐量时遭遇BrokenPipeError错误求助
Nvidia Triton Server计算吞吐量时遭遇BrokenPipeError错误求助
大家好,我现在在测试Nvidia Triton Server的吞吐量,需求是让客户端先发送10000个请求到服务器,等所有请求都发送完成后,服务器再开始顺序处理这些请求,最后用10000/(end_time - start_time)来计算吞吐量。目前代码里暂时设置的是发送5000个请求,但不管数量多少,运行过程中一直碰到BrokenPipeError: [Errno 32] Broken pipe错误。
我注意到代码里的count变量应该加锁,但感觉这个不是导致错误的原因。我用的是nvcr.io/nvidia/tritonserver:24.02-py3镜像,下面是我的代码(已经去掉了导入语句,方便阅读):
Model.py
class TritonPythonModel: def initialize(self, args): """ Initialize MobileNet model and ImageNet labels. """ self.model = ResNet50(weights="imagenet") self.model.predict(np.zeros((1, 224, 224, 3), dtype=np.float32)) # Warm-up labels_path = os.path.join(os.path.dirname(__file__), "imagenet_labels.json") with open(labels_path) as f: class_idx = json.load(f) self.labels = {int(k): v[1] for k, v in class_idx.items()} print("SLEEPING...") time.sleep(1) print("STARTING...") def execute(self, requests): responses = [] for request in requests: try: input_tensor = pb_utils.get_input_tensor_by_name(request, "INPUT_IMAGE") imgs = input_tensor.as_numpy() # Ensure the images are float32 (if not already) if imgs.dtype != np.float32: imgs = imgs.astype(np.float32) preds = self.model.predict(imgs, verbose=0) # Process each image in the batch batch_labels = [] for b in range(imgs.shape[0]): top5 = np.argsort(preds[b])[::-1][:5] labels = [self.labels[i].encode("utf-8") for i in top5] batch_labels.append(labels) output_array = np.array(batch_labels, dtype=object) output_tensor = pb_utils.Tensor("OUTPUT_CLASS", output_array) responses.append(pb_utils.InferenceResponse(output_tensors=[output_tensor])) except Exception as e: print(e) responses.append(pb_utils.InferenceResponse(error=pb_utils.TritonError(str(e)))) return responses
Client.py
BASE_FILE_PATH = os.path.dirname(os.path.abspath('__file__')) CSV_HEADER = ['timestamp', 'latency'] image_pool = preload_images_parallel(pool_size=100, max_workers=100) image_pool_size = len(image_pool) def get_image(index): """Returns the image at index (circular access if needed).""" return image_pool[index % image_pool_size] count = 0 def infer_request(url, image_input, model_name): global count """Perform inference on a Triton server and return latency""" client = httpclient.InferenceServerClient(url=url, network_timeout=2**31 - 1, connection_timeout=2**31 - 1) start_t = time.time() try: response = client.infer(model_name=model_name, inputs=[image_input]) count += 1 end_t = time.time() duration_ms = (end_t - start_t) * 1000 print(f"{duration_ms} ms") if count % 200 == 0: print(f"Processed {count} latencies", flush=True) return duration_ms except InferenceServerException as e: print(f"Inference failed: {e}") return None def send_10k_reqs(shared_data): NUM_REQS = 5000 RPS = 250 success = 0 start_time = time.time() with ThreadPoolExecutor() as executor: futures = [] # Use tqdm to wrap the loop and display progress for i in tqdm(range(NUM_REQS), desc="Sending requests"): image = get_image(i) image_input = httpclient.InferInput("INPUT_IMAGE", [1, 224, 224, 3], datatype="FP32") image_input.set_data_from_numpy(image, binary_data=True) with shared_data['lock']: url = shared_data['url'] model_name = shared_data['model_name'] future = executor.submit(infer_request, url, image_input, model_name) futures.append(future) # Uncomment below if you need to pace requests: sleep_time = np.random.exponential(1/RPS) # time.sleep(sleep_time) for future in futures: result = future.result() end_time = time.time() total_time = end_time - start_time throughput = NUM_REQS / total_time print(f"Total Time: {total_time}, Throughput: {throughput}") return total_time, throughput def main(): parser = argparse.ArgumentParser(description="ML Inference Throughput & Latency Benchmark") parser.add_argument("-d", "--duration", type=int, required=False, help="Workload duration (in seconds)") parser.add_argument("-o", "--output", type=str, required=False, help="CSV file to store results") parser.add_argument("-m", "--mode", type=str, required=True, choices=["test_rps", "throughput"], help="Mode: 'test_rps' for latency benchmarking across RPS, 'throughput' for sending 10K requests") args = parser.parse_args() server = "localhost" model_id = "resnet50_python" shared_data = { 'url': f'{server}:8000', 'model_name': f'{model_id}', 'lock': threading.Lock() } if args.mode == "throughput": total_time, throughput = send_10k_reqs(shared_data) with open('output.csv', 'a', newline='') as csvfile: csvwriter = csv.writer(csvfile) csvwriter.writerow(["Total Time (sec)", "Throughput (req/sec)"]) csvwriter.writerow([total_time, throughput]) print("Throughput benchmarking completed.") if __name__ == "__main__": main()
希望大家能帮我看看这个BrokenPipeError是怎么回事,怎么解决才能完成吞吐量测试?
备注:内容来源于stack exchange,提问作者Tanay Joshi
相关产品推荐
相关产品推荐

