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

文章详情

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

PaddleNLP 模型并行基础组件解析:paddlenlp.ops.distributed.parallel 源码级指南

PaddleNLP 模型并行基础组件解析:paddlenlp.ops.distributed.parallel 源码级指南 人工智能大模型预训练微调LoRARLHF强化学习分布式训练【免费下载链接】PaddleNLPEasy-to-use and powerful LLM and SLM library with awesome model zoo.项目地址https://gitcode.com/gh_mirrors/pa/PaddleNLP点击查看免费下载导读paddlenlp.ops.distributed.parallel是 PaddleNLP 面向大规模模型训练提供的轻量级模型并行Model Parallelism基础组件模块集中实现了设备守卫装饰器guard、词表并行嵌入ParallelEmbedding、列并行线性层ColumnParallelLiner与行并行线性层RowParallelLiner四个核心原语。本篇文章将以 parallel.py 的源码为主线逐一讲解每个 API 的参数含义、切分策略、通信行为与底层实现原理并结合模块内配套的 RNG 状态跟踪与拓扑工具帮助你理解 PaddleNLP 中 Megatron 式张量并行Tensor Parallelism是如何被搭建出来的以及如何在自定义模型中使用这些原语完成权重切分与聚合。一、模块定位PaddleNLP 自研的轻量级张量并行原语该模块由 API 文档 paddlenlp.ops.distributed.parallel 通过automodule指令自动生成参考内容其文档主题即指向paddlenlp/ops/distributed/parallel.py这个 Python 源码文件。因此理解这份文档本质就是理解该源码文件中的全部公开接口。从模块导出来看parallel.py通过__all__明确声明了四个公开 API__all__ [ guard, ParallelEmbedding, ColumnParallelLiner, RowParallelLiner, ]而 paddlenlp/ops/distributed/init.py 又通过from .parallel import *与from .utils import *将parallel与utils两个子模块的接口汇总导出其中utils子模块包含utils/topo.py定义Topology类与GroupInfo命名元组用于描述 dp / pp / sharding / mp / sep 五维并行拓扑utils/random.py定义RNGStatesTracker与全局单例RNG_STATE_TRACKER用于在模型并行下隔离随机数状态。从代码结构上看这套原语与paddle.distributed.fleet.meta_parallel提供的VocabParallelEmbedding、ColumnParallelLinear、RowParallelLinear定位一致属于同一类按维度切分权重 通信聚合的张量并行基础设施。在 PaddleNLP 的 GPT、Llama、Qwen 等模型的 pipeline parallel 建模文件如 gpt/modeling.py以及 PEFT Prefix 模型中可以看到fleet.meta_parallel系实现的成熟用法而本模块则是仓库内一套独立、轻量、自包含的并行原语实现二者在接口设计上高度同构适合作为理解张量并行原理的切入点。模块开头对 paddle 版本做了防御性处理若当前 paddle 不包含paddle.distributed.fleet会抛出warnings.warn(paddle.distributed is not contains in you paddle!)警告而非直接崩溃方便在单机非分布式环境下也能安全 import。二、guard(device)设备守卫装饰器guard是一个装饰器工厂接收一个设备字符串如gpu:0返回一个装饰器用于将某个nn.Layer子类的初始化和前向计算强制约束在指定设备上执行def guard(device): def decorator(Layer): class WrapperClass(Layer): def __init__(self, *args, **kw): with paddle.static.device_guard(device): print(Init {} on {}.format(Layer.__name__, device)) super().__init__(*args, **kw) def forward(self, *args, **kw): with paddle.static.device_guard(device): print(Forward {} on {}.format(Layer.__name__, device)) return super().forward(*args, **kw) return WrapperClass return decorator要点如下机制通过动态生成WrapperClass继承被装饰的 Layer在__init__和forward中分别包上paddle.static.device_guard(device)上下文从而控制参数创建与前向计算发生的设备适用场景适合在多设备环境下将不同子层显式调度到不同卡上是朴素流水线/设备摆放的一种轻量手段注意事项实现中使用了print输出初始化与前向的设备信息属于调试用途其行为依赖静态图模式的device_guard语义在纯动态图场景下的生效范围需要结合具体 paddle 版本验证。三、ParallelEmbedding词表并行嵌入层ParallelEmbedding将词表按模型并行维度切分每个 rank 只保存词表的1/world_size分片输入 id 通过start_index偏移定位到本 rank 的分片区间。3.1 构造参数参数类型含义默认值num_embeddingsint词表大小embedding 字典大小决定输入 id 的最大取值上界必填embedding_dimint每个 embedding 向量的维度必填rankint当前分片的 rank决定词表分片的起始索引必填world_sizeinttrainer 数量即模型并行组大小必填weight_attrTensor / ParamAttr指定权重参数属性含初始化方式None使用默认参数属性namestr参数命名一般无需用户设置None3.2 切分逻辑与校验assert ( num_embeddings % self.world_size 0 ), The length of the vocabulary must be divisible by the parallelism degree of MP per_part_size num_embeddings // self.world_size self.vocab_start_index self.rank * per_part_size self._size [per_part_size, embedding_dim]硬性约束num_embeddings必须能被world_size整除否则直接断言失败这是词表均匀切分的前提分片规则每个 rank 持有的权重形状为[per_part_size, embedding_dim]其中per_part_size num_embeddings / world_size起始索引vocab_start_index rank * per_part_size即当前 rank 负责的词表区间为[rank * per_part_size, (rank 1) * per_part_size)。3.3 随机状态隔离与分布式标记当启用模型并行world_size 1且处于动态图模式时参数创建被包进 RNG 状态跟踪器以保证各 rank 参数初始化的一致性if self.is_mp and paddle.in_dynamic_mode(): with get_rng_state_tracker().rng_state(): self.weight self.create_parameter(...)随后代码将权重标记为分布式张量并同步到 startup 与 main 两个 Program 的 global block 中self.weight.is_distributed True if self.is_mp else False startup_block paddle.static.default_startup_program().global_block() main_block paddle.static.default_main_program().global_block() startup_block.vars[self.weight.name].is_distributed True main_block.vars[self.weight.name].is_distributed True这里is_distributed True是告诉底层分布式框架该参数是已经按 MP 切分好的分片在后续的保存、加载与自动并行处理中不应再被二次切分。3.4 前向计算def forward(self, x): if self.is_mp: output_parallel paddle.distributed.collective._c_lookup_table( self.weight, x, start_indexself.vocab_start_index, nameself._name ) output paddle.distributed.collective._mp_allreduce( output_parallel, groupNone, use_calc_streamTrue, use_model_parallelTrue ) else: output paddle.nn.functional.embedding( x, weightself.weight, padding_idxNone, sparseFalse, nameself._name ) return output输入x包含 id 信息的 Tensor数据类型应为int32或int64且 id 取值应在[0, weight.shape[0]]范围内此处指原始未切分词表大小num_embeddingsMP 分支先通过底层_c_lookup_table在本地分片上做查表用start_index完成 id 偏移再通过_mp_allreduce做模型并行组内的 AllReduce把各 rank 的局部 embedding 求和还原为完整输出这是词表并行中各 rank 查到不同 id 的向量后相加的经典 Megatron 做法非 MP 分支退化为标准的paddle.nn.functional.embedding常规查表。四、ColumnParallelLiner按列切分的并行线性层ColumnParallelLiner对应 Megatron 的 Column Parallel Linear权重按输出维度列axis1切分即把[in, out]的权重切成num_partitions份[in, out/num_partitions]的列分片。4.1 构造参数参数类型含义默认值size(int, int)原始线性层形状(输入维度, 输出维度)必填num_partitionsint模型并行组内的切分数1gather_outbool是否对输出做聚合concatTrueparam_attrTensor / ParamAttr权重参数属性含初始化方式Nonebias_attrTensor / ParamAttr偏置参数属性Nonenamestr参数命名None4.2 rank 推导与形状校验if paddle.in_dynamic_mode(): rank paddle.distributed.get_rank() else: assert fleet._role_maker, To use paddle.distributed.split, you must call fleet.init() firstly. rank fleet.worker_index() inner_rank rank % num_partitions动态图下直接取全局 rank静态图下要求先调用fleet.init()再通过fleet.worker_index()获取inner_rank rank % num_partitions将全局 rank 映射到模型并行组内的局部编号硬性约束size[1] % num_partitions 0即输出维度必须能被切分数整除。4.3 权重构建与命名self.per_part_size size[1] // num_partitions linear_size (size[0], self.per_part_size) if not name: name fc_by_col_rank_%d % inner_rank else: name name _by_col_rank_%d % inner_rank每个 rank 上的权重形状为[size[0], size[1] / num_partitions]命名自动追加_by_col_rank_{inner_rank}后缀以便区分不同分片。切分完成后权重同样被标记为is_distributed True并同步到 startup / main Program由于列切分时偏置也随列一同切分因此self.linear.bias也会被标记为分布式if self.linear._bias_attr: startup_block.vars[self.linear.bias.name].is_distributed True main_block.vars[self.linear.bias.name].is_distributed True self.bias self.linear.bias4.4 前向计算group None x paddle.distributed.collective._c_identity(x, groupgroup) output_parallel self.linear(x) if self.gather_out is False: return output_parallel return paddle.distributed.collective._c_concat(output_parallel, groupgroup)前向流程分三步_c_identity对输入做身份拷贝保证输入在各 rank 上完整可用本地self.linear(x)每个 rank 用自己那列权重计算局部输出若gather_outTrue通过_c_concat沿最后一维拼接各 rank 的局部输出还原完整输出若为False则直接返回并行输出供下游的 Row 切分层作为并行输入消费省去一次通信。返回 Tensor 为映射后的结果其数据类型可为 int 或 float。五、RowParallelLiner按行切分的并行线性层RowParallelLiner对应 Megatron 的 Row Parallel Linear权重按输入维度行axis0切分即把[in, out]的权重切成num_partitions份[in/num_partitions, out]的行分片前向时对局部输出做 AllReduce 求和。5.1 构造参数参数类型含义默认值size(int, int)原始线性层形状(输入维度, 输出维度)必填num_partitionsint模型并行组内的切分数1input_is_parallelbool输入是否已经是并行切分状态Falseparam_attrTensor / ParamAttr权重参数属性含初始化方式Nonebias_attrTensor / ParamAttr偏置参数属性Nonenamestr参数命名None5.2 权重构建与偏置处理assert ( size[0] % num_partitions 0 ), Number of rows of the weight for linear ({}) must be divisible by num_partitions ({}).format(...) self.per_part_size size[0] // num_partitions linear_size (self.per_part_size, size[1])硬性约束size[0] % num_partitions 0即输入维度必须能被切分数整除每个 rank 上的权重形状为[size[0] / num_partitions, size[1]]命名自动追加_by_row_rank_{inner_rank}后缀关键设计行切分时每个 rank 只持有部分输入行偏置不能在本地直接加上必须等 AllReduce 汇总后再加。因此源码中构建内部 Linear 时显式传bias_attrFalse并在注释中说明# NOTE(wangxi): row split, bias need add after allreduceself.linear paddle.nn.Linear( num_rows, num_cols, weight_attrparam_attr, bias_attrFalse, # row split, bias need add after allreduce namename, )若用户传入的bias_attr不为False则在类外部单独创建形状为[num_cols]的完整偏置参数if bias_attr is not False: self.bias self.create_parameter(shape[num_cols], attrbias_attr, dtypeself._dtype, is_biasTrue) else: self.bias None5.3 前向计算group None if self.input_is_parallel: assert x.shape[-1] self.per_part_size, ( The width ({}) of the input x must be equal to the height ({}) of the weight. Maybe you should split the input x using paddle.split..format(x.shape[-1], self.per_part_size) ) else: # split last dim x paddle.distributed.collective._c_split(x, groupgroup) output_parallel self.linear(x) output paddle.distributed.collective._mp_allreduce( output_parallel, groupgroup, use_calc_streamTrue, use_model_parallelTrue ) output output self.bias if self.bias is not None else output return output前向流程分四步若input_is_parallelTrue要求输入最后一维宽度恰好等于per_part_size即输入已由上游 Column 层切分好否则断言报错并提示用paddle.split手动切分若input_is_parallelFalse先通过_c_split沿最后一维切分输入使各 rank 只拿到与自身行分片对应的输入片段本地线性计算后用_mp_allreduce在模型并行组内求和恢复完整输出最后再加上完整的偏置向量。六、典型组合Transformer 前馈网络的列切分 行切分范式在 Megatron 式张量并行中Transformer 的 FFN 通常采用第一个线性层按列切分无通信聚合第二个线性层按行切分的组合从而将两次全连接之间的中间激活也切分到各 rank 上只保留最终的 AllReduce。用本模块原语可以这样组织from paddlenlp.ops.distributed import ColumnParallelLiner, RowParallelLiner hidden_size 4096 ffn_size 16384 mp_size 8 # 模型并行度需满足 ffn_size % mp_size 0 # 第一个线性层按列切分gather_outFalse 直接输出并行激活 fc1 ColumnParallelLiner( size(hidden_size, ffn_size), num_partitionsmp_size, gather_outFalse, # 不聚合把切分后的激活直接交给下一个层 ) # 第二个线性层按行切分input_is_parallelTrue 消费并行激活 fc2 RowParallelLiner( size(ffn_size, hidden_size), num_partitionsmp_size, input_is_parallelTrue, # 输入已是切分状态无需再切 )对应到仓库中的成熟实现可以在 PEFT Prefix 模型 prefix_model.py 中看到完全一致的模式fleet.meta_parallel.VocabParallelEmbedding负责词表并行嵌入fleet.meta_parallel.ColumnParallelLinear负责列切分fleet.meta_parallel.RowParallelLinear负责行切分——这与本模块三个类的职责一一对应可作为对照学习素材。七、配套工具RNG 状态跟踪与并行拓扑paddlenlp.ops.distributed.utils提供了两个与本模块协同工作的基础工具7.1 RNGStatesTracker模型并行下的随机状态隔离utils/random.py 实现了一个简单的 RNG 状态跟踪器add(name, seed)以给定 seed 登记一个命名的随机状态重复 seed 或重复 name 会抛出ValueErrorrng_state(nameMODEL_PARALLEL_RNG)上下文管理器进入时切换 CUDA RNG 状态到指定命名状态退出时保存并还原模块级单例RNG_STATE_TRACKER与工厂函数get_rng_state_tracker()。ParallelEmbedding在动态图 MP 模式下正是通过get_rng_state_tracker().rng_state()包裹参数创建从而保证各 rank 使用一致的初始化 seed实现权重分片间的初始化对齐。7.2 Topology五维并行拓扑描述utils/topo.py 定义了Topology类以device_rank、world_size与五个维度度数dp_degree、pp_degree、sharding_degree、mp_degree、sep_degree为输入通过order默认[dp, pp, sharding, mp, sep]指定维度排列顺序利用 numpy 多维数组索引推导出每个维度对应的GroupInfosize / rank / world 命名元组并额外计算data_info数据并行与参数分片合并后的数据组与data_inner_times。其构造函数中assert set(order) {dp, pp, sharding, mp, sep}保证了维度集合的合法性__repr__则可打印全部组的拓扑信息方便调试。八、使用前提与注意事项安装与导入前提parallel.py依赖paddle.distributed.fleet相关能力。若当前 paddle 发行版不包含分布式模块import 时会触发warnings.warn警告相关并行 API 将不可用可整除性约束ParallelEmbedding要求num_embeddings % world_size 0ColumnParallelLiner要求size[1] % num_partitions 0RowParallelLiner要求size[0] % num_partitions 0使用前需保证维度可整除静态图前置条件静态图模式下使用ColumnParallelLiner/RowParallelLiner必须先调用fleet.init()否则fleet._role_maker断言失败行切分的偏置语义Row 切分层内部 Linear 不带 bias完整 bias 在 AllReduce 之后手动相加切勿在切分前自行拼接带 bias 的 Lineargather_out与input_is_parallel的配对使用列层gather_outFalse产出并行激活行层需同步设置input_is_parallelTrue才能省去中间通信二者必须配套分布式标记三个类都会将切分后的权重及列层的偏置标记为is_distributed True并同步到 startup / main Program保证统一保存加载时按分片处理。九、总结paddlenlp.ops.distributed.parallel用不足 330 行的代码完整复现了张量并行的三大核心原语词表并行嵌入、列并行线性与行并行线性并配套设备守卫装饰器、RNG 状态跟踪器与五维拓扑工具。理解这份源码等于掌握了 PaddleNLP 模型并行建模的最小骨架列层负责把权重和中间激活横向铺开行层负责把部分结果汇总还原AllReduce 与 split / concat 通信点清晰可见。对于希望将自定义模型改造成张量并行、或深入理解fleet.meta_parallel同类实现的读者本文所述的 parallel.py 是一份不可多得的入门级参考实现。赞分享人工智能大模型预训练微调LoRARLHF强化学习分布式训练【免费下载链接】PaddleNLPEasy-to-use and powerful LLM and SLM library with awesome model zoo.项目地址https://gitcode.com/gh_mirrors/pa/PaddleNLP点击查看免费下载相关推荐PaddleNLP中DeepseekR1-MTP模型的并行解码技术解析PaddleNLP中DeepseekR1 MTP模型的并行解码技术解析 引言大模型推理的性能瓶颈与突破 在大语言模型LLM的实际部署中推理性能往往是制约大模型NLP预训练微调模型推理服务模型量化分布式训练RLHF人工智能智能OCR革新如何用开源工具实现跨语言无缝沟通智能OCR革新如何用开源工具实现跨语言无缝沟通 在全球化数字时代语言障碍成为阻碍信息获取的最大挑战之一。Translumo作为一款先进的实时屏幕翻译工具通人工智能大模型预训练微调LoRARLHF强化学习分布式训练模型推理服务推理引擎模型量化模型压缩本地部署NLPPaddleNLP 分词器基础设施全解PretrainedTokenizerBase 与 BatchEncoding 源码级剖析PaddleNLP 分词器基础设施全解PretrainedTokenizerBase 与 BatchEncoding 源码级剖析 PaddleNLP 是面向人工智能大模型预训练微调LoRARLHF强化学习分布式训练模型推理服务推理引擎模型量化模型压缩本地部署NLP创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表