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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.13 19:39:50