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_count | Optional[PositiveFloat] | 期望分配的 CPU 核数,可为小数;调度时内部换算为mcpu(毫核) |
gpu_count | Optional[NonNegativeInt] | GPU 数量;0表示不申请 GPU(会显式移除资源请求中的gpu键) |
memory | 字符串,匹配正则^[0-9]+(KB\|KIB\|MB\|MIB\|GB\|GIB\|TB\|TIB\|PB\|PIB)$ | 内存大小必须带单位后缀,例如"16GB"、"8GiB";二进制单位(KiB/MiB…)按 2 的幂换算 |
pool_resources | Optional[Dict[str, PositiveInt]] | 用于 ZenML 资源池的自定义资源键(如tpu、vcpus);当gpu_count/cpu_count/memory同时设置时,后者优先覆盖gpu/mcpu/memory_mb |
preemptible | bool,默认True | 资源是否可被抢占,仅在使用 ZenML 资源池时生效 |
从源码结构看,ResourceSettings通过merged_requested_resources()方法把上述字段归一化为调度器可理解的资源请求映射:gpu_count→gpu、cpu_count→mcpu(ceil(cpu * 1000))、memory→memory_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 集成),或用
CommandStep跑torchrun。无需外部启动器或 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:跑torchrun的CommandStep
如果你已经习惯用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 协调——分配RANK、WORLD_SIZE、LOCAL_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
| Launcher | worker 启动在 | 额外需要 | 适用场景 |
|---|---|---|---|
TorchX(dist.ddp) | Kubernetes、Slurm、local | K8s 上需要 Volcano 做 gang scheduling | 在裸 Kubernetes 上用纯torch.distributed做多节点 |
Ray(ray job submit) | Ray 集群 / 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_on。command 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),仅供参考