多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

PaddleNLP 分布式工具集 paddlenlp.ops.distributed.utils 深度解析:RNG 状态追踪与并行拓扑建模

PaddleNLP 分布式工具集 paddlenlp.ops.distributed.utils 深度解析:RNG 状态追踪与并行拓扑建模 人工智能大模型预训练微调LoRARLHF强化学习分布式训练【免费下载链接】PaddleNLPEasy-to-use and powerful LLM and SLM library with awesome model zoo.项目地址https://gitcode.com/gh_mirrors/pa/PaddleNLP点击查看免费下载导读paddlenlp.ops.distributed.utils是 PaddleNLP 中面向大规模分布式训练的底层工具子包它由两个核心模块构成random.pyRNG 状态追踪器与topo.py并行拓扑建模。前者负责在模型并行MP场景下隔离并复现各卡之间的随机数生成序列保证权重初始化与 Dropout 行为可复现后者负责把全局设备秩device_rank映射为 dp / pp / sharding / mp / sep 五维并行组为流水线、张量并行等训练编排提供组信息。读完本文你将掌握这两个工具的设计动机、源码实现细节、在 PaddleNLP 各模型如 Llama、Qwen、GPT中的实际用法以及它们与 Paddle 框架fleet.meta_parallel的协同关系。一、子包定位与文档骨架本文所依据的文档为 docs/zh/source/paddlenlp.ops.distributed.utils.rst它通过 Sphinxautomodule指令自动生成该子包的 API 文档并把两个子模块纳入 toctreepaddlenlp.ops.distributed.utils.randomRNG 状态追踪器paddlenlp.ops.distributed.utils.topo并行拓扑描述。在仓库中的实际代码路径为paddlenlp/ops/distributed/utils/其中__init__.py对外导出两个公开符号from .topo import Topology from .random import get_rng_state_tracker __all__ [ Topology, get_rng_state_tracker, ]从源码结构可以推断这个子包是一个轻量但关键的分布式基础层它不直接实现通信算子通信由 parallel.py 中的ParallelEmbedding、ColumnParallelLiner、RowParallelLiner等完成而是为这些并行算子以及上层模型提供随机性隔离与拓扑分组两项基础设施。二、RNG 状态追踪器让模型并行下的随机序列可复现2.1 为什么需要 RNG 追踪器在模型并行训练中同一个模型的不同切分如张量并行的不同 rank需要各自初始化自己的参数切片。如果所有 rank 共享同一个随机数生成器不同 rank 上初始化出的权重彼此错位、互相依赖切分后的算子无法保证数学等价性训练也无法复现。因此需要一种机制每个模型并行分片拥有独立且可复现的随机序列同时又不破坏全局随机流如数据加载、Dropout的原始顺序。2.2 核心类RNGStatesTracker源码位于 paddlenlp/ops/distributed/utils/random.py其核心是一个按名字 → CUDA RNG 状态保存的字典class RNGStatesTracker: def __init__(self): # Map from name to the rng state. self.states_ {} self.seeds_ set()模块级还定义了默认状态名常量与全局单例MODEL_PARALLEL_RNG model_parallel_rng RNG_STATE_TRACKER RNGStatesTracker()get_rng_state_tracker()直接返回这个全局单例供全仓库复用从源码结构看这是一个进程内共享的全局追踪器。2.3add(name, seed)登记一个可复现的随机状态def add(self, name, seed): if seed in self.seeds_: raise ValueError(seed {} already exists.format(seed)) self.seeds_.add(seed) if name in self.states_: raise ValueError(state {} already exists.format(name)) orig_rng_state paddle.get_cuda_rng_state() paddle.seed(seed) self.states_[name] paddle.get_cuda_rng_state() paddle.set_cuda_rng_state(orig_rng_state)要点seed 唯一性校验同一个 seed 不能注册两次seeds_集合避免两个状态共享同一随机流name 唯一性校验同一个 name 不能重复注册states_字典状态捕获方式先保存当前 CUDA RNG 状态 → 用paddle.seed(seed)重置生成器 → 抓取此时的 RNG 状态存入字典 → 恢复原始状态。这样注册动作本身不会污染调用方的随机流。2.4rng_state(name)进入/退出该随机状态rng_state是一个contextlib.contextmanager配合with语句使用contextlib.contextmanager def rng_state(self, nameMODEL_PARALLEL_RNG): if name not in self.states_: raise ValueError(state {} does not exist.format(name)) orig_cuda_rng_state paddle.get_cuda_rng_state() paddle.set_cuda_rng_state(self.states_[name]) try: yield finally: self.states_[name] paddle.get_cuda_rng_state() paddle.set_cuda_rng_state(orig_cuda_rng_state)工作流程保存当前全局 CUDA RNG 状态 → 切换到已登记的状态 → 执行随机数生成如参数初始化→ 在finally中把演进后的状态写回字典保证下次进入时序列继续而不是从头重复→ 恢复调用方原本的 RNG 状态。默认状态名为MODEL_PARALLEL_RNG即模型并行专用随机流。2.5 在模型初始化中的真实用法以 Llama 为例paddlenlp/transformers/llama/modeling.py 中_init_weights在tensor_parallel_degree 1时取得追踪器并对分布式权重切片执行if self.config.tensor_parallel_degree 1: rng_tracker get_rng_state_tracker().rng_state ... if layer.weight.is_distributed: with rng_tracker(): layer.weight.set_value( paddle.tensor.normal( mean0.0, stdself.config.initializer_range, shapelayer.weight.shape, ) )同样地Llama 的LlamaLMHead在词汇表被切分vocab_size ! config.vocab_size时也用with get_rng_state_tracker().rng_state():包裹create_parameter见同一文件第 1946 行附近。这样每个张量并行分片的参数初始化都走独立的 RNG 流序列一致、切片互不干扰。同样的模式遍布 Qwen、Qwen2、Qwen3、GPT、Gemma、Mixtral、DeepSeek-V2、ChatGLMv2、Jamba、GLM 等模型见paddlenlp/transformers/下对应modeling.py的get_rng_state_tracker引用说明这是 PaddleNLP 全部主流大模型的统一参数初始化约定。2.6 与 Paddle 框架fleet.meta_parallel的协同值得注意的是模型文件中使用的get_rng_state_tracker有两个来源一部分直接from paddle.distributed.fleet.meta_parallel import get_rng_state_trackerPaddle 框架自带本子包paddlenlp/ops/distributed/utils/__init__.py导出的get_rng_state_tracker()则是仓库内的全局单例封装。二者通过 trainer 桥接在 paddlenlp/trainer/trainer.py 中训练器会调用fleet.meta_parallel.get_rng_state_tracker().set_states_tracker(...)与get_states_tracker()来同步 RNG 状态第 2461、3152 行附近配合set_seed(self.args.seed)第 340 行附近完成全局种子设置。这构成了框架级 tracker 与仓库级单例共存、由 trainer 统一协调的层次。2.7 在 Dropout/重计算场景的应用除了参数初始化RNG 追踪器也用于需要确定性随机行为的算子。例如 paddlenlp/transformers/refined_recompute.py、paddlenlp/peft/lora/utils.py 及多个 MoE 模型Qwen2-MoE、Qwen3-MoE、Mixtral都引用了rng_state。从代码结构可以推断这些场景用同一个模型并行随机流来保证 Dropout mask、MoE 路由噪声等在张量并行下的一致性——这是 Megatron 风格混合并行训练中保证可复现性的关键做法。三、Topology把全局秩映射成五维并行组3.1 设计动机大规模训练通常同时叠加多种并行策略数据并行dp、流水线并行pp、参数分片sharding、张量/模型并行mp、序列并行sep。每个设备需要快速回答两个问题我属于哪个 dp 组组里有哪些 rank以及我在流水线的第几段是否是最后一段。Topology类就是回答这些问题的统一工具。3.2 构造函数与参数源码位于 paddlenlp/ops/distributed/utils/topo.pyclass Topology: def __init__( self, device_rank, world_size, dp_degreeNone, pp_degree1, sharding_degree1, mp_degree1, sep_degree1, order[dp, pp, sharding, mp, sep], ):参数类型默认值含义device_rankint必填当前设备在全局的秩0 ~ world_size-1world_sizeint必填全局设备总数dp_degreeintNone数据并行度pp_degreeint1流水线并行度sharding_degreeint1参数分片并行度mp_degreeint1张量/模型并行度sep_degreeint1序列并行度orderlist[dp, pp, sharding, mp, sep]五维轴的排列顺序注意assert set(order) {dp, pp, sharding, mp, sep}order必须是这五个维度的全排列否则直接抛错。这保证了维度命名始终与内部映射一致。3.3 分组信息的载体GroupInfoGroupInfo namedtuple(GroupInfo, [size, rank, world])每个组信息由三部分组成size组内设备数、rank当前设备在组内的秩、world组内全部 rank 列表。Topology最终产出的所有组对象都是该 namedtuple 的实例。3.4 分组计算的核心算法构造过程分四步第一步构造五维秩张量。按order顺序把五个并行度排成 shape用np.arange(...).reshape(shape)生成一个五维数组arr其扁平索引恰好就是全局 rankshape [degree_map[key] for key in self.order] arr np.arange(0, dp_degree * pp_degree * sharding_degree * mp_degree * sep_degree).reshape(shape) ranks [rank[0] for rank in np.where(arr device_rank)]np.where(arr device_rank)返回当前设备在五个轴上的坐标ranks。第二步对每个轴切出该维为全切片、其余维固定的子数组。对第i个轴构造索引tuple(ranks[:i] [slice(None)] ranks[(i 1):])即在第i维取全部、其他维取当前坐标从而得到该并行维对应的组内 rank 列表worlds [] for i in range(len(ranks)): indexes tuple(ranks[:i] [slice(None)] ranks[(i 1):]) worlds.append(arr[indexes])第三步按order把结果写入五个*_info字段。遍历order把worlds[i]的长度、当前坐标与 rank 列表分别存入dp_info/pp_info/sharding_info/mp_info/sep_infofor i, key in enumerate(self.order): if key dp: self.dp_info GroupInfo(sizelen(worlds[i]), rankranks[i], worldworlds[i].tolist()) ...同时计算self.is_last self.pp_info.rank self.pp_info.size - 1用于判断当前设备是否为流水线最后一段决定是否输出最终 logits、是否计算 loss 等。第四步合成data_info。数据并行与 sharding 在谁消费同一份数据的语义上等价因此把 dp×sharding 两维压平成一个数据组data_arr np.arange(0, dp_degree * sharding_degree).reshape([dp_degree, sharding_degree]) for i, key in enumerate(self.order): if key ! dp and key ! sharding: data_arr np.expand_dims(data_arr, axisi).repeat(degree_map[key], axisi) self.data_info GroupInfo( sizeint(self.dp_info.size * self.sharding_info.size), rankint(self.dp_info.rank * self.sharding_info.size self.sharding_info.rank), worlddata_arr.reshape(-1).tolist(), )随后有一行内部一致性断言assert self.data_info.world[device_rank] self.data_info.rank若映射出错会在构造期立刻暴露data_inner_times self.world.size // self.data_info.size表示每个数据组内部被多少个非数据维的副本复用用于估算数据被重复消费的次数。3.5 拓扑信息的可读输出__repr__以多行文本一次性输出全部六组信息def __repr__(self): return fdp_info:\n\t {self.dp_info}, \npp_info:\n\t {self.pp_info}, \nsharding_info:\n\t {self.sharding_info}, \nmp_info:\n\t {self.mp_info}, \nsep_info:\n\t {self.sep_info}, \ndata_info:\n\t {self.data_info}, \norder:\n\t {self.order}这对调试非常有用在任意 rank 上打印Topology对象即可一目了然地看到自己所属的每个并行组。3.6 在仓库中的使用位置从代码引用看Topology主要在训练与并行编排层被使用典型位置包括paddlenlp/trainer/trainer.py 与 paddlenlp/trainer/trainer_utils.py用于初始化混合并行训练时的分组paddlenlp/trainer/auto_trainer.py 与 paddlenlp/trl/dpo_auto_trainer.py、paddlenlp/trl/sft_auto_trainer.py 等自动训练器paddlenlp/transformers/model_utils.py、configuration_utils.py 等模型装配环节paddlenlp/rl/utils/reshard_utils.py强化学习训练中的参数重分片场景。结合 llm/utils/argument.py 中的并行度配置tensor_parallel_degree、pipeline_parallel_degree、sharding_parallel_degree、sequence_parallel_degree、data_parallel_degree可以推断Topology正是把这些高层训练参数翻译成具体通信组的关键桥梁。四、与并行算子层的配套关系paddlenlp.ops.distributed的完整结构为paddlenlp/ops/distributed/ ├── __init__.py ├── parallel.py # guard / ParallelEmbedding / ColumnParallelLiner / RowParallelLiner └── utils/ ├── __init__.py # 导出 Topology、get_rng_state_tracker ├── random.py # RNGStatesTracker / RNG_STATE_TRACKER / get_rng_state_tracker └── topo.py # Topology / GroupInfoparallel.py 中的并行算子会直接消费 utils 提供的基础设施例如ParallelEmbedding在is_mpworld_size 1且动态图模式下用get_rng_state_tracker().rng_state()包裹create_parameter第 92-96 行保证词汇表切分后每片 embedding 的初始化随机流独立可复现ColumnParallelLiner/RowParallelLiner则负责列切分 / 行切分线性层的通信_c_identity、_c_concat、_c_split、_mp_allreduce。而Topology提供的组信息如mp_info.world正是这些通信算子构造通信组时的依据。三者共同构成了一个拓扑分组 → 随机性隔离 → 切分算子的完整并行基础层。五、小结与实践建议随机性隔离凡是在张量并行/模型并行下需要切分参数或执行随机算子的地方都应通过get_rng_state_tracker().rng_state()默认MODEL_PARALLEL_RNG状态包裹初始化逻辑以保证不同分片的随机序列独立且整体可复现拓扑查询构建混合并行训练时用Topology(device_rank, world_size, dp_degree..., pp_degree..., sharding_degree..., mp_degree..., sep_degree..., order[...])一次性得到全部并行组信息pp_info判断流水线位置is_lastdata_info判断数据消费分组一致性保证Topology内部自带data_info.world[device_rank] data_info.rank断言注册 RNG 状态时对 seed 与 name 均做唯一性校验错误配置会在构造期快速失败而不是在训练中途产生难查的随机数错乱。对于想要深入源码的读者建议按以下顺序阅读先看 utils/random.py 与 utils/topo.py 两个核心实现再看 parallel.py 中的并行算子如何消费它们最后以 paddlenlp/transformers/llama/modeling.py 的_init_weights与 paddlenlp/trainer/trainer.py 的种子同步逻辑作为端到端用例闭环验证。赞分享人工智能大模型预训练微调LoRARLHF强化学习分布式训练【免费下载链接】PaddleNLPEasy-to-use and powerful LLM and SLM library with awesome model zoo.项目地址https://gitcode.com/gh_mirrors/pa/PaddleNLP点击查看免费下载相关推荐PaddleNLP 分布式随机数状态追踪器RNGStatesTracker源码解析与模型并行实践PaddleNLP 分布式随机数状态追踪器RNGStatesTracker源码解析与模型并行实践 本篇指南以 paddlenlp.ops.distribut人工智能大模型预训练微调LoRARLHF强化学习分布式训练模型推理服务推理引擎模型量化模型压缩本地部署NLP攻克CAD拓扑追踪难题BRepTools_History.Merge()深度解析与工业级应用攻克CAD拓扑追踪难题BRepTools_History.Merge 深度解析与工业级应用 引言拓扑操作的历史追踪痛点 在三维建模与CAD应用开发中你是否3D建模图形学Apache APISIX skywalking 插件基于 SkyWalking 的分布式追踪与拓扑分析接入指南Apache APISIX skywalking 插件基于 SkyWalking 的分布式追踪与拓扑分析接入指南 导读 本文围绕 Apache APISIX后端微服务云原生创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表