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/preprocessing.py b/bench/inception_v3/preprocessing.py index f0eabd8..233d8ea 100644 --- a/bench/inception_v3/preprocessing.py +++ b/bench/inception_v3/preprocessing.py @@ -1,15 +1,11 @@ #image 전처리 library -import tensorflow as tf import numpy as np 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/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/inference_request_workload_manager.py b/bench/inference_request_workload_manager.py new file mode 100644 index 0000000..7fbee05 --- /dev/null +++ b/bench/inference_request_workload_manager.py @@ -0,0 +1,66 @@ +from numpy import random +import pickle as pk + +# (modle name, requests per second) +inference_request_info = [ + ('mobilenet_v1', 20), + ('mobilenet_v2', 2), + ('inception_v3', 2), + ('yolo_v5', 1) +] + +file_name = 'workload4' +req_time_num = 20 + +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: + 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]) + + workload_info = {} + workload_info['total_request_num'] = total_req_num + workload_info['requests'] = requests + + 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) + + +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_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) 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/preprocessing.py b/bench/mobilenet_v1/preprocessing.py index 49164c0..66c07d0 100644 --- a/bench/mobilenet_v1/preprocessing.py +++ b/bench/mobilenet_v1/preprocessing.py @@ -1,15 +1,11 @@ #image 전처리 library -import tensorflow as tf import numpy as np 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_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/preprocessing.py b/bench/mobilenet_v2/preprocessing.py index 49164c0..66c07d0 100644 --- a/bench/mobilenet_v2/preprocessing.py +++ b/bench/mobilenet_v2/preprocessing.py @@ -1,15 +1,11 @@ #image 전처리 library -import tensorflow as tf import numpy as np 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/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/preprocessing_data_manager.py b/bench/preprocessing_data_manager.py new file mode 100644 index 0000000..cf7bf7b --- /dev/null +++ b/bench/preprocessing_data_manager.py @@ -0,0 +1,32 @@ +import tensorflow as tf +import json +import importlib +import os.path + + +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 get_file_path(filename): + return os.path.join(os.path.dirname(__file__), filename) + +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(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(get_file_path(data_source))) + preprocessed_datas.update({model: data}) + + return preprocessed_datas \ No newline at end of file diff --git a/bench/round_robin_scheduler.py b/bench/round_robin_scheduler.py new file mode 100644 index 0000000..5a98368 --- /dev/null +++ b/bench/round_robin_scheduler.py @@ -0,0 +1,458 @@ +# 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 +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) + +scheduler_args = parser.parse_args() +edges_to_inference = scheduler_args.edge + +# grpc or rest +request_type = 'rest' +# workload file name +workload_file_name = scheduler_args.workload +# 서버의 포트 정보 +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', 'mobilenet_v2', 'inception_v3'] + } + } + + +# --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}') + + +# 각 모델의 preprocessing 데이터 저장 +preprocessed_datas = preprocessing_data_manager.regist_preprocessed_datas(request_type) + + +# 딕셔너리에 모델별로 엣지장비이름 등록 -> 들어오는 요청에 따라 어느 장비에 보낼 차례인지 확인 할 수 있는 딕셔너리 생성 +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, thread_id): + 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') + request_data = preprocessed_datas.get(model) + + req_time = time.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() + + + # 결과 분석을 위한 기록 + 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') + +# 결과 분석을 위한 자료구조들 +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} # 요청을 보낸 시각과 응답을 받은 시간을 기록 + + +# 총 추론시간 측정 +start_inference_time = time.time() + +# 요청을 각 장비에 전달, 여러요청을 동시에 다룰 수 있도록 쓰레드 이용 +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 + + 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초동안 원하는 양만큼 전송하는지 확인하는 코드 + + + # 매초마다 처리되지 않은 요청량과 이전의 보냈던 요청처리량, 현재 시점에 보낸 요청처리량을 출력 + 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 + +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() + + +# 시나리오 전체에 대한 추론결과를 출력 (총처리량, 평균 최소 최대 처리시간 등) +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() + + +# 시나리오에서 매초 처리한 처리량, 평균 최소 최대 처리시간 등 출력 +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 (measurement including request processing time, network delay)') +print() + + +# 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 + + 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) + 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 + 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) + + idle_time_percent = [] + for (s, e) in idle_time_info: + idle_time_percent.append(f'{((e - s) / total_inference_time) * 100}%') + + 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_idle_time_by_edge() \ No newline at end of file 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/preprocessing.py b/bench/yolo_v5/preprocessing.py index d555ad1..a0fe44f 100644 --- a/bench/yolo_v5/preprocessing.py +++ b/bench/yolo_v5/preprocessing.py @@ -1,16 +1,11 @@ #image 전처리 library -import tensorflow as tf import numpy as np -from PIL import Image 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 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: 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); + } +} + 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() +}