Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -6,4 +6,9 @@
*.swp
.vscode/
.vs/
.idea/
.idea/

# OS / Python cache
.DS_Store
__pycache__/
*.pyc
31 changes: 9 additions & 22 deletions GPU-Virtual-Service/gpu-remoting/dispatcher.py
Original file line number Diff line number Diff line change
@@ -1,24 +1,15 @@
import socket
import threading
import json
import redis
from ctypes import Structure, c_char, c_int, c_size_t, sizeof
from cffi import FFI
import struct
import logging

from runtime_config import create_redis_connection, load_runtime_config

logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s', datefmt='%Y-%m-%d %H:%M:%S')

def create_redis_connection(redis_config):

r = redis.Redis(
host=redis_config['RedisConfig']['host'],
port=redis_config['RedisConfig']['port'],
db=redis_config['RedisConfig']['db'],
password=redis_config['RedisConfig']['password']
)
return r
logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s', datefmt='%Y-%m-%d %H:%M:%S')

def sort_gpu_info(gpu_Info_properties,conf):
if conf == 1:
Expand All @@ -37,9 +28,8 @@ def sort_gpu_info(gpu_Info_properties,conf):


def find_available_gpus(gpuInfoProperties, cnt):
with open('config.json', 'r') as f:
config = json.load(f)
r = create_redis_connection(config)
config = load_runtime_config()
r = create_redis_connection()
avail_gpus = list()
avail_gpus_properties = list()
if config['DispatcherMethod']['method'] == 1:
Expand Down Expand Up @@ -131,8 +121,7 @@ def find_available_gpus(gpuInfoProperties, cnt):

def handle_message(conn, addr, message, gpuInfoProperties):
if message.startswith("TypeA:"):
with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()
if config['DispatcherMethod']['method'] != 1:
gpuInfoProperties = handle_redis()
logging.info(f"Received message from {addr}: {message}")
Expand Down Expand Up @@ -170,8 +159,7 @@ def handle_client(conn, addr, gpuInfo):
def handle_redis():
gpu_info_properties = list()
# gpu_properties = list()
with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()
conf = config['DispatcherMethod']['method']
if conf == 1:#简单轮询调度
logging.info("RoundRobin scheduling")
Expand All @@ -181,7 +169,7 @@ def handle_redis():
logging.info("Utilization scheduling")
else:
logging.info("Free Memory and Utilization scheduling")
r = create_redis_connection(config)
r = create_redis_connection()
serv_ip = config['ServerConfig']['serverIp_']
cursor = 0
while True:
Expand All @@ -195,8 +183,7 @@ def handle_redis():
return gpu_info_properties

def main():
with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()
HOST = config['DispatcherConfig']['dpcIp_']
PORT = config['DispatcherConfig']['dpcPort_']

Expand All @@ -222,4 +209,4 @@ def main():
server_socket.close()

if __name__ == '__main__':
main()
main()
31 changes: 9 additions & 22 deletions GPU-Virtual-Service/gpu-remoting/dispatcher_.py
Original file line number Diff line number Diff line change
@@ -1,24 +1,15 @@
import socket
import threading
import json
import redis
from ctypes import Structure, c_char, c_int, c_size_t, sizeof
from cffi import FFI
import struct
import logging

from runtime_config import create_redis_connection, load_runtime_config

logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s', datefmt='%Y-%m-%d %H:%M:%S')

def create_redis_connection(redis_config):

r = redis.Redis(
host=redis_config['RedisConfig']['host'],
port=redis_config['RedisConfig']['port'],
db=redis_config['RedisConfig']['db'],
password=redis_config['RedisConfig']['password']
)
return r
logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s', datefmt='%Y-%m-%d %H:%M:%S')

def sort_gpu_info(gpu_Info_properties,conf):
if conf == 1:
Expand All @@ -40,9 +31,8 @@ def sort_gpu_info(gpu_Info_properties,conf):


def find_available_gpus(gpuInfoProperties, cnt, priority):
with open('config.json', 'r') as f:
config = json.load(f)
r = create_redis_connection(config)
config = load_runtime_config()
r = create_redis_connection()
avail_gpus = list()
avail_gpus_properties = list()
if config['DispatcherMethod']['method'] == 1:
Expand Down Expand Up @@ -138,8 +128,7 @@ def find_available_gpus(gpuInfoProperties, cnt, priority):

def handle_message(conn, addr, message, gpuInfoProperties):
if message.startswith("TypeA:"):
with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()
if config['DispatcherMethod']['method'] != 1:
gpuInfoProperties = handle_redis()
logging.info(f"Received message from {addr}: {message}")
Expand Down Expand Up @@ -178,8 +167,7 @@ def handle_client(conn, addr, gpuInfo):
def handle_redis():
gpu_info_properties = list()
# gpu_properties = list()
with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()
conf = config['DispatcherMethod']['method']
if conf == 1:#简单轮询调度
logging.info("RoundRobin scheduling")
Expand All @@ -189,7 +177,7 @@ def handle_redis():
logging.info("Utilization scheduling")
else:
logging.info("Free Memory and Utilization scheduling")
r = create_redis_connection(config)
r = create_redis_connection()
serv_ip = config['ServerConfig']['serverIp_']
cursor = 0
while True:
Expand All @@ -203,8 +191,7 @@ def handle_redis():
return gpu_info_properties

def main():
with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()
HOST = config['DispatcherConfig']['dpcIp_']
PORT = config['DispatcherConfig']['dpcPort_']

Expand All @@ -230,4 +217,4 @@ def main():
server_socket.close()

if __name__ == '__main__':
main()
main()
46 changes: 14 additions & 32 deletions GPU-Virtual-Service/gpu-remoting/monitor.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
import socket
import threading
import json
import redis
import pynvml
import os
import numpy as np
Expand All @@ -12,22 +11,14 @@
import struct
import logging

from runtime_config import create_redis_connection, load_runtime_config

# 删除 CUDA_VISIBLE_DEVICES 环境变量
if 'CUDA_VISIBLE_DEVICES' in os.environ:
del os.environ['CUDA_VISIBLE_DEVICES']

logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s', datefmt='%Y-%m-%d %H:%M:%S')

def create_redis_connection(redis_config):

r = redis.Redis(
host=redis_config['host'],
port=redis_config['port'],
db=redis_config['db'],
password=redis_config['password']
)
return r


# 加载 CUDA 运行时库
cuda = ctypes.CDLL("/usr/local/cuda/lib64/libcudart.so")
Expand All @@ -48,8 +39,7 @@ def get_device_properties_bytes(device_id=0):
#在初始化的时候调用,初始化redis里的GPU信息
def get_gpu_info():
# global r
with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()
logging.debug("Get GPU info")
gpuInfo = list()
compressed_datas = list()
Expand All @@ -60,14 +50,11 @@ def get_gpu_info():
dev_cnt = pynvml.nvmlDeviceGetCount()
# torch.cuda.device_count()
logging.info(f"Total {dev_cnt} GPU(s) found")
with open('config.json', 'r') as f:
config = json.load(f)

IP_addr = config['ServerConfig']['serverIp_']
Port = config['ServerConfig']['serverPort_']

logging.debug(f"IP_addr: {IP_addr}, Port: {Port}")
r = create_redis_connection(config['RedisConfig'])
r = create_redis_connection()

# 打印每个GPU的内存信息和利用率
for i in range(dev_cnt):
Expand Down Expand Up @@ -117,8 +104,7 @@ def get_gpu_info():
#更新redis中GPU信息
def update_gpu_info():
# global r
with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()
logging.info("Updating GPU info...")
gpuInfo = list()
# 初始化NVML,这对于使用pynvml库是必需的
Expand All @@ -130,7 +116,7 @@ def update_gpu_info():
# IP_addr = '192.168.0.208'
Ip_addr = config['ServerConfig']['serverIp_']

r = create_redis_connection(config['RedisConfig'])
r = create_redis_connection()


# 更新每个GPU的内存信息和利用率
Expand Down Expand Up @@ -164,9 +150,8 @@ def update_gpu_info():
def reduce_gpu_job(data):
# 解析data字符串,提取GPU ID
gpu_ids = data.split(':')[1].split(',')
with open('config.json', 'r') as f:
config = json.load(f)
r = create_redis_connection(config['RedisConfig'])
config = load_runtime_config()
r = create_redis_connection()

Ip_addr = config['ServerConfig']['serverIp_']
for gpu_id in gpu_ids:
Expand All @@ -184,9 +169,8 @@ def reduce_gpu_job(data):
def handle_new_job(data):
# 解析data字符串,提取GPU ID
gpu_ids = data.split(':')[1].split(',')
with open('config.json', 'r') as f:
config = json.load(f)
r = create_redis_connection(config['RedisConfig'])
config = load_runtime_config()
r = create_redis_connection()
Ip_addr = config['ServerConfig']['serverIp_']
for gpu_id in gpu_ids:
key = f"{Ip_addr}:{gpu_id}"
Expand All @@ -202,15 +186,14 @@ def handle_new_job(data):
logging.info(f"GPU {gpu_id} not found")

def collect_gpu_utilization():
with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()
# 初始化pynvml
pynvml.nvmlInit()
dev_cnt = pynvml.nvmlDeviceGetCount()

IP_addr = config['ServerConfig']['serverIp_']
Port = config['ServerConfig']['serverPort_']
r = create_redis_connection(config['RedisConfig'])
r = create_redis_connection()

# 创建一个字典来存储每个GPU的utilization数据
gpu_utilizations = {i: [] for i in range(dev_cnt)}
Expand Down Expand Up @@ -311,8 +294,7 @@ def schedule_update():
def main():
#一开始先初始化redis里的服务器信息,然后启动TCP连接,监听Server请求
get_gpu_info()
with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()

HOST = config['MonitorConfig']['monitorIp_']
PORT = config['MonitorConfig']['monitorPort_']
Expand Down Expand Up @@ -349,4 +331,4 @@ def main():


if __name__ == '__main__':
main()
main()
8 changes: 4 additions & 4 deletions GPU-Virtual-Service/gpu-remoting/proxy.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@
import lz4.frame
import struct

from runtime_config import load_runtime_config

logging.basicConfig(level=logging.DEBUG, format='%(asctime)s [%(levelname)s] %(message)s', datefmt='%Y-%m-%d %H:%M:%S')

clientId2gpuInfos = {}
Expand Down Expand Up @@ -215,9 +217,7 @@ def handle_client_worker(conn, addr, proxy, proxy_lock):
logging.info(f"Client#{clientID} released all GPU resources")

def main():

with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()

dpcIp = config['DispatcherConfig']['dpcIp_']
dpcPort = config['DispatcherConfig']['dpcPort_']
Expand Down Expand Up @@ -250,4 +250,4 @@ def main():

if __name__ == "__main__":
main()


8 changes: 4 additions & 4 deletions GPU-Virtual-Service/gpu-remoting/proxy_.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@
import lz4.frame
import struct

from runtime_config import load_runtime_config

logging.basicConfig(level=logging.DEBUG, format='%(asctime)s [%(levelname)s] %(message)s', datefmt='%Y-%m-%d %H:%M:%S')

clientId2gpuInfos = {}
Expand Down Expand Up @@ -216,9 +218,7 @@ def handle_client_worker(conn, addr, proxy, proxy_lock):
logging.info(f"Client#{clientID} released all GPU resources")

def main():

with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()

dpcIp = config['DispatcherConfig']['dpcIp_']
dpcPort = config['DispatcherConfig']['dpcPort_']
Expand Down Expand Up @@ -251,4 +251,4 @@ def main():

if __name__ == "__main__":
main()


7 changes: 3 additions & 4 deletions GPU-Virtual-Service/gpu-remoting/proxy_msg.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import time
from scheduler.msg_queue import *

from runtime_config import load_runtime_config


logging.basicConfig(level=logging.DEBUG, format='%(asctime)s [%(levelname)s] %(message)s', datefmt='%Y-%m-%d %H:%M:%S')
Expand Down Expand Up @@ -321,9 +322,7 @@ def handle_client_worker(conn, addr):


def main():

with open('config.json', 'r') as f:
config = json.load(f)
config = load_runtime_config()

listenIp = config['ClientConfig']['proxyIp_']
listenPort = config['ClientConfig']['proxyPort_']
Expand All @@ -347,4 +346,4 @@ def main():

if __name__ == "__main__":
main()


Loading
Loading