ZenML GPU 分布式训练实战:从 ResourceSettings 申请 GPU 到 Accelerate 与多节点启动器
2026/9/18 6:10:54 网站建设 项目流程

ZenML GPU 分布式训练实战:从 ResourceSettings 申请 GPU 到 Accelerate 与多节点启动器

【免费下载链接】zenmlZenML 🙏: One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenml

当你需要比笔记本更强的算力时,ZenML 提供了一条完整的 GPU 训练路径:在单个@step上用ResourceSettings申请 CPU/GPU/内存,用DockerSettings构建带 CUDA 运行时的容器镜像,用 🤗 Accelerate 集成把训练 fan-out 到多卡,最终通过CommandStep包装 TorchX、Ray 等分布式启动器完成多节点训练。读完本篇,你可以掌握 ZenML 中 GPU 资源声明、CUDA 镜像准备、多卡/多节点训练的三种单机模式与两种多节点模式的完整配置方法,并理解每种模式背后在 src/zenml/config/resource_settings.py 与 src/zenml/integrations/huggingface/steps/accelerate_runner.py 中的实际实现。

1. 为单个 Step 申请额外资源(CPU / GPU / 内存)

如果你的 orchestrator 支持,可以直接在 ZenML 的@step上预留 CPU、GPU 和内存:

from zenml import step from zenml.config import ResourceSettings @step(settings={ "resources": ResourceSettings(cpu_count=8, gpu_count=2, memory="16GB") }) def training_step(...): ... # heavy training logic

各参数在源码 src/zenml/config/resource_settings.py 中的定义与约束如下:

字段类型 / 取值说明
cpu_countOptional[PositiveFloat]期望分配的 CPU 核数,可为小数;调度时内部换算为mcpu(毫核)
gpu_countOptional[NonNegativeInt]GPU 数量;0表示不申请 GPU(会显式移除资源请求中的gpu键)
memory字符串,匹配正则^[0-9]+(KB\|KIB\|MB\|MIB\|GB\|GIB\|TB\|TIB\|PB\|PIB)$内存大小必须带单位后缀,例如"16GB""8GiB";二进制单位(KiB/MiB…)按 2 的幂换算
pool_resourcesOptional[Dict[str, PositiveInt]]用于 ZenML 资源池的自定义资源键(如tpuvcpus);当gpu_count/cpu_count/memory同时设置时,后者优先覆盖gpu/mcpu/memory_mb
preemptiblebool,默认True资源是否可被抢占,仅在使用 ZenML 资源池时生效

从源码结构看,ResourceSettings通过merged_requested_resources()方法把上述字段归一化为调度器可理解的资源请求映射:gpu_countgpucpu_countmcpuceil(cpu * 1000))、memorymemory_mb,再叠加pool_resources中的自定义键。这正是各类 orchestrator/step operator 读取资源需求时看到的最终形态。

两点使用建议(原文档提示):

  • 请查阅你所用 orchestrator 的文档,部分 orchestrator(如 SkyPilot)不直接用ResourceSettings,而是暴露自己的专用 settings;
  • 如果你的 orchestrator 无法满足这些资源要求,可以考虑把该 step 卸载(off-load)到一个专用的 step operator。

2. 构建 CUDA 容器镜像

仅申请 GPU 是不够的——Docker 镜像本身必须携带 CUDA 运行时,GPU 才会真正可见。

from zenml import pipeline from zenml.config import DockerSettings docker = DockerSettings( parent_image="pytorch/pytorch:2.1.0-cuda12.1-cudnn8-runtime", python_package_installer_args={"system": None}, requirements=["zenml", "torchvision"] ) @pipeline(settings={"docker": docker}) def my_gpu_pipeline(...): ...

建议优先使用 TensorFlow/PyTorch 官方 CUDA 镜像,或 AWS、GCP、Azure 提供的预构建镜像。

可选:在 Step 开始时清空 CUDA 缓存

如果你要压榨 GPU 的每一 MB 显存,可以在每个 step 开头清空 CUDA 缓存:

import gc, torch def cleanup_memory(): while gc.collect(): torch.cuda.empty_cache()

在 GPU 密集型 step 的入口处调用cleanup_memory()即可。

3. 用 🤗 Accelerate 做单步多卡/多机训练

ZenML 与 Hugging Face Accelerate 启动器集成。用run_with_accelerate装饰你的训练 step,即可把它 fan-out 到多张 GPU 或多台机器:

from zenml import step, pipeline from zenml.integrations.huggingface.steps import run_with_accelerate @run_with_accelerate(num_processes=4, multi_gpu=True) @step def training_step(...): ... # your distributed training code @pipeline def dist_pipeline(...): training_step(...)

常用参数:

  • num_processes:要启动的进程总数(每个 GPU 一个)
  • multi_gpu=True:启用多 GPU 模式
  • cpu=True:强制使用 CPU 训练
  • mixed_precision"fp16"/"bf16"/"no"

⚠️重要限制:Accelerate 装饰的 step 必须用关键字参数调用,且不能在 pipeline 定义内被二次包装。

限制并非文档口头约定,而是有源码强制保证。在 src/zenml/integrations/huggingface/steps/accelerate_runner.py 中:

  • inner()入口检测到位置参数即抛ValueError("Accelerated steps do not support positional arguments.")
  • 装饰器内部通过get_pipeline_context()检查:若当前已处于 pipeline 上下文中做函数式调用,直接抛RuntimeError,要求“把装饰器直接应用到 step 上”,并给出允许/禁止的写法示例;
  • 运行时实现是:把 step 的 entrypoint 经create_cli_wrapped_script(entrypoint, flavor="accelerate")包装成一个 CLI 脚本,把@step调用时传入的关键字参数转换为--arg value命令行参数,再用 Accelerate 的launch_command_parser()解析、合并装饰器参数后调用launch_command(args);step 返回值用 cloudpickle 序列化到output_path,训练结束后由包装层pickle.load还原,从而保证 step 的返回值正常回流给 ZenML。

准备容器

使用与第 2 节相同的 CUDA 镜像,并在 requirements 中追加 Accelerate:

DockerSettings( parent_image="pytorch/pytorch:2.1.0-cuda12.1-cudnn8-runtime", python_package_installer_args={"system": None}, requirements=["zenml", "accelerate", "torchvision"] )

4. 跨进程/跨节点分布式训练:单机多卡与多节点两条路线

第 3 节的 Accelerate 已经把单个 step 在单机上 fan-out 到多张 GPU。更一般地,分布式训练分两种场景,对应不同的工具选择:

  • 单机多卡(single node, multiple GPUs):用原生 step(自管多进程,或 Accelerate 集成),或用CommandSteptorchrun。无需外部启动器或 gang scheduler。
  • 多节点(multiple nodes):把一个分布式launcher包进 CommandStep:ZenML 负责管理 run,launcher 负责 worker 进程。

4.1 单机多卡的三种模式

所有进程都生活在 step 自己的容器里,因此不需要外部启动器或 gang scheduler。以下三种模式当前都可用,选与你现有训练习惯最匹配的一种。

模式 1:自己拉起进程的原生 step

ZenML 开箱支持多进程训练——用ResourceSettings申请 GPU,然后由你自己的代码为每张卡起一个进程(例如torch.multiprocessing.spawn并内部搭 DDP)。无需额外库,同时保留 ZenML 的输入/输出与日志追踪:

import torch.distributed as dist import torch.multiprocessing as mp from zenml import pipeline, step from zenml.config import ResourceSettings def _worker(rank: int, world_size: int) -> None: dist.init_process_group("nccl", rank=rank, world_size=world_size) # ... your DDP training ... dist.destroy_process_group() @step( runtime="isolated", # run in its own container with the GPUs, not inline settings={"resources": ResourceSettings(gpu_count=4)}, ) def train() -> None: mp.spawn(_worker, args=(4,), nprocs=4) # one process per GPU @pipeline(dynamic=True) def training() -> None: train()

runtime="isolated"让 orchestrator 把 step 放进一个按 GPU 资源请求规格配置的新容器中运行,而不是内联在编排进程里——这正是重型训练 step 需要的运行方式。由于 isolated 容器使用 pipeline 镜像,请确保该镜像已包含 CUDA 和 torch(见第 2 节)。isolated 运行是动态 pipeline 特性,因此 pipeline 需要声明dynamic=True

模式 2:Accelerate 集成

如果不想自己搭进程组,就用run_with_accelerate(见第 3 节)处理 fan-out:

from zenml import step from zenml.config import ResourceSettings from zenml.integrations.huggingface.steps import run_with_accelerate @run_with_accelerate(num_processes=4, multi_gpu=True) @step(settings={"resources": ResourceSettings(gpu_count=4)}) def train(...): ... # your training code
模式 3:跑torchrunCommandStep

如果你已经习惯用torchrun(或任意 launcher CLI)驱动训练,直接把它包进 command step。launcher 和它的工作进程共享同一个容器,因此一个镜像同时携带zenml+torch+train.py即可:

from zenml import CommandStep, pipeline from zenml.config import DockerSettings, ResourceSettings train = CommandStep( command=["torchrun", "--standalone", "--nproc-per-node=4", "train.py"], step_operator="gmi-k8s", settings={ "docker": DockerSettings( parent_image="pytorch/pytorch:2.4.0-cuda12.1-cudnn9-runtime", requirements=["zenml"], ), "resources": ResourceSettings(gpu_count=4), }, ) @pipeline(dynamic=True, depends_on=[train]) def training() -> None: train()

三种模式的权衡:模式 1 和 2 保留 ZenML 的 artifact 与日志追踪;模式 3 把训练当作不透明命令(日志落在后端,没有 inputs/outputs),但可以让你原封不动地复用现有torchrun入口。同一份原生train.py(见 4.2 节)即可用于模式 3。

CommandStep的能力边界由源码 src/zenml/steps/command_step.py 明确定义:它的entrypoint()只是subprocess.run(command, check=True)——命令退出码 0 即 step 成功,非零即失败;构造函数会强制 command 非空,并拒绝声明 inputs/outputs 或配置 hook 的用法,同时默认enable_cache=False

4.2 多节点:把 launcher 包进CommandStep

一旦训练跨越多台机器,最干净的做法是让专门的 launcher 拥有 worker gang,让 ZenML 拥有 run。

为什么由 launcher(而不是 ZenML)启动 worker

torch.distributed/torchrun只做 rank 之间的rendezvous 协调——分配RANKWORLD_SIZELOCAL_RANK并把进程连起来,但负责在别的节点上开机器或起进程。必须有一个能供给并 gang-schedule 这 N 个 worker 进程的 launcher:TorchX、Ray、Slurm 等。

ZenML 不重写这部分,而是给你一个清晰的分工缝:

  • CommandStep在 step operator 的容器里执行一条不透明命令
  • 把这条命令指向 launcher,由 launcher 负责调度并启动 worker gang;
  • ZenML 记录 run,并通过 step operator 的submit/get_status/cancel生命周期跟踪launcher进程。launcher pod 会阻塞到整个作业结束,因此 ZenML step 与作业同生共死。

骨架永远相同,只是命令不同:

from zenml import CommandStep, pipeline from zenml.config import DockerSettings train = CommandStep( command=[...launcher CLI...], # TorchX / Ray step_operator="<your-step-operator>", # where the launcher itself runs settings={"docker": DockerSettings(requirements=["<launcher-package>"])}, ) # A command step runs on a step operator. In a dynamic pipeline, only steps # named in depends_on get their own image built — list it here so ZenML builds # an image with the launcher installed (otherwise the step falls back to the # orchestrator image). dynamic=True also unlocks resource pools. @pipeline(dynamic=True, depends_on=[train]) def training() -> None: train()

这里的depends_on语义可在 src/zenml/pipelines/dynamic/pipeline_definition.py 中得到印证:DynamicPipeline专门接收depends_on参数并做去重校验(重复 step 会抛RuntimeError,提示可用step.with_options(...)传多份配置),且动态 pipeline 的默认执行模式为STOP_ON_FAILURE

选择哪个 launcher
Launcherworker 启动在额外需要适用场景
TorchXdist.ddpKubernetes、Slurm、localK8s 上需要 Volcano 做 gang scheduling在裸 Kubernetes 上用纯torch.distributed做多节点
Rayray job submitRay 集群 / KubeRay一个运行中的 Ray 集群你已经在跑 Ray;Ray Train 会替你搭好torch.distributed
完整示例:TorchX + Volcano on Kubernetes

训练脚本是原生torch.distributed——读取 launcher 注入的 rank/world-size,不含任何 ZenML 或 launcher 特定代码:

# train.py import os import torch import torch.distributed as dist from torch.nn.parallel import DistributedDataParallel as DDP def main() -> None: dist.init_process_group("nccl") # rendezvous via env vars local_rank = int(os.environ["LOCAL_RANK"]) torch.cuda.set_device(local_rank) model = DDP(MyModel().cuda(local_rank), device_ids=[local_rank]) # ... your normal training loop ... dist.destroy_process_group() if __name__ == "__main__": main()

pipeline 把torchx run包进CommandStep。TorchX 的dist.ddpbuiltin 使用 torchelastic,并在 Volcano 上 gang-schedule worker:

# pipeline.py from zenml import CommandStep, pipeline from zenml.config import DockerSettings from zenml.integrations.kubernetes.flavors import KubernetesStepOperatorSettings NNODES, NPROC = 2, 8 # 2 nodes x 8 GPUs train = CommandStep( command=[ "torchx", "run", "-s", "kubernetes", "-cfg", "queue=default", "--wait", "--log", "dist.ddp", "-j", f"{NNODES}x{NPROC}", "--gpu", str(NPROC), "--image", "<registry>/ddp-worker:latest", # the CUDA worker image "--script", "train.py", "--env", "EPOCHS=5", ], step_operator="gmi-k8s", settings={ # ZenML builds a slim launcher image (its base + zenml + torchx). "docker": DockerSettings(requirements=["torchx"]), "step_operator": KubernetesStepOperatorSettings( service_account_name="torchx-launcher" ), }, ) @pipeline(dynamic=True, depends_on=[train]) def training() -> None: train() if __name__ == "__main__": training()

你只需手工构建worker镜像(CUDA + torch + 你的脚本——worker 不是 ZenML step,不需要zenml):

# worker image -> <registry>/ddp-worker:latest FROM pytorch/pytorch:2.4.0-cuda12.1-cudnn9-runtime WORKDIR /app COPY train.py .
同样的模式换成 Ray

只有命令不同。Ray(提交到已有集群,让 Ray Train 拥有torch.distributed):

train = CommandStep( command=["ray", "job", "submit", "--address", "http://ray-head:8265", "--", "python", "train_ray.py"], step_operator="gmi-k8s", settings={"docker": DockerSettings(requirements=["ray[client]"])}, )
关键注意事项
  • 把 command step 写进depends_oncommand step 运行在 step operator 上,而在动态 pipeline 中只有depends_on列出的 step 才会构建专属镜像——否则 step 会回落到 orchestrator 镜像,其中没有安装 launcher。@pipeline(dynamic=True, depends_on=[train])会构建正确的镜像;dynamic=True同时解锁 resource pools。(上面的单机普通 step 不需要depends_on——它们直接用 pipeline 镜像。)
  • 镜像必须携带 launcher(和zenml)。CommandStep在 step operator 上通过 ZenML 的 entrypoint 运行,因此 launcher 镜像需要同时有zenml和 launcher 包。可以像上文一样用requirements=[...]让 ZenML 安装,也可以自己烘焙镜像并用DockerSettings(skip_build=True, parent_image=...)——自定义parent_image必须已含zenml。传给 launcher 的worker镜像(例如--image)只需要你的训练栈,不需要zenml
  • 日志在 launcher 的后端里。command step 的日志不由 ZenML 追踪——worker 日志留在 launcher 放它们的地方(pod 日志、Ray dashboard 等)。参见 command steps 的限制说明。
  • 两个容量管理器、两个 job。launcher 的 gang scheduler(如 Volcano)以 all-or-nothing 方式预留worker容量;ZenML resource pools(如果用了)管的是launcherstep 本身。两者互不重叠。

5. 故障排查与技巧

问题快速修复
GPU 未被使用在容器内验证 CUDA toolkit(nvcc --version),检查驱动兼容性
清缓存后仍 OOM减小 batch size、使用梯度累积,或申请更多显存
Accelerate 挂起确保节点间端口开放;显式传main_process_port

参考资料

  • 本文对应教程原文:docs/book/user-guide/tutorial/distributed-training.md
  • 资源设置实现:src/zenml/config/resource_settings.py
  • Accelerate step 包装实现:src/zenml/integrations/huggingface/steps/accelerate_runner.py
  • CommandStep 实现(不透明命令执行、无 inputs/outputs、无 hook):src/zenml/steps/command_step.py
  • 动态 pipeline 的depends_on校验逻辑:src/zenml/pipelines/dynamic/pipeline_definition.py
  • Command Steps 完整说明与限制:docs/book/how-to/steps-pipelines/command_steps.md
  • 动态 pipeline 指南:docs/book/how-to/steps-pipelines/dynamic_pipelines.md
  • ZenML Pro Resource Pools:docs/book/getting-started/zenml-pro/resource-pools.md

【免费下载链接】zenmlZenML 🙏: One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenml

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

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

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

立即咨询