资讯动态

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

发布时间:2026/9/16 20:37:27 来源:尧图企业网站定制
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/skypilotSkyPilot 的 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的子类如RayCodeGen、SlurmCodeGen负责生成作业驱动脚本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 ExecutorSlurm 执行器在每个 Slurm 计算节点上通过如下命令被调用slurm.py 模块文档srun python -m sky.skylet.executor.slurm --scriptuser_script --log-dirpath ...在 task_codegen.py 中SlurmCodeGen.build_task_runner_cmd()实际构造的srun命令为unset $(env | awk -F /^SLURM_/ $1 !~ /^SLURM_CONF/ {print $1}) \ srun --exportALL --quiet --unbuffered --kill-on-bad-exit --jobidSLURM_JOB_ID \ --job-namesky-job_id --ntasks-per-node1 [--container-remap-root --container-namename: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/envuv 安装产生、不可执行遮蔽系统env导致execvp失败的 Slurm 怪癖。该模块专门处理 Slurm 特有的三类问题文档原文通过SLURM_PROCID与集群 IP 映射确定节点身份通过共享 NFS 上的信号文件协调 setup/run 两个阶段将日志写入每个节点唯一的日志文件并实时流式输出。完整的命令行参数main()函数slurm.py通过argparse定义以下参数可直接在自有集群中手动调用参数是否必选默认值说明--script二选一—用户脚本内联、Shell 引用形式脚本过长时改用--script-path--script-path二选一—脚本文件路径内联超长时的备选方案由 CodeGen 判断命令长度后自动切换--env-vars否{}JSON 编码的环境变量字典--log-dir是—日志文件目录--cluster-num-nodes是—集群节点总数--cluster-ips是—集群节点 IP 的逗号分隔列表--task-name否None单节点集群日志前缀使用的任务名--is-setup否flagFalse是否为 setup 命令影响日志前缀与文件名--cluster-home-dir否—集群共享文件系统 home 目录容器内~为容器本地路径需显式传入共享路径用于跨节点协调--alloc-signal-file否—资源分配完成信号文件路径--setup-done-signal-file否—setup 完成信号文件路径此外main()断言--script与--script-path至少提供一个否则直接失败。节点身份识别SLURM_PROCID IP 映射执行器首先从环境变量读取任务等级slurm.pyrank 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 fworker{node_idx}_get_ip_address()刻意不用hostname -I因为 Docker bridge 网段172.17.x.x会排在最前导致 IP 错配改用gethostbyname与_get_job_node_ips()保持一致。若 IP 不在列表中则抛RuntimeError避免静默错位。清理 step 级 Slurm 环境变量执行器在用户脚本前统一前置一段环境清理命令slurm.pyunset ${!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_TASK1而srun会把其中很多当作输入默认值导致用户脚本内嵌套的srun被静默约束成执行器 step 的形状每任务 1 CPU、沿用执行器的 CPU binding而非完整的 job allocation。保留job 级变量SLURM_JOB_ID、SLURM_JOB_NODELIST、SLURM_GPUS_ON_NODE等保证srun --overlap --jobid$SLURM_JOB_ID这类模式仍能作用于完整分配。日志分流与唯一命名由于所有节点的~/sky_logs目录共享在同一个文件系统上每个节点必须使用唯一文件名否则会互相覆盖slurm.pysetup 阶段setup-node_name.log如setup-head.log、setup-worker1.log。源码 TODO 注明这与其它云上统一的setup.log命名不一致但 Slurm 场景下必须如此单节点集群run.log多节点集群rank-node_name.log如0-head.log、1-worker1.log。日志通过run_bash_command_with_log()来自 sky/skylet/log_lib.py写入文件并实时流式输出同时附加带颜色的前缀前缀随场景区分slurm.pysetup / head(setup pid{pid})setup / worker(setup pid{pid}, ip1.2.3.4)单节点(task_name, pid{pid})多节点 head(head, rank0, pid{pid})多节点 worker(worker1, rank1, pid{pid}, ip1.2.3.4)其中{pid}占位符由run_with_log实际填充。信号文件协调 setup / run 阶段Slurm 执行器通过共享文件系统上的两个信号文件把资源分配与setup 完成两个事件在 Job Driver 与各节点 Executor 之间同步task_codegen.pyalloc_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 后端的工作目录或改用其它协调机制。整体时序为Job Driver 在后台线程中启动 run 阶段的srun--exclusive抢占分配等待alloc_signal_filerank 0 的 Executor 在获得分配后touch分配信号文件slurm.pyDriver 检测到分配完成后同时监控后台线程存活若srun提前失败则直接报FAILED_SETUP退出如有 setup 命令则再以--overlap --nodessetup_nodes启动 setup 的srun--overlap避免与已占用的分配互相阻塞死锁setup 成功后在驱动侧touchsetup 完成信号文件各节点 Executor 轮询等待setup_done_signal_file出现100ms 间隔后才真正运行用户脚本slurm.py运行结束Driver 清理两个信号文件并回收退出码。注入 SKYPILOT 环境变量对于非 setup 的 run 阶段执行器会向用户进程注入三个关键环境变量slurm.pyenv_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_idrank 0 先清空并创建目录防止残留文件提前满足屏障其余节点轮询等待目录出现每个节点写完自己的 done 文件后调用_wait_for_all_ranks()轮询每个对等节点的 done 文件500ms 间隔该函数永不抛异常——屏障存在的唯一目的是保持本 task 的 cgroup 开放直到对等任务结束协调失败不得改变用户脚本的退出码。_wait_for_all_ranks()slurm.py对文件系统错误做了精细处理按文件名逐个探测ENOENT 视为尚未完成与文件系统故障区分仅当文件系统持续报错超过BARRIER_ERROR_TIMEOUT_SECONDS 300秒才放弃该值设计上要超过 NFS 客户端默认acdirmax60s的属性缓存窗口若目录整体消失外部删除任何 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 函数。对应实现是RayCodeGentask_codegen.py它通过 Ray 的pgPlacement Groupray.remote()完成资源预留与任务分发节点身份、日志、同步都由 Ray 运行时自身承载因此 SkyPilot 的 Executor 模块只服务 Slurm 这一需要显式跨节点编排的执行模型。与 Job Driver 的完整联动流程综合文档与源码一次 Slurm 任务的完整生命周期如下生成SlurmCodeGentask_codegen.py把任务 YAML 编译为sky_job_id驱动脚本注册SIGTERM处理器_cancel_slurm_job_steps通过squeue -s -j jobid找到名为sky-job_id的 step 并scancel用于失败时取消运行线程的srun。投递驱动脚本在 head 节点执行置任务状态为PENDING。分配后台线程发起 run 阶段srun--nodesN --cpus-per-taskceil(CPU) --mem0 --gpus-per-nodegpus --exclusive--exclusive保证抢占整节点。同步rank 0 Executor touch 分配信号 → Driver 确认分配 →如有 setupDriver 以--overlap启动 setupsrun→ 成功后再 touch setup 完成信号 → 各节点 Executor 解除等待。执行各节点 Executor 清理 step 级SLURM_*变量、注入SKYPILOT_*变量、按节点唯一命名写日志并流式输出。收尾成功时走 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 --nodes2 --ntasks-per-node1 --cpus-per-task4 \ python -m sky.skylet.executor.slurm \ --scriptecho hello from rank $SLURM_PROCID; hostname \ --env-vars{FOO:bar} \ --log-dir$HOME/sky_logs \ --cluster-num-nodes2 \ --cluster-ips$(scontrol show hostnames $SLURM_JOB_NODELIST | tr \n , | sed s/,$//) \ --cluster-home-dir$HOME观察$HOME/sky_logs下生成的0-head.log与1-worker1.log以及日志前缀中的节点名、rank 与 IP。这即是对 slurm.py 核心逻辑IP 映射 → 环境清理 → 日志分流 → 屏障等待的一次端到端验证。进一步可阅读以下源码深入执行器本体sky/skylet/executor/slurm.py驱动脚本生成sky/backends/task_codegen.pySlurmCodeGen与 L301RayCodeGen日志工具sky/skylet/log_lib.pyrun_bash_command_with_log环境变量与路径常量sky/skylet/constants.py/DSMLparameter /DSMLinvoke /DSMLtool_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),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价