83 行
2.6 KiB
Python
83 行
2.6 KiB
Python
import yaml
|
|
import logging
|
|
from apscheduler.schedulers.background import BackgroundScheduler
|
|
from .collector import GPUCollector
|
|
|
|
# 配置日志
|
|
logging.basicConfig(level=logging.INFO)
|
|
logger = logging.getLogger("GPUScheduler")
|
|
|
|
class GPUScheduler:
|
|
"""
|
|
负责定时触发远程采集任务并将结果分发的调度类
|
|
"""
|
|
def __init__(self, config_path, socketio=None):
|
|
self.config_path = config_path
|
|
self.socketio = socketio # 可选:如果传入 socketio 实例,则直接推送数据到前端
|
|
self.scheduler = BackgroundScheduler()
|
|
self.last_data = {} # 存储各服务器最后一次采集的数据 {server_alias: [gpu_info]}
|
|
|
|
def load_config(self):
|
|
"""加载本地配置文件"""
|
|
try:
|
|
with open(self.config_path, 'r', encoding='utf-8') as f:
|
|
return yaml.safe_load(f)
|
|
except Exception as e:
|
|
logger.error(f"Failed to load config file: {e}")
|
|
return None
|
|
|
|
def collect_all_servers(self):
|
|
"""遍历所有服务器并采集数据"""
|
|
config = self.load_config()
|
|
if not config:
|
|
return
|
|
|
|
servers = config.get('servers', [])
|
|
all_results = {}
|
|
|
|
for s_conf in servers:
|
|
alias = s_conf.get('alias', 'Unknown')
|
|
logger.info(f"Collecting data from {alias}...")
|
|
|
|
collector = GPUCollector(s_conf)
|
|
data = collector.fetch_gpu_data()
|
|
|
|
if data is not None:
|
|
all_results[alias] = data
|
|
else:
|
|
all_results[alias] = None # 标记为采集失败
|
|
|
|
# 更新内存状态
|
|
self.last_data = all_results
|
|
|
|
# 如果配置了 socketio,则实时推送给前端
|
|
if self.socketio:
|
|
self.socketio.emit('gpu_update', all_results)
|
|
logger.info("GPU data broadcasted via SocketIO")
|
|
|
|
def start(self):
|
|
"""启动定时任务"""
|
|
config = self.load_config()
|
|
if not config:
|
|
logger.error("Cannot start scheduler: config file missing or invalid")
|
|
return False
|
|
|
|
interval = config.get('settings', {}).get('interval', 5)
|
|
|
|
# 添加定时任务
|
|
self.scheduler.add_job(
|
|
self.collect_all_servers,
|
|
'interval',
|
|
seconds=interval,
|
|
id='gpu_collection_job'
|
|
)
|
|
|
|
self.scheduler.start()
|
|
logger.info(f"Scheduler started. Interval: {interval}s")
|
|
return True
|
|
|
|
def stop(self):
|
|
"""停止调度器"""
|
|
self.scheduler.shutdown()
|
|
logger.info("Scheduler stopped")
|