From 75db787a43d1c862c59fade16058ebb1c67186c7 Mon Sep 17 00:00:00 2001 From: kh3654po Date: Mon, 15 May 2023 01:23:11 +0900 Subject: [PATCH 01/13] =?UTF-8?q?=ED=8F=AC=EC=95=84=EC=86=A1=20=EB=B6=84?= =?UTF-8?q?=ED=8F=AC=EB=A5=BC=20=EB=94=B0=EB=A5=B4=EB=8F=84=EB=A1=9D=20?= =?UTF-8?q?=EC=B6=94=EB=A1=A0=20=EC=9A=94=EC=B2=AD=20=EC=8B=9C=EB=82=98?= =?UTF-8?q?=EB=A6=AC=EC=98=A4=EB=A5=BC=20=EC=9E=91=EC=84=B1=ED=95=98?= =?UTF-8?q?=EB=8A=94=20=EB=AA=A8=EB=93=88?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bench/inference_request_workload_manager.py | 47 +++++++++++++++++++++ 1 file changed, 47 insertions(+) create mode 100644 bench/inference_request_workload_manager.py diff --git a/bench/inference_request_workload_manager.py b/bench/inference_request_workload_manager.py new file mode 100644 index 0000000..a11776a --- /dev/null +++ b/bench/inference_request_workload_manager.py @@ -0,0 +1,47 @@ +from numpy import random +import pickle as pk + +# (modle name, requests per second) +inference_request_info = [ + ('mobilenet_v1', 10), + ('mobilenet_v2', 2), + ('inception_v3', 2), + ('yolo_v5', 1) +] + +file_name = 'inference_request_workload.pickle' + + +def create_inference_request_workload(req_time): + requests = [[] for _ in range(req_time)] + for (model_name, req_per_sec) in inference_request_info: + workloads = random.poisson(lam=req_per_sec, size=req_time) + for idx in range(req_time): + requests[idx].extend([model_name for _ in range(workloads[idx])]) + + for idx in range(req_time): + random.shuffle(requests[idx]) + + return requests + + +def save_workload_to_file(file_name, workloads): + with open(file_name, 'wb') as f: + pk.dump(workloads, f) + + +def load_workload_from_file(file_name): + with open(file_name, 'rb') as f: + loaded_workloads = pk.load(f) + return loaded_workloads + + +# workloads = create_inference_request_workload(100) + +# save_workload_to_file(file_name, workloads) +# loaded_workloads = load_workload(file_name) + + +# for i in loaded_workloads: +# print(i) +# print() \ No newline at end of file From 8aba62f86669cb28d2209d89a5d79d4561b46d2b Mon Sep 17 00:00:00 2001 From: kh3654po Date: Mon, 15 May 2023 01:24:36 +0900 Subject: [PATCH 02/13] =?UTF-8?q?=EC=B6=94=EB=A1=A0=EC=9A=94=EC=B2=AD?= =?UTF-8?q?=EB=93=A4=EC=9D=84=20=EB=9D=BC=EC=9A=B4=EB=93=9C=EB=A1=9C?= =?UTF-8?q?=EB=B9=88=20=EB=B0=A9=EC=8B=9D=EC=9C=BC=EB=A1=9C=20=EC=97=A3?= =?UTF-8?q?=EC=A7=80=EC=9E=A5=EB=B9=84=EC=97=90=EA=B2=8C=20=EC=A0=84?= =?UTF-8?q?=EC=86=A1=ED=95=98=EB=8A=94=20=EC=8A=A4=EC=BC=80=EC=A4=84?= =?UTF-8?q?=EB=9F=AC=20=EA=B5=AC=ED=98=84?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bench/round_robin_scheduler.py | 183 +++++++++++++++++++++++++++++++++ 1 file changed, 183 insertions(+) create mode 100644 bench/round_robin_scheduler.py diff --git a/bench/round_robin_scheduler.py b/bench/round_robin_scheduler.py new file mode 100644 index 0000000..23f6673 --- /dev/null +++ b/bench/round_robin_scheduler.py @@ -0,0 +1,183 @@ +# from module import put_data_into_sheet +# put_data_into_sheet.put_data(variables.rest_spreadsheet_id, result, variables.num_tasks) + +import argparse +import roundrobin +import time +from threading import Thread +import importlib +import inference_request_workload_manager + + +parser = argparse.ArgumentParser() +parser.add_argument('--edge', default=None, type=str) + +args = parser.parse_args() +edges_to_inference = args.edge + +# grpc or rest +request_type = 'rest' +# workload file name +workload_file_name = 'inference_request_workload.pickle' +# 서버의 포트 정보 +grpc_port = 8500 +rest_port = 8501 + +# 각 장비의 ip와 로드된 모델들을 설정해주어야함. +edges_info = {'nvidia-xavier2': {'ip_addr': 'nvidia-xavier2', + 'models': ['mobilenet_v1', 'mobilenet_v2', 'inception_v3', 'yolo_v5'] + }, + 'nvidia-tx2': {'ip_addr': 'nvidia-tx2', + 'models': ['mobilenet_v1', 'mobilenet_v2', 'inception_v3', 'yolo_v5'] + }, + 'nvidia-nano1': {'ip_addr': 'nvidia-nano1', + 'models': ['mobilenet_v1'] + } + } + + +# --edge 옵션이 없을 시 등록되어 있는 모든 장비들에 추론 요청, 요청장비들은 edges_info에 등록되어 있어야함. 입력 형식은 'a, b, ...' +registered_edges = list(edges_info.keys()) + +if request_type not in ['grpc', 'rest']: + print(f'request_type must in [grpc, rest] / current type : {request_type}') + exit(1) + +if edges_to_inference is None: + edges_to_inference = registered_edges +else: + edges_to_inference = edges_to_inference.split(',') + +for edge in edges_to_inference: + if edge not in registered_edges: + print(f'--edge arg must be in {registered_edges}') + exit(1) + +print(f'Edges to inference: {edges_to_inference}') + + +# 추론 요청 할 장비들에서 요청 가능한 모델들 +def get_enable_models_to_inference(): + models_to_inference = [] + + for edge_name in edges_to_inference: + edge_info = edges_info.get(edge_name) + models = edge_info.get('models') + models_to_inference.extend(models) + + return set(models_to_inference) + +models_to_inference = get_enable_models_to_inference() +print(f'Models to inference: {models_to_inference}') + + +# 추론요청 할 각 모델의 모듈을 저장. +def regist_enable_model_modules(models_to_infer): + model_modules = {} + + for model in models_to_infer: + module = None + if request_type == 'grpc': + module = importlib.import_module(f"{model}.grpc_bench") + else: + module = importlib.import_module(f"{model}.rest_bench") + model_modules.update({model: module}) + + return model_modules + +model_modules = regist_enable_model_modules(models_to_inference) +print(f'Model modules: {model_modules}') + + +# 딕셔너리에 모델별로 엣지장비이름 등록 -> 들어오는 요청에 따라 어느 장비에 보낼 차례인지 확인 할 수 있는 딕셔너리 생성 +def regist_edges_to_model(): + model_edge_info = {} + + for edge in edges_to_inference: + edge_info = edges_info.get(edge) + for model in edge_info.get('models'): + if model not in model_edge_info.keys(): + model_edge_info[model] = [] + + model_edge_info[model].append((edge, 1)) + + for model in model_edge_info.keys(): + dataset = model_edge_info.get(model) + model_edge_info[model] = roundrobin.smooth(dataset) + + return model_edge_info + +model_edge_info = regist_edges_to_model() +print(f'model-edge dataset: {model_edge_info}') + + +# 모델로 엣지장비 이름 얻는 함수, 모델을 키값으로 하여 해당 모델이 로드된 엣지장비들을 라운드로빈으로 함수를 호출할 때마다 하나씩 얻어옴 +def get_edge_by_model_rr(model): + if model in model_edge_info.keys(): + return model_edge_info.get(model)() + else: + return None + + +# 추론을 요청하는 함수, 인자로는 추론을 요청할 엣지 장비, 모델. 엣지장비와 모델은 위의 edges_info에 등록되어 있어야함 +def model_request(edge, model, idx): + if edge not in edges_to_inference: + print(f'edge must be in {edges_to_inference}/ input value: {edge}') + return + + if model not in models_to_inference: + print(f'model must be in {models_to_inference}/ input value: {model}') + return + + edge_info = edges_info.get(edge) + edge_ip_addr = edge_info.get('ip_addr') + + if request_type == 'grpc': + edge_ip_addr = f'{edge_ip_addr}:{grpc_port}' + result = model_modules[model].run_bench(1, edge_ip_addr, 0) + else: + edge_ip_addr = f'http://{edge_ip_addr}:{rest_port}/' + result = model_modules[model].run_bench(1, edge_ip_addr) + + inference_times[idx].extend(result) + + +# 들어오는 요청들 +requests_list = inference_request_workload_manager.load_workload_from_file(workload_file_name) + + +# 요청을 각 장비에 전달, 여러요청을 동시에 다룰 수 있도록 쓰레드 이용 +threads = [] +inference_times = [[] for _ in range(len(requests_list))] +idx = -1 + +start_inference_time = time.time() + +for cur_reqs in requests_list: + idx += 1 + request_sleep_time = 1 / len(cur_reqs) # 요청들을 1초에 나눠서 보내기 위한 슬립시간 + for req in cur_reqs: + edge_to_infer = get_edge_by_model_rr(req) + if edge_to_infer is None: + print(f'{req} can\'t be inference') + continue + + th = Thread(target=model_request, args=(edge_to_infer, req, idx)) + th.start() + threads.append(th) + time.sleep(request_sleep_time) + +for th in threads: + th.join() + +end_inference_time = time.time() +total_inference_time = end_inference_time - start_inference_time + +print() +print('total inference time', total_inference_time) +print('----------------------') +print('inference time info (each argument is info about requests per sec)') +print() +for time_info in inference_times: + print(time_info) + print() From 818b82f86340918f89ee8aba8049e9208328587d Mon Sep 17 00:00:00 2001 From: kh3654po Date: Thu, 1 Jun 2023 15:23:12 +0900 Subject: [PATCH 03/13] =?UTF-8?q?=EC=B6=94=EB=A1=A0=EC=9A=94=EC=B2=AD=20?= =?UTF-8?q?=EC=9B=8C=ED=81=AC=EB=A1=9C=EB=93=9C=20=EC=83=9D=EC=84=B1=20?= =?UTF-8?q?=EC=BD=94=EB=93=9C=20=EC=88=98=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bench/inference_request_workload_manager.py | 31 +++++++++++---------- 1 file changed, 16 insertions(+), 15 deletions(-) diff --git a/bench/inference_request_workload_manager.py b/bench/inference_request_workload_manager.py index a11776a..479ba9f 100644 --- a/bench/inference_request_workload_manager.py +++ b/bench/inference_request_workload_manager.py @@ -3,45 +3,46 @@ # (modle name, requests per second) inference_request_info = [ - ('mobilenet_v1', 10), + ('mobilenet_v1', 20), ('mobilenet_v2', 2), ('inception_v3', 2), ('yolo_v5', 1) ] -file_name = 'inference_request_workload.pickle' - +file_name = 'workload4' +req_time_num = 20 def create_inference_request_workload(req_time): requests = [[] for _ in range(req_time)] + total_req_num = 0 for (model_name, req_per_sec) in inference_request_info: workloads = random.poisson(lam=req_per_sec, size=req_time) for idx in range(req_time): requests[idx].extend([model_name for _ in range(workloads[idx])]) + total_req_num += sum(workloads) for idx in range(req_time): random.shuffle(requests[idx]) - return requests + workload_info = {} + workload_info['total_request_num'] = total_req_num + workload_info['requests'] = requests + + return workload_info -def save_workload_to_file(file_name, workloads): +def save_workload_info_to_file(file_name, workloads): with open(file_name, 'wb') as f: pk.dump(workloads, f) -def load_workload_from_file(file_name): +def load_workload_info_from_file(file_name): with open(file_name, 'rb') as f: loaded_workloads = pk.load(f) return loaded_workloads -# workloads = create_inference_request_workload(100) - -# save_workload_to_file(file_name, workloads) -# loaded_workloads = load_workload(file_name) - - -# for i in loaded_workloads: -# print(i) -# print() \ No newline at end of file +# workloads = create_inference_request_workload(req_time_num) +# save_workload_info_to_file(file_name, workloads) +# loaded_workload_info = load_workload_info_from_file(file_name) +# print(loaded_workload_info) From 11b3a8c70942e02c5bfe3c49201acfcf200e208b Mon Sep 17 00:00:00 2001 From: kh3654po Date: Thu, 1 Jun 2023 15:24:01 +0900 Subject: [PATCH 04/13] =?UTF-8?q?=EA=B0=81=20=EB=AA=A8=EB=8D=B8=EB=B3=84?= =?UTF-8?q?=20=EC=B6=94=EB=A1=A0=EC=9A=94=EC=B2=AD=EC=8B=9C=20=EB=B3=B4?= =?UTF-8?q?=EB=82=BC=20=EB=8D=B0=EC=9D=B4=ED=84=B0=EB=A5=BC=20=EB=AF=B8?= =?UTF-8?q?=EB=A6=AC=20=EC=A0=84=EC=B2=98=EB=A6=AC=20=ED=9B=84=20=EC=A0=80?= =?UTF-8?q?=EC=9E=A5=ED=95=98=EB=8A=94=20=EB=AA=A8=EB=93=88?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bench/preprocessing_data_manager.py | 27 +++++++++++++++++++++++++++ 1 file changed, 27 insertions(+) create mode 100644 bench/preprocessing_data_manager.py diff --git a/bench/preprocessing_data_manager.py b/bench/preprocessing_data_manager.py new file mode 100644 index 0000000..795c0c8 --- /dev/null +++ b/bench/preprocessing_data_manager.py @@ -0,0 +1,27 @@ +import tensorflow as tf +import json +import importlib + + +data_source_info = {'mobilenet_v1': '../dataset/imagenet/imagenet_1000_raw/n01843383_1.JPEG', + 'mobilenet_v2': '../dataset/imagenet/imagenet_1000_raw/n01843383_1.JPEG', + 'inception_v3': '../dataset/imagenet/imagenet_1000_raw/n01843383_1.JPEG', + 'yolo_v5': '../dataset/coco_2017/coco/images/val2017/000000089761.jpg', + } + +def regist_preprocessed_datas(request_type): + preprocessed_datas = {} + + for model in data_source_info.keys(): + if request_type == 'rest': + preprocessing_module = importlib.import_module(f"{model}.preprocessing") + data_source = data_source_info.get(model) + data = json.dumps({"instances": preprocessing_module.run_preprocessing(data_source).tolist()}) + preprocessed_datas.update({model: data}) + elif request_type == 'grpc': + preprocessing_module = importlib.import_module(f"{model}.preprocessing") + data_source = data_source_info.get(model) + data = tf.make_tensor_proto(preprocessing_module.run_preprocessing(data_source)) + preprocessed_datas.update({model: data}) + + return preprocessed_datas \ No newline at end of file From e2cd4c45bd22e6b23f7d827ab883b4a3d4f7cd54 Mon Sep 17 00:00:00 2001 From: kh3654po Date: Thu, 1 Jun 2023 15:27:31 +0900 Subject: [PATCH 05/13] =?UTF-8?q?=EC=B6=94=EB=A1=A0=EC=9A=94=EC=B2=AD?= =?UTF-8?q?=EC=8B=9C=20=EC=A0=84=EC=B2=98=EB=A6=AC=EB=90=9C=20=EB=8D=B0?= =?UTF-8?q?=EC=9D=B4=ED=84=B0=EB=A5=BC=20=EC=A0=84=EC=86=A1=ED=95=98?= =?UTF-8?q?=EB=8A=94=EA=B2=83=EC=9C=BC=EB=A1=9C=20=EB=B3=80=EA=B2=BD?= =?UTF-8?q?=ED=96=88=EC=9C=BC=EB=AF=80=EB=A1=9C=20=EC=A0=84=EC=B2=98?= =?UTF-8?q?=EB=A6=AC=20=EC=BD=94=EB=93=9C=20=EC=82=AD=EC=A0=9C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bench/inception_v3/grpc_bench.py | 10 +--------- bench/inception_v3/rest_bench.py | 10 +--------- bench/mobilenet_v1/grpc_bench.py | 10 +--------- bench/mobilenet_v1/rest_bench.py | 10 +--------- bench/mobilenet_v2/grpc_bench.py | 10 +--------- bench/mobilenet_v2/rest_bench.py | 10 +--------- bench/yolo_v5/grpc_bench.py | 10 +--------- bench/yolo_v5/rest_bench.py | 10 +--------- 8 files changed, 8 insertions(+), 72 deletions(-) diff --git a/bench/inception_v3/grpc_bench.py b/bench/inception_v3/grpc_bench.py index d762422..5d73836 100644 --- a/bench/inception_v3/grpc_bench.py +++ b/bench/inception_v3/grpc_bench.py @@ -5,21 +5,13 @@ #tf log setting import os os.environ['TF_CPP_MIN_LOG_LEVEL'] = '3' -import tensorflow as tf -import numpy as np - -#preprocessing library -from inception_v3 import preprocessing #병렬처리 library import concurrent.futures -def run_bench(num_tasks, server_address, use_https): +def run_bench(num_tasks, server_address, use_https, data): model_name = "inception_v3" - image_file_path = "../../dataset/imagenet/imagenet_1000_raw/n01843383_1.JPEG" - data = tf.make_tensor_proto(preprocessing.run_preprocessing(image_file_path)) - stub = module_grpc.create_grpc_stub(server_address, use_https) # gRPC 요청 생성 diff --git a/bench/inception_v3/rest_bench.py b/bench/inception_v3/rest_bench.py index 672df73..03919c3 100644 --- a/bench/inception_v3/rest_bench.py +++ b/bench/inception_v3/rest_bench.py @@ -1,19 +1,11 @@ -#preprocessing library -from inception_v3 import preprocessing -import numpy as np - #REST 요청 관련 library from module import module_rest -import json #병렬처리 library import concurrent.futures -def run_bench(num_tasks, server_address): +def run_bench(num_tasks, server_address, data): model_name = "inception_v3" - image_file_path = "../../../dataset/imagenet/imagenet_1000_raw/n01843383_1.JPEG" - - data = json.dumps({"instances": preprocessing.run_preprocessing(image_file_path).tolist()}) # REST 요청 병렬 처리 with concurrent.futures.ThreadPoolExecutor(max_workers=num_tasks) as executor: diff --git a/bench/mobilenet_v1/grpc_bench.py b/bench/mobilenet_v1/grpc_bench.py index 9a11370..9c4efd3 100644 --- a/bench/mobilenet_v1/grpc_bench.py +++ b/bench/mobilenet_v1/grpc_bench.py @@ -5,21 +5,13 @@ #tf log setting import os os.environ['TF_CPP_MIN_LOG_LEVEL'] = '3' -import tensorflow as tf -import numpy as np - -#preprocessing library -from mobilenet_v1 import preprocessing #병렬처리 library import concurrent.futures -def run_bench(num_tasks, server_address, use_https): +def run_bench(num_tasks, server_address, use_https, data): model_name = "mobilenet_v1" - image_file_path = "../../dataset/imagenet/imagenet_1000_raw/n01843383_1.JPEG" - data = tf.make_tensor_proto(preprocessing.run_preprocessing(image_file_path)) - stub = module_grpc.create_grpc_stub(server_address, use_https) # gRPC 요청 생성 diff --git a/bench/mobilenet_v1/rest_bench.py b/bench/mobilenet_v1/rest_bench.py index b890f30..2b28993 100644 --- a/bench/mobilenet_v1/rest_bench.py +++ b/bench/mobilenet_v1/rest_bench.py @@ -1,19 +1,11 @@ -#preprocessing library -from mobilenet_v1 import preprocessing -import numpy as np - #REST 요청 관련 library from module import module_rest -import json #병렬처리 library import concurrent.futures -def run_bench(num_tasks, server_address): +def run_bench(num_tasks, server_address, data): model_name = "mobilenet_v1" - image_file_path = "../../../dataset/imagenet/imagenet_1000_raw/n01843383_1.JPEG" - - data = json.dumps({"instances": preprocessing.run_preprocessing(image_file_path).tolist()}) # REST 요청 병렬 처리 with concurrent.futures.ThreadPoolExecutor(max_workers=num_tasks) as executor: diff --git a/bench/mobilenet_v2/grpc_bench.py b/bench/mobilenet_v2/grpc_bench.py index 6b64325..306ab8e 100644 --- a/bench/mobilenet_v2/grpc_bench.py +++ b/bench/mobilenet_v2/grpc_bench.py @@ -5,21 +5,13 @@ #tf log setting import os os.environ['TF_CPP_MIN_LOG_LEVEL'] = '3' -import tensorflow as tf -import numpy as np - -#preprocessing library -from mobilenet_v2 import preprocessing #병렬처리 library import concurrent.futures -def run_bench(num_tasks, server_address, use_https): +def run_bench(num_tasks, server_address, use_https, data): model_name = "mobilenet_v2" - image_file_path = "../../dataset/imagenet/imagenet_1000_raw/n01843383_1.JPEG" - data = tf.make_tensor_proto(preprocessing.run_preprocessing(image_file_path)) - stub = module_grpc.create_grpc_stub(server_address, use_https) # gRPC 요청 생성 diff --git a/bench/mobilenet_v2/rest_bench.py b/bench/mobilenet_v2/rest_bench.py index 33facbc..9bcee65 100644 --- a/bench/mobilenet_v2/rest_bench.py +++ b/bench/mobilenet_v2/rest_bench.py @@ -1,20 +1,12 @@ -#preprocessing library -from mobilenet_v2 import preprocessing -import numpy as np - #REST 요청 관련 library from module import module_rest -import json #병렬처리 library import concurrent.futures -def run_bench(num_tasks, server_address): +def run_bench(num_tasks, server_address, data): model_name = "mobilenet_v2" - image_file_path = "../../../dataset/imagenet/imagenet_1000_raw/n01843383_1.JPEG" - data = json.dumps({"instances": preprocessing.run_preprocessing(image_file_path).tolist()}) - # REST 요청 병렬 처리 with concurrent.futures.ThreadPoolExecutor(max_workers=num_tasks) as executor: futures = [executor.submit(lambda: module_rest.predict(server_address, model_name, data)) for _ in range(num_tasks)] diff --git a/bench/yolo_v5/grpc_bench.py b/bench/yolo_v5/grpc_bench.py index 4f47e85..edad901 100644 --- a/bench/yolo_v5/grpc_bench.py +++ b/bench/yolo_v5/grpc_bench.py @@ -5,21 +5,13 @@ #tf log setting import os os.environ['TF_CPP_MIN_LOG_LEVEL'] = '3' -import tensorflow as tf -import numpy as np - -#preprocessing library -from yolo_v5 import preprocessing #병렬처리 library import concurrent.futures -def run_bench(num_tasks, server_address, use_https): +def run_bench(num_tasks, server_address, use_https, data): model_name = "yolo_v5" - image_file_path = "../../dataset/imagenet/imagenet_1000_raw/n01843383_1.JPEG" - data = tf.make_tensor_proto(preprocessing.run_preprocessing(image_file_path)) - stub = module_grpc.create_grpc_stub(server_address, use_https) # gRPC 요청 생성 diff --git a/bench/yolo_v5/rest_bench.py b/bench/yolo_v5/rest_bench.py index 5110a79..91fe11d 100644 --- a/bench/yolo_v5/rest_bench.py +++ b/bench/yolo_v5/rest_bench.py @@ -1,19 +1,11 @@ -#preprocessing library -from mobilenet_v1 import preprocessing -import numpy as np - #REST 요청 관련 library from module import module_rest -import json #병렬처리 library import concurrent.futures -def run_bench(num_tasks, server_address): +def run_bench(num_tasks, server_address, data): model_name = "yolo_v5" - image_file_path = "../../../dataset/imagenet/imagenet_1000_raw/n01843383_1.JPEG" - - data = json.dumps({"instances": preprocessing.run_preprocessing(image_file_path).tolist()}) # REST 요청 병렬 처리 with concurrent.futures.ThreadPoolExecutor(max_workers=num_tasks) as executor: From 9147915db4e8ccafc53119566244c5021886c5b0 Mon Sep 17 00:00:00 2001 From: kh3654po Date: Thu, 1 Jun 2023 15:28:10 +0900 Subject: [PATCH 06/13] =?UTF-8?q?=EC=A0=84=EC=B2=98=EB=A6=AC=20=EB=AA=A8?= =?UTF-8?q?=EB=93=88=EC=97=90=EC=84=9C=20=EB=B6=88=ED=95=84=EC=9A=94?= =?UTF-8?q?=ED=95=98=EA=B2=8C=20import=EB=90=9C=20=EB=AA=A8=EB=93=88=20?= =?UTF-8?q?=EC=82=AD=EC=A0=9C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bench/inception_v3/preprocessing.py | 1 - bench/mobilenet_v1/preprocessing.py | 1 - bench/mobilenet_v2/preprocessing.py | 1 - bench/yolo_v5/preprocessing.py | 2 -- 4 files changed, 5 deletions(-) diff --git a/bench/inception_v3/preprocessing.py b/bench/inception_v3/preprocessing.py index f0eabd8..a8be0d5 100644 --- a/bench/inception_v3/preprocessing.py +++ b/bench/inception_v3/preprocessing.py @@ -1,5 +1,4 @@ #image 전처리 library -import tensorflow as tf import numpy as np from PIL import Image import os diff --git a/bench/mobilenet_v1/preprocessing.py b/bench/mobilenet_v1/preprocessing.py index 49164c0..427f7c0 100644 --- a/bench/mobilenet_v1/preprocessing.py +++ b/bench/mobilenet_v1/preprocessing.py @@ -1,5 +1,4 @@ #image 전처리 library -import tensorflow as tf import numpy as np from PIL import Image import os diff --git a/bench/mobilenet_v2/preprocessing.py b/bench/mobilenet_v2/preprocessing.py index 49164c0..427f7c0 100644 --- a/bench/mobilenet_v2/preprocessing.py +++ b/bench/mobilenet_v2/preprocessing.py @@ -1,5 +1,4 @@ #image 전처리 library -import tensorflow as tf import numpy as np from PIL import Image import os diff --git a/bench/yolo_v5/preprocessing.py b/bench/yolo_v5/preprocessing.py index d555ad1..45b3ed5 100644 --- a/bench/yolo_v5/preprocessing.py +++ b/bench/yolo_v5/preprocessing.py @@ -1,7 +1,5 @@ #image 전처리 library -import tensorflow as tf import numpy as np -from PIL import Image import os import cv2 From aba59829a4a583393b30f7dd69a6d398d4cab357 Mon Sep 17 00:00:00 2001 From: kh3654po Date: Thu, 1 Jun 2023 15:30:33 +0900 Subject: [PATCH 07/13] =?UTF-8?q?=EC=9A=94=EC=B2=AD=EC=9D=84=20=EC=B2=98?= =?UTF-8?q?=EB=A6=AC=20=ED=9B=84=20=EA=B0=81=20=EC=9E=A5=EB=B9=84=EC=9D=98?= =?UTF-8?q?=20idle=20=EC=8B=9C=EA=B0=84=EA=B3=BC=20=EC=B6=94=EB=A1=A0?= =?UTF-8?q?=EC=8B=9C=EA=B0=84=EC=9D=84=20=EC=B6=9C=EB=A0=A5=ED=95=98?= =?UTF-8?q?=EB=8F=84=EB=A1=9D=20=EB=B3=80=EA=B2=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bench/round_robin_scheduler.py | 208 +++++++++++++++++++++++++++++---- 1 file changed, 184 insertions(+), 24 deletions(-) diff --git a/bench/round_robin_scheduler.py b/bench/round_robin_scheduler.py index 23f6673..baab066 100644 --- a/bench/round_robin_scheduler.py +++ b/bench/round_robin_scheduler.py @@ -5,20 +5,23 @@ import roundrobin import time from threading import Thread +import concurrent.futures import importlib import inference_request_workload_manager +import preprocessing_data_manager parser = argparse.ArgumentParser() parser.add_argument('--edge', default=None, type=str) +parser.add_argument('--workload', type=str, required=True) -args = parser.parse_args() -edges_to_inference = args.edge +scheduler_args = parser.parse_args() +edges_to_inference = scheduler_args.edge # grpc or rest request_type = 'rest' # workload file name -workload_file_name = 'inference_request_workload.pickle' +workload_file_name = scheduler_args.workload # 서버의 포트 정보 grpc_port = 8500 rest_port = 8501 @@ -26,13 +29,13 @@ # 각 장비의 ip와 로드된 모델들을 설정해주어야함. edges_info = {'nvidia-xavier2': {'ip_addr': 'nvidia-xavier2', 'models': ['mobilenet_v1', 'mobilenet_v2', 'inception_v3', 'yolo_v5'] - }, + }, 'nvidia-tx2': {'ip_addr': 'nvidia-tx2', 'models': ['mobilenet_v1', 'mobilenet_v2', 'inception_v3', 'yolo_v5'] - }, + }, 'nvidia-nano1': {'ip_addr': 'nvidia-nano1', - 'models': ['mobilenet_v1'] - } + 'models': ['mobilenet_v1', 'mobilenet_v2', 'inception_v3'] + } } @@ -89,6 +92,11 @@ def regist_enable_model_modules(models_to_infer): print(f'Model modules: {model_modules}') +# 각 모델의 preprocessing 데이터 저장 +preprocessed_datas = preprocessing_data_manager.regist_preprocessed_datas(request_type) +print(f'Preprocessed datas: {preprocessed_datas}') + + # 딕셔너리에 모델별로 엣지장비이름 등록 -> 들어오는 요청에 따라 어느 장비에 보낼 차례인지 확인 할 수 있는 딕셔너리 생성 def regist_edges_to_model(): model_edge_info = {} @@ -120,7 +128,7 @@ def get_edge_by_model_rr(model): # 추론을 요청하는 함수, 인자로는 추론을 요청할 엣지 장비, 모델. 엣지장비와 모델은 위의 edges_info에 등록되어 있어야함 -def model_request(edge, model, idx): +def model_request(edge, model, idx, thread_id): if edge not in edges_to_inference: print(f'edge must be in {edges_to_inference}/ input value: {edge}') return @@ -131,53 +139,205 @@ def model_request(edge, model, idx): edge_info = edges_info.get(edge) edge_ip_addr = edge_info.get('ip_addr') + request_data = preprocessed_datas.get(model) + req_time = time.time() + print(f'[{thread_id}] req start: {req_time}') if request_type == 'grpc': edge_ip_addr = f'{edge_ip_addr}:{grpc_port}' - result = model_modules[model].run_bench(1, edge_ip_addr, 0) + result = model_modules[model].run_bench(1, edge_ip_addr, 0, request_data) else: edge_ip_addr = f'http://{edge_ip_addr}:{rest_port}/' - result = model_modules[model].run_bench(1, edge_ip_addr) + result = model_modules[model].run_bench(1, edge_ip_addr, request_data) - inference_times[idx].extend(result) + res_time = time.time() + print(f'[{thread_id}] req end: {res_time}') + #print(f'[{thread_id}] total: {res_time - req_time}') + inference_times[idx].extend(result) # 추론시간 기록 + # inference_times.extend(result) # 추론시간 기록 + edge_req_time_info[edge].append((req_time, res_time)) # 들어오는 요청들 -requests_list = inference_request_workload_manager.load_workload_from_file(workload_file_name) +requests_info = inference_request_workload_manager.load_workload_info_from_file(workload_file_name) +total_req_num = requests_info.get('total_request_num') +requests_list = requests_info.get('requests') # 요청을 각 장비에 전달, 여러요청을 동시에 다룰 수 있도록 쓰레드 이용 threads = [] +# inference_times = [] inference_times = [[] for _ in range(len(requests_list))] -idx = -1 + +edge_req_time_info = {} +for edge in edges_to_inference: + edge_req_time_info[edge] = [] start_inference_time = time.time() -for cur_reqs in requests_list: - idx += 1 +# cur_progress = 0 +# for cur_reqs_idx, cur_reqs in enumerate(requests_list): +# request_sleep_time = 1 / len(cur_reqs) # 요청들을 1초에 나눠서 보내기 위한 슬립시간 + +# # cur_reqs_start_time = time.time() +# for req_idx, req in enumerate(cur_reqs): +# edge_to_infer = get_edge_by_model_rr(req) +# if edge_to_infer is None: +# print(f'{req} can\'t be inference') +# continue + +# # thread_start_time = time.time() + +# th = Thread(target=model_request, args=(edge_to_infer, req, cur_reqs_idx, cur_progress)) +# th.start() +# threads.append(th) + +# # thread_complete_time = time.time() - thread_start_time +# # if thread_complete_time < request_sleep_time: +# # time.sleep(request_sleep_time - thread_complete_time) +# cur_progress += 1 +# print(f'progress: {cur_progress}/{total_req_num}', end='\r') + +# # cur_reqs_end_time = time.time() +# # cur_reqs_total_time = cur_reqs_end_time - cur_reqs_start_time +# # if cur_reqs_total_time < 1: +# # time.sleep(1 - cur_reqs_total_time) + +# print(f'waiting to complete... (it takes about {cur_reqs_idx+1} seconds)') + +# for th in threads: +# th.join() + +cur_progress = 0 +executor = concurrent.futures.ThreadPoolExecutor(1000) +for cur_reqs_idx, cur_reqs in enumerate(requests_list): request_sleep_time = 1 / len(cur_reqs) # 요청들을 1초에 나눠서 보내기 위한 슬립시간 - for req in cur_reqs: + + cur_reqs_start_time = time.time() + for req_idx, req in enumerate(cur_reqs): edge_to_infer = get_edge_by_model_rr(req) if edge_to_infer is None: print(f'{req} can\'t be inference') continue + + # thread_start_time = time.time() + + executor.submit(model_request, edge_to_infer, req, cur_reqs_idx, cur_progress) + + cur_progress += 1 + print(f'progress: {cur_progress}/{total_req_num}', end='\r') - th = Thread(target=model_request, args=(edge_to_infer, req, idx)) - th.start() - threads.append(th) - time.sleep(request_sleep_time) + cur_reqs_end_time = time.time() + cur_reqs_total_time = cur_reqs_end_time - cur_reqs_start_time + if cur_reqs_total_time < 1: + time.sleep(1 - cur_reqs_total_time) + +print(f'waiting to complete... (it takes about {cur_reqs_idx+1} seconds)') + + + +# cur_progress = 0 +# s = sched.scheduler() +# for cur_reqs_idx, cur_reqs in enumerate(requests_list): +# request_sleep_time = 1 / len(cur_reqs) # 요청들을 1초에 나눠서 보내기 위한 슬립시간 +# for req_idx, req in enumerate(cur_reqs): +# edge_to_infer = get_edge_by_model_rr(req) +# if edge_to_infer is None: +# print(f'{req} can\'t be inference') +# continue + +# s.enter((request_sleep_time * req_idx) + cur_reqs_idx, 1, model_request, (edge_to_infer, req, cur_reqs_idx, cur_progress,)) +# cur_progress += 1 +# print(f'progress: {cur_progress}/{total_req_num}', end='\r') -for th in threads: - th.join() +# print(f'waiting to complete... (it takes about {cur_reqs_idx-1} seconds)') + +# s.run() + +print('complete!') end_inference_time = time.time() total_inference_time = end_inference_time - start_inference_time print() +print('total request num: ', total_req_num) print('total inference time', total_inference_time) +print() print('----------------------') print('inference time info (each argument is info about requests per sec)') print() -for time_info in inference_times: - print(time_info) + +for i, time_info in enumerate(inference_times, start=1): + time_info.sort() + time_info_len = len(time_info) + + print(f'[{i}] reqeust num per sec: ', time_info_len) + print(f'[{i}] avg:', sum(time_info)/time_info_len) + print(f'[{i}] min:', time_info[0]) + print(f'[{i}] max:', time_info[-1]) + + print(f'[{i}] 25%:', time_info[int(time_info_len/4)]) + print(f'[{i}] 50%:', time_info[int(time_info_len/2)]) + print(f'[{i}] 75%:', time_info[int((time_info_len*3)/4)]) + + print() + +# inference_times.sort() +# inference_times_len = len(inference_times) +# print(f'reqeust num per sec: ', inference_times_len) +# print(f'avg:', sum(inference_times)/inference_times_len) +# print(f'min:', inference_times[0]) +# print(f'max:', inference_times[-1]) + +# print(f'25%:', inference_times[int(inference_times_len/4)]) +# print(f'50%:', inference_times[int(inference_times_len/2)]) +# print(f'75%:', inference_times[int((inference_times_len*3)/4)]) + +print('----------------------') +print('idle time by edge') +print() + + +for edge in edge_req_time_info: + req_time_info = edge_req_time_info.get(edge) + req_time_info.sort() + + cur_req_time = 0 + cur_res_time = 0 + + idle_time = [] + + for (req_time, res_time) in req_time_info: + if cur_req_time == 0: + idle_time.append(req_time - start_inference_time) + cur_req_time = req_time + cur_res_time = res_time + elif req_time >= cur_req_time and req_time <= cur_res_time: + if res_time > cur_res_time: + cur_res_time = res_time + else: + idle_time.append(req_time - cur_res_time) + cur_req_time = req_time + cur_res_time = res_time + if cur_res_time != 0: + idle_time.append(end_inference_time - cur_res_time) + + idle_time.sort() + idle_time_len = len(idle_time) + total_idle_time = sum(idle_time) + + if idle_time_len == 0: + print(f'{edge}\'s idle time is 0') + print() + continue + + print(f'[{edge}] total:', total_idle_time) + print(f'[{edge}] avg:', total_idle_time/idle_time_len) + print(f'[{edge}] min:', idle_time[0]) + print(f'[{edge}] max:', idle_time[-1]) + + print(f'[{edge}] 25%:', idle_time[int(idle_time_len/4)]) + print(f'[{edge}] 50%:', idle_time[int(idle_time_len/2)]) + print(f'[{edge}] 75%:', idle_time[int((idle_time_len*3)/4)]) + print() From 732646f77908a9a06f827c7f23f2d65f1cdd3951 Mon Sep 17 00:00:00 2001 From: kh3654po Date: Thu, 1 Jun 2023 16:06:39 +0900 Subject: [PATCH 08/13] =?UTF-8?q?=EC=A0=84=EC=B2=98=EB=A6=AC=20=ED=95=A8?= =?UTF-8?q?=EC=88=98=EC=97=90=EC=84=9C=20=EC=86=8C=EC=8A=A4=EB=8D=B0?= =?UTF-8?q?=EC=9D=B4=ED=84=B0=EC=9D=98=20=EC=A0=88=EB=8C=80=EA=B2=BD?= =?UTF-8?q?=EB=A1=9C=EB=A5=BC=20=EB=B0=9B=EB=8F=84=EB=A1=9D=20=EC=88=98?= =?UTF-8?q?=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bench/inception_v3/preprocessing.py | 5 +---- bench/mobilenet_v1/preprocessing.py | 5 +---- bench/mobilenet_v2/preprocessing.py | 5 +---- bench/preprocessing_data_manager.py | 9 +++++++-- bench/yolo_v5/preprocessing.py | 5 +---- 5 files changed, 11 insertions(+), 18 deletions(-) diff --git a/bench/inception_v3/preprocessing.py b/bench/inception_v3/preprocessing.py index a8be0d5..233d8ea 100644 --- a/bench/inception_v3/preprocessing.py +++ b/bench/inception_v3/preprocessing.py @@ -3,12 +3,9 @@ from PIL import Image import os -def get_file_path(filename): - return os.path.join(os.path.dirname(__file__), filename) - # 이미지 로드 및 전처리 (for inception) def run_preprocessing(image_file_path): - img = Image.open(get_file_path(image_file_path)) + img = Image.open(image_file_path) img = img.resize((299, 299)) img_array = np.array(img) img_array = (img_array - np.mean(img_array)) / np.std(img_array) diff --git a/bench/mobilenet_v1/preprocessing.py b/bench/mobilenet_v1/preprocessing.py index 427f7c0..66c07d0 100644 --- a/bench/mobilenet_v1/preprocessing.py +++ b/bench/mobilenet_v1/preprocessing.py @@ -3,12 +3,9 @@ from PIL import Image import os -def get_file_path(filename): - return os.path.join(os.path.dirname(__file__), filename) - # 이미지 로드 및 전처리 (for mobilenet) def run_preprocessing(image_file_path): - img = Image.open(get_file_path(image_file_path)) + img = Image.open(image_file_path) img = img.resize((224, 224)) img_array = np.array(img) img_array = img_array.astype('float32') / 255.0 diff --git a/bench/mobilenet_v2/preprocessing.py b/bench/mobilenet_v2/preprocessing.py index 427f7c0..66c07d0 100644 --- a/bench/mobilenet_v2/preprocessing.py +++ b/bench/mobilenet_v2/preprocessing.py @@ -3,12 +3,9 @@ from PIL import Image import os -def get_file_path(filename): - return os.path.join(os.path.dirname(__file__), filename) - # 이미지 로드 및 전처리 (for mobilenet) def run_preprocessing(image_file_path): - img = Image.open(get_file_path(image_file_path)) + img = Image.open(image_file_path) img = img.resize((224, 224)) img_array = np.array(img) img_array = img_array.astype('float32') / 255.0 diff --git a/bench/preprocessing_data_manager.py b/bench/preprocessing_data_manager.py index 795c0c8..cf7bf7b 100644 --- a/bench/preprocessing_data_manager.py +++ b/bench/preprocessing_data_manager.py @@ -1,6 +1,7 @@ import tensorflow as tf import json import importlib +import os.path data_source_info = {'mobilenet_v1': '../dataset/imagenet/imagenet_1000_raw/n01843383_1.JPEG', @@ -9,6 +10,9 @@ 'yolo_v5': '../dataset/coco_2017/coco/images/val2017/000000089761.jpg', } +def get_file_path(filename): + return os.path.join(os.path.dirname(__file__), filename) + def regist_preprocessed_datas(request_type): preprocessed_datas = {} @@ -16,12 +20,13 @@ def regist_preprocessed_datas(request_type): if request_type == 'rest': preprocessing_module = importlib.import_module(f"{model}.preprocessing") data_source = data_source_info.get(model) - data = json.dumps({"instances": preprocessing_module.run_preprocessing(data_source).tolist()}) + + data = json.dumps({"instances": preprocessing_module.run_preprocessing(get_file_path(data_source)).tolist()}) preprocessed_datas.update({model: data}) elif request_type == 'grpc': preprocessing_module = importlib.import_module(f"{model}.preprocessing") data_source = data_source_info.get(model) - data = tf.make_tensor_proto(preprocessing_module.run_preprocessing(data_source)) + data = tf.make_tensor_proto(preprocessing_module.run_preprocessing(get_file_path(data_source))) preprocessed_datas.update({model: data}) return preprocessed_datas \ No newline at end of file diff --git a/bench/yolo_v5/preprocessing.py b/bench/yolo_v5/preprocessing.py index 45b3ed5..a0fe44f 100644 --- a/bench/yolo_v5/preprocessing.py +++ b/bench/yolo_v5/preprocessing.py @@ -3,12 +3,9 @@ import os import cv2 -def get_file_path(filename): - return os.path.join(os.path.dirname(__file__), filename) - # 이미지 로드 및 전처리 (for yolo) def run_preprocessing(image_file_path): - img = cv2.imread(get_file_path(image_file_path)) + img = cv2.imread(image_file_path) img = cv2.cvtColor(img, cv2.COLOR_BGR2RGB) img = cv2.resize(img, (640, 640)) img = img.astype('float32') / 255.0 From 0108318e812678b4fa732966bfe2e93a09950506 Mon Sep 17 00:00:00 2001 From: kh3654po Date: Thu, 1 Jun 2023 16:12:20 +0900 Subject: [PATCH 09/13] =?UTF-8?q?=EB=AA=A8=EB=93=A0=20=EC=9A=94=EC=B2=AD?= =?UTF-8?q?=EC=9D=84=20=EC=B2=98=EB=A6=AC=20=ED=9B=84=20=EA=B2=B0=EA=B3=BC?= =?UTF-8?q?=EB=A5=BC=20=EC=B6=9C=EB=A0=A5=ED=95=98=EB=8F=84=EB=A1=9D=20?= =?UTF-8?q?=EC=88=98=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bench/round_robin_scheduler.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/bench/round_robin_scheduler.py b/bench/round_robin_scheduler.py index baab066..a5d0d70 100644 --- a/bench/round_robin_scheduler.py +++ b/bench/round_robin_scheduler.py @@ -94,7 +94,6 @@ def regist_enable_model_modules(models_to_infer): # 각 모델의 preprocessing 데이터 저장 preprocessed_datas = preprocessing_data_manager.regist_preprocessed_datas(request_type) -print(f'Preprocessed datas: {preprocessed_datas}') # 딕셔너리에 모델별로 엣지장비이름 등록 -> 들어오는 요청에 따라 어느 장비에 보낼 차례인지 확인 할 수 있는 딕셔너리 생성 @@ -222,7 +221,7 @@ def model_request(edge, model, idx, thread_id): # thread_start_time = time.time() - executor.submit(model_request, edge_to_infer, req, cur_reqs_idx, cur_progress) + threads.append(executor.submit(model_request, edge_to_infer, req, cur_reqs_idx, cur_progress)) cur_progress += 1 print(f'progress: {cur_progress}/{total_req_num}', end='\r') @@ -234,7 +233,8 @@ def model_request(edge, model, idx, thread_id): print(f'waiting to complete... (it takes about {cur_reqs_idx+1} seconds)') - +for th in concurrent.futures.as_completed(threads): + continue # cur_progress = 0 # s = sched.scheduler() From 2a56743afa539efe4798261d4e93b0744b3a877d Mon Sep 17 00:00:00 2001 From: kh3654po Date: Fri, 14 Jul 2023 13:42:08 +0900 Subject: [PATCH 10/13] =?UTF-8?q?=EC=8B=9C=EB=82=98=EB=A6=AC=EC=98=A4?= =?UTF-8?q?=EB=A5=BC=20=EC=83=9D=EC=84=B1=ED=95=A0=20=EB=95=8C=20=EB=A7=A4?= =?UTF-8?q?=EC=B4=88=20=EC=9A=94=EC=B2=AD=EB=9F=89=EC=9D=B4=20=EC=9D=BC?= =?UTF-8?q?=EC=A0=95=ED=95=98=EB=8F=84=EB=A1=9D=20=EC=83=9D=EC=84=B1?= =?UTF-8?q?=ED=95=98=EB=8A=94=20=EA=B8=B0=EB=8A=A5=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bench/inference_request_workload_manager.py | 22 +++++++++++++++++++-- 1 file changed, 20 insertions(+), 2 deletions(-) diff --git a/bench/inference_request_workload_manager.py b/bench/inference_request_workload_manager.py index 479ba9f..7fbee05 100644 --- a/bench/inference_request_workload_manager.py +++ b/bench/inference_request_workload_manager.py @@ -12,7 +12,7 @@ file_name = 'workload4' req_time_num = 20 -def create_inference_request_workload(req_time): +def create_inference_request_workload_by_poisson(req_time): requests = [[] for _ in range(req_time)] total_req_num = 0 for (model_name, req_per_sec) in inference_request_info: @@ -31,6 +31,23 @@ def create_inference_request_workload(req_time): return workload_info +def create_inference_request_workload_regularly(req_time): + requests = [[] for _ in range(req_time)] + total_req_num = 0 + for (model_name, req_per_sec) in inference_request_info: + for idx in range(req_time): + requests[idx].extend([model_name for _ in range(req_per_sec)]) + total_req_num += req_per_sec + + for idx in range(req_time): + random.shuffle(requests[idx]) + + workload_info = {} + workload_info['total_request_num'] = total_req_num + workload_info['requests'] = requests + + return workload_info + def save_workload_info_to_file(file_name, workloads): with open(file_name, 'wb') as f: pk.dump(workloads, f) @@ -42,7 +59,8 @@ def load_workload_info_from_file(file_name): return loaded_workloads -# workloads = create_inference_request_workload(req_time_num) +# workloads = create_inference_request_workload_by_poisson(req_time_num) +# workloads = create_inference_request_workload_regularly(req_time_num) # save_workload_info_to_file(file_name, workloads) # loaded_workload_info = load_workload_info_from_file(file_name) # print(loaded_workload_info) From 2b31071ef354fef71f8b01e1dc3c9958b52d6af0 Mon Sep 17 00:00:00 2001 From: kh3654po Date: Fri, 14 Jul 2023 13:45:55 +0900 Subject: [PATCH 11/13] =?UTF-8?q?round=20robin=EC=9C=BC=EB=A1=9C=20?= =?UTF-8?q?=EC=8A=A4=EC=BC=80=EC=A4=84=EB=A7=81=ED=95=98=EB=8A=94=20?= =?UTF-8?q?=EC=BD=94=EB=93=9C=20=EB=A6=AC=ED=8C=A9=ED=86=A0=EB=A7=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bench/round_robin_scheduler.py | 387 +++++++++++++++++++++------------ 1 file changed, 251 insertions(+), 136 deletions(-) diff --git a/bench/round_robin_scheduler.py b/bench/round_robin_scheduler.py index a5d0d70..5a98368 100644 --- a/bench/round_robin_scheduler.py +++ b/bench/round_robin_scheduler.py @@ -4,7 +4,6 @@ import argparse import roundrobin import time -from threading import Thread import concurrent.futures import importlib import inference_request_workload_manager @@ -55,7 +54,7 @@ if edge not in registered_edges: print(f'--edge arg must be in {registered_edges}') exit(1) - + print(f'Edges to inference: {edges_to_inference}') @@ -96,7 +95,7 @@ def regist_enable_model_modules(models_to_infer): preprocessed_datas = preprocessing_data_manager.regist_preprocessed_datas(request_type) -# 딕셔너리에 모델별로 엣지장비이름 등록 -> 들어오는 요청에 따라 어느 장비에 보낼 차례인지 확인 할 수 있는 딕셔너리 생성 +# 딕셔너리에 모델별로 엣지장비이름 등록 -> 들어오는 요청에 따라 어느 장비에 보낼 차례인지 확인 할 수 있는 딕셔너리 생성 def regist_edges_to_model(): model_edge_info = {} @@ -141,124 +140,99 @@ def model_request(edge, model, idx, thread_id): request_data = preprocessed_datas.get(model) req_time = time.time() - print(f'[{thread_id}] req start: {req_time}') if request_type == 'grpc': edge_ip_addr = f'{edge_ip_addr}:{grpc_port}' result = model_modules[model].run_bench(1, edge_ip_addr, 0, request_data) else: edge_ip_addr = f'http://{edge_ip_addr}:{rest_port}/' result = model_modules[model].run_bench(1, edge_ip_addr, request_data) - res_time = time.time() - print(f'[{thread_id}] req end: {res_time}') - #print(f'[{thread_id}] total: {res_time - req_time}') - inference_times[idx].extend(result) # 추론시간 기록 - # inference_times.extend(result) # 추론시간 기록 - edge_req_time_info[edge].append((req_time, res_time)) + + + # 결과 분석을 위한 기록 + inference_times_per_model_in_edge[edge][model].extend(result) + inference_times_per_sec[idx].extend(result) + inference_times_per_edge_and_sec[idx][edge].extend(result) + inference_times_per_model_in_edge_and_sec[idx][edge][model].extend(result) + models_per_edge_and_sec[idx][edge].append(model) + inference_times_all.extend(result) + req_and_res_time_per_edge[edge].append((req_time, res_time)) # 들어오는 요청들 requests_info = inference_request_workload_manager.load_workload_info_from_file(workload_file_name) total_req_num = requests_info.get('total_request_num') -requests_list = requests_info.get('requests') +requests_list = requests_info.get('requests') +# 결과 분석을 위한 자료구조들 +inference_times_all = [] # 모든 요청에 대한 추론시간을 기록 +inference_times_per_sec = [[] for _ in range(len(requests_list))] # 매초 보낸 요청이 처리되는데 걸린 추론시간을 기록 +inference_times_per_edge_and_sec = [{edge: [] for edge in edges_to_inference} for _ in range(len(requests_list))] # +inference_times_per_model_in_edge = {edge: {model: [] for model in models_to_inference} for edge in edges_to_inference} # 각 엣지에서 각 모델이 추론하는데 걸린 시간을 기록 +inference_times_per_model_in_edge_and_sec = [{edge: {model: [] for model in models_to_inference} for edge in edges_to_inference} for _ in range(len(requests_list))] # 매초 보낸 요청이 각 엣지 각 모델에서 추론되는데 걸린 시간을 기록 +models_per_edge_and_sec = [{edge: [] for edge in edges_to_inference} for _ in range(len(requests_list))] # 매초 보내는 요청이 어떤 엣지에 대한 요청인지를 기록 +req_and_res_time_per_edge = {edge: [] for edge in edges_to_inference} # 요청을 보낸 시각과 응답을 받은 시간을 기록 -# 요청을 각 장비에 전달, 여러요청을 동시에 다룰 수 있도록 쓰레드 이용 -threads = [] -# inference_times = [] -inference_times = [[] for _ in range(len(requests_list))] - -edge_req_time_info = {} -for edge in edges_to_inference: - edge_req_time_info[edge] = [] +# 총 추론시간 측정 start_inference_time = time.time() -# cur_progress = 0 -# for cur_reqs_idx, cur_reqs in enumerate(requests_list): -# request_sleep_time = 1 / len(cur_reqs) # 요청들을 1초에 나눠서 보내기 위한 슬립시간 - -# # cur_reqs_start_time = time.time() -# for req_idx, req in enumerate(cur_reqs): -# edge_to_infer = get_edge_by_model_rr(req) -# if edge_to_infer is None: -# print(f'{req} can\'t be inference') -# continue - -# # thread_start_time = time.time() - -# th = Thread(target=model_request, args=(edge_to_infer, req, cur_reqs_idx, cur_progress)) -# th.start() -# threads.append(th) - -# # thread_complete_time = time.time() - thread_start_time -# # if thread_complete_time < request_sleep_time: -# # time.sleep(request_sleep_time - thread_complete_time) -# cur_progress += 1 -# print(f'progress: {cur_progress}/{total_req_num}', end='\r') - -# # cur_reqs_end_time = time.time() -# # cur_reqs_total_time = cur_reqs_end_time - cur_reqs_start_time -# # if cur_reqs_total_time < 1: -# # time.sleep(1 - cur_reqs_total_time) - -# print(f'waiting to complete... (it takes about {cur_reqs_idx+1} seconds)') - -# for th in threads: -# th.join() - +# 요청을 각 장비에 전달, 여러요청을 동시에 다룰 수 있도록 쓰레드 이용 +threads = [] cur_progress = 0 executor = concurrent.futures.ThreadPoolExecutor(1000) +pre_not_complete_reqs = 0 for cur_reqs_idx, cur_reqs in enumerate(requests_list): + if len(cur_reqs) <= 0: + continue + request_sleep_time = 1 / len(cur_reqs) # 요청들을 1초에 나눠서 보내기 위한 슬립시간 + + # 시나리오에 따라 요청을 전송, 각 요청은 쓰레드로 병렬처리 cur_reqs_start_time = time.time() for req_idx, req in enumerate(cur_reqs): edge_to_infer = get_edge_by_model_rr(req) if edge_to_infer is None: print(f'{req} can\'t be inference') continue - - # thread_start_time = time.time() - + threads.append(executor.submit(model_request, edge_to_infer, req, cur_reqs_idx, cur_progress)) - + time.sleep(request_sleep_time) + cur_progress += 1 print(f'progress: {cur_progress}/{total_req_num}', end='\r') + print(f'{cur_reqs_idx+1}초대 모든 요청전송 완료시간 : {time.time() - cur_reqs_start_time}') # 시나리오가 실제로 1초동안 원하는 양만큼 전송하는지 확인하는 코드 - cur_reqs_end_time = time.time() - cur_reqs_total_time = cur_reqs_end_time - cur_reqs_start_time - if cur_reqs_total_time < 1: - time.sleep(1 - cur_reqs_total_time) -print(f'waiting to complete... (it takes about {cur_reqs_idx+1} seconds)') + # 매초마다 처리되지 않은 요청량과 이전의 보냈던 요청처리량, 현재 시점에 보낸 요청처리량을 출력 + completed_req_num = 0 + real_req_num = 0 + real_req_num = 0 + for i in range(cur_reqs_idx+1): + real_req_num += len(requests_list[i]) + completed_req_num += len(inference_times_per_sec[i]) + if i == cur_reqs_idx-1 and cur_reqs_idx != 0: + pre_completed_req_num = completed_req_num + print(f'{cur_reqs_idx+1}초대 완료되지 않은 요청량: {real_req_num - completed_req_num}/{real_req_num}', f'{len(cur_reqs)+pre_completed_req_num-completed_req_num}/{len(cur_reqs)}') + print(f'{cur_reqs_idx}초대 완료된 요청량(완료되지않은 요청량): {pre_completed_req_num}/{real_req_num-len(cur_reqs)}({real_req_num-len(cur_reqs)-pre_completed_req_num})') + print(f'{cur_reqs_idx+1}초대 처리량: {pre_not_complete_reqs+len(cur_reqs) - (real_req_num - completed_req_num)}') + pre_not_complete_reqs = real_req_num - completed_req_num + +print(f'waiting to complete...') for th in concurrent.futures.as_completed(threads): continue -# cur_progress = 0 -# s = sched.scheduler() -# for cur_reqs_idx, cur_reqs in enumerate(requests_list): -# request_sleep_time = 1 / len(cur_reqs) # 요청들을 1초에 나눠서 보내기 위한 슬립시간 -# for req_idx, req in enumerate(cur_reqs): -# edge_to_infer = get_edge_by_model_rr(req) -# if edge_to_infer is None: -# print(f'{req} can\'t be inference') -# continue - -# s.enter((request_sleep_time * req_idx) + cur_reqs_idx, 1, model_request, (edge_to_infer, req, cur_reqs_idx, cur_progress,)) -# cur_progress += 1 -# print(f'progress: {cur_progress}/{total_req_num}', end='\r') - -# print(f'waiting to complete... (it takes about {cur_reqs_idx-1} seconds)') - -# s.run() - print('complete!') +# 총 추론시간 측정 end_inference_time = time.time() total_inference_time = end_inference_time - start_inference_time + +# 테스트 결과를 출력 + print() print('total request num: ', total_req_num) print('total inference time', total_inference_time) @@ -267,77 +241,218 @@ def model_request(edge, model, idx, thread_id): print('inference time info (each argument is info about requests per sec)') print() -for i, time_info in enumerate(inference_times, start=1): - time_info.sort() - time_info_len = len(time_info) - print(f'[{i}] reqeust num per sec: ', time_info_len) - print(f'[{i}] avg:', sum(time_info)/time_info_len) - print(f'[{i}] min:', time_info[0]) - print(f'[{i}] max:', time_info[-1]) +# 시나리오 전체에 대한 추론결과를 출력 (총처리량, 평균 최소 최대 처리시간 등) +def print_overall_inference_result(): + inference_times_all.sort() + inference_times_len = len(inference_times_all) + print(f'reqeust num per sec: ', inference_times_len) + print(f'avg:', sum(inference_times_all)/inference_times_len) + print(f'min:', inference_times_all[0]) + print(f'max:', inference_times_all[-1]) + + print(f'25%:', inference_times_all[int(inference_times_len/4)]) + print(f'50%:', inference_times_all[int(inference_times_len/2)]) + print(f'75%:', inference_times_all[int((inference_times_len*3)/4)]) - print(f'[{i}] 25%:', time_info[int(time_info_len/4)]) - print(f'[{i}] 50%:', time_info[int(time_info_len/2)]) - print(f'[{i}] 75%:', time_info[int((time_info_len*3)/4)]) - print() -# inference_times.sort() -# inference_times_len = len(inference_times) -# print(f'reqeust num per sec: ', inference_times_len) -# print(f'avg:', sum(inference_times)/inference_times_len) -# print(f'min:', inference_times[0]) -# print(f'max:', inference_times[-1]) -# print(f'25%:', inference_times[int(inference_times_len/4)]) -# print(f'50%:', inference_times[int(inference_times_len/2)]) -# print(f'75%:', inference_times[int((inference_times_len*3)/4)]) +# 시나리오에서 매초 처리한 처리량, 평균 최소 최대 처리시간 등 출력 +def print_inference_result_per_sec(): + for i, time_info in enumerate(inference_times_per_sec, start=1): + time_info.sort() + time_info_len = len(time_info) + + print(f'[{i}] reqeust num per sec: ', time_info_len) + if time_info_len == 0: + print() + continue + print(f'[{i}] avg:', sum(time_info)/time_info_len) + print(f'[{i}] min:', time_info[0]) + print(f'[{i}] max:', time_info[-1]) + + print(f'[{i}] 25%:', time_info[int(time_info_len/4)]) + print(f'[{i}] 50%:', time_info[int(time_info_len/2)]) + print(f'[{i}] 75%:', time_info[int((time_info_len*3)/4)]) + + print() + + +# 리스트안의 각 요소의 개수를 세는 함수 +def count_each_element_in_list(list_to_count: list): + element_names = list(set(list_to_count)) + element_names.sort() + result = [] + for name in element_names: + result.append(f'{name}: {list_to_count.count(name)}') + + return result + + +# 엣지마다 요청이 들어온 모델을 기준으로 처리량, 평균 최소 최대 처리시간 등 출력 +def print_inference_result_per_model_in_edge(): + for edge in inference_times_per_model_in_edge: + models = inference_times_per_model_in_edge[edge] + print(edge) + for model in models: + infer_time_per_model = models[model] + infer_time_per_model.sort() + infer_time_per_model_len = len(infer_time_per_model) + total_infer_time_per_model = sum(infer_time_per_model) + + if infer_time_per_model_len == 0: + print(f'{model} is not infered') + continue + print(f'[{model}] request num:', infer_time_per_model_len) + print(f'[{model}] total:', total_infer_time_per_model) + print(f'[{model}] avg:', total_infer_time_per_model/infer_time_per_model_len) + print(f'[{model}] min:', infer_time_per_model[0]) + print(f'[{model}] max:', infer_time_per_model[-1]) + + print(f'[{model}] 25%:', infer_time_per_model[int(infer_time_per_model_len/4)]) + print(f'[{model}] 50%:', infer_time_per_model[int(infer_time_per_model_len/2)]) + print(f'[{model}] 75%:', infer_time_per_model[int((infer_time_per_model_len*3)/4)]) + print() + + +# 엣지에서 매초 처리한 처리량, 평균 최소 최대 처리시간 등 출력 +def print_inference_result_per_edge_and_sec(): + for i, time_info in enumerate(inference_times_per_edge_and_sec, start=1): + for edge in time_info.keys(): + times = time_info.get(edge) + + times.sort() + times_len = len(times) + + models = models_per_edge_and_sec[i-1].get(edge) + requests = count_each_element_in_list(models) + + print(f'[{i}초대 {edge} 장비] reqeust num per sec: ', times_len) + if times_len == 0: + print() + continue + print(f'[{i}초대 {edge} 장비] requests: ', requests) + print(f'[{i}초대 {edge} 장비] avg:', sum(times)/times_len) + print(f'[{i}초대 {edge} 장비] min:', times[0]) + print(f'[{i}초대 {edge} 장비] max:', times[-1]) + + print(f'[{i}초대 {edge} 장비] 25%:', times[int(times_len/4)]) + print(f'[{i}초대 {edge} 장비] 50%:', times[int(times_len/2)]) + print(f'[{i}초대 {edge} 장비] 75%:', times[int((times_len*3)/4)]) + + print() + + +# 매초 엣지마다 요청이 들어온 모델을 기준으로 처리량, 평균 최소 최대 처리시간 등 출력 +def print_inference_result_per_model_in_edge_and_sec(): + for i, time_info in enumerate(inference_times_per_model_in_edge_and_sec, start=1): + for edge in time_info.keys(): + times_per_model_info = time_info.get(edge) + times = [] + + for model in times_per_model_info.keys(): + times_per_model = times_per_model_info.get(model) + times_per_model.sort() + times_per_model_len = len(times_per_model) + + print(f'[{i}초대 {edge} 장비 {model} 모델] request num per sec:', times_per_model_len) + if times_per_model_len == 0: + print() + continue + print(f'[{i}초대 {edge} 장비 {model} 모델] avg:', sum(times_per_model)/times_per_model_len) + print(f'[{i}초대 {edge} 장비 {model} 모델] min:', times_per_model[0]) + print(f'[{i}초대 {edge} 장비 {model} 모델] max:', times_per_model[-1]) + times.extend(times_per_model) + print() + + times.sort() + times_len = len(times) + + models = models_per_edge_and_sec[i-1].get(edge) + requests = count_each_element_in_list(models) + + print(f'[{i}초대 {edge} 장비] reqeust num per sec: ', times_len) + if times_len == 0: + print() + continue + print(f'[{i}초대 {edge} 장비] requests: ', requests) + print(f'[{i}초대 {edge} 장비] avg:', sum(times)/times_len) + print(f'[{i}초대 {edge} 장비] min:', times[0]) + print(f'[{i}초대 {edge} 장비] max:', times[-1]) + + print(f'[{i}초대 {edge} 장비] 25%:', times[int(times_len/4)]) + print(f'[{i}초대 {edge} 장비] 50%:', times[int(times_len/2)]) + print(f'[{i}초대 {edge} 장비] 75%:', times[int((times_len*3)/4)]) + + print() + +print_overall_inference_result() +print_inference_result_per_sec() +print_inference_result_per_model_in_edge() +print_inference_result_per_edge_and_sec() +print_inference_result_per_model_in_edge_and_sec() + print('----------------------') -print('idle time by edge') +print('idle time by edge (measurement including request processing time, network delay)') print() -for edge in edge_req_time_info: - req_time_info = edge_req_time_info.get(edge) - req_time_info.sort() +# idle time 출력, 실제로 장비의 유휴시간이 아닌 네트워크 지연을 포함 +def print_idle_time_by_edge(): + for edge in req_and_res_time_per_edge: + req_time_info = req_and_res_time_per_edge.get(edge) + req_time_info.sort() - cur_req_time = 0 - cur_res_time = 0 + cur_req_time = 0 + cur_res_time = 0 - idle_time = [] + idle_time = [] + idle_time_info = [] - for (req_time, res_time) in req_time_info: - if cur_req_time == 0: - idle_time.append(req_time - start_inference_time) - cur_req_time = req_time - cur_res_time = res_time - elif req_time >= cur_req_time and req_time <= cur_res_time: - if res_time > cur_res_time: + for (req_time, res_time) in req_time_info: + if cur_req_time == 0: + idle_time.append(req_time - start_inference_time) + idle_time_info.append((start_inference_time - start_inference_time, req_time - start_inference_time)) + cur_req_time = req_time cur_res_time = res_time - else: - idle_time.append(req_time - cur_res_time) - cur_req_time = req_time - cur_res_time = res_time - if cur_res_time != 0: - idle_time.append(end_inference_time - cur_res_time) - - idle_time.sort() - idle_time_len = len(idle_time) - total_idle_time = sum(idle_time) - - if idle_time_len == 0: - print(f'{edge}\'s idle time is 0') - print() - continue + elif req_time >= cur_req_time and req_time <= cur_res_time: + if res_time > cur_res_time: + cur_res_time = res_time + else: + idle_time.append(req_time - cur_res_time) + idle_time_info.append((cur_res_time - start_inference_time, req_time - start_inference_time)) + cur_req_time = req_time + cur_res_time = res_time + if cur_res_time != 0: + idle_time.append(end_inference_time - cur_res_time) + idle_time_info.append((cur_res_time - start_inference_time, end_inference_time - start_inference_time)) + + idle_time.sort() + idle_time_len = len(idle_time) + total_idle_time = sum(idle_time) + + requests_per_edge = [] + for ed in models_per_edge_and_sec: + requests_per_edge.extend(ed.get(edge)) + requests = count_each_element_in_list(requests_per_edge) + print(f'[{edge}] requests num:', len(req_time_info)) + if idle_time_len == 0: + print() + continue + print(f'[{edge}] requests:', requests) + print(f'[{edge}] total idle time:', total_idle_time) - print(f'[{edge}] total:', total_idle_time) - print(f'[{edge}] avg:', total_idle_time/idle_time_len) - print(f'[{edge}] min:', idle_time[0]) - print(f'[{edge}] max:', idle_time[-1]) + idle_time_percent = [] + for (s, e) in idle_time_info: + idle_time_percent.append(f'{((e - s) / total_inference_time) * 100}%') - print(f'[{edge}] 25%:', idle_time[int(idle_time_len/4)]) - print(f'[{edge}] 50%:', idle_time[int(idle_time_len/2)]) - print(f'[{edge}] 75%:', idle_time[int((idle_time_len*3)/4)]) + print('total inference time:', total_inference_time) + print('idle time info:', idle_time_info) + print('idle time:', idle_time) + print('idle time percent:', idle_time_percent) - print() + print() + +print_idle_time_by_edge() \ No newline at end of file From dd07358872ced3ae09f0ff16b6b701ca181907df Mon Sep 17 00:00:00 2001 From: kh3654po Date: Wed, 26 Jul 2023 11:09:59 +0900 Subject: [PATCH 12/13] =?UTF-8?q?gpu=EB=A5=BC=20=EC=82=AC=EC=9A=A9?= =?UTF-8?q?=ED=95=B4=20=ED=96=89=EB=A0=AC=EC=97=B0=EC=82=B0=ED=95=98?= =?UTF-8?q?=EB=8A=94=20c,=20cuda=20=EC=BD=94=EB=93=9C=20=EA=B5=AC=ED=98=84?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cuda/c/README.md | 14 ++++++++ cuda/c/maxmul.c | 84 ++++++++++++++++++++++++++++++++++++++++++++++++ cuda/c/maxmul.cu | 56 ++++++++++++++++++++++++++++++++ 3 files changed, 154 insertions(+) create mode 100644 cuda/c/README.md create mode 100644 cuda/c/maxmul.c create mode 100644 cuda/c/maxmul.cu diff --git a/cuda/c/README.md b/cuda/c/README.md new file mode 100644 index 0000000..d4b9386 --- /dev/null +++ b/cuda/c/README.md @@ -0,0 +1,14 @@ +1. 먼저 maxmul.cu를 공유 라이브러리로 컴파일한다. +``` +nvcc --ptxas-options=-v --compiler-options '-fPIC' -o libmaxmul.so --shared maxmul.cu +``` + +2. 위 라이브러리를 포함하여 maxmul.c를 컴파일한다. +``` +gcc -o maxmul maxmul.c -L. -lpthread -lmaxmul +``` + +3. 실행 +``` +./maxmul +``` \ No newline at end of file diff --git a/cuda/c/maxmul.c b/cuda/c/maxmul.c new file mode 100644 index 0000000..65a627e --- /dev/null +++ b/cuda/c/maxmul.c @@ -0,0 +1,84 @@ +#include +#include +#include +#include +#include +#include + +#define SIZE 600 // 행렬의 크기 +#define NUM_THREADS 1000 // 사용할 스레드 수 +#define REQUEST_TIME 10 + +void maxmul(int *A, int* B, int *C, int size); + +// 쓰레드 인자로 쓸 구조체 정의 +struct thread_args { + int thread_id; + struct timeval st; +}; + +// 입력 행렬 A와 B +int A[SIZE*SIZE]; +int B[SIZE*SIZE]; + +// 스레드 함수 +void *multiply(void *arg) { + struct thread_args* args = (struct thread_args*)arg; + int thread_id = args->thread_id; + struct timeval st = args->st; + + // 고유한 행렬 생성 + int matrix_A[SIZE*SIZE]; + int matrix_B[SIZE*SIZE]; + int matrix_C[SIZE*SIZE]; + + // 행렬 A와 B 복사 + for (int i = 0; i < SIZE*SIZE; i++) { + matrix_A[i] = A[i]; + matrix_B[i] = B[i]; + } + + // 행렬 곱셈 수행 + struct timeval start_time, end_time; + gettimeofday(&start_time, NULL); // 시작 시간 기록 + maxmul(matrix_A, matrix_B, matrix_C, SIZE); + gettimeofday(&end_time, NULL); // 종료 시간 기록 + + // 처리 시간 출력 + double execution_time = (double)(end_time.tv_sec - start_time.tv_sec) + (double)(end_time.tv_usec - start_time.tv_usec) / 1000000; + printf("Thread %d, matrix_result[0]=%d, start time: %.4f, end time: %.4f, processing time: %.4f seconds\n", thread_id, matrix_C[0], ((double)start_time.tv_sec + (double)start_time.tv_usec / 1000000) - ((double)st.tv_sec + (double)st.tv_usec / 1000000), ((double)end_time.tv_sec + (double)end_time.tv_usec / 1000000) - ((double)st.tv_sec + (double)st.tv_usec / 1000000), execution_time); + + pthread_exit(NULL); +} + +int main() { + pthread_t threads[NUM_THREADS]; + struct thread_args thread_args[NUM_THREADS]; + int num_thread_per_sec = NUM_THREADS / REQUEST_TIME; + float usleep_time = ((float)1 / (float)num_thread_per_sec) * 1000000; + + // 행렬 A와 B를 초기화 + for (int i = 0; i < SIZE*SIZE; i++) { + A[i] = rand()%100; + B[i] = rand()%100; + } + + struct timeval start_time; + gettimeofday(&start_time, NULL); // 시작 시간 기록 + + // 스레드 생성 및 행렬 곱셈 수행 + for (int i = 0; i < NUM_THREADS; i++) { + thread_args[i].thread_id = i; + thread_args[i].st = start_time; + pthread_create(&threads[i], NULL, multiply, (void *)&thread_args[i]); + + usleep((int)usleep_time); + } + + // 모든 스레드의 종료를 기다림 + for (int i = 0; i < NUM_THREADS; i++) { + pthread_join(threads[i], NULL); + } + + return 0; +} \ No newline at end of file diff --git a/cuda/c/maxmul.cu b/cuda/c/maxmul.cu new file mode 100644 index 0000000..3a21c6c --- /dev/null +++ b/cuda/c/maxmul.cu @@ -0,0 +1,56 @@ +// maxmul.cu + +#include +#include + +__global__ void vecmul(int *A, int* B, int *C, int size) { + // Row and Column indexes: + int row = blockIdx.y * blockDim.y + threadIdx.y; + int col = blockIdx.x * blockDim.x + threadIdx.x; + + // Are they below the maximum? + if (col < size && row < size) { + int result = 0; + for (int ix = 0; ix < size; ix++) { + result += A[row * size + ix] * B[ix * size + col]; + } + C[row * size + col] = result; + } +} + +// maxmul 함수의 구현 +extern "C" { + void maxmul(int *A, int* B, int *C, int size) { + + int total = size * size; + // Allocate device memory: + int* gpu_A; + int* gpu_B; + int* gpu_C; + int msize = total * sizeof(int); + cudaMalloc((void**)&gpu_A, msize); + cudaMemcpy(gpu_A, A, msize, cudaMemcpyHostToDevice); + cudaMalloc((void**)&gpu_B, msize); + cudaMemcpy(gpu_B, B, msize, cudaMemcpyHostToDevice); + cudaMalloc((void**)&gpu_C, msize); + + // Blocks & grids: + int block_size = 32; + int grid_size = (size + block_size - 1) / block_size; + dim3 blocks(block_size, block_size); + dim3 grid(grid_size, grid_size); + + // Call the kernel: + vecmul<<>>(gpu_A, gpu_B, gpu_C, size); + cudaDeviceSynchronize(); + + // Get the result Matrix: + cudaMemcpy(C, gpu_C, msize, cudaMemcpyDeviceToHost); + + // Free device matrices + cudaFree(gpu_A); + cudaFree(gpu_B); + cudaFree(gpu_C); + } +} + From 054d4ee466ce0bad1a0259c8ae2a98a7550771f9 Mon Sep 17 00:00:00 2001 From: kh3654po Date: Wed, 26 Jul 2023 11:10:41 +0900 Subject: [PATCH 13/13] =?UTF-8?q?gpu=EB=A5=BC=20=EC=82=AC=EC=9A=A9?= =?UTF-8?q?=ED=95=B4=20=ED=96=89=EB=A0=AC=EC=97=B0=EC=82=B0=ED=95=98?= =?UTF-8?q?=EB=8A=94=20go,=20cuda=20=EC=BD=94=EB=93=9C=20=EA=B5=AC?= =?UTF-8?q?=ED=98=84?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cuda/go/README.md | 9 +++++++ cuda/go/maxmul.cu | 56 ++++++++++++++++++++++++++++++++++++++++++ cuda/go/maxmul.go | 62 +++++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 127 insertions(+) create mode 100644 cuda/go/README.md create mode 100644 cuda/go/maxmul.cu create mode 100644 cuda/go/maxmul.go diff --git a/cuda/go/README.md b/cuda/go/README.md new file mode 100644 index 0000000..92441cd --- /dev/null +++ b/cuda/go/README.md @@ -0,0 +1,9 @@ +1. 먼저 maxmul.cu를 공유 라이브러리로 컴파일한다. +``` +nvcc --ptxas-options=-v --compiler-options '-fPIC' -o libmaxmul.so --shared maxmul.cu +``` + +2. maxmul.go 파일 실행 +``` +go run maxmul.go +``` diff --git a/cuda/go/maxmul.cu b/cuda/go/maxmul.cu new file mode 100644 index 0000000..3a21c6c --- /dev/null +++ b/cuda/go/maxmul.cu @@ -0,0 +1,56 @@ +// maxmul.cu + +#include +#include + +__global__ void vecmul(int *A, int* B, int *C, int size) { + // Row and Column indexes: + int row = blockIdx.y * blockDim.y + threadIdx.y; + int col = blockIdx.x * blockDim.x + threadIdx.x; + + // Are they below the maximum? + if (col < size && row < size) { + int result = 0; + for (int ix = 0; ix < size; ix++) { + result += A[row * size + ix] * B[ix * size + col]; + } + C[row * size + col] = result; + } +} + +// maxmul 함수의 구현 +extern "C" { + void maxmul(int *A, int* B, int *C, int size) { + + int total = size * size; + // Allocate device memory: + int* gpu_A; + int* gpu_B; + int* gpu_C; + int msize = total * sizeof(int); + cudaMalloc((void**)&gpu_A, msize); + cudaMemcpy(gpu_A, A, msize, cudaMemcpyHostToDevice); + cudaMalloc((void**)&gpu_B, msize); + cudaMemcpy(gpu_B, B, msize, cudaMemcpyHostToDevice); + cudaMalloc((void**)&gpu_C, msize); + + // Blocks & grids: + int block_size = 32; + int grid_size = (size + block_size - 1) / block_size; + dim3 blocks(block_size, block_size); + dim3 grid(grid_size, grid_size); + + // Call the kernel: + vecmul<<>>(gpu_A, gpu_B, gpu_C, size); + cudaDeviceSynchronize(); + + // Get the result Matrix: + cudaMemcpy(C, gpu_C, msize, cudaMemcpyDeviceToHost); + + // Free device matrices + cudaFree(gpu_A); + cudaFree(gpu_B); + cudaFree(gpu_C); + } +} + diff --git a/cuda/go/maxmul.go b/cuda/go/maxmul.go new file mode 100644 index 0000000..09e7b97 --- /dev/null +++ b/cuda/go/maxmul.go @@ -0,0 +1,62 @@ +package main + +/* +#cgo LDFLAGS: -L. -lmaxmul +void maxmul(int *A, int* B, int *C, int size); +*/ +import "C" + +import ( + "fmt" + "math/rand" + "sync" + "time" +) + +func Maxmul(a []C.int, b []C.int, size int, wg *sync.WaitGroup, reqIdx int, standardTime time.Time) { + wg.Add(1) + defer wg.Done() + + var startTime = time.Since(standardTime) + var totalSize = size * size + var c = make([]C.int, totalSize) + + C.maxmul(&a[0], &b[0], &c[0], C.int(size)) + + var endTime = time.Since(standardTime) + var totalTime = endTime - startTime + fmt.Printf("[%d] c[0]=%d, start time: %s, end time: %s, processing time: %s\n", reqIdx, c[0], startTime, endTime, totalTime) +} + +func main() { + var waitGroup sync.WaitGroup + + var reqNumPerSec = 100 + var reqSec = 10 + var totalReq = reqNumPerSec * reqSec + var rowSize int = 600 + var totalSize = rowSize * rowSize + + var aa = make([]C.int, totalSize) + var bb = make([]C.int, totalSize) + + for i := 0; i < totalSize; i++ { + aa[i] = C.int(rand.Intn(100)) + bb[i] = C.int(rand.Intn(100)) + } + + var sleepTime = float64(1) / float64(reqNumPerSec) + var sleepDuration = time.Second + if sleepTime < 1 { + sleepTime *= 1000 + sleepDuration = time.Millisecond + } + + var standardTime = time.Now() + for i := 0; i < totalReq; i++ { + go Maxmul(aa, bb, rowSize, &waitGroup, i, standardTime) + time.Sleep(sleepDuration * time.Duration(sleepTime)) + } + + waitGroup.Wait() +}