- 人工智能
- 机器学习
- 分布式训练
- 图计算
- 后端
【免费下载链接】angel
A Flexible and Powerful Parameter Server for large-scale machine learning
模型分区是参数服务器的核心本质功能:为了支撑宽模型、深模型等超大模型,Angel 必须把模型参数切分成多个部分,分布到不同的 PSServer 节点上,并提供便捷的访问服务。本文围绕 docs/design/model_partitioner.md 展开,深入讲解 Angel 默认的 RangePartitioner 划分算法、整体设计链路、以及通过
blockRow/blockCol和自定义Partitioner接口实现灵活分区的两种方法,并结合仓库源码揭示分区从生成、分配到 PSServer 上实例化的完整调用链。读完本文,你将掌握:如何判断模型是否需要自定义分区、如何编写并注入一个自定义 Partitioner、以及分区行为会如何影响 PS 负载均衡、单点瓶颈和请求路由。
为什么参数服务器需要模型分区
无论是宽模型(特征维度极大,例如亿级特征)还是深模型(网络层次极深、参数规模庞大),单个节点的内存都无法容纳完整的模型参数。因此,Angel 需要将模型划分为多个部分(Partition),存储在不同的 PSServer 节点上,并提供统一的访问服务——这就是参数服务器的本质与最基本功能之一。
模型划分方式是参数服务器设计中非常值得关注的通用性工程问题:不同的划分方式会直接导致计算性能上的差异。一个好的模型划分应尽量满足以下三个目标:
- 保证 PS 负载均衡:避免某个 PSServer 成为热点,造成整体吞吐瓶颈;
- 降低 PS 单点性能瓶颈:让请求均匀分散,降低单个节点的压力;
- 关联的数据放在同一个 PS 上:如果一次计算需要访问的数据分布在多个节点,会引入大量跨节点通信,因此把有关联的数据聚合到同一节点可以显著减少网络开销。
但在实际应用中,不同算法、甚至同一算法的不同实现,对划分方式的要求往往并不一致(例如某些行被高频访问、某些分区大小需要不同、某些分区需要强关联地放在同一台 PS 上)。为此,Angel 采取"默认 + 自定义"的双轨策略:提供一套默认划分算法满足一般需求,同时开放自定义分区能力满足特殊需求。
Angel 模型分区整体设计
Angel 的模型分区整体设计如下:
该设计图清晰展示了三个层级:最外层是用户入口PSModel,其内部持有MatrixContext(矩阵上下文),MatrixContext集成了分区器Partitioner(文档中写作默认实现PSPartitioner),分区器负责输出一串Partition1、Partition2、Partition3……;每个分区由左上角(startRow, startCol)与右下角(endRow, endCol)唯一确定。整体设计有以下要点:
- 对用户来说,入口类是
PSModel:应尽量通过它操作模型,屏蔽底层分区细节; - 默认分区算法类是
RangePartitioner(文档中亦写作PSPartitioner,两者指同一默认实现),但它可以被替换为自定义实现; - 分区(Partition)是最小单位:每个分区对应 PSServer 上的一个具体 Model Shard(模型分片)。
从源码结构看,这一设计在仓库中有完整对应:PSModel(PSModel.scala)内部持有matrixCtx: MatrixContext;MatrixContext(MatrixContext.java)则维护矩阵的行数、列数、块大小、partitionerClass等分区元信息;Partitioner接口及RangePartitioner、HashPartitioner等实现位于 partitioner 包。
需要强调的一点是:PSServer 会在 Load Model(加载模型)阶段,根据传入的 Partitioner 进行 Model Shard 的初始化。因此 Model Partitioner 的设计,与运行时的真实模型数据,会共同决定 PS Server 上的模型分布和访问行为。
分区从生成到 PSServer 实例化的完整链路
结合源码可以还原分区的完整生命周期。当应用注册矩阵时,Master 侧的AMMatrixMetaManager(AMMatrixMetaManager.java)会:
- 通过反射实例化
MatrixContext中指定的 partitionerClass,并调用其init(matrixContext, context)初始化(AMMatrixMetaManager.java#L193-L202); - 调用
partitioner.getPartitions()生成分区列表(若从 HDFS 加载已保存模型,则改为读取模型 meta 文件中的分区信息)(AMMatrixMetaManager.java#L218-L236); - 调用
partitioner.assignPartToServer(partId)把每个分区指派给具体的 PSServer,并写入PartitionMeta的storedPs列表(首个 PS 为 Master,其余为 Slave 副本)(AMMatrixMetaManager.java#L301-L309); - 最终把分区元信息同步给各 PSServer,PSServer 在初始化
ServerMatrix时遍历PartitionMeta,为每个分区创建对应的ServerPartition(即 Model Shard)并置为READ_AND_WRITE状态(ServerMatrix.java#L128-L136)。
PartitionMeta(PartitionMeta.java)记录了分区的PartitionKey(含 partId、matrixId、startRow/endRow、startCol/endCol)以及该分区所在的所有 PS 列表,其中storedPs.get(0)是该分区的 Master PS——客户端只能通过 Master PS 对该分区执行 get/put 操作。
请求路由如何利用分区信息
分区信息不仅决定了模型如何存储,还决定了应用层请求如何被切分和路由。当 Worker 上的 PSAgent 发起一次行读取或索引读取时,会基于MatrixMeta中的PartitionKey[]把请求拆分为多个分区请求,再分别发往各分区所在的 PS。例如UserRequestAdapter中按分区拆分请求(UserRequestAdapter.java#L213-L249),而 Range 类型分区对应的RangeRouterUtils则利用各分区getEndCol()边界,对排序后的索引序列进行线性切分,快速定位每个 key 所属的分区(RangeRouterUtils.java#L62-L118)。这正是 RangePartitioner"按 index 区间划分、请求时可快速计算落点"优势的源码体现。
默认分区算法:RangePartitioner
RangePartitioner使用模型的 index 下标范围来划分模型:把模型的 index 范围切分成一个个互不重合的区间。这样做的好处是,在划分应用层请求时可以快速计算出需要的模型部分落在哪个区间上。
在 Angel 中,模型使用 Matrix(矩阵)来表示。当应用没有指定矩阵划分参数或方法时,Angel 使用默认的划分算法,其遵循以下原则:
- 尽量将一个模型平均分配到所有 PS 节点上;
- 对于非常小的模型,尽量将它们放在一个 PS 节点上(避免过度切分带来的管理开销);
- 对于多行的模型,尽量将同一行放在一个 PS 节点上(保证按行访问时的局部性)。
RangePartitioner 的源码实现细节
仓库中默认实现的源码为 RangePartitioner.java,其关键逻辑如下:
init阶段从配置中读取三项全局参数(RangePartitioner.java#L52-L70):默认分区大小angel.model.partitioner.partition.size、最大分区数angel.model.partitioner.max.partition.number、每台 PS 的分区数angel.model.partitioner.partition.number.perserver,并结合angel.ps.number计算maxPartNum;getPartitions阶段依据行数row、列数col(或 index 区间 start/end)、块大小等参数计算blockRow与blockCol,然后双重循环按行块、列块切分,生成一系列PartitionMeta(RangePartitioner.java#L72-L139);assignPartToServer的默认实现是partId % serverNum,即按分区 ID 轮流分配到各 PS,天然实现负载均衡(RangePartitioner.java#L142-L146)。
代码中的blockRow/blockCol计算逻辑体现了"小模型放一个 PS、多行模型同行同节点"的原则:当行数大于 PS 数时,blockRow取min(row/serverNum, ...)使每个块尽可能包含多行;而当行数很少时blockRow = row,即整行不拆分。
相关的全局配置参数
RangePartitioner行为受以下全局配置影响(定义见 AngelConf.java#L500-L510):
| 配置项 | 默认值 | 说明 |
|---|---|---|
angel.model.partitioner.partition.size | 500000 | 默认分区大小(元素个数),用于估算块大小 |
angel.model.partitioner.max.partition.number | 500 | 全局最大分区数上限 |
angel.model.partitioner.partition.number.perserver | 5 | 每台 PS 上的分区数(>0 时参与限制总分区数) |
angel.ps.number | 1 | PS 节点个数,参与负载均衡与分区数计算 |
这些配置可在 angel-ps/conf/angel-default.xml 基础上通过angel-site.xml或提交参数进行覆盖。
仓库中其他内置分区器
除RangePartitioner外,partitioner 包 还提供了另外两个内置实现,可作为理解"分区器可替换"的补充:
HashPartitioner(HashPartitioner.java):按分区总数切分,partId直接映射到partId % psNum的 PS 上;注意MatrixContext.check()中规定稠密类型矩阵(如T_DOUBLE_DENSE)不能使用 HashPartitioner(MatrixContext.java#L723-L732);ColumnRangePartitioner(ColumnRangePartitioner.java):继承RangePartitioner,在getPartitions()中强制把blockRow设为整个行数,即按整行进行范围切分,保证每一行完整地落在同一个分区。
需要说明的是,MatrixContext中partitionerClass的默认值是HashPartitioner(MatrixContext.java#L107),而在initPartitioner()中若用户未显式设置则回退为RangePartitioner(MatrixContext.java#L571-L577)。从这一实现可以推断:文档中"默认分区算法类 RangePartitioner(可替换)"的表述对应的是未指定分区方式时的默认行为,实际选择哪个分区器由MatrixContext中最终生效的partitionerClass决定。
自定义模型分区:两种方式
为了实现更复杂的模型分区方式、适配复杂算法,Angel 允许用户自定义矩阵分区,并提供了两种方法,在方便性与灵活性之间取得平衡。
方式一(简单):定义 PSModel 时传入 blockRow 和 blockCol
第一种方式最直接:在定义PSModel时传入blockRow和blockCol参数,控制每个分区的行列块大小:
val sketch = PSModelTDoubleVector从PSModel的构造函数签名可以看到,blockRow和blockCol的默认值都是 -1(即不指定,交由 RangePartitioner 自动计算):
class PSModel( val modelName: String, val row: Int, val col: Long, val blockRow: Int = -1, val blockCol: Long = -1, val validIndexNum: Long = -1, var needSave: Boolean = true)(implicit ctx: TaskContext)(见 PSModel.scala#L50-L57)。构造函数内部会把blockRow/blockCol写入matrixCtx = new MatrixContext(modelName, row, col, validIndexNum, blockRow, blockCol)(PSModel.scala#L62),MatrixContext对应的setMaxRowNumInBlock/setMaxColNumInBlock会把这些值传给默认的 RangePartitioner 使用。
通过这种方式,用户可以轻松控制矩阵分区的大小,满足一般需求,但它的灵活性有限:分区间大小必须一致,无法表达"不同分区不同大小"或"按访问频率差异化切分"等复杂策略。
方式二(高阶):实现自定义 Partitioner 接口
当算法非常复杂时,往往会出现"奇奇怪怪"的需求,例如:
- 矩阵各部分访问频率不一致(某些行被高频访问);
- 不同分区需要不同的大小;
- 需要把某些存在关联的分区放到同一个 PS 上;
- ……
为了满足这类特殊需求,Angel 抽象了Partitioner分区接口——RangePartitioner只是其中一个默认实现,用户可以通过实现该接口来定义自己的模型分区方式,并注入到PSModel中,从而改变模型的分区行为。
仓库中Partitioner接口的当前定义如下(Partitioner.java):
public interface Partitioner { /** * Init matrix partitioner * * @param mContext matrix context * @param context AMContext */ void init(MatrixContext mContext, AMContext context); /** * Generate the partitions for the matrix * * @return the partitions for the matrix */ List<PartitionMeta> getPartitions(); /** * Assign a matrix partition to a parameter server * * @param partId matrix partition id * @return parameter server index */ int assignPartToServer(int partId); /** * Get partition type * * @return partition type */ PartitionType getPartitionType(); }注意:与文档中的接口定义相比,仓库当前版本的init签名略有演进(第二参数由Configuration conf变为AMContext context),并新增了getPartitionType()方法返回分区类型(RANGE_PARTITION/HASH_PARTITION等),说明接口在持续演进,读者以当前仓库源码为准。
用户需要重点实现的两个方法是:
getPartitions:获取该矩阵的分区列表。返回值List<PartitionMeta>中每个PartitionMeta都包含matrixId、partId、startRow、endRow、startCol、endCol等完整的分区坐标信息;assignPartToServer:决定一个分区(通过partId标识)被分配给哪个 PSServer(返回 PS 索引)。
实战示例:按访问频率自定义分区
下面通过一个具体例子,演示如何通过实现Partitioner接口来自定义矩阵划分方式。假设使用场景如下:
- 模型是一个
3 * 10,000,000维的矩阵,PS 个数为 8; - 其中第一行访问得非常频繁,需要被更细粒度地切分以分散热点。
为此,我们制定如下分区策略:
- 第一行切分成 4 个分区,分布到 4 台 PS 上;
- 其他两行各切分成 2 个分区,每行分布到 2 台 PS 上。
图例如下:
由于默认分区方式要求分区大小相对均匀,无法实现这种大小不一的分区,因此我们定制一个分区类CustomizedPartitioner:
public class CustomizedPartitioner implements Partitioner { …… @Override public List<MLProtos.Partition> getPartitions() { List<MLProtos.Partition> partitions = new ArrayList<MLProtos.Partition>(6); int row = mContext.getRowNum(); int col = mContext.getColNum(); int blockCol = col / 4; int partitionId = 0; // Split the first row to 4 partitions for (int i = 0; i < 4; i++) { if (i < 3) { partitions.add(MLProtos.Partition.newBuilder().setMatrixId(mContext.getId()) .setPartitionId(partitionId++).setStartRow(0).setEndRow(1).setStartCol(i * blockCol) .setEndCol((i + 1) * blockCol).build()); } else { partitions.add(MLProtos.Partition.newBuilder().setMatrixId(mContext.getId()) .setPartitionId(partitionId++).setStartRow(0).setEndRow(1).setStartCol(i * blockCol) .setEndCol(col).build()); } } blockCol = col / 2; // Split other row to 2 partitions for (int rowIndex = 1; rowIndex < row; rowIndex++) { partitions.add( MLProtos.Partition.newBuilder().setMatrixId(mContext.getId()).setPartitionId(partitionId++) .setStartRow(rowIndex).setEndRow(rowIndex + 1).setStartCol(0).setEndCol(blockCol) .build()); partitions.add( MLProtos.Partition.newBuilder().setMatrixId(mContext.getId()).setPartitionId(partitionId++) .setStartRow(rowIndex).setEndRow(rowIndex + 1).setStartCol(blockCol).setEndCol(col) .build()); } return partitions; } …… }实现要点说明:
- 第一行(
startRow=0, endRow=1)按blockCol = col / 4切成 4 段,前 3 段结束列为(i+1) * blockCol,最后一段结束列取完整的col(保证不丢列); - 其余行(
rowIndex = 1, 2)按blockCol = col / 2各切成 2 段; - 通过
partitionId++为每个分区分配递增且唯一的partitionId; - 示例代码使用
MLProtos.Partition.newBuilder()构建分区。从仓库源码看,该接口已演变为返回PartitionMeta(Partitioner.java#L41),因此在实际实现时,getPartitions()应按当前版本返回List<PartitionMeta>,分区对象通过new PartitionMeta(matrixId, partId, startRow, endRow, startCol, endCol)构造(PartitionMeta.java#L57-L60),核心的分区坐标语义与示例完全一致; - 示例中还省略了
assignPartToServer的实现,可参考RangePartitioner的默认实现partId % serverNum轮询分配,或根据业务需要自定义(例如把访问频繁的第一行分区分散到多台 PS,把有关联的分区放到同一台 PS)。
实现了CustomizedPartitioner后,将其注入到PSModel的MatrixContext之中,即可让模型使用自定义分区:
psModel.matrixCtx.setPartitioner(new CustomizedPartitioner());从源码看,MatrixContext中与分区器相关的注入/查询方法为setPartitionerClass(Class<? extends Partitioner>)与getPartitionerClass()(MatrixContext.java#L375-L377、L463-L465),实际使用时可通过matrixCtx.setPartitionerClass(CustomizedPartitioner.class)注入分区器类,Master 侧会通过反射实例化并调用(AMMatrixMetaManager.java#L193-L202)。
此外,MatrixContext还支持通过addPart(PartContext)/setParts(List<PartContext>)直接手工指定分区列表(MatrixContext.java#L618-L638),这可以看作自定义分区的另一种数据驱动方式:当parts非空时,AMMatrixMetaManager会直接按PartContext构造PartitionMeta而不再调用partitioner.getPartitions()(AMMatrixMetaManager.java#L224-L235)。
分区策略的选择建议与注意事项
结合默认算法与自定义机制,可以给出如下实践建议:
- 优先使用默认 RangePartitioner:大多数模型(尤其是行数不多、特征维度大的宽模型)都能从"index 区间切分 +
partId % serverNum轮询分配"中获得良好的负载均衡与请求路由效率,无需额外定制; - 分区粒度要匹配访问模式:如果某些行/区间访问明显更频繁,建议像示例那样将这些区域切得更细、分散到更多 PS 上,降低单点瓶颈;若模型很小,则可利用默认规则整体放在一台 PS 上,减少分片管理开销;
- 关注分区与行局部性:需要按整行操作的模型,可参考
ColumnRangePartitioner的思路保证整行不跨分区;需要按行聚合的计算,尽量让同一行落在同一台 PS 上,避免跨节点通信; - 注意稠密矩阵与 HashPartitioner 的限制:稠密类型矩阵不能使用
HashPartitioner,自定义分区器也要避免与所选RowType冲突(MatrixContext.java#L723-L732); - 分区数不宜过大:每个分区在 PSServer 上对应一个 ServerPartition 实例,分区过多会带来元信息与状态管理开销,建议利用
angel.model.partitioner.max.partition.number、angel.model.partitioner.partition.number.perserver等配置把分区规模控制在合理范围。
总结
模型分区是 Angel 参数服务器中最基本、最关键的功能之一,它决定了超大模型在 PS 集群上的分布方式,并直接影响负载均衡、单点瓶颈与请求路由效率。Angel 通过默认的RangePartitioner(按 index 区间切分、按partId % serverNum轮询分配)满足了绝大多数场景;同时开放了Partitioner接口与blockRow/blockCol快捷参数,让复杂算法可以按访问频率、分区大小、数据关联性等维度自由定制分区行为。从PSModel入口,到MatrixContext携带分区信息,再到 Master 侧生成分区、分配 PS、最终在 PSServer 上实例化 ServerPartition,这条完整链路贯穿了 Angel 分布式模型管理的始终。理解并善用分区机制,是编写高效、可扩展的 Angel 机器学习算法的基本功。
- 人工智能
- 机器学习
- 分布式训练
- 图计算
- 后端
【免费下载链接】angel
A Flexible and Powerful Parameter Server for large-scale machine learning
相关推荐
MXNet 自定义模型分区器(Custom Partitioner)实战:基于动态库的图分区扩展开发指南
MXNet 自定义模型分区器(Custom Partitioner)实战:基于动态库的图分区扩展开发指南 导读 本指南以 MXNet 仓库中的 example/
深度学习人工智能机器学习分布式训练Angel:面向超大规模机器学习的灵活高性能参数服务器平台
Angel:面向超大规模机器学习的灵活高性能参数服务器平台 Angel 是一个以 模型为中心 设计理念构建的高性能分布式机器学习与图计算平台,核心实现是一套通用
人工智能机器学习分布式训练图计算后端米家API终极实战:5个智能家居自动化场景深度解析
米家API终极实战:5个智能家居自动化场景深度解析 米家API是Python开发者控制小米智能家居设备的完整解决方案,提供了从账号认证到设备控制的全套工具链。本
智能家居物联网MCP 服务AI 技能
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考