You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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作业状态关联节点:

  1. 初始化时,用sinfo获取所有空闲节点,按算力权重降序生成可用节点列表。
  2. 提交任务到指定节点时,用srun/sbatch绑定节点,同时记录作业ID与对应节点。
  3. 定期查询作业状态,当作业完成(不在队列或状态为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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 12:15:57