SkyPilot Task Executors 深度解析:Slurm 分布式任务执行器的工作机制与源码实现
2026/9/16 20:37:19 网站建设 项目流程

SkyPilot Task Executors 深度解析:Slurm 分布式任务执行器的工作机制与源码实现

【免费下载链接】skypilotThe AI Compute Platform for frontier teams. SkyPilot turns fragmented AI compute into one AI supercomputer, so frontier AI teams build custom intelligence faster.项目地址: https://gitcode.com/GitHub_Trending/sk/skypilot

SkyPilot 的 Task Executors 模块负责在每个集群节点上运行用户的训练脚本,是"Code Generator → Job Driver → Task Executor"三层执行流水线的最后一环。本文以该模块的官方文档为主体,结合 slurm.py 与 task_codegen.py 的源码实现,完整讲解 Slurm 场景下任务执行器的启动方式、命令行参数、节点身份识别、环境隔离、日志分流与跨节点同步屏障,并解释为何 Ray 后端不需要独立的 Executor 模块。读完本文,你将掌握 SkyPilot 在 Slurm 集群(含容器环境)上编排分布式任务的全部底层机制,并可直接用srun python -m sky.skylet.executor.slurm ...在自有 Slurm 集群上复现这套执行流程。

三个核心概念

根据 Task Executors 文档,该模块围绕三个相互协作的组件展开:

  • Code Generator(代码生成器)TaskCodeGen的子类(如RayCodeGenSlurmCodeGen),负责生成作业驱动脚本(Job Driver Script),源码位于 sky/backends/task_codegen.py。
  • Job Driver(作业驱动):生成出的 Python 脚本(~/.sky/sky_app/sky_job_<id>),运行在集群的 head 节点上,负责编排跨所有节点的分布式执行。
  • Task Executor(任务执行器):运行在每个集群节点上的模块,负责环境准备(setup)、日志记录(logging),以及与 Job Driver 之间的协调(coordination)。

三层组件呈严格的流水线关系:Code Generator 在用户侧把 YAML/命令行翻译成一段 Python 驱动脚本;驱动脚本被投递到 head 节点后,通过分布式调度器(Slurm 的srun或 Ray 的ray.remote())把真正的用户脚本分派到每个节点;而 Task Executor 就是在每个节点上真正执行用户脚本的那段代码。

整体架构

文档给出了如下三层架构图(原文 ASCII 图),直观展示了数据流方向:

┌─────────────────────────────────────────────────────────┐ │ Code Generator │ │ (RayCodeGen / SlurmCodeGen in task_codegen.py) │ │ Generates the job driver script │ └─────────────────────────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────┐ │ Job Driver │ │ (~/.sky/sky_app/sky_job_<id> - runs on head node) │ └─────────────────────────────────────────────────────────┘ │ ┌───────────────┼───────────────┐ ▼ ▼ ▼ ┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐ │ Task Executor │ │ Task Executor │ │ Task Executor │ │ (head) │ │ (worker1) │ │ (worker2) │ └─────────────────┘ └─────────────────┘ └─────────────────┘

值得注意的是,每个集群节点上都会运行一个 Task Executor 实例,head 与各 worker 通过共享文件系统上的信号文件完成状态同步(详见下文"信号文件协调机制")。从源码结构看(sky/skylet/executor/ 目录仅含slurm.py__init__.py),当前仓库实际落地的 Executor 实现只有 Slurm 一个,Ray 场景则将执行逻辑内联进驱动脚本,两者在架构上互补。

Executors:两种执行模型

slurm.py— Slurm Task Executor

Slurm 执行器在每个 Slurm 计算节点上通过如下命令被调用(slurm.py 模块文档):

srun python -m sky.skylet.executor.slurm --script=<user_script> --log-dir=<path> ...

在 task_codegen.py 中,SlurmCodeGen.build_task_runner_cmd()实际构造的srun命令为:

unset $(env | awk -F= '/^SLURM_/ && $1 !~ /^SLURM_CONF/ {print $1}') && \ srun --export=ALL --quiet --unbuffered --kill-on-bad-exit --jobid=<SLURM_JOB_ID> \ --job-name=sky-<job_id> --ntasks-per-node=1 [--container-remap-root --container-name=<name>:exec] \ <额外标志> /bin/bash -c '<python> -m sky.skylet.executor.slurm <runner_args>'

其中 Python 解释器由常量SKY_SLURM_PYTHON_CMD决定(定义于 constants.py,会先取消继承的PYTHONPATH以避免环境串扰)。代码中特别使用/usr/bin/env显式定位解释器,以规避$HOME/.local/bin/env(uv 安装产生、不可执行)遮蔽系统env导致execvp失败的 Slurm 怪癖。

该模块专门处理 Slurm 特有的三类问题(文档原文):

  1. 通过SLURM_PROCID与集群 IP 映射确定节点身份;
  2. 通过共享 NFS 上的信号文件协调 setup/run 两个阶段;
  3. 将日志写入每个节点唯一的日志文件并实时流式输出。
完整的命令行参数

main()函数(slurm.py)通过argparse定义以下参数,可直接在自有集群中手动调用:

参数是否必选默认值说明
--script二选一用户脚本(内联、Shell 引用形式);脚本过长时改用--script-path
--script-path二选一脚本文件路径(内联超长时的备选方案,由 CodeGen 判断命令长度后自动切换)
--env-vars'{}'JSON 编码的环境变量字典
--log-dir日志文件目录
--cluster-num-nodes集群节点总数
--cluster-ips集群节点 IP 的逗号分隔列表
--task-nameNone单节点集群日志前缀使用的任务名
--is-setup否(flag)False是否为 setup 命令(影响日志前缀与文件名)
--cluster-home-dir集群共享文件系统 home 目录(容器内~为容器本地路径,需显式传入共享路径用于跨节点协调)
--alloc-signal-file资源分配完成信号文件路径
--setup-done-signal-filesetup 完成信号文件路径

此外main()断言--script--script-path至少提供一个,否则直接失败。

节点身份识别:SLURM_PROCID + IP 映射

执行器首先从环境变量读取任务等级(slurm.py):

rank = int(os.environ['SLURM_PROCID']) # 任务 rank,注意不是节点索引 num_nodes = int(os.environ.get('SLURM_NNODES', 1))

随后模仿 Ray 的cluster_ips_to_node_id思路,把--cluster-ips拆分后用socket.gethostbyname(socket.gethostname())反查本机 IP,再在列表中定位节点索引:

ip_addr = _get_ip_address() node_idx = cluster_ips.index(ip_addr) node_name = 'head' if node_idx == 0 else f'worker{node_idx}'

_get_ip_address()刻意不用hostname -I,因为 Docker bridge 网段(172.17.x.x)会排在最前导致 IP 错配;改用gethostbyname_get_job_node_ips()保持一致。若 IP 不在列表中则抛RuntimeError,避免静默错位。

清理 step 级 Slurm 环境变量

执行器在用户脚本前统一前置一段环境清理命令(slurm.py):

unset "${!SLURM_STEP_@}" "${!SLURM_CPU_BIND@}" "${!SLURM_MEM_BIND@}" \ SLURM_CPUS_PER_TASK SLURM_TRES_PER_TASK SLURM_CPUS_ON_NODE \ SLURM_NTASKS SLURM_NPROCS SLURM_NTASKS_PER_NODE SLURM_TASKS_PER_NODE \ SLURM_DISTRIBUTION SLURM_SRUN_COMM_HOST SLURM_SRUN_COMM_PORT \ SLURM_LAUNCH_NODE_IPADDR SLURM_TASK_PID \ SLURM_PROCID SLURM_LOCALID SLURM_NODEID SLURM_GTIDS SLURM_STEPID

原因(源码注释结合 SchedMD Bug 14298 描述)在于:Slurm 会为执行器自身所在的 job step 填充 step 级SLURM_*变量(如SLURM_CPUS_PER_TASK=1),而srun会把其中很多当作输入默认值,导致用户脚本内嵌套的srun被静默约束成执行器 step 的形状(每任务 1 CPU、沿用执行器的 CPU binding),而非完整的 job allocation。保留job 级变量(SLURM_JOB_IDSLURM_JOB_NODELISTSLURM_GPUS_ON_NODE等),保证srun --overlap --jobid=$SLURM_JOB_ID这类模式仍能作用于完整分配。

日志分流与唯一命名

由于所有节点的~/sky_logs目录共享在同一个文件系统上,每个节点必须使用唯一文件名,否则会互相覆盖(slurm.py):

  • setup 阶段:setup-<node_name>.log(如setup-head.logsetup-worker1.log)。源码 TODO 注明这与其它云上统一的setup.log命名不一致,但 Slurm 场景下必须如此;
  • 单节点集群:run.log
  • 多节点集群:<rank>-<node_name>.log(如0-head.log1-worker1.log)。

日志通过run_bash_command_with_log()(来自 sky/skylet/log_lib.py)写入文件并实时流式输出,同时附加带颜色的前缀,前缀随场景区分(slurm.py):

  • setup / head:(setup pid={pid})
  • setup / worker:(setup pid={pid}, ip=1.2.3.4)
  • 单节点:(<task_name>, pid={pid})
  • 多节点 head:(head, rank=0, pid={pid})
  • 多节点 worker:(worker1, rank=1, pid={pid}, ip=1.2.3.4)

其中{pid}占位符由run_with_log实际填充。

信号文件协调 setup / run 阶段

Slurm 执行器通过共享文件系统上的两个信号文件,把"资源分配"与"setup 完成"两个事件在 Job Driver 与各节点 Executor 之间同步(task_codegen.py):

alloc_signal_file = f'~/.sky_alloc_{slurm_job_id}_{job_id}' setup_done_signal_file = f'~/.sky_setup_done_{slurm_job_id}_{job_id}'

信号文件存储于 home 目录,源码注释明确指出这依赖 home 目录挂在共享 NFS 上;若要支持非 NFS home,需让用户指定 NFS 后端的工作目录或改用其它协调机制。

整体时序为:

  1. Job Driver 在后台线程中启动 run 阶段的srun--exclusive抢占分配),等待alloc_signal_file
  2. rank 0 的 Executor 在获得分配后touch分配信号文件(slurm.py);
  3. Driver 检测到分配完成后(同时监控后台线程存活,若srun提前失败则直接报FAILED_SETUP退出),如有 setup 命令则再以--overlap --nodes=<setup_nodes>启动 setup 的srun--overlap避免与已占用的分配互相阻塞死锁);
  4. setup 成功后在驱动侧touchsetup 完成信号文件;
  5. 各节点 Executor 轮询等待setup_done_signal_file出现(100ms 间隔)后才真正运行用户脚本(slurm.py);
  6. 运行结束,Driver 清理两个信号文件并回收退出码。
注入 SKYPILOT 环境变量

对于非 setup 的 run 阶段,执行器会向用户进程注入三个关键环境变量(slurm.py):

env_vars['SKYPILOT_NODE_RANK'] = str(rank) # 本节点 rank env_vars['SKYPILOT_NUM_NODES'] = str(num_nodes) # 总节点数 env_vars['SKYPILOT_NODE_IPS'] = _get_job_node_ips() # 全部节点 IP(换行分隔)

_get_job_node_ips()hostlist.expand_hostlist()展开压缩格式的SLURM_JOB_NODELIST(如"node[1-3,5]"node1\nnode2...),再逐个gethostbyname解析为 IP。注释说明这里刻意不用scontrol show hostnames,因为scontrol及 Slurm CLI 在容器内可能不存在。相关常量名定义于 sky/skylet/constants.py,如SKYPILOT_NUM_GPUS_PER_NODE则由 CodeGen 在驱动侧注入。

run-done 跨节点同步屏障

这是多节点 Slurm 任务正确性最关键的细节。当任务成功、节点数大于 1、且 Slurm 使用 proctrack/cgroup 时,每个节点在退出前必须等待所有对等节点完成(slurm.py)。

背景(源码注释):proctrack/cgroup 会在某个 task 的主进程退出时,杀掉该 task 的 cgroup 内的所有进程。若一个节点提前退出,即使其它节点仍在运行(如 Ray worker 作为子进程),也会被误杀。因此失败的任务必须立即退出,以便srun --kill-on-bad-exit终止其余任务;而成功的任务则要等待全部对等任务完成。

实现方式:

  • 屏障目录为共享 home 下的.sky_run_done_<job_id>_<step_id>
  • rank 0 先清空并创建目录(防止残留文件提前满足屏障),其余节点轮询等待目录出现;
  • 每个节点写完自己的 done 文件后调用_wait_for_all_ranks(),轮询每个对等节点的 done 文件(500ms 间隔);
  • 该函数"永不抛异常"——屏障存在的唯一目的是保持本 task 的 cgroup 开放直到对等任务结束,协调失败不得改变用户脚本的退出码。

_wait_for_all_ranks()(slurm.py)对文件系统错误做了精细处理:按文件名逐个探测(ENOENT 视为"尚未完成",与文件系统故障区分);仅当文件系统持续报错超过BARRIER_ERROR_TIMEOUT_SECONDS = 300秒才放弃(该值设计上要超过 NFS 客户端默认acdirmax=60s的属性缓存窗口);若目录整体消失(外部删除,任何 rank 都无法上报)也按错误处理,避免无限等待。源码 TODO 同时指出:若有对等节点存活却不写 done 文件,其余 rank 会无限等待,仅靠--kill-on-bad-exit覆盖不了该场景,后续需要引入对端 Slurm task 状态之类的活性信号。

容器环境下,_is_proctrack_cgroup_enabled()会从显式传入的共享 home 目录(--cluster-home-dir)读取.sky_proctrack_type文件(常量定义见 constants.py),因为容器内~解析为容器本地路径/root/;文件缺失时保守地默认启用 cgroup 屏障。

Ray:无需独立 Executor

文档明确说明:Ray 直接使用ray.remote()把任务分派到 worker 节点,执行逻辑内联在生成的驱动脚本中,而不需要独立模块——因为 Ray 可以直接执行 Python 函数。对应实现是RayCodeGen(task_codegen.py):它通过 Ray 的pg(Placement Group)+ray.remote()完成资源预留与任务分发,节点身份、日志、同步都由 Ray 运行时自身承载,因此 SkyPilot 的 Executor 模块只服务 Slurm 这一需要显式跨节点编排的执行模型。

与 Job Driver 的完整联动流程

综合文档与源码,一次 Slurm 任务的完整生命周期如下:

  1. 生成SlurmCodeGen(task_codegen.py)把任务 YAML 编译为sky_job_<id>驱动脚本,注册SIGTERM处理器(_cancel_slurm_job_steps通过squeue -s -j <jobid>找到名为sky-<job_id>的 step 并scancel,用于失败时取消运行线程的srun)。
  2. 投递:驱动脚本在 head 节点执行,置任务状态为PENDING
  3. 分配:后台线程发起 run 阶段srun--nodes=N --cpus-per-task=<ceil(CPU)> --mem=0 --gpus-per-node=<gpus> --exclusive),--exclusive保证抢占整节点。
  4. 同步:rank 0 Executor touch 分配信号 → Driver 确认分配 →(如有 setup)Driver 以--overlap启动 setupsrun→ 成功后再 touch setup 完成信号 → 各节点 Executor 解除等待。
  5. 执行:各节点 Executor 清理 step 级SLURM_*变量、注入SKYPILOT_*变量、按节点唯一命名写日志并流式输出。
  6. 收尾:成功时走 run-done 屏障保持 cgroup 直到全体完成;失败则立即退出交由--kill-on-bad-exit清理;驱动脚本最终通过job_lib.set_exit_codes/set_status上报状态。

边界情况与工程取舍

从源码可以提炼出若干值得借鉴的工程细节:

  • 命令长度兜底build_task_runner_cmd()backend_utils.is_command_length_over_limit()判断脚本是否超长,超长则写入临时文件并用--script-path传递,运行后清理临时文件(task_codegen.py)。
  • 嵌套 srun 约束srun继承父 step 的SLURM_*变量会约束内层分配,驱动侧与执行器侧各做一层unset SLURM_*(但保留SLURM_CONF/SLURM_CONF_SERVER,否则slurmctld无法定位导致 DNS SRV 查找失败)。
  • 日志命名对齐问题:Slurm 下 setup 日志必须带节点名,与其它云setup.log不一致,源码 TODO 建议未来把该命名推广到所有云。
  • 屏障容错:屏障绝不改变用户脚本退出码;目录被外部删除、文件系统持续硬错误都有超时兜底;正常但未完成的对端节点仍在 TODO 中,等待更完善的活性检测方案。

如何在你的 Slurm 集群上复现

无需修改仓库代码,即可在自有 Slurm 集群手动验证执行器的两个核心行为。多节点示例(两个节点):

# 在 sbatch 分配内、两个节点各启动一个 Executor srun --nodes=2 --ntasks-per-node=1 --cpus-per-task=4 \ python -m sky.skylet.executor.slurm \ --script='echo "hello from rank $SLURM_PROCID"; hostname' \ --env-vars='{"FOO":"bar"}' \ --log-dir="$HOME/sky_logs" \ --cluster-num-nodes=2 \ --cluster-ips="$(scontrol show hostnames $SLURM_JOB_NODELIST | tr '\n' ',' | sed 's/,$//')" \ --cluster-home-dir="$HOME"

观察$HOME/sky_logs下生成的0-head.log1-worker1.log,以及日志前缀中的节点名、rank 与 IP。这即是对 slurm.py 核心逻辑(IP 映射 → 环境清理 → 日志分流 → 屏障等待)的一次端到端验证。

进一步可阅读以下源码深入:

  • 执行器本体:sky/skylet/executor/slurm.py
  • 驱动脚本生成:sky/backends/task_codegen.py(SlurmCodeGen)与 L301(RayCodeGen
  • 日志工具:sky/skylet/log_lib.py(run_bash_command_with_log
  • 环境变量与路径常量:sky/skylet/constants.py

</|DSML|parameter> </|DSML|invoke> </|DSML|tool_calls>

【免费下载链接】skypilotThe AI Compute Platform for frontier teams. SkyPilot turns fragmented AI compute into one AI supercomputer, so frontier AI teams build custom intelligence faster.项目地址: https://gitcode.com/GitHub_Trending/sk/skypilot

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询