CANN ops-transformer DistributeBarrier 算子全解析:NPU 通信域全卡同步屏障的原理与 aclnn 调用实战
【免费下载链接】ops-transformer本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。项目地址: https://gitcode.com/cann/ops-transformer
DistributeBarrier 是 CANN ops-transformer 项目(mc2/distribute_barrier 目录)中提供的一个轻量级同步算子:在指定 HCCL 通信域内完成所有卡的全卡同步(barrier),其唯一的业务输入xRef只用于构建 Tensor 依赖、不参与任何计算。本文以 mc2/distribute_barrier/README.md 为主体,结合 aclnnDistributeBarrier 接口文档、aclnnDistributeBarrierV2 接口文档 以及 op_api、op_host、op_kernel 源码,系统讲解该算子的产品支持、两段式接口、参数语义、底层调用链与约束,并给出可直接参考的调用示例,帮助读者在 MoE 分布式并行等需要屏蔽快慢卡性能波动的场景中正确接入该算子。
一、算子定位:为什么需要“无计算”的全卡同步
在超大规模分布式训练/推理网络中,各卡由于负载差异、通信抖动等原因会出现明显的快慢不一致(straggler 问题)。DistributeBarrier 算子专门用于在通信域内完成全卡同步:调用后所有卡必须全部到达该屏障点,才能继续执行后续任务。其设计上有两个鲜明特点:
- 零业务语义:
xRef仅用于构建 Tensor 依赖,接口内部不对xRef做任何读写操作; - 纯同步语义:算子本身不做数据搬运与计算,只负责“等待与放行”。
从源码看,算子的 Host 侧定义也印证了这一点:distribute_barrier_def.cpp 将x_ref同时声明为输入与输出(输出与输入 shape/数据类型完全一致),并通过MC2().HcclGroup({"group"})将该算子纳入 MC2 通信框架、绑定到名为group的 HCCL 通信域上。
典型使用场景:在需要进行全卡同步的网络模型中调用该算子,可以屏蔽快慢卡引入的性能波动问题,协助分析性能;也可以连续多次调用,在关键节点之间人为插入同步点。在实际 MoE 分布式流水线(Dispatch → Barrier → Combine)中,Barrier 被置于专家分发与合并之间,用于确保各 EP 卡完成调度后再进行结果合并。
二、产品支持情况
根据 README.md 与两份接口文档,该算子在不同产品上的支持情况如下:
| 产品 | 是否支持 |
|---|---|
| Ascend 950DT | √ |
| Atlas A3 训练系列产品 / Atlas A3 推理系列产品 | √ |
| Atlas A2 训练系列产品 / Atlas A2 推理系列产品 | × |
| Atlas 200I/500 A2 推理产品 | × |
| Atlas 推理系列产品 | × |
| Atlas 训练系列产品 | × |
即该算子仅在 Ascend 950DT 与 Atlas A3 系列产品上可用。这一结论在算子定义中得到印证:distribute_barrier_def.cpp 仅为ascend910_93(A3 架构)与ascend950两个平台注册了 AICore 配置,分别指向 kernel 目录下的distribute_barrier_a3与distribute_barrier_apt两个实现(见 arch22/distribute_barrier_a3.cpp 与 arch35/distribute_barrier_apt.cpp)。
三、接口总览:V1 与 V2 两段式接口
算子对外提供两套 aclnn 接口:
| 接口 | 说明 |
|---|---|
aclnnDistributeBarrier/aclnnDistributeBarrierGetWorkspaceSize | 基础版本,入参为xRef、group、worldSize |
aclnnDistributeBarrierV2/aclnnDistributeBarrierV2GetWorkspaceSize | 增强版本,在 V1 基础上新增timeOutOptional与elasticInfoOptional两个可选参数 |
V2 与 V1 的差别(见 aclnnDistributeBarrierV2.md):
- 新增
elasticInfoOptional参数,用于支持 EP 通信域动态缩容; - 新增
timeOutOptional参数,用于设置超时时间。
从实现上看,V1 与 V2 共享同一套底层逻辑:aclnn_distribute_barrier.cpp 中aclnnDistributeBarrierGetWorkspaceSize直接以nullptr替代timeOut与elasticInfo调用aclnnDistributeBarrierGetWorkspaceSizeBase;而 aclnn_distribute_barrier_v2.cpp 则原样透传用户传入的可选参数,二者的执行阶段(第二阶段接口)则完全共用aclnnDistributeBarrierBase。
四、两段式接口与函数原型
与 CANN 其他 aclnn 算子一致,DistributeBarrier 采用两段式接口设计:必须先调用GetWorkspaceSize接口获取计算所需 workspace 大小以及包含算子计算流程的执行器(executor),再调用执行接口真正下发计算。
4.1 aclnnDistributeBarrier(V1)
aclnnStatus aclnnDistributeBarrierGetWorkspaceSize( aclTensor* xRef, const char* group, int64_t worldSize, uint64_t* workspaceSize, aclOpExecutor** executor)aclnnStatus aclnnDistributeBarrier( void *workspace, uint64_t workspaceSize, aclOpExecutor *executor, aclrtStream stream)4.2 aclnnDistributeBarrierV2(V2)
aclnnStatus aclnnDistributeBarrierV2GetWorkspaceSize( const aclTensor *xRef, const aclTensor *timeOutOptional, const aclTensor *elasticInfoOptional, const char *group, int64_t worldSize, uint64_t *workspaceSize, aclOpExecutor **executor)aclnnStatus aclnnDistributeBarrierV2( void *workspace, uint64_t workspaceSize, aclOpExecutor *executor, aclrtStream stream)4.3 第一阶段接口(GetWorkspaceSize)参数说明
| 参数名 | 输入/输出 | 描述 | 使用说明 | 数据类型 | 数据格式 | 维度(shape) | 非连续 Tensor |
|---|---|---|---|---|---|---|---|
| xRef | 输入 | 无业务语义,仅用于构建输入 Tensor 依赖,接口内不做任何操作 | 无 | BFLOAT16、FLOAT16、FLOAT32、BOOL、INT8、INT16、INT32、INT64、UINT8、UINT16、UINT32、UINT64、FLOAT8_E5M2、FLOAT8_E4M3FN、FLOAT4_E1M2、FLOAT4_E2M1、HIFLOAT8、INT4 | ND | 0-8(V2 中 INT4 支持 2-3) | √ |
| timeOutOptional | 输入 | 超时时间设置,如果在此时间内所在卡未完成全卡同步,则认为该卡存在超时异常 | 可选:传有效数据或空指针 | INT32 | ND | 1 | √ |
| elasticInfoOptional | 输入 | EP 通信域动态缩容信息 | 可选:传有效数据或空指针,空指针表示不开启动态缩容 | INT32 | ND(支持非连续 Tensor) | 1 | √ |
| group | 输入 | 通信域名称,进行所有卡同步的通信域 | 支持长度 [1,127] | STRING | - | - | - |
| worldSize | 输入 | 通信域大小 | - | INT64 | - | - | - |
| workspaceSize | 输出 | 返回需要在 Device 侧申请的 workspace 大小 | - | UINT64 | - | - | - |
| executor | 输出 | 返回 op 执行器,包含算子计算流程 | - | aclOpExecutor* | - | - | - |
4.4 第二阶段接口(执行)参数说明
| 参数名 | 输入/输出 | 描述 |
|---|---|---|
| workspace | 输入 | 在 Device 侧申请的 workspace 内存地址 |
| workspaceSize | 输入 | Device 侧申请的 workspace 大小,由第一阶段接口返回 |
| executor | 输入 | op 执行器,包含算子计算流程 |
| stream | 输入 | 指定执行任务的 Stream |
4.5 返回值与错误码
两个接口均返回aclnnStatus状态码。第一阶段接口完成入参校验,出现以下场景时报错:
| 返回值 | 错误码 | 描述 |
|---|---|---|
| ACLNN_ERR_PARAM_NULLPTR | 161001 | V1:输入的必选参数 Tensor 是空指针;V2:传入的 xRef、group 或 worldSize 是空指针 |
| ACLNN_ERR_PARAM_INVALID | 161002 | V2 独有:xRef、timeOut、elasticInfo、group、worldSize 的数据类型/数据格式不在支持范围内,或 shape 不匹配 |
| ACLNN_ERR_INNER_TILING_ERROR | 561002 | V1 独有:参数的取值不在支持的范围内 |
4.6 平台相关的参数细节
- Atlas A3 训练/推理系列产品:不支持 FLOAT8_E5M2、FLOAT8_E4M3FN、FLOAT4_E1M2、FLOAT4_E2M1、HIFLOAT8、INT4 类型;
epWorldSize取值支持 [2, 384]; - Ascend 950DT:
timeOutOptional参数中的超时时间单位为微秒(us),建议配置 5000000us,根据实际环境不同超时时间下限可能不同;epWorldSize取值支持 [2, 1024]。
五、elasticInfoOptional 动态缩容信息详解
elasticInfoOptional是 V2 接口支持 EP 通信域动态缩容的关键参数,其语义如下:
- Atlas A2 系列产品:不支持,传空指针;
- Atlas A3 系列产品:传入 1D Tensor,shape 为
4 + 2 * epWorldSize,INT32 类型。其中前 4 位为缩容配置,后 2 * epWorldSize 为 rank 映射表。
在 V2 的调用示例中,elasticInfo的构造方式如下(EP_WORLD_SIZE = 2时 shape 为{4 + 2*2} = {8},示例中实际按更大规模构造):
std::vector<int32_t> elasticInfoHostData{ isElastic, rankNumAfterElastic, sharedExpertRankNumAfterElastic, moeExpertNumAfterElastic, 0, 1, -1, -1, -1, -1, 2, 3, 0, 1, 6, 7, -1, -1, -1, -1 };即前 4 个元素描述缩容后的配置(是否缩容、缩容后的 rank 数、共享专家 rank 数、MoE 专家数),后续元素构成新旧 rank 的映射表,-1表示该位置无有效映射。Tiling 阶段会严格校验其维度:distribute_barrier_tiling.cpp中的CheckElasticInfo要求dim0 == ELASTIC_METAINFO_OFFSET + RANK_LIST_NUM * worldSize,其中ELASTIC_METAINFO_OFFSET对应前 4 位配置、RANK_LIST_NUM对应每个 rank 的映射项数,与文档中4 + 2 * epWorldSize的表述一致。
参数一致性约束:开启elasticInfoOptional时,必须确保aclnnMoeDistributeDispatchV3与aclnnMoeDistributeCombineV3(或aclnnMoeDistributeCombineAddRmsNormV2)也开启该参数,并且传入的elasticInfo取值保持一致,否则会导致 EP 通信域内的 rank 视图不一致。
六、约束说明
- 通信域使用约束:一个模型中的
aclnnDistributeBarrier/aclnnDistributeBarrierV2需要使用单独通信域,该通信域中不允许有其他算子。示例代码中专门通过HcclCommInitAll创建了独立的commsEpBarrier通信域(hcclEpBarrierComm),并在调用时通过HcclGetCommName取出通信域名称传入接口,正是对这一约束的落实; - 通信方式约束:Ascend 950DT 仅支持 UB Memory 通信;
- 确定性计算:
aclnnDistributeBarrier与aclnnDistributeBarrierV2均为默认确定性实现; - 使用场景:在需要进行全卡同步的网络模型中调用该算子,可屏蔽快慢卡引入的性能波动问题,协助分析性能;可以连续调用,入图时需将上个算子的输入、下个算子的输出作为入参传入接口。
七、底层实现剖析:从 aclnn 到 Kernel 的调用链
理解底层调用链有助于在排查问题、评估性能时做出正确判断。DistributeBarrier 的完整调用链如下:
7.1 op_api 层:平台分流的入口
distribute_barrier_base.cpp 是 V1/V2 共用的核心实现,主要做了三件事:
- 入参校验:
BarrierCheckNullStatus检查xRef与group非空,BarrierCheckParams校验group字符串长度在(0, HCCL_GROUP_NAME_MAX)范围内; - 平台分流:通过
GetCurrentPlatformInfo().GetCurNpuArch() == DAV_3510判断是否为 Ascend 950DT:- 非 950(即 Atlas A3)走
aclnnInnerDistributeBarrierGetWorkspaceSize/aclnnInnerDistributeBarrier; - 950 平台则先通过
Mc2Aclnn::Mc2Context::GetMc2ContextTensor基于 group 获取 MC2 上下文 Tensor,再调用aclnnInnerDistributeBarrierExtendGetWorkspaceSize/aclnnInnerDistributeBarrierExtend;
- 非 950(即 Atlas A3)走
- 服务类型设置:执行阶段在 950 平台上会调用
NnopbaseSetHcclServerType(executor, NNOPBASE_HCCL_SERVER_TYPE_MTE),将 HCCL 服务类型设为 MTE,保证同步原语以正确的硬件引擎执行。
7.2 op_host 层:算子注册、InferShape 与 Tiling
- 算子原型(distribute_barrier_def.cpp):声明输入
x_ref(必选,18 种数据类型,ND 格式)、time_out(可选,INT32)、elastic_info(可选,INT32),输出x_ref,属性group(必选字符串)与world_size(必选整型);为ascend910_93与ascend950分别注册 Kernel,并声明MC2().HcclGroup({"group"}); - InferShape(distribute_barrier_infershape.cpp):输出 shape 与 dtype 直接拷贝自输入
xRef,即输出只是对输入依赖的“透传”,进一步印证算子不产生任何计算; - Tiling(distribute_barrier_tiling.cpp):负责参数校验与运行时配置生成,关键逻辑包括:
CheckAndSetAttrs:校验worldSize必须在[MIN_WORLD_SIZE, MAX_WORLD_SIZE]范围内,其中 Ascend 950 与 A3 的MAX_WORLD_SIZE取值不同,对应文档中epWorldSize的 [2, 1024] 与 [2, 384] 限制;CheckTimeOut:校验timeOut为 1D、INT32、非 FRACTAL_NZ 格式、dim0 == 1;CheckElasticInfo:校验elasticInfo为 1D、INT32、dim0 == 4 + 2 * worldSize;SetHcommCfg:通过Mc2CcTilingConfig配置通信算法AlltoAll=level0:fullmesh;level1:pairwise,并设置AIV_ENGINE通信引擎——注释明确指出“通过不拉起 AICPU,提高算子退出性能”,这是该算子低开销设计的关键;SetWorkSpace:为系统预留固定大小 workspace(SYSTEM_NEED_WORKSPACE);- 计算
numBlocks并设置ScheduleMode(1)(batch mode,所有核同时启动),同时将totalUbSize、aivNum、isInputTimeOut、isInputElasticInfo等写入 TilingData,供 Kernel 使用。
7.3 op_kernel 层:架构差异化实现
Kernel 侧针对两种架构分别实现:arch22/distribute_barrier_a3.cpp(Atlas A3)与 arch35/distribute_barrier_apt.cpp(Ascend 950DT)。TilingData 定义在 distribute_barrier_tiling.h,包含worldSize、rankId、aivNum、totalUbSize、isInputTimeOut、isInputElasticInfo、isMc2Context等字段,Kernel 据此完成全卡同步的集合通信等待逻辑。
八、调用示例:以 MoE Dispatch → Barrier → Combine 为背景
仓库提供了完整的可运行样例 examples/test_aclnn_distribute_barrier.cpp(V1)与 V2 接口的示例代码(见 aclnnDistributeBarrierV2.md 调用示例章节)。样例以 2 卡 EP(EP_WORLD_SIZE=2、TP_WORLD_SIZE=1)的 MoE 专家并行为例,演示了 Dispatch、Barrier、Combine 三个算子的串联调用。其中 Barrier 的调用核心逻辑如下。
8.1 创建独立的 Barrier 通信域
HcclComm commsEpBarrier[TP_WORLD_SIZE][EP_WORLD_SIZE]; for (int32_t tpId = 0; tpId < TP_WORLD_SIZE; tpId++) { ret = HcclCommInitAll(EP_WORLD_SIZE, devicesEp[tpId], commsEpBarrier[tpId]); CHECK_RET(ret == ACL_SUCCESS, LOG_PRINT("[ERROR] HcclCommInitAll epBarrier %d failed, ret %d\n", tpId, ret); return ret); }每个 rank 通过HcclGetCommName(args.hcclEpBarrierComm, hcomEpBarrierName)取得通信域名称字符串,供后续传入接口。
8.2 两段式调用(V1 路径)
// 第一阶段:获取 workspace 大小与 executor ret = aclnnDistributeBarrierGetWorkspaceSize(expandX, hcomEpBarrierName, EP_WORLD_SIZE, &barrierWorkspaceSize, &barrierExecutor); CHECK_RET(ret == ACL_SUCCESS, LOG_PRINT("[ERROR] aclnnDistributeBarrierGetWorkspaceSize failed. ret = %d \n", ret); return ret); // 根据 workspaceSize 申请 device 内存 if (barrierWorkspaceSize > 0) { ret = aclrtMalloc(&barrierWorkspaceAddr, barrierWorkspaceSize, ACL_MEM_MALLOC_HUGE_FIRST); CHECK_RET(ret == ACL_SUCCESS, LOG_PRINT("[ERROR] aclrtMalloc workspace failed. ret = %d \n", ret); return ret); } // 第二阶段:下发执行 ret = aclnnDistributeBarrier(barrierWorkspaceAddr, barrierWorkspaceSize, barrierExecutor, args.barrierStream); CHECK_RET(ret == ACL_SUCCESS, LOG_PRINT("[ERROR] aclnnDistributeBarrier failed. ret = %d \n", ret); return ret); // 同步等待任务执行结束 ret = aclrtSynchronizeStreamWithTimeout(args.barrierStream, 10000);其中expandX是 Dispatch 阶段的输出 Tensor,作为 Barrier 的xRef传入,从而在计算图上建立 Dispatch → Barrier → Combine 的依赖关系——这正是“xRef 仅用于构建 Tensor 依赖”的典型应用方式。
8.3 带超时与动态缩容的 V2 调用
V2 的调用差别仅在于第一阶段多传两个可选 Tensor(详见 aclnn_distribute_barrier_v2.cpp 的函数原型):
// 创建 timeOut 与 elasticInfo Tensor(1D INT32) std::vector<int64_t> elasticInfoShape{4 + EP_WORLD_SIZE * 2}; std::vector<int64_t> timeOutShape{1}; std::vector<int32_t> timeOutHostData(timeOutShapeSize, 1000000); // 超时 1000000us // 第一阶段:warm up 时可传空指针(不开启超时与缩容) ret = aclnnDistributeBarrierV2GetWorkspaceSize(expandX, nullptr, nullptr, hcomEpBarrierName, EP_WORLD_SIZE, &barrierWorkspaceSize, &barrierExecutor); // 正式调用时传入 timeOut 与 elasticInfo ret = aclnnDistributeBarrierV2GetWorkspaceSize(expandX, timeOut, elasticInfo, hcomEpBarrierName, EP_WORLD_SIZE, &barrierWorkspaceSize, &barrierExecutor); // 第二阶段:执行 ret = aclnnDistributeBarrierV2(barrierWorkspaceAddr, barrierWorkspaceSize, barrierExecutor, args.barrierStream); ret = aclrtSynchronizeStreamWithTimeout(args.barrierStream, 10000);注意:timeOutOptional与elasticInfoOptional都是可选参数,可选择传入有效数据或填空指针。在 950DT 上timeOut的单位为 us,文档建议配置 5000000us;示例中 warm up 阶段传空指针、正式阶段传有效数据,以规避首次建链开销对超时统计的影响。
8.4 单测验证
仓库为该算子提供了完整的单测覆盖,可作调用与行为验证的参考:
- op_api 层:tests/ut/op_api/test_aclnn_distribute_barrier.cpp;
- op_host 层:tests/ut/op_host/test_distribute_barrier_infershape.cpp 与 tests/ut/op_host/arch22/test_distribute_barrier_tiling.cpp;
- op_kernel 层:tests/ut/op_kernel/test_distribute_barrier.cpp。
九、使用建议与注意事项
- 独立通信域是硬约束:Barrier 必须独占一个 HCCL 通信域,通信域内不允许再挂载其他算子,务必在模型初始化阶段单独创建;
- 依赖构建方式:入图场景下,将上一个算子的输出 Tensor 作为
xRef传入,即可在图中建立正确的执行依赖;连续调用时同理,将上个 Barrier 的输出作为下个 Barrier 的输入; - 超时参数按平台配置:950DT 上
timeOut单位为 us、建议 5000000us,且不同环境超时下限可能不同,需根据实际集群环境调整; - 动态缩容的一致性:开启
elasticInfoOptional时,Dispatch、Barrier、Combine 三处传入的elasticInfo必须完全一致,否则会导致 EP 通信域内 rank 视图错乱; - 性能开销可控:从 Tiling 实现看,算子通过 AIV 引擎通信、不拉起 AICPU(见 distribute_barrier_tiling.cpp 中
SetCommEngine(mc2tiling::AIV_ENGINE)的注释),并采用 batch mode 使所有核同时启动,整体同步开销被控制在较低水平,适合在分析性能、屏蔽快慢卡波动时高频插入。
十、总结
DistributeBarrier 是一个“小而专”的同步算子:功能上只做通信域内全卡屏障,设计上通过xRef依赖注入、独立通信域、AIV 引擎通信与平台分流(A3 / 950DT)实现了低开销、可入图、可连续调用的同步能力;V2 版本进一步引入超时检测与 EP 通信域动态缩容,使其在 MoE 专家并行这类对容错与可观测性要求较高的场景中更加实用。读者可结合本仓库的 README、V1 接口文档、V2 接口文档 与 示例代码 快速上手,并按上文约束完成工程接入。
【免费下载链接】ops-transformer本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。项目地址: https://gitcode.com/cann/ops-transformer
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考