
人工智能深度学习分布式训练【免费下载链接】accelerate A simple way to launch, train, and use PyTorch models on almost any device and distributed configuration, automatic mixed precision (including fp8), and easy-to-configure FSDP and DeepSpeed support项目地址https://gitcode.com/gh_mirrors/ac/accelerate点击查看免费下载在分布式训练中多个进程在 GPU 上并行执行不同进程的运行速度天然存在差异有些进程先完成有些后完成有些操作如打印日志、上传模型、保存 checkpoint只需要执行一次且不能在其它进程尚未完成前抢先执行。Hugging Face Accelerate 提供了一整套进程执行编排 API帮助你精确控制哪些代码在哪个进程上运行、何时运行确保所有设备上的进程始终保持同步。本文基于 Accelerate 仓库中 docs/source/basic_tutorials/execution.md 教程系统讲解只在单个进程上执行在指定进程上执行与延迟执行全局同步屏障三类能力并结合 src/accelerate/accelerator.py 与 src/accelerate/state.py 的源码实现与仓库内真实示例深入剖析其底层原理。为什么需要显式编排进程执行在分布式训练系统中所有进程都在独立运行同一份脚本但进程之间天然存在执行节奏差异部分代码只需在一台机器节点上执行一次例如打印一条日志、只在本地主进程上展示一个进度条部分代码需要在整个训练集群的所有进程中只执行一次例如把最终模型上传到模型仓库、合并权重或清理缓存部分代码需要保证所有进程都到达某个执行点之后才能继续例如保存模型前必须确认每个进程的训练都已完成。Accelerate 围绕Accelerator对象和底层的PartialState单例状态提供了一套配套工具。核心判断依据是进程索引process_index是全局进程编号跨所有机器local_process_index是当前机器上的本地进程编号每台机器从 0 开始。在 src/accelerate/state.py 中is_main_process与is_local_main_process两个属性的实现为property def is_main_process(self) - bool: Returns whether the current process is the main process return ( self.process_index 0 if self.distributed_type ! DistributedType.MEGATRON_LM else self.is_last_process ) property def is_local_main_process(self) - bool: Returns whether the current process is the main process on the local node return ( self.local_process_index 0 if self.distributed_type ! DistributedType.MEGATRON_LM else self.is_last_process )可以看到在标准分布式环境中主进程即process_index 0全局或local_process_index 0本机而在 Megatron-LM 集成场景下则通过is_last_process判定。这些属性同样作为Accelerator的属性暴露见 src/accelerate/accelerator.py用法完全一致。只在单进程上执行在本地主进程上执行一次is_local_main_processaccelerator.is_local_main_process用于标记只在一台服务器上执行一次的代码典型场景是日志打印、展示单个进度条、设置各库的日志等级等。例如在 examples/by_feature/deepspeed_with_config_support.py 中只在本地主进程上设置详细日志级别其余进程静默if accelerator.is_local_main_process: datasets.utils.logging.set_verbosity_warning() transformers.utils.logging.set_verbosity_info() else: datasets.utils.logging.set_verbosity_error() transformers.utils.logging.set_verbosity_error()训练循环中最常见的用法是控制进度条只显示一次文档原例from tqdm.auto import tqdm progress_bar tqdm(range(args.max_train_steps), disablenot accelerator.is_local_main_process)也可以把任意语句包在if accelerator.is_local_main_process:中if accelerator.is_local_main_process: print(Accelerate is the best)[!TIP] 对于没有包在accelerator.is_local_main_process中的零散print语句可以直接把print替换为 Accelerate 的Accelerator.print方法它内部通过if self.is_local_main_process: print(*args, **kwargs)实现每台服务器只打印一次。用装饰器限定本地主进程函数on_local_main_process如果是一段每台服务器只应执行一次的完整函数使用accelerator.on_local_main_process装饰器accelerator.on_local_main_process def do_my_thing(): Something done once per server do_thing_once_per_server()在所有进程中只执行一次is_main_process当需要跨所有机器、所有进程只执行一次时典型场景是上传最终模型到模型仓库、创建输出目录使用accelerator.is_main_processif accelerator.is_main_process: repo.push_to_hub()仓库中的真实示例同样遵循这一模式。examples/by_feature/deepspeed_with_config_support.py 在accelerator.is_main_process中创建 Hub 仓库或创建输出目录随后立即调用accelerator.wait_for_everyone()让其它进程等主进程完成后再继续读取数据集if accelerator.is_main_process: if args.push_to_hub: api HfApi(tokenargs.hub_token) ... elif args.output_dir is not None: os.makedirs(args.output_dir, exist_okTrue) accelerator.wait_for_everyone()对应地用accelerator.on_main_process装饰器限定整个函数accelerator.on_main_process def do_my_thing(): Something done once per server do_thing_once()两种主进程的差异对比判断方式判定条件执行次数适用场景is_local_main_process/on_local_main_process每台机器上local_process_index 0的进程每台机器各执行一次N 台机器共执行 N 次日志、进度条、本机数据准备is_main_process/on_main_process全局process_index 0的进程整个集群只执行一次上传模型、创建共享目录、合并权重装饰器的底层实现原理从源码看Accelerator.on_main_process/on_local_main_process都只是把实际逻辑转发给PartialState的对应方法见 src/accelerate/accelerator.py。PartialState的核心实现位于 src/accelerate/state.pydef on_main_process(self, function: Callable[..., Any] | None None): ... if self.is_main_process or not self.use_distributed: return function return do_nothing def on_local_main_process(self, function: Callable[..., Any] | None None): ... if self.is_local_main_process or not self.use_distributed: return function return do_nothing即当条件满足时返回原函数照常执行条件不满足时返回模块级的空操作函数do_nothing(*args, **kwargs)定义于 src/accelerate/state.py。注意两点实现细节当没有使用分布式单卡 / CPU时use_distributed为假函数会无条件执行因此这些装饰器在单进程环境下退化为普通函数无需改动代码即可兼容do_nothing会吞掉返回值返回None若被装饰函数有返回值且依赖该返回值需改用is_main_process判断语句而非装饰器。在指定进程上执行除了主进程Accelerate 还支持把函数精确调度到任意进程索引上执行。在指定全局进程上执行on_process使用accelerator.on_process(process_index...)并指定全局进程索引accelerator.on_process(process_index0) def do_my_thing(): Something done on process index 0 do_thing_on_index_zero()在指定本地进程上执行on_local_process使用accelerator.on_local_process(local_process_index...)并指定本地进程索引函数会在每台机器的对应本地进程上各执行一次accelerator.on_local_process(local_process_idx0) def do_my_thing(): Something done on process index 0 on each server do_thing_on_index_zero_on_each_server()从源码实现看Accelerator.on_process与Accelerator.on_local_process支持两种调用形式直接accelerator.on_process(process_index2)传参或accelerator.on_process无参形式配合partial机制延迟绑定索引。内部的判定同样委托给PartialState.on_process/PartialState.on_local_process通过比较当前process_index或local_process_index与目标索引决定是否替换为do_nothing。例如假设 4 个进程accelerator.on_process(process_index2)装饰的函数只会在全局编号为 2 的进程上输出 Printed on process 2而假设 2 台服务器各 4 个进程accelerator.on_local_process(local_process_index2)会在每台服务器的本地 2 号进程上各执行一次。此外仓库还提供同族的扩展能力可一并了解见 src/accelerate/accelerator.pyaccelerator.on_last_process只在最后一个进程process_index num_processes - 1上执行with accelerator.main_process_first():与with accelerator.local_main_process_first():上下文管理器让主进程先进入代码块其它进程在主进程退出后再进入见 src/accelerate/state.py。延迟执行全局同步屏障wait_for_everyone何时需要使用同步屏障在多个 GPU 同时运行同一脚本时有些进程会比其他进程更快地执行到后续代码。如果在某些进程还在训练时就去保存模型会得到不完整、不可用的 checkpoint。因此保存模型 / 做最终汇总这类操作前必须保证所有进程都已到达同一执行点。使用方法在代码中插入accelerator.wait_for_everyone()accelerator.wait_for_everyone()该调用会阻塞所有先到达的进程直到剩余进程全部抵达同一位置在单 GPU 或 CPU 单进程环境下该调用不产生任何效果因此可以无副作用地保留在代码中。一个直观示例来自 src/accelerate/accelerator.py 的 docstring# Assuming two GPU processes import time from accelerate import Accelerator accelerator Accelerator() if accelerator.is_main_process: time.sleep(2) else: print(Im waiting for the main process to finish its sleep...) accelerator.wait_for_everyone() # Should print on every process at the same time print(Everyone is here)运行时会看到非主进程先打印等待主进程完成睡眠随后wait_for_everyone()使两个进程对齐最后两个进程几乎同时打印 Everyone is here。底层实现barrier 与 rendezvousAccelerator.wait_for_everyone()src/accelerate/accelerator.py转发给PartialState().wait_for_everyone()后者根据分布式类型选择不同的底层同步原语src/accelerate/state.pyif self.distributed_type in ( DistributedType.MULTI_GPU, DistributedType.MULTI_MLU, DistributedType.MULTI_SDAA, DistributedType.MULTI_MUSA, DistributedType.MULTI_NPU, DistributedType.MULTI_XPU, DistributedType.MULTI_CPU, DistributedType.MULTI_HPU, DistributedType.MULTI_NEURON, DistributedType.DEEPSPEED, DistributedType.FSDP, ): torch.distributed.barrier(device_ids[self.local_process_index]) elif self.distributed_type DistributedType.XLA: xm.rendezvous(accelerate.utils.wait_for_everyone)也就是说对绝大多数分布式后端多 GPU / 多 CPU / DeepSpeed / FSDP 等底层是torch.distributed.barrier集合通信屏障并以本地进程索引作为设备号对 TPUXLA 分布式类型则调用xm.rendezvous(accelerate.utils.wait_for_everyone)实现进程汇合。⚠️ 使用警告必须确保所有进程最终都会到达wait_for_everyone()这一行。若某个进程因为条件分支提前退出或走不到该调用已到达的进程将永久挂起等待。因此不要把wait_for_everyone()放进只由if accelerator.is_main_process:保护的代码块内部而应放在其外部、所有进程都会执行到的位置。仓库内的典型应用模式wait_for_everyone()在仓库中被广泛用于训练与保存的关键节点以下模式值得直接借鉴保存 checkpoint 前先wait_for_everyone()确保所有进程训练完成再在is_main_process中保存模型典型写法见 examples/by_feature/deepspeed_with_config_support.py创建共享资源后主进程创建输出目录 / 上传仓库后立即wait_for_everyone()防止其它进程提前读取尚未就绪的资源同上文件 L320多轮保存 / 上传混合场景在 examples/by_feature/megatron_lm_gpt_pretraining.py 中可以看到保存 checkpoint 前等待 → 主进程保存 → 再次等待 → 主进程上传模型的完整节奏编排分布式推理examples/inference/distributed/distributed_image_generation.py 在distributed_state.is_main_process主进程分发输入后调用wait_for_everyone()保证各进程同步推进。此外wait_for_everyone也被封装进main_process_first()/local_main_process_first()上下文管理器_goes_first实现见 src/accelerate/state.py非主进程在进入代码块前先等待主进程退出后再等待一次实现主进程优先执行、其余进程随后跟上的语义Accelerator.save_state与profiler上下文等内部实现也会调用wait_for_everyone()来保证 checkpoint 与性能剖析输出的完整性。小结与最佳实践进程执行编排是分布式训练健壮性的基石。Accelerate 用四组 API 覆盖了全部场景本机一次accelerator.is_local_main_process语句/accelerator.on_local_main_process函数/accelerator.print单条打印全局一次accelerator.is_main_process语句/accelerator.on_main_process函数指定进程accelerator.on_process(process_index...)全局索引/accelerator.on_local_process(local_process_index...)本地索引全局同步accelerator.wait_for_everyone()底层为torch.distributed.barrierTPU 上为xm.rendezvous单进程环境自动退化为空操作。实战建议总结打印日志、进度条统一使用is_local_main_process或Accelerator.print避免 N 份重复输出上传模型、创建共享目录、保存最终权重使用is_main_process保存模型前一定先调用wait_for_everyone()且该调用必须放在所有进程都会经过的路径上防止进程挂起熟悉Accelerator.save_state与accelerate.utils.ProjectConfiguration等配套能力详见 docs/source/package_reference/accelerator.md它们内部已封装好同步逻辑需要轻量级进程控制时可以直接使用PartialState不经过Accelerator它在 src/accelerate/state.py 中作为单例维护全部进程状态并通过共享状态字典在多次实例化间保持一致。掌握这些 API你就能在多卡、多机、TPU、DeepSpeed、FSDP 等任意分布式配置下精确控制进程节奏写出不会因为竞态而崩溃的分布式训练与推理代码。赞分享人工智能深度学习分布式训练【免费下载链接】accelerate A simple way to launch, train, and use PyTorch models on almost any device and distributed configuration, automatic mixed precision (including fp8), and easy-to-configure FSDP and DeepSpeed support项目地址https://gitcode.com/gh_mirrors/ac/accelerate点击查看免费下载相关推荐n8n执行策略并行与串行执行控制n8n执行策略并行与串行执行控制 在现代工作流自动化平台中执行策略的选择直接影响系统性能和任务处理效率。n8n作为一款结合代码灵活性和无代码高效性的工作流自工作流自动化人工智能AI Agent后端前端低代码Ray Actor 任务执行顺序全解析同步单线程、Async/Threaded Actor 与乱序执行机制Ray Actor 任务执行顺序全解析同步单线程、Async/Threaded Actor 与乱序执行机制 导读 在 Ray 分布式运行时中Actor 是承人工智能分布式训练强化学习任务调度模型推理服务后端Playwright CLI让智能助手成为你的浏览器操作专家Playwright CLI让智能助手成为你的浏览器操作专家 凌晨三点我还在为那个该死的电商测试脚本头疼。页面元素定位器像捉迷藏一样难以捉摸异步等待让代码AI 技能浏览器控制GUI 自动化测试上一篇联发科设备逆向工程终极指南MTKClient深度解析与实战应用下一篇WVP-GB28181-Pro把摄像头接进国标28181体系的开源视频平台创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考