
PyTorch Elastic 多进程管理使用torch.distributed.elastic.multiprocessing启动与编排多 Worker【免费下载链接】pytorchTensors and Dynamic neural networks in Python with strong GPU acceleration项目地址: https://gitcode.com/GitHub_Trending/py/pytorchtorch.distributed.elastic.multiprocessingTorchElastic 的 multiprocessing 子模块是一个用于启动并管理n份 worker 子进程的库既可以按「函数」Callable方式用torch.multiprocessingspawn/fork 出进程也可以按「二进制可执行程序」如/usr/bin/echo、python train.py方式用subprocess.Popen拉起子进程。它是torchrun/TorchElastic Agent 在单机上真正拉起训练 worker 的底层实现也是自定义「一个入口点、多个副本」并行任务时的标准工具。读完本文你将掌握start_processes的完整参数语义、PContext 进程上下文家族的使用方法、标准输出/错误stdout/stderr的 redirect/tee 与日志目录结构、进程返回值与失败错误文件error.json的传播机制并理解函数型与二进制型两种启动方式的差异与适用场景。本文以 docs/source/elastic/multiprocessing.md 为骨架展开其全部 API 引用均落在 torch/distributed/elastic/multiprocessing 模块内入口模块init.py核心实现 api.py。一、模块定位两种入口点、一套统一上下文模块 docstring见init.py明确了它的职责函数入口点Callable底层使用torch.multiprocessing即 Pythonmultiprocessing通过spawn/fork/forkserver方式启动进程返回MultiprocessContext二进制入口点str底层使用subprocess.Popen创建子进程返回SubprocessContext两者都是父类PContext的具体实现。这与torch.multiprocessing.start_processes的约定保持一致start_processes的返回值是一个进程上下文PContext。无论用哪种方式启动调用方都通过同一套 PContext 接口start()、wait()、close()、pids()管理全部 worker从而在 TorchElastic 的 agent 代码中做到「调度逻辑与启动机制解耦」。从源码结构看整个模块还包含四个支撑文件分别承担不同职责文件职责api.pyPContext/MultiprocessContext/SubprocessContext、RunProcsResult、LogsSpecs/DefaultLogsSpecs/LogsDest、Std等核心类errors/基于文件的跨进程错误传播record装饰器、ProcessFailure、error.json写入redirects.py在文件描述符层面对 stdout/stderr 的重定向dup2实现tail_log.pyTailLog模拟 unixtail -f将日志文件内容实时输出到控制台实现 teesubprocess_handler/二进制 worker 的子进程处理器进程信号发送、关闭等二、start_processes启动多个 Worker 的唯一入口start_processes是模块对外暴露的唯一「启动入口」定义在init.pydef start_processes( name: str, entrypoint: Callable | str, args: dict[int, tuple], envs: dict[int, dict[str, str]], logs_specs: LogsSpecs, log_line_prefixes: dict[int, str] | None None, start_method: str spawn, numa_options: NumaOptions | None None, duplicate_stdout_filters: list[str] | None None, duplicate_stderr_filters: list[str] | None None, ) - PContext:2.1 参数语义核心参数的完整语义来自函数 docstringinit.pyname人类可读的短名称描述这些进程的用途作为 tee 输出到控制台时的行头[{name}{local_rank}]:。entrypoint要么是Callable函数要么是str命令/二进制路径。进程副本数量由args的条目数决定。args/envs以 local rank副本序号为 key、分别映射每个副本的入参与环境变量。两者的 key 集合必须一致且完整覆盖{0, 1, ..., nprocs-1}即「所有 local rank 都必须被覆盖」。logs_specsLogsSpecs实例统一定义log_dir、redirects、tee详见第五节。log_line_prefixes为每个 local rank 指定日志行的自定义前缀agent 常用它注入role/rank/local_rank/hostname模板见 local_elastic_agent.py。start_methodmultiprocessing 启动方式spawn/fork/forkserver对二进制入口点无效忽略。numa_optionsNumaOptions用于给 worker 绑定 NUMA 节点传入后函数入口点会被_maybe_wrap_with_numa_binding包装见 api.py。duplicate_stdout_filters/duplicate_stderr_filters非空时把被tee选中的 rank 的 stdout/stderr 中命中任意一个过滤子串的行额外聚合写到filtered_stdout.log/filtered_stderr.log。2.2 用法一以函数方式启动两个 trainerdocstring 给出了可以直接复制的经典用法init.pyfrom torch.distributed.elastic.multiprocessing import Std, start_processes def trainer(a, b, c): pass # train # 运行两个 trainer # LOCAL_RANK0 - trainer(1,2,3) # LOCAL_RANK1 - trainer(4,5,6) ctx start_processes( nametrainer, entrypointtrainer, args{0: (1, 2, 3), 1: (4, 5, 6)}, envs{0: {LOCAL_RANK: 0}, 1: {LOCAL_RANK: 1}}, logs_specsDefaultLogsSpecs(log_dir/tmp/foobar, redirectsStd.ALL), tee{0: Std.ERR}, # 仅把 local rank 0 的 stderr 同步打印到控制台 ) # 等待所有 trainer 结束 ctx.wait()注意示例中的redirects/tee需通过logs_specsDefaultLogsSpecs传入tee语义与 unixtee命令一致——既写入日志文件、又打印到控制台而redirects则只写文件、不打印。2.3 用法二以二进制方式启动 echo当entrypoint为字符串时等价于直接执行该命令init.pyctx start_processes( nameecho, entrypointecho, # 二进制 / 命令行 args{0: hello, 1: world}, redirects{1: Std.OUT}, logs_specsDefaultLogsSpecs(log_dir/tmp/foobar), ) # 实际执行echo helloecho world /tmp/foobar/1/stdout.log2.4 三个容易踩的约束从 docstring 与实现init.py可以提炼出以下硬性规则二进制入口点时args只能传字符串若传其他类型如1,2,3、[1,2,3]会被强制str()化后再执行。例如args{0: (1, 2, 3), 1: ([1, 2, 3],)}实际运行的是echo 1 2 3与echo [1, 2, 3]。args与envs的 key 集合必须都恰好等于{0,...,nprocs-1}。启动前会调用_validate_full_rankapi.py校验key 不齐会抛出RuntimeError: local rank mapping mismatch。例如envs{0:{}}而args有两个条目就是非法的。函数入口点失败时错误文件默认会被记录无需额外标注二进制入口点失败要写出error.json则入口点主函数必须显式标注torch.distributed.elastic.multiprocessing.errors.record。三、Std位掩码redirects 与 tee 的表达方式Std定义在 api.py本质是一个IntFlag枚举class Std(IntFlag): NONE 0 OUT 1 ERR 2 ALL OUT | ERR # 3它既可以全局作用于所有 rank传单个Std值也可以按 rank 选择性作用传dict[int, Std]未出现的 rank 默认Std.NONE。redirects与tee在LogsSpecs/start_processes的语义如下init.pyredirects把指定的 std 流重定向写入log_dir下的日志文件tee把指定的 std 流同时写日志文件并打印到控制台redirect print若不想让 worker 输出刷屏应使用redirects而非tee。3.1Std.from_str与to_mapStd.from_strapi.py支持两种字符串输入Agent 常把配置从命令行以字符串形式传下来Std.from_str(0) # - Std.NONE Std.from_str(1) # - Std.OUT Std.from_str(0:3,1:0,2:1) # - {0: Std.ALL, 1: Std.NONE, 2: Std.OUT}to_mapapi.py把「单个值或局部映射」统一扩展成覆盖所有 local rank 的映射to_map(Std.OUT, local_world_size2) # {0: Std.OUT, 1: Std.OUT} to_map({1: Std.OUT}, local_world_size2) # {0: Std.NONE, 1: Std.OUT}四、Process ContextPContext家族文档列出的四类进程上下文类集中在 api.py 中。PContext是抽象基类api.py名字有意与torch.multiprocessing.ProcessContext区分。4.1PContext统一生命周期接口start()api.py在主线程调用时会按环境变量TORCHELASTIC_SIGNALS_TO_HANDLE默认SIGTERM,SIGINT,SIGHUP,SIGQUIT注册信号处理函数_terminate_process_handler——收到终止信号会抛出SignalExceptionapi.py该异常不应被吞掉否则进程永不退出。随后调用_start()真正拉起 worker并启动TailLog线程。wait(timeout-1, period1)api.py每隔period秒轮询一次等待所有进程结束。timeout0等价于一次 poll非阻塞查询timeout0表示无限等待超时返回None。close(death_sigNone, timeout30)api.py用death_sigunix 默认SIGTERM、Windows 默认CTRL_C_EVENT终止所有进程并清理资源超时后升级为强杀信号unixSIGKILL、WindowsCTRL_C_EVENT见_get_kill_signal。实现中还带SIGKILL后的有界 join避免 worker 卡在不可中断内核态如 NCCL/GPU 集合通信挂死时把 agent 自身拖死。pids()返回{local_rank: pid}映射。close/wait与信号配合的推荐写法docstring 内嵌示例api.pypc start_processes(...) try: pc.wait(1) # ... 做其他工作 except SignalException as e: pc.shutdown(e.sigval, timeout30) # 优雅退出超时再强杀4.2MultiprocessContext函数型 workerMultiprocessContextapi.py用于entrypoint为函数的情形构造时按start_method为每个 local rank 建立mp.SimpleQueue_ret_vals用于回传返回值_start()内部调用mp.start_processes(fn_wrap, ..., joinFalse, daemonFalse, start_method...)api.py。_wrap先注入该 rank 的环境变量再进入 stdout/stderr 重定向上下文并用record(fn)(*args_)包装执行最后把返回值put进队列api.py_poll()依赖「所有进程全部结束」与「任意进程失败」两种终态通过ProcessContext.join 轮询返回队列避免大返回值导致管道死锁api.py。4.3SubprocessContext二进制型 workerSubprocessContextapi.py用于entrypoint为命令字符串的情形_start()为每个 local rank 通过get_subprocess_handler创建一个SubprocessHandlerPopen封装_poll()逐个proc.poll()采集退出码任意副本失败或全部结束后即收敛结果all-or-nothing 策略并对仍在运行的副本执行close()二进制 worker没有返回值因此成功时RunProcsResult.return_values会被填充为None占位以与MultiprocessContext保持一致的 API 形态api.py。4.4RunProcsResult运行结果容器RunProcsResult是一个 dataclassapi.py由PContext的轮询/等待返回注意以下几点约束字段说明return_values: dict[int, Any]各 rank 的返回值仅函数型启动时填充按 local rank 索引failures: dict[int, ProcessFailure]各 rank 的失败信息ProcessFailure包含 local_rank、pid、exitcode、error_file、message、timestamp 等stdouts/stderrs各 rank 的 stdout.log / stderr.log 路径未重定向则为空串is_failed()len(failures) 0即为失败五、LogsSpecs、DefaultLogsSpecs与LogsDest日志目录规划文档列出的日志三件套同样定义在 api.py。5.1LogsSpecs抽象基类LogsSpecsapi.py定义日志处理与重定向的抽象协议构造参数log_dir日志根目录、redirects、tee均可传单个Std或{local_rank: Std}映射抽象方法reify(envs) - LogsDest根据各 rank 的环境变量为每个 rank 计算出日志文件目标路径抽象属性root_log_dir。5.2DefaultLogsSpecs默认实现DefaultLogsSpecsapi.py是现成可用、也是 agent 默认采用的实现行为如下log_dir不存在则自动创建为文件时抛NotADirectoryError未指定时自动tempfile.mkdtemp(prefixtorchelastic_)日志目录按「运行 尝试 rank」分层组织reify()依据环境变量TORCHELASTIC_RUN_ID默认test_run_id与TORCHELASTIC_RESTART_COUNT默认0生成目录并在每次 restart 前用shutil.rmtree清理旧 attempt 目录log_dir/rdzv_run_id/attempt_attempt/ ├── rank/stdout.log # redirects OUT 时 ├── rank/stderr.log # redirects ERR 时 ├── rank/error.json # 失败时记录错误 ├── filtered_stdout.log # duplicate_stdout_filters 非空时 └── filtered_stderr.log # duplicate_stderr_filters 非空时实现要点api.pytee 的实现方式先把 tee 目标并入 redirects先落盘再靠TailLog把文件内容同步输出到控制台——因此LogsSpecs内部要求 stdouts/stderrs 一定是 tee_stdouts/tee_stderrs 的超集每个 rank 的error.json路径会以环境变量TORCHELASTIC_ERROR_FILE注入 worker若配置了local_ranks_filter不在集合内的 rank 即使配置了 tee 也只落盘不 tail、未重定向的流则导向os.devnull从而实现「只打印选定 rank 的日志到控制台」。5.3LogsDestLogsDestapi.py是reify()的返回值 dataclass按日志类型保存{local_rank: 文件路径}的映射stdouts/stderrs重定向的日志文件路径未重定向为空串tee_stdouts/tee_stderrs被 tail 输出的文件error_files{local_rank: error.json}filtered_stdout/filtered_stderr按过滤子串聚合的日志文件路径。六、Tee 的底层实现TailLog与文件描述符级重定向6.1TailLog无文件等待式 tailTailLogtail_log.py为每个日志文件起一个线程模拟tail -f日志文件尚未创建时也会优雅等待轮询interval_sec默认 0.1 秒因此可以在 worker 真正写出日志前就启动默认输出行头为[{name}{local_rank}]:可用log_line_prefixes覆盖tail_log.pystop()通过 Event 通知各 tail 线程退出并回收线程池由于跨文件缓冲日志行不保证严格按墙钟顺序打印——官方建议业务日志自带时间戳。6.2redirects.py重定向的是文件描述符不只是sys.stdoutredirect_stdout/redirect_stderrredirects.py基于 POSIXdup2实现因此同时覆盖 Python 层与 C 层输出with redirect_stdout(/tmp/stdout.log): print(python stdouts are redirected) libc ctypes.CDLL(libc.so.6) libc.printf(bc stdouts are also redirected) # 也会被重定向 os.system(echo system stdouts are also redirected) # 也会被重定向 print(stdout restored)平台限制源码中可查macOS 目前不支持该重定向机制加载 libc 时直接告警并返回None重定向退化为nullcontextWindows 则需要通过_dup2/SetStdHandle/CRTFILE*多层重定向redirects.py。因此get_std_cmapi.py在IS_WINDOWS or IS_MACOS时返回空上下文仅 unix 平台真正执行 fd 级重定向。七、跨进程错误传播error.json、record与ProcessFailure多进程场景下异常发生在 worker 进程agent 无法用 try-catch 直接捕获。TorchElastic 采用基于文件的跨进程错误传播详见 errors/init.py任何被record装饰器包裹的入口点在发生未捕获异常时会连同完整 traceback 写入环境变量TORCHELASTIC_ERROR_FILE指向的 JSON 文件父进程agent为每个子进程设置该环境变量并在子进程失败后聚合各error.json父进程选择timestamp 最小最先发生的错误作为 root cause 向上传播。ErrorHandler.record_exception写入的 JSON 结构见 error_handler.py核心字段为嵌套的message含py_callstack与timestamp。ProcessFailure会尝试解析该文件来构造message/timestamp若error.json不存在如被信号杀死则会据退出码生成信息如Signal -15 (SIGTERM) received by PID ...见 errors/init.py。ChildFailedErrorerrors/init.py允许record包裹的父函数把子进程根异常原样上抛不包裹父 traceback适合「父进程只是简单保姆进程、真正的计算都在子进程」的场景。八、与 TorchElastic Agent 的集成一切从_start_workers开始该模块并非孤立 API它是 TorchElasticLocalElasticAgent拉起 worker 的直接底层。在 local_elastic_agent.py 的_start_workers中agent 逐 rank 组装好envs注入LOCAL_RANK、TORCHELASTIC_ERROR_FILE等与args含宏替换macros.substitute后一次性调用self._pcontext start_processes( namespec.role, entrypointspec.entrypoint, argsargs, envsenvs, logs_specsself._logs_specs, log_line_prefixeslog_line_prefixes, start_methodself._start_method, numa_optionsspec.numa_options, duplicate_stdout_filtersspec.duplicate_stdout_filters, duplicate_stderr_filtersspec.duplicate_stderr_filters, )由此可见torchrun/agent 对日志目录、redirect/tee、错误文件、NUMA 绑定、日志前缀等一系列编排能力最终都会收敛到start_processes与LogsSpecs这套 API 上。需要自定义训练启动器或调试 worker 拉起/退出/日志行为时这套接口就是直接的切入点。九、实战小结与最佳实践进程数量由args决定且args与envs的 key 必须恰好是{0,...,nprocs-1}——参数不齐会在进程启动前快速失败想让 worker 输出进文件用redirects想既进文件又打印到控制台用tee只想打印个别 rank可用{rank: Std}映射或local_ranks_filter函数型入口点自动获得返回值收集经mp.SimpleQueue回传与record错误记录二进制型入口点没有返回值RunProcsResult.return_values为None占位错误文件需入口程序自身标注record等待结果优先使用ctx.wait(timeout)返回值RunProcsResult注意轮询/等待是 all-or-nothing 语义全部成功或任一失败即收敛优雅退出务必处理主进程收到的SIGTERM/SIGINT等信号SignalException随后调用close(death_sig, timeout)超时后框架会自动升级为SIGKILL日志目录建议直接使用DefaultLogsSpecs它会自动创建目录并按rdzv_run_id/attempt_n/rank/组织便于失败定位与多轮重试排障。围绕本文涉及的 API 与语义还可进一步查阅 docs/source/elastic/multiprocessing.md 及模块内 api.py、tail_log.py、redirects.py 与 errors/ 的源码注释与 docstring以获得逐类、逐参数的权威细节。【免费下载链接】pytorchTensors and Dynamic neural networks in Python with strong GPU acceleration项目地址: https://gitcode.com/GitHub_Trending/py/pytorch创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考