资讯动态

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

发布时间:2026/9/18 6:11:04 来源:尧图企业网站定制
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_count8, gpu_count2, memory16GB) }) 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 资源池的自定义资源键如tpu、vcpus当gpu_count/cpu_count/memory同时设置时后者优先覆盖gpu/mcpu/memory_mbpreemptiblebool默认True资源是否可被抢占仅在使用 ZenML 资源池时生效从源码结构看ResourceSettings通过merged_requested_resources()方法把上述字段归一化为调度器可理解的资源请求映射gpu_count→gpu、cpu_count→mcpuceil(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_imagepytorch/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_processes4, multi_gpuTrue) step def training_step(...): ... # your distributed training code pipeline def dist_pipeline(...): training_step(...)常用参数num_processes要启动的进程总数每个 GPU 一个multi_gpuTrue启用多 GPU 模式cpuTrue强制使用 CPU 训练mixed_precisionfp16/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, flavoraccelerate)包装成一个 CLI 脚本把step调用时传入的关键字参数转换为--arg value命令行参数再用 Accelerate 的launch_command_parser()解析、合并装饰器参数后调用launch_command(args)step 返回值用 cloudpickle 序列化到output_path训练结束后由包装层pickle.load还原从而保证 step 的返回值正常回流给 ZenML。准备容器使用与第 2 节相同的 CUDA 镜像并在 requirements 中追加 AccelerateDockerSettings( parent_imagepytorch/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包进 CommandStepZenML 负责管理 runlauncher 负责 worker 进程。4.1 单机多卡的三种模式所有进程都生活在 step 自己的容器里因此不需要外部启动器或 gang scheduler。以下三种模式当前都可用选与你现有训练习惯最匹配的一种。模式 1自己拉起进程的原生 stepZenML 开箱支持多进程训练——用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, rankrank, world_sizeworld_size) # ... your DDP training ... dist.destroy_process_group() step( runtimeisolated, # run in its own container with the GPUs, not inline settings{resources: ResourceSettings(gpu_count4)}, ) def train() - None: mp.spawn(_worker, args(4,), nprocs4) # one process per GPU pipeline(dynamicTrue) def training() - None: train()runtimeisolated让 orchestrator 把 step 放进一个按 GPU 资源请求规格配置的新容器中运行而不是内联在编排进程里——这正是重型训练 step 需要的运行方式。由于 isolated 容器使用 pipeline 镜像请确保该镜像已包含 CUDA 和 torch见第 2 节。isolated 运行是动态 pipeline 特性因此 pipeline 需要声明dynamicTrue。模式 2Accelerate 集成如果不想自己搭进程组就用run_with_accelerate见第 3 节处理 fan-outfrom zenml import step from zenml.config import ResourceSettings from zenml.integrations.huggingface.steps import run_with_accelerate run_with_accelerate(num_processes4, multi_gpuTrue) step(settings{resources: ResourceSettings(gpu_count4)}) def train(...): ... # your training code模式 3跑torchrun的CommandStep如果你已经习惯用torchrun或任意 launcher CLI驱动训练直接把它包进 command step。launcher 和它的工作进程共享同一个容器因此一个镜像同时携带zenmltorchtrain.py即可from zenml import CommandStep, pipeline from zenml.config import DockerSettings, ResourceSettings train CommandStep( command[torchrun, --standalone, --nproc-per-node4, train.py], step_operatorgmi-k8s, settings{ docker: DockerSettings( parent_imagepytorch/pytorch:2.4.0-cuda12.1-cudnn9-runtime, requirements[zenml], ), resources: ResourceSettings(gpu_count4), }, ) pipeline(dynamicTrue, 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, checkTrue)——命令退出码 0 即 step 成功非零即失败构造函数会强制 command 非空并拒绝声明 inputs/outputs 或配置 hook 的用法同时默认enable_cacheFalse。4.2 多节点把 launcher 包进CommandStep一旦训练跨越多台机器最干净的做法是让专门的 launcher 拥有 worker gang让 ZenML 拥有 run。为什么由 launcher而不是 ZenML启动 workertorch.distributed/torchrun只做 rank 之间的rendezvous 协调——分配RANK、WORLD_SIZE、LOCAL_RANK并把进程连起来但不负责在别的节点上开机器或起进程。必须有一个能供给并 gang-schedule 这 N 个 worker 进程的 launcherTorchX、Ray、Slurm 等。ZenML 不重写这部分而是给你一个清晰的分工缝CommandStep在 step operator 的容器里执行一条不透明命令把这条命令指向 launcher由 launcher 负责调度并启动 worker gangZenML 记录 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_operatoryour-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). dynamicTrue also unlocks resource pools. pipeline(dynamicTrue, 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。选择哪个 launcherLauncherworker 启动在额外需要适用场景TorchXdist.ddpKubernetes、Slurm、localK8s 上需要 Volcano 做 gang scheduling在裸 Kubernetes 上用纯torch.distributed做多节点Rayray job submitRay 集群 / KubeRay一个运行中的 Ray 集群你已经在跑 RayRay 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, queuedefault, --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, EPOCHS5, ], step_operatorgmi-k8s, settings{ # ZenML builds a slim launcher image (its base zenml torchx). docker: DockerSettings(requirements[torchx]), step_operator: KubernetesStepOperatorSettings( service_account_nametorchx-launcher ), }, ) pipeline(dynamicTrue, 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.distributedtrain CommandStep( command[ray, job, submit, --address, http://ray-head:8265, --, python, train_ray.py], step_operatorgmi-k8s, settings{docker: DockerSettings(requirements[ray[client]])}, )关键注意事项把 command step 写进depends_on。command step 运行在 step operator 上而在动态 pipeline 中只有depends_on列出的 step 才会构建专属镜像——否则 step 会回落到 orchestrator 镜像其中没有安装 launcher。pipeline(dynamicTrue, depends_on[train])会构建正确的镜像dynamicTrue同时解锁 resource pools。上面的单机普通 step 不需要depends_on——它们直接用 pipeline 镜像。镜像必须携带 launcher和zenml。CommandStep在 step operator 上通过 ZenML 的 entrypoint 运行因此 launcher 镜像需要同时有zenml和 launcher 包。可以像上文一样用requirements[...]让 ZenML 安装也可以自己烘焙镜像并用DockerSettings(skip_buildTrue, 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 toolkitnvcc --version检查驱动兼容性清缓存后仍 OOM减小 batch size、使用梯度累积或申请更多显存Accelerate 挂起确保节点间端口开放显式传main_process_port参考资料本文对应教程原文docs/book/user-guide/tutorial/distributed-training.md资源设置实现src/zenml/config/resource_settings.pyAccelerate step 包装实现src/zenml/integrations/huggingface/steps/accelerate_runner.pyCommandStep 实现不透明命令执行、无 inputs/outputs、无 hooksrc/zenml/steps/command_step.py动态 pipeline 的depends_on校验逻辑src/zenml/pipelines/dynamic/pipeline_definition.pyCommand Steps 完整说明与限制docs/book/how-to/steps-pipelines/command_steps.md动态 pipeline 指南docs/book/how-to/steps-pipelines/dynamic_pipelines.mdZenML Pro Resource Poolsdocs/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),仅供参考

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

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

免费获取报价