如何在SLURM集群中按算力优先级分发Python计算任务?
基于SLURM的动态节点调度实现方案
一、SLURM基础命令:获取节点状态与算力信息
首先用SLURM自带命令获取节点核心数据:
获取空闲节点与算力参数:
执行sinfo筛选空闲节点并提取CPU、GPU信息:sinfo -o "%N %c %G %t" | grep "idle"输出字段说明:
%N:节点名%c:CPU核心总数%G:GPU信息(格式如gpu:tesla_v100:2,代表2块V100显卡)%t:节点状态(idle为空闲,alloc为已占用)
自定义算力权重:
可根据硬件配置计算节点优先级,比如设定CPU核心数权重为1,V100显卡权重为10,节点总权重=CPU核心数 + GPU数量×10,确保算力高的节点排序靠前。
二、动态维护空闲节点列表的核心逻辑
无需手动跟踪节点占用状态,通过SLURM作业状态关联节点:
- 初始化时,用
sinfo获取所有空闲节点,按算力权重降序生成可用节点列表。 - 提交任务到指定节点时,用
srun/sbatch绑定节点,同时记录作业ID与对应节点。 - 定期查询作业状态,当作业完成(不在队列或状态为
COMPLETED),将对应节点放回可用列表并重排序。
三、Python代码实现示例
以下是贴合你伪代码的实际实现,结合SLURM命令调用:
import subprocess import time import re def calculate_node_weight(node_line): """计算节点算力权重:CPU核心数 + GPU数量×10(可自定义权重规则)""" node_name, cpu_cores, gpu_info, _ = node_line.split() cpu = int(cpu_cores) # 提取GPU数量,适配SLURM的GPU字段格式 gpu_match = re.search(r':(\d+)$', gpu_info) gpu_count = int(gpu_match.group(1)) if gpu_match else 0 return cpu + gpu_count * 10 def get_sorted_idle_nodes(): """获取当前空闲节点,按算力权重降序排序""" result = subprocess.run( ['sinfo', '-o', '%N %c %G %t'], capture_output=True, text=True ) # 跳过表头,处理每行节点数据 node_lines = result.stdout.strip().split('\n')[1:] idle_nodes = [line.strip() for line in node_lines if 'idle' in line] # 按算力权重排序后提取节点名 idle_nodes.sort(key=calculate_node_weight, reverse=True) return [line.split()[0] for line in idle_nodes] def submit_task_to_node(task_input, node_name): """提交任务到指定节点,返回作业ID""" # 替换为你的计算函数执行脚本/命令 cmd = [ 'srun', '--nodes=1', '--nodelist=' + node_name, '--job-name=sim_task', './run_expensive_function.sh', str(task_input) ] # 后台提交任务,获取SLURM返回的作业ID result = subprocess.run( cmd + ['--wait=no'], capture_output=True, text=True ) job_id_match = re.search(r'job (\d+)', result.stderr) return job_id_match.group(1) if job_id_match else None def check_job_completed(job_id): """检查作业是否已完成""" result = subprocess.run( ['squeue', '-j', job_id, '-o', '%T'], capture_output=True, text=True ) # 作业不在队列中或状态为COMPLETED则判定为完成 if not result.stdout.strip(): return True status = result.stdout.strip().split('\n')[1] return status == 'COMPLETED' # 初始化任务与运行中作业记录 input_tasks = [f'input_{i}' for i in range(10)] # 替换为你的实际任务输入 running_jobs = {} # 键:作业ID,值:节点名 while input_tasks or running_jobs: # 回收已完成任务的节点 completed_job_ids = [] for job_id, node_name in running_jobs.items(): if check_job_completed(job_id): completed_job_ids.append(job_id) for job_id in completed_job_ids: freed_node = running_jobs.pop(job_id) print(f"节点 {freed_node} 已释放") # 获取最新空闲节点列表 idle_nodes = get_sorted_idle_nodes() # 分配任务到空闲节点 while input_tasks and idle_nodes: current_task = input_tasks.pop(0) target_node = idle_nodes.pop(0) job_id = submit_task_to_node(current_task, target_node) if job_id: running_jobs[job_id] = target_node print(f"任务 {current_task} 提交到节点 {target_node},作业ID {job_id}") else: # 提交失败,将任务放回队列重试 input_tasks.insert(0, current_task) print(f"任务 {current_task} 提交到节点 {target_node} 失败,稍后重试") # 轮询间隔,可根据任务时长调整 time.sleep(10)
四、关键注意事项
- 节点权限:确认你有提交任务到指定节点的权限,部分集群可能限制直接指定节点,需联系管理员确认。
- 权重规则调整:
calculate_node_weight函数可根据集群硬件配置自定义,比如给A100显卡设置更高权重。 - 错误处理:示例简化了错误逻辑,实际使用中需添加节点不可用、任务提交失败的重试机制。
- 轮询频率:
time.sleep(10)可根据任务平均时长调整,避免过于频繁查询SLURM造成集群负载。
内容的提问来源于stack exchange,提问作者GTOgod
相关产品推荐
相关产品推荐

