已开启
[Requirement|需求建议]: 集合通信Executor统一算法结构方案 #607
卢杨创建于 8月12日
8月12日 修改了issue 的描述
8月12日 关联了pull request:[WIP][Feature] executor refactor rfc
8月12日 将 luyang20 设为负责人
8月12日 修改了issue 的描述
8月12日 修改了issue 的描述
8月27日 添加了label:tech-debt
8月28日 修改了issue 的描述
8月28日 修改了issue 的描述
8月28日 关联了pull request:feat(oxc): add experimental oxc core, build integration
8月31日 关联了pull request:feat(oxc): add OXC skeleton framework with full headers and stub implementations
8月31日 关联了pull request:feat(oxc): add OXC skeleton framework with full headers and stub implementations
9月1日 修改了issue 的描述
29 天前 修改了issue 的描述
29 天前 修改了issue 的描述
29 天前 修改了issue 的描述


RFC: 集合通信Executor统一算法结构方案
概要
本RFC提出一套以
AlgoDesc为静态算法描述、以OpsExecutor为通用解释器、以 Template 为单层执行单元、以 CommPlanner 为通信计划生成器的 HCCL 重构方案。方案通过递归算法树描述 Sequence、Parallel 和 OmniPipe 等组合,通过ranksForInputData/ranksForOutputData传递逻辑数据归属,并建立Input → CCL Buffer → ... → CCL Buffer → Output的统一内存与数据流模型。本重构以
experimental/ops/op_common/recursive_executor/为落地目录,采用插件式零侵入接入方式对接到 src 原流程:recursive_executor 代码以 OBJECT 库形式编入libhccl.so,通过REGISTER_ADAPTOR_EXECUTOR把AdaptorExecutor(继承 src 的InsCollAlgBase)注册进 src 的CollAlgExecRegistryV2;src 的 Selector 对 4 级拓扑选择 recursive_executor 算法名后,原流程Selector → HcclExecOp → GetAlgExec → CalcAlgHierarchyInfo/CalcRes → Orchestrate自然调度到 recursive_executor 执行器,src 侧除 Selector 的 4 级拓扑分支外零改动。背景与动机
现状问题
当前HCCL的代码架构随着多层拓扑、并行切分和复合算子增加,同一种 Mesh 或 NHR 通信过程会在多个算法类中重复出现,存在以下问题:
直接触发原因:四层组网拓扑适配
上述问题被激化的直接触发因素是四层组网拓扑的新增需求。
HCCL 当前支持 2~3 层网络拓扑:
每种拓扑层级组合下,每个算子(AllGather/AllReduce/Broadcast/ReduceScatter/Scatter)都需要对应的 Executor + Template 实现。四层组网拓扑在此基础上新增第 4 层(跨 super-pod 的 OCS 层):
在旧架构下,适配四层拓扑意味着每个算子新增 4 层编排执行器,且与 3 层代码大量重复(4 层的前三阶段与 3 层完全相同,仅多了第 4 层),工作量线性膨胀:6 个算子 ×(Sequence + Parallel + 可能的 OmniPipe 变体)≈ 12-18 个新执行器类,每个 300-500 行。
这正是重构的核心驱动力:需要一个能抽象描述任意层级组合的统一数据结构,使新增拓扑层级只需在算法表中注册算法执行策略,而非复制整套执行器代码,即可实现新增算法。
在重构后的架构中,4 层 AllGather 的算法只是在 3 层算法上再追加一层
AllGather_Mesh1DOcs:flowchart LR L0["Mesh1D<br/>Level 0"] --> L1["NHR<br/>Level 1"] --> L2["NHR<br/>Level 2"] --> L3["<span style='color:red'>Mesh1DOcs<br/>Level 3 ← 新增</span>"]代码量对比
当前 Executor 有 6 个命令字(AllGather、AllReduce、Broadcast、Reduce、ReduceScatter、Scatter),每个命令字可能有 5 种常见编排方式(Solo、Sequence、Concurrent、Parallel、OmniPipe),随着定制机型扩展和通信维度扩展,算子编排方式会出现倍数增长。以
broadcast_parallel执行器为例,算法需要 4 个 Template 并行在 inter 和 intra 维度计算,由于缺少数据抽象,OrchestrateLoop函数超过 150 行,并生成每个模板的TemplateDataParams对象。初步设计可将所有 Executor 合并为 1 个通用执行器(约 800 行),大幅降低维护成本。execPolicy和dataSplitRatio(~5 行)重构目标与非目标
目标:
AlgoDesc + AlgoExecDesc + TemplateExecDesc描述算法结构。约束:
src/common/hcomm_dlsym/的符号表 + dlsym。术语表
架构与接口契约
整体架构
重构后的核心链路只有四层:
flowchart LR Algo["AlgoDesc<br/>算法结构定义"] Executor["OpsExecutor<br/>算法执行"] Template["Template<br/>算法模板"] CommPlanner["CommPlanner<br/>生成通信计划"] Comm["通信与本地计算操作"] Algo --> Executor Executor --> Template Template --> CommPlanner CommPlanner --> Template Template --> Comm四层分别回答不同问题:
AlgoDescOpsExecutor核心设计思想是把静态算法结构和动态数据状态分离:
AlgoDesc在算法选择完成后保持不变;AlgoExecDataDesc随 Loop、Sequence 阶段和 Parallel 子片动态变化。Executor 递归解释算法树时不修改AlgoDesc,只为每个 Child 派生一份AlgoExecDataDesc。对外接口
本次重构为算子开发者提供以下接口,开发者基于这些接口即可新增算法而无需编写执行器类(详细流程见 4. 新增算法指南):
1. REGISTER_ALG — 算法与执行器一步注册
#define REGISTER_ALG(cmdType, algName, algoDesc)一步完成算法入
AlgSelector和执行器入CollAlgExecRegistryV2,两表以同一算法名关联。cmdTypeHcclCMDTypeHCCL_CMD_ALLGATHER),决定执行器在 src 注册表中的分类槽位algNamestd::string"AicpuAllGatherSequenceXxxMesh"),Selector 返回此名称,Executor 据此查找算法树algoDescAlgoDesc2. AlgoDesc / AlgoExecDesc / TemplateExecDesc — 算法树描述结构
struct TemplateExecDesc { TemplateDesc templateDesc; // 算子类型 + 算法类型 int subCommIndex; // 子通信域索引 int netLayer = -1; // 网络层索引,-1=遍历所有层取首个匹配 }; struct AlgoExecDesc { HcclAlgExecPolicy execPolicy; // 执行策略:SEQUENCE / PARALLEL / OMNIPIPE std::vector<VariantType> children; // 子节点(叶子=TemplateExecDesc,非叶子=嵌套AlgoExecDesc) std::vector<u32> dataSplitRatio; // 并行数据切分比例,元素个数须与 children 一致 }; // AlgoDesc 即 HcclAlgorithm,描述完整算法入口 class HcclAlgorithm { HcclCMDType hcclCmdType; // 算子类型 HcclAlgEngineType engineType; // 引擎类型 std::shared_ptr<TopoMatchBase> topoMatch; // 拓扑匹配器 AlgoExecDesc algoExecDesc; // 执行树 std::string algName; // 算法名 };开发者用这三层数据结构组装算法树:
TemplateExecDesc描述单层 Template,AlgoExecDesc描述执行策略与子节点列表,AlgoDesc描述完整算法入口。3. AicpuBaseTemplate — Template 基类
class AicpuBaseTemplate : public BaseTemplate { public: // Template 执行入口(由 Engine 调用),编排 PreCopy → RunAlgorithm → SendAll → PostCopy HcclResult KernelRun(const TemplateDataParams &tempAlgParams, TemplateResource &templateResource, std::vector<u32> &ranksForOutputData) override; protected: // [纯虚] 子类必须实现:调用 CommPlanner 生成收发描述列表 virtual HcclResult RunAlgorithm(std::vector<TxRxSlicesList> &txRxSlicesLists, std::vector<u32> &ranksForOutputData) = 0; // [可选重写] 统一执行 SendRecv;默认走 WRITE 方向 virtual HcclResult SendAll(const std::vector<TxRxSlicesList> &txRxSlicesLists, TemplateResource &templateResource, const std::vector<ThreadHandle> &threads); // [可选重写] 通信后本地处理(ccl buffer → output) virtual HcclResult PostCopy(const std::vector<ThreadHandle> &threads); // [可选重写] 计算所需线程数与 notify 数 virtual HcclResult GetRes(AlgResourceRequest &res) const; };开发者继承此类,实现
RunAlgorithm()调用 CommPlanner 生成通信描述列表,按需重写SendAll()/PostCopy()/GetRes()。基类提供PreCopy → RunAlgorithm → SendAll → PostCopy的执行骨架。4. GetTemplate — Template 工厂
std::unique_ptr<BaseTemplate> GetTemplate(const TemplateDesc &templateDesc, const std::vector<u32> &ranks, u32 myRank);templateDescTemplateDesc{HCCL_CMD_ALLGATHER, HCCL_ALGO_TYPE_FULLMESH})ranksstd::vector<u32>myRanku32std::unique_ptr<BaseTemplate>通过
REGISTER_TEMPLATE(cmdType, algType, TemplateClass)宏将 Template 类注册进全局工厂表,GetTemplate按TemplateDesc查表创建实例。新增 Template 时调用此宏即可,无需修改GetTemplate本身。依赖接口
本重构依赖以下 src / HCOMM 接口,不引入新外部依赖:
InsCollAlgBasesrc/ops/op_common/executor/executor_v2_base.hAdaptorExecutor继承此类,实现三个纯虚接口(Prepare/GetRes/KernelRun),桥接 src 执行框架OpParam/AlgResourceRequest/AlgResourceCtxSerializablesrc/ops/op_common/...CollAlgExecRegistryV2src/ops/op_common/executor/registry/coll_alg_v2_exec_registry.hREGISTER_ALG通过此表将AdaptorExecutor注册到 srchcomm_dlsym符号表src/common/hcomm_dlsym/cann/hcomm的编译期硬依赖param.algName路由机制影响分析
OpsExecutor的递归编排相比特化执行器的 inline 调用引入额外开销,小消息场景敏感;统一的 Slot 布局可能使 Scatter 算子 CCL Buffer 占用上升。设计稳定后需跑 benchmark 对比。REGISTER_ALG)注册到商用代码运行,不改变现有构建和发布流程。兼容性考虑
本次重构不改变算法选择策略、不改变公开 API,也不改变 src 执行框架。迁移只发生在算法被选中之后:已选算法名 →
GetAlgExec返回 recursive_executor 执行器 →AlgoDesc树 → 通用 Executor 解释 → 新 Template/CommPlanner 执行。1. 代码合入路径:experimental/ops/op_common/recursive_executor/
根据
experimental/README.md规范,重构代码落在experimental/ops/op_common/recursive_executor/(位于experimental/ops/op_common/下,结构与src一致,含executor/、template/、topo/、inc/;algorithm/目录待 Phase 1 算法注册补全时创建)。2. 渐进式接入
灰度接入分三阶段:
experimental/ops/op_common/recursive_executor/已搭建骨架(TopoMatchFourLevel+AdaptorExecutor+OpsExecutor骨架 +AllGatherMesh/NhrTemplate+Mesh/NhrCommPlanner),算法注册(algorithm/all_gather.cc)待补全。详细设计
1. 数据结构
1.1 AlgoDesc 三层描述结构
AlgoDesc由三个层次组成:classDiagram class AlgoDesc { HcclCMDType hcclCmdType CommEngine engineType shared_ptr~TopoMatchBaseV2~ topoMatch AlgAttrs algAttrs AlgoExecDesc algoExecDesc string algName + GetExecutor(OpParam&) OpsExecutor + Dump() } class TopoMatchBaseV2 { + MatchTopo(topoInfo, algHierarchyInfo, algAttrs) } class AlgoExecDesc { AlgExecPolicy execPolicy vector~VariantType~ children vector~u32~ dataSplitRatio } class TemplateExecDesc { TemplateDesc templateDesc int subCommIndex int netLayer } class TemplateDesc { HcclCMDType hcclCmdType HcclAlgoType algType } AlgoDesc *-- AlgoExecDesc AlgoDesc o-- TopoMatchBaseV2 AlgoExecDesc *-- AlgoExecDesc AlgoExecDesc *-- TemplateExecDesc TemplateExecDesc *-- TemplateDescAlgoDesc:描述一个完整集合通信算法,携带topoMatch(拓扑匹配器)和算法树,通过GetExecutor()创建通用执行器。这是注册进算法表的最小单元。TopoMatchBaseV2:拓扑匹配器,负责把通信域拓扑拆成逐层子通信域(AlgHierarchyInfoForAllLevel),供 Selector/Executor 共用。AlgoExecDesc:组合节点,描述 Children 采用何种执行策略(SEQUENCE/PARALLEL/OMNIPIPE)。TemplateExecDesc:算法模板,描述在哪个子通信域(subCommIndex)上执行哪种 Template,netLayer用于跨层模板(默认 -1)。TemplateDesc:描述算子语义与拓扑算法类型。对应代码(
experimental/ops/op_common/recursive_executor/inc/algo_desc.h):enum class AlgExecPolicy { SEQUENCE, PARALLEL, OMNIPIPE }; struct TemplateDesc { HcclCMDType hcclCmdType; HcclAlgoType algType; }; enum SubCommIndexType : int { SUB_COMM_INDEX_0 = 0, SUB_COMM_INDEX_1 = 1, SUB_COMM_INDEX_2 = 2, SUB_COMM_INDEX_3 = 3, SUB_COMM_INDEX_4 = 4, SUB_COMM_INDEX_5 = 5, }; struct TemplateExecDesc { TemplateDesc templateDesc; int subCommIndex; int netLayer = -1; }; struct AlgoExecDesc; using VariantType = std::variant<TemplateExecDesc, std::shared_ptr<AlgoExecDesc>>; struct AlgoExecDesc { AlgExecPolicy execPolicy = AlgExecPolicy::SEQUENCE; std::vector<VariantType> children; std::vector<u32> dataSplitRatio; }; class AlgoDesc { public: std::unique_ptr<OpsExecutor> GetExecutor(OpParam& param); void Dump(); HcclCMDType hcclCmdType; CommEngine engineType; std::shared_ptr<TopoMatchBaseV2> topoMatch; AlgAttrs algAttrs; AlgoExecDesc algoExecDesc; std::string algName; };AlgoExecDesc::children是递归 Variant:每个 Child 要么是一个算法模板(TemplateExecDesc),要么是一棵子树(shared_ptr<AlgoExecDesc>)。Executor 遍历时使用std::get_if区分两种类型:对TemplateExecDesc走RunTemplateDesc,对shared_ptr<AlgoExecDesc>递归调用OrchestrateLoop。1.2 执行策略
现状 Executor 盘点
当前
src/ops/下共有 53 个 Executor 类(统计口径:直接继承InsCollAlgBase/ExecutorBase的类共 63 个,剔除 Send/Recv/BatchSendRecv 等 10 个点对点类),按编排方式和算子维度分类如下(下表为主要类别,未穷举 AllGatherV/ReduceScatterV/Aiv 等变体):分析上述 53 个 Executor 的编排方式,可归纳为三种本质不同的 Children 组织模式:
SEQUENCE。PARALLEL,通过dataSplitRatio描述切分比例。OMNIPIPE,作为 Parallel 的流水线扩展。SEQUENCEPARALLELdataSplitRatio切给多个 Children,各 Child 处理不同数据子片OMNIPIPE现状中的 Sole、Sequence、TwoShot 三种 Executor 都映射到
SEQUENCE:Sole 是单算法模板的 Sequence,TwoShot 是两个算法模板的 Sequence,区别仅在 children 数量和类型。Parallel 和 Concurrent 都映射到PARALLEL:Concurrent 是 Parallel 的特例(各 Child 使用不同子通信域)。OrderPreserved 作为算法语义约束,由 Template 层保证,不影响 Executor 编排策略。映射等价性论证:旧架构中 Parallel 与 Concurrent 在数据流上都属于"按比例切分数据 + 在不同子通信域上并发执行",差异仅在于切分比例的来源——
InsV2AllGatherParallelExecutor在CalcCostCoeff中使用固定 50/50 比例(constexpr float ratio = 0.5f)把数据切给 mesh 与 NHR 两个轴;InsV2AllGatherConcurrentExecutor通过GetParallelDataSplit按端口数比例(splitData = portNum0 / (portNum0 + portNum1))切分,并在 mesh + CLOS 两个子通信域并发下发。两者都切分数据,不存在"同时下发但不切分"的语义。重构用dataSplitRatio统一表达切分比例——固定比例或端口比例都可在算法构建阶段算好填入,因此合并为PARALLEL不改变运行行为。这三种策略是正交的:串行描述时序依赖,并行描述数据切分,流水描述执行重叠。任意复杂算法都可以通过这三种策略的递归组合表达。新增编排模式只需扩展枚举值和 Executor 的
OrchestrateLoop分支,不影响已有策略的执行逻辑。拓扑层级(
subCommIndex)和执行顺序(execPolicy)是两个独立维度:前者回答"在哪个 Rank 集合上执行",后者回答"Children 之间如何组织"。subCommIndex到实际通信域的映射:subCommIndex是逐层子通信域表AlgHierarchyInfoForAllLevel::infos[]的下标,该表由TopoMatchBaseV2::MatchTopo在CalcAlgHierarchyInfo阶段填充(infos[i]为第 i 层子通信域的 Rank 分组)。Template 运行时通过同一下标取本层通信资源:algHierarchyInfo_.infos[subCommIndex].at(0)取本层 Rank 列表,GenTemplateRes以channelTable_.at(subCommIndex)/subThreads_.at(subCommIndex)取本层 channel 与线程(资源由CalcRes阶段CalcTemplateChannelRes按subCommIndex逐层申请、经resCtx.channels[level]序列化、InitRes时RestoreChannelMap重建,见 3.6)。约束:subCommIndex必须小于拓扑层数,越界直接报错。1.3 算法注册表(AlgSelector)
算法注册表是重构架构的核心枢纽——它将"算法"从代码逻辑转化为可查询的数据结构,实现 Selector、Executor 和 Template 三层的彻底解耦。
设计动机
算法注册表的核心思想是算法即数据:每个算法用一棵
AlgoDesc树完整描述,预构造后注册进全局表。Selector 只负责返回算法名,Executor 只负责解释执行,两者都不包含算法定义本身。新增算法只需在算法文件里追加一个REGISTER_ALG宏。当前阶段尚未完成算法注册机制的落地,Selector 仍需侵入式修改 src 代码(新增 4 级拓扑分支);后续算法注册机制完善后,Selector 与 Executor 均可实现零代码侵入。数据结构
class AlgSelector { public: static AlgSelector& Instance(); HcclResult Register(const std::string& algName, AlgoDesc algo); bool GetAlgorithm(const std::string& algName, AlgoDesc& algo) const; private: AlgSelector() = default; std::map<std::string, AlgoDesc> algMap_; mutable std::mutex mu_; }; // 注册宏:静态初始化期把算法对象登记进注册表 #define REGISTER_ALGORITHM(algName, algo) \ static bool g_reg_##algName = \ AlgSelector::Instance().Register(#algName, algo)GetAlgorithm按名字返回AlgoDesc的拷贝(topoMatch为shared_ptr,共享同一匹配器)。注册表在 Host 库和 Device 内核中各自静态初始化,双端都可通过算法名重建算法定义,无需序列化算法树。算法注册示例(4 级 AllGather)
experimental/ops/op_common/recursive_executor/algorithm/all_gather.cc展示了完整注册流程:先组装算法树,再注册算法 + 注册执行器(REGISTER_ALG见第 3 节):// 4级串行:Mesh(layer3) -> NHR(layer2) -> NHR(layer1) -> Mesh(layer0) static AlgoExecDesc MakeAllGather4LevelAlgoExecDesc() { TemplateDesc meshDesc{HcclCMDType::HCCL_CMD_ALLGATHER, HcclAlgoType::HCCL_ALGO_TYPE_FULLMESH}; TemplateDesc nhrDesc{HcclCMDType::HCCL_CMD_ALLGATHER, HcclAlgoType::HCCL_ALGO_TYPE_NHR}; AlgoExecDesc desc; desc.execPolicy = AlgExecPolicy::SEQUENCE; desc.children = { TemplateExecDesc{meshDesc, SUB_COMM_INDEX_3}, TemplateExecDesc{nhrDesc, SUB_COMM_INDEX_2}, TemplateExecDesc{nhrDesc, SUB_COMM_INDEX_1}, TemplateExecDesc{meshDesc, SUB_COMM_INDEX_0}, }; desc.dataSplitRatio = {1, 1, 1, 1}; return desc; } static AlgoDesc MakeAllGather4LevelAlgo() { AlgoDesc algo; algo.hcclCmdType = HcclCMDType::HCCL_CMD_ALLGATHER; algo.engineType = CommEngine::COMM_ENGINE_AICPU; algo.topoMatch = std::make_shared<TopoMatchFourLevel>(); algo.algoExecDesc = MakeAllGather4LevelAlgoExecDesc(); algo.algName = "AicpuAllGatherSequenceMeshNHRNHRMesh"; return algo; } // 注册算法到 AlgSelector + 注册执行器到 CollAlgExecRegistryV2 REGISTER_ALG(HcclCMDType::HCCL_CMD_ALLGATHER, AicpuAllGatherSequenceMeshNHRNHRMesh, MakeAllGather4LevelAlgo());新增算法只需三步:编写工厂函数(组装
AlgoExecDesc树)、复用或新增TopoMatchBaseV2匹配器、追加一行REGISTER_ALG。Selector 与算法表的交互
Selector 的职责不变——根据拓扑层级、数据量、Rank 数等条件选择最优算法。变化在于:对 4 级拓扑直接返回 recursive_executor 算法名字符串(
"AicpuAllGatherSequenceMeshNHRNHRMesh")。后续链路完全复用 src 的字符串路由机制(param.algName),无需 Selector 直接持有算法对象。2. 关键逻辑
2.1 算法组装示例
快速组装算法的关键不是新增 Executor 类,而是复用三种积木:算法模板(
TemplateExecDesc)、执行策略(SEQUENCE/PARALLEL节点)、递归(把一个AlgoExecDesc作为另一个节点的 Child)。最小算法——单层 Mesh AllGather,一棵只有一个算法模板的 Sequence 树:
AlgoExecDesc root { .execPolicy = SEQUENCE, .children = { TemplateExecDesc{allGatherMeshDesc, 0} }, .dataSplitRatio = {1} };两层 Sequence——先 Level 0 执行 Mesh,再 Level 1 执行 NHR:
AlgoExecDesc root { .execPolicy = SEQUENCE, .children = { TemplateExecDesc{allGatherMeshDesc, 0}, TemplateExecDesc{allGatherNhrDesc, 1} }, .dataSplitRatio = {1, 1} };Concurrent——当前重构将 Concurrent 表示为一个
PARALLEL节点,两个算法模板处理不同数据子片、使用不同子通信域并发提交:AlgoExecDesc root { .execPolicy = PARALLEL, .children = { TemplateExecDesc{allGatherMeshDesc, 0}, TemplateExecDesc{allGatherNhrDesc, 1} }, .dataSplitRatio = {1, 1} };嵌套 Parallel——两个 Parallel 阶段串行执行,第一阶段数据不同子片分别沿两个维度扩散,第二阶段交换维度顺序:
auto phase0 = AlgoExecDesc{ PARALLEL, {meshLevel0, nhrLevel1}, {1, 1} }; auto phase1 = AlgoExecDesc{ PARALLEL, {nhrLevel1, meshLevel0}, {1, 1} }; AlgoExecDesc root{ SEQUENCE, {shared(phase0), shared(phase1)}, {1, 1} };AllReduce TwoShot——算法语义是
AllReduce = ReduceScatter → AllGather,直接组合两个已有算法模板,无需专用大 Template:AlgoExecDesc allReduceTwoShot { .execPolicy = SEQUENCE, .children = { TemplateExecDesc{reduceScatterMeshDesc, 0}, TemplateExecDesc{allGatherMeshDesc, 0} }, .dataSplitRatio = {1, 1} };多层 AllReduce 同样只是扩展 Sequence,核心顺序是"沿拓扑逐层 ReduceScatter,再按相反顺序逐层 AllGather":
AlgoExecDesc root { .execPolicy = SEQUENCE, .children = { TemplateExecDesc{reduceScatterMeshDesc, 0}, TemplateExecDesc{reduceScatterNhrDesc, 1}, TemplateExecDesc{allGatherNhrDesc, 1}, TemplateExecDesc{allGatherMeshDesc, 0} }, .dataSplitRatio = {1, 1, 1, 1} };由此可见,3 层→4 层拓扑只需在算法表追加一个算法模板,Executor 无需任何修改。
2.2 Executor 递归编排
功能流程
Executor 是一个与引擎类型、算法命令、Topo 无关的通用对象。
OpsExecutor::Orchestrate实际执行流程如下:flowchart TB Start["executor->Orchestrate(resCtx)"] InitRes["InitRes<br/>从 resCtx 恢复 cclBuffer/线程/channel 表"] Prepare["PrepareOrchestrate<br/>计算 dataCount / maxProcCntPerLoop / loopTimes / dataStride"] Loop{"for loopIdx < loopTimes"} InitDesc["InitAlgoExecDataDesc<br/>初始化 dataOffset/sliceCount/tailCount/ranksForInputDataGroup"] Orche["OrchestrateLoop<br/>递归编排算法树"] Next["offsetCount += processCount"] Done["结束"] Start --> InitRes --> Prepare --> Loop Loop -->|"是"| InitDesc --> Orche --> Next --> Loop Loop -->|"否"| DoneOpsExecutor是AlgoDesc::GetExecutor()创建的通用执行器,构造时从OpParam采集 input/output/root/dataType 等运行时信息。静态结构与动态状态
Executor 同时持有两类信息:
AlgoDescAlgoExecDataDescAlgoExecDataDesc是"数据在算法树某个节点入口处的状态快照",核心字段包括:inputBufferType/outputBufferType/cclBufferType(本阶段输入/输出/CCL Buffer 来源)、dataOffset(Loop 在用户内存中的起始偏移)、sliceOffset/sliceCount(Parallel 子片偏移与数量)、dataStride/scratchStride(用户内存/CCL Buffer 相邻 Slot 间距)、ranksForInputDataGroup/ranksForOutputDataGroup(当前 Buffer 中各 Slot 的 Owner)。递归编排流程
OrchestrateLoop是 Executor 的核心递归函数,统一处理 Sequence、Parallel 和嵌套组合:flowchart TB Entry["OrchestrateLoop(algoExecDesc, algoExecDataDesc)"] Init["初始化 children AlgoExecDataDesc<br/>(复用或从 Parent 拷贝)"] Sync1{"PARALLEL 且 children>1?"} PreSync["PreSyncBySubCommMask<br/>并行前同步"] Loop["遍历 children 节点"] Split{"PARALLEL?"} PSplit["UpdateDataSplitParallel<br/>按 dataSplitRatio 切分数据/归属"] SSplit["UpdateDataSplitSequence<br/>传播 ranksForInput/inputBufferType"] IsLeaf{"递归到 Template?"} RunT["RunTemplateDesc<br/>GenTemplateRes + GenTemplateDataParams + KernelRun"] Recurse["OrchestrateLoop<br/>递归子树"] Sync2{"SEQUENCE 且 children>1?"} PreSync1["PreSyncSingleSubComm 串行前同步"] PostSync1["PostSyncSingleSubComm 串行后同步"] Merge["MergeChildrenOutput<br/>合并子节点输出归属"] Sync3{"PARALLEL 且 children>1?"} PostSync["PostSyncBySubCommMask<br/>并行后同步"] Entry --> Init --> Sync1 Sync1 -->|"是"| PreSync --> Loop Sync1 -->|"否"| Loop Loop --> Split Split -->|"是"| PSplit Split -->|"否"| SSplit PSplit --> IsLeaf SSplit --> IsLeaf IsLeaf -->|"是"| Sync2 IsLeaf -->|"否"| Recurse Sync2 -->|"是"| PreSync1 --> RunT --> PostSync1 --> Next{"还有 child?"} Sync2 -->|"否"| RunT --> Next Recurse --> Next Next -->|"是"| Loop Next -->|"否"| Merge --> Sync3 Sync3 -->|"是"| PostSyncSequence 状态传播
Sequence 不仅表示调用顺序,还定义两个状态传递:
Executor 决定每个 Child 的输出位置:非最后 Child 输出到
HCCL_BUFFER,最后 Child 输出到 Parent 目标 Buffer。因此一个三阶段 Sequence 的数据状态为:ranksForOutputData不是调试信息,而是 Sequence 正确连接的必要状态。注意 Buffer 推导(INPUT/CCL/OUTPUT)只决定数据放在哪个 Buffer,与归属组数无关;归属组数由下述消费规则决定。多组归属的消费规则:
UpdateDataSplitSequence把前一 Child 的全部输出归属(可能多组)原样赋给下一 Child,具体如何消费取决于下一 Child 的类型:GenTemplateDataParams强制ranksForInputDataGroup.size() == 1——Template 只能消费单组归属。因此多组归属必须先被一个 PARALLEL 节点消化、或合并归为单组,否则执行期报错。UpdateDataSplitParallel要求输入组数等于子节点数(否则报错),按下标一一对应把 N 组归属分给 N 个 PARALLEL 子节点——多组归属必须匹配后续的 PARALLEL 结构。对应 2.1 节嵌套示例
root = SEQUENCE[phase0(PARALLEL), phase1(PARALLEL)]:phase0 的两个 Child(mesh/nhr)若输出归属不同而保留 2 组,phase1 作为 2 子节点的 PARALLEL 节点恰好按下标一一消费这 2 组归属。设计算法树时需保证多组归属的语义顺序与后续 PARALLEL 子节点的通信域排列一致,必要时调整 Child 排列或插入合并节点。Parallel 数据切分与归属传播
dataSplitRatio是比例(ratio),不是绝对计数。例如{2, 1}表示 Child 0 和 Child 1 按 2:1 比例分配 Parent 的sliceCount。具体计算:childSlice[i] = floor(parentSlice * ratio[i] / sum(ratio)),最后一个 Child 承接整除余量以保证数据不丢失。以{2, 1}且parentSlice = 10为例:Child 0 得floor(10 * 2 / 3) = 6,Child 1 得10 - 6 = 4。余量分配的影响:余量(最多
childrenSize - 1个元素)固定集中到最后一个 Child。数据量远大于 Rank 数时可忽略;但在对称 Concurrent 场景(如{1, 1}且parentSlice为奇数)下,最后一个子通信域会比其它子通信域多处理一个 Slot,造成轻微负载不均。该策略当前不可配置;对均衡敏感的场景,建议在算法构建阶段按实际端口/带宽比例设置dataSplitRatio(使切分与各通信域能力匹配),余量影响随数据量增大自然稀释,必要时可评估扩展余量分散策略。Parallel 的每个 Child 从 Parent 继承大部分状态,但重新计算
sliceCount和sliceOffset:按上述比例切分,Offset 累加前一 Child 覆盖范围。Tail 只传给最后一个 Child。Parallel 切分改变的是 Slot 内部 Offset 和 Count,不改变
stride。归属传播规则:Parent 只有一组输入归属时,每个 Child 处理同一组 Owner 的不同子片;Parent 有多组归属时,各 Child 取得自己的那一组。Children 执行完成后,
MergeChildrenOutput若判定所有 Child 输出归属相同则 Parent 保留一组,否则保留多组供后续消费。OmniPipe 编排(OrchestrateOmniPipeLoop)
OmniPipe(跨层流水)是 2D 网格上"慢轴/快轴按 Step 交替通信以重叠传输时间"的编排方式。旧架构为每个算子维护一个专用执行器(如
InsV2AllGatherOmniPipeExecutor/InsV2AllGatherOmniPipe2dExecutor),内部硬编码 3 层拓扑、逐轴切片与多线程调度。在重构架构中,OmniPipe 不再需要专用执行器:它只是AlgoExecDesc的一种execPolicy(AlgExecPolicy::OMNIPIPE),由通用OpsExecutor解释,两个轴对应算法树的 2 个 Child,复用同一套 Template/CommPlanner 与同步原语。OmniPipe 算法表达
OMNIPIPE节点与PARALLEL不同:它要求恰好 2 个 Child(OmniPipeUpdateEqBWAndReorder对children.size() != 2直接报错),两个 Child 分别代表慢轴 X 和快轴 Y,各 Child 既可以是TemplateExecDesc叶子,也可以是子树(如每个轴各是一棵 Sequence 树):AlgoExecDesc root { .execPolicy = OMNIPIPE, .children = { // 轴 X(慢轴):沿某层子通信域的 AllGather 子树 AlgoExecDesc { SEQUENCE, { meshLevel0, nhrLevel1 }, {1, 1} }, // 轴 Y(快轴):沿另一层子通信域的 AllGather 子树 AlgoExecDesc { SEQUENCE, { nhrLevel1, meshLevel0 }, {1, 1} } }, .dataSplitRatio = {1, 1} };当前约束:
OMNIPIPE策略仅支持ALLREDUCE/ALLGATHER两个命令(OpsExecutor::Orchestrate对其它命令直接报错)。OmniPipe 执行流程
OMNIPIPE顶层走与SEQUENCE/PARALLEL不同的分支(OpsExecutor::Orchestrate):其中
PreCopy以ranksForInputData = {myRank_}把本 Rank 的 Input 拷贝进 CCL Buffer,随后把inputBufferType/outputBufferType置为HCCL_BUFFER,编排结束后PostCopy以ranksForOutputData = [0..rankSize-1]把全量结果写回 Output。InitRes阶段先调用OmniPipeUpdateEqBWAndReorder做一次拓扑预处理,结果按AlgoExecDesc*缓存进omniPipeXYdataMap_,供后续OrchestrateOmniPipeLoop查询:OMIN_MESH_BW=56,CLOS 层OMIN_CLOS_BW/(eqRankSize-1)(OMIN_CLOS_BW=112)。xEqBw ≤ yEqBw)。CalcBandwidth2D(xB, yB, xRankSize, yRankSize, OMIN_MAX_STEP_NUM, steps, scale)按带宽比计算流水步数steps(≤ 5)与scale缩放,连同bandwidthRatio = yB/xB、xEqRankSize/yEqRankSize一起存入OmniPipeXYdata。OrchestrateOmniPipeLoop的编排骨架:flowchart TB Start["OrchestrateOmniPipeLoop(desc, dataDesc)"] Get["查询 omniPipeXYdataMap_[desc]<br/>得到 steps/scale/bandwidthRatio/xEqRankSize/yEqRankSize"] Slice["OmniPipeCalcExecData<br/>为 X/Y 两轴各生成 steps 份 AlgoExecDataDesc<br/>(CalcOmniPipeDataSlice 逐步切分 + 修正逐步 ranksForInputDataGroup)"] Loop{"step i < steps?"} Pre["PreSyncBySubCommMask"] RunX["执行 X 轴 Child (RunTemplateDesc / 递归)"] RunY["执行 Y 轴 Child (RunTemplateDesc / 递归)"] Post["PostSyncBySubCommMask"] Next["i++"] Done["结束"] Start --> Get --> Slice --> Loop Loop -->|"是"| Pre --> RunX --> RunY --> Post --> Next --> Loop Loop -->|"否"| DoneOmniPipeUpdateDataSlice调CalcOmniPipeDataSlice(bandwidthRatio, xRankSize, yRankSize, steps, scale, sliceCount, xSliceCount, ySliceCount)(executor/omnipipe_utils.cc),按"第一步快轴满载、慢轴scale/bw起步,中间步按growth = (xRankSize-1)/bw递推,倒数第二步收口,最后一步斜对角切分"的递推公式算出每步 X/Y 轴各自要处理的sliceCount与sliceOffset,并保证xSliceCount是yRankSize-1的倍数、ySliceCount是xRankSize-1的倍数。ranksForInputDataGroup由CalcPeerAxisRanksForOutput用对侧轴的子通信域重算,保证每步通信的 Peer Owner 集合正确。PreSyncBySubCommMask/PostSyncBySubCommMask做并行前/后同步,与PARALLEL复用同一套线程/notify 机制,无需旧的ntfIdxCtrlToTempXY_等专用通道映射。与旧 OmniPipe 执行器的对比
OMNIPIPE只是execPolicy,通用OpsExecutor解释BuildSubCommAndTempMap硬编码 3 层(level0/1/2)OmniPipeSliceInfo(dataSliceLevel0/1/2)+CalcAGOmniPipeSliceInfo/CalcRSOmniPipeSliceInfo等算子专属函数CalcOmniPipeDataSlice(2D 精简版),去算子/引擎特化分支tempMainThreadsXY_/tempMainThreadsZ_+ 专用 notify 索引PreSyncBySubCommMask/PostSyncBySubCommMask通用同步inputOmniPipeSliceStride等推导ranksForInputDataGroup显式契约 +CalcPeerAxisRanksForOutput对 src 的接入
OMNIPIPE的接入与其他 recursive_executor 算法完全一致,无需额外的 Selector 分支或注册宏:只要在算法表中用REGISTER_ALG注册一棵execPolicy = OMNIPIPE的算法树,第 3 节的param.algName → GetAlgExec → AdaptorExecutor → OpsExecutor链路自动生效。这也正是兼容性考虑第 2 节渐进接入中"Phase 3:全算子覆盖,OmniPipe 流水"这一最后阶段的核心支撑。2.3 ranksForInputData:数据流的主线
定义
ranksForInputData表示当前 Buffer 中按物理布局顺序存在的逻辑 Slot Owner:它表达三个信息:当前有多少个有效 Slot、每个 Slot 的 Owner 是谁、这些 Slot 按什么逻辑顺序排列。它不表示当前 Template 要与哪些 Peer 通信。
CommPlanner 前后变化
CommPlanner 接收
ranksForInputData并计算ranksForOutputData:AllReduce TwoShot 中的归属变化
以 4 Rank 单层 Mesh 为例:
这正是
SEQUENCE[ReduceScatter, AllGather]能够成立的数据契约。2.4 Template 与 CommPlanner
功能流程
Template 和 CommPlanner 不是两套算法实现,而是"执行框架"和"通信计划"的关系:
flowchart TB Params["DataParams + TemplateResource"] Template["Template KernelRun"] Pre["PreCopy<br/>Input → CCL(必要时)"] Plan["CommPlanner<br/>生成通信计划"] Desc["TxRxSlicesList[]<br/>ranksForOutputData"] Execute["Template 执行通信与 LocalReduce"] Post["PostCopy<br/>CCL → Output(必要时)"] Params --> Template Template --> Pre --> Plan --> Desc --> Execute --> PostTemplate 的稳定执行骨架为
PreCopy → RunAlgorithm(CommPlanner) → SendAll/ReduceAll → PostCopy。各阶段按 Buffer 状态退化:输入已在 CCL Buffer 时 PreCopy 跳过;输出给下一 Sequence 阶段时 PostCopy 跳过。分工
CommPlanner 不感知引擎,只返回通信描述列表。这样可复用同一 CommPlanner,同时保留不同 Template 对拷贝、归约和执行方式的差异。
2.5 内存布局与数据流
Slot 排列
CCL Buffer 的内存布局对所有算子类型一致——均按 Rank ID 顺序排列为一维 Slot 数组,Slot i 存放 Rank i 的数据,相邻 Slot 间距为
scratchStride。flowchart LR subgraph CCL["CCL Buffer"] C0["Slot0<br/>owner=0"] C1["Slot1<br/>owner=1"] C2["Slot2<br/>owner=2"] C3["SlotN-1<br/>owner=N-1"] C0 -.->|"scratchStride"| C1 C1 -.->|"scratchStride"| C2 C2 -.->|"scratchStride"| C3 end对于 ReduceScatter,Slot[myAlgRank] 兼任归约目标位(Final Slot),其余 Slot 充当临时接收位(Temp Slot)。通信阶段为纯搬运(
reduceOp强制为RESERVED),归约操作在 PostCopy 阶段统一执行:将各 Temp Slot 的数据 LocalReduce 到 Final Slot(即 Slot[myAlgRank])。当数据量过大时,Executor 按
maxProcCntPerLoop分批处理(多 Loop 循环),CCL Buffer 每个 Loop 只需容纳分批数据。统一内存流向
flowchart LR Input["用户 Input"] -->|"首阶段 PreCopy"| C0["CCL 逻辑快照 0"] C0 -->|"CommPlanner 通信"| C1["CCL 逻辑快照 1"] C1 -->|"CommPlanner 通信"| C2["CCL 逻辑快照 2"] C2 -->|"末阶段 PostCopy"| Output["用户 Output"]各 CCL 快照通常是同一块物理 CCL Buffer 在不同时间的逻辑快照。Sequence 的中间 Child 输出到 CCL Buffer,最后 Child 输出到 Parent 目标 Buffer。
ranksForInputData/ranksForOutputData是连接前后阶段的契约:前一阶段输出归属 = 后一阶段输入归属,PostCopy 按 Owner 将数据写入 Output 的对应 Slot。3. 对接到 src 原流程
本重构不修改、不替换 src 的执行框架,而是作为插件挂到 src 已有的算法路由上。核心思路一句话:recursive_executor 只是一个实现了
InsCollAlgBase接口的"新执行器",通过注册表混入 src 原流程,src 只多了一条 Selector 分支来选中它。接入点在三个层面:编译期(OBJECT 库并入)、选择期(Selector 返回 recursive_executor 算法名)、执行期(
AdaptorExecutor桥接InsCollAlgBase→OpsExecutor)。3.1 编译期接入:hccl_oxc OBJECT 库
experimental/ops/op_common/recursive_executor/CMakeLists.txt将 recursive_executor 编译为 OBJECT 库hccl_oxc,链接进libhccl.so:add_library(hccl_oxc OBJECT ${RE_CORE_SRC}) set_target_properties(hccl_oxc PROPERTIES POSITION_INDEPENDENT_CODE ON) target_compile_definitions(hccl_oxc PRIVATE _GLIBCXX_USE_CXX11_ABI=0) target_include_directories(hccl_oxc PUBLIC ${RE_INCLUDE_LIST}) target_link_libraries(hccl_oxc PUBLIC hccl_compat) if(TARGET hccl) target_link_libraries(hccl PRIVATE hccl_oxc) endif()关键点:
RE_INCLUDE_LIST直接引用src/ops/op_common/...等路径,recursive_executor 与 src 共享同一套OpParam/AlgResourceRequest/AlgResourceCtxSerializable/InsCollAlgBase类型,天然类型一致,无需包装层。src/common/hcomm_dlsym/的符号表 + dlsym(alg_param.h引入的hccl_res_dl.h等即 dlsym 封装),不引入对cann/hcomm的编译期硬依赖。scatter_aicpu_kernel目标时,把RE_CORE_SRC注入该 AICPU 内核,保证 Host/Device 两端都有 recursive_executor 的注册表与执行器。3.2 选择期:Selector 的 4 级拓扑分支
src 的
AllGatherAutoSelector::SelectAicpuAlgo(src/ops/all_gather/selector/all_gather_auto_selector.cc)在多级拓扑分支里新增了 4 级处理:if (topoInfo->topoLevelNums > 1) { // recursive_executor 4 级拓扑算法 if (topoInfo->topoLevelNums == 4) { selectAlgName = "AicpuAllGatherSequenceMeshNHRNHRMesh"; HCCL_INFO("[AllGatherAutoSelector] topoLevelNums=%u, select recursive_executor algorithm [%s]", topoInfo->topoLevelNums, selectAlgName.c_str()); return SelectorStatus::MATCH; } if (...) { // ... 原 3 级逻辑保持不变selectAlgName,此后与 src 其它算法走完全相同的路由,不感知 recursive_executor 存在。3.3 执行期:AdaptorExecutor 桥接层
InsCollAlgBase(src/ops/op_common/executor/executor_v2_base.h)是 src 所有 V2 执行器的统一抽象,src 只通过三个纯虚接口驱动执行器:class InsCollAlgBase { public: virtual HcclResult CalcAlgHierarchyInfo( HcclComm comm, TopoInfoWithNetLayerDetails* topoInfo, AlgHierarchyInfoForAllLevel& algHierarchyInfo) = 0; virtual HcclResult CalcRes( HcclComm comm, const OpParam& param, const TopoInfoWithNetLayerDetails* topoInfo, const AlgHierarchyInfoForAllLevel& algHierarchyInfo, AlgResourceRequest& resourceRequest) = 0; virtual HcclResult Orchestrate(const OpParam& param, const AlgResourceCtxSerializable& resCtx) = 0; };AdaptorExecutorBase(experimental/ops/op_common/recursive_executor/executor/adaptor_executor.h)继承该接口,把三个接口全部转发给 recursive_executor 的OpsExecutor:class AdaptorExecutorBase : public InsCollAlgBase { public: AdaptorExecutorBase() = default; ~AdaptorExecutorBase() override = default; HcclResult CalcAlgHierarchyInfo( HcclComm comm, TopoInfoWithNetLayerDetails* topoInfo, AlgHierarchyInfoForAllLevel& algHierarchyInfo) override; HcclResult CalcRes( HcclComm comm, const OpParam& param, const TopoInfoWithNetLayerDetails* topoInfo, const AlgHierarchyInfoForAllLevel& algHierarchyInfo, AlgResourceRequest& resourceRequest) override; HcclResult Orchestrate(const OpParam& param, const AlgResourceCtxSerializable& resCtx) override; std::string Describe() const override; protected: std::string algName_; // 子类构造时绑定算法名 std::unique_ptr<OpsExecutor> executor_; // 延迟构造的通用执行器 }; template <const char *AlgName> class AdaptorExecutorImpl : public AdaptorExecutorBase { public: AdaptorExecutorImpl() : AdaptorExecutorBase() { algName_ = AlgName; } ~AdaptorExecutorImpl() override = default; };三个接口的转发实现(
experimental/ops/op_common/recursive_executor/executor/adaptor_executor.cc):// 1. 拓扑匹配:从 AlgSelector 取算法对象的 topoMatch 完成匹配(传 algAttrs) HcclResult AdaptorExecutorBase::CalcAlgHierarchyInfo(...) { AlgoDesc alg; if (!AlgSelector::Instance().GetAlgorithm(algName_, alg)) { ... } return alg.topoMatch->MatchTopo(topoInfo, algHierarchyInfo, alg.algAttrs); } // 2. 资源计算:按 param.algName 构造 OpsExecutor,并补一次拓扑匹配 HcclResult AdaptorExecutorBase::CalcRes(HcclComm comm, const OpParam& param, ...) { if (!executor_) { AlgoDesc alg; AlgSelector::Instance().GetAlgorithm(param.algName, alg); OpParam& mutableParam = const_cast<OpParam&>(param); executor_ = alg.GetExecutor(mutableParam); // new OpsExecutor executor_->CalcAlgHierarchyInfo(comm, topoInfo, algHierarchyInfo); } return executor_->CalcRes(comm, resourceRequest); } // 3. 执行:复用或重建 OpsExecutor 后编排 HcclResult AdaptorExecutorBase::Orchestrate(const OpParam& param, const AlgResourceCtxSerializable& resCtx) { if (!executor_) { AlgoDesc algo; AlgSelector::Instance().GetAlgorithm(param.algName, algo); executor_ = algo.GetExecutor(const_cast<OpParam&>(param)); } return executor_->Orchestrate(const_cast<AlgResourceCtxSerializable&>(resCtx)); }要点:
CalcRes/Orchestrate均按param.algName从AlgSelector取回算法定义,与 src 通过param.algName路由执行器的机制完全一致,双端(Host 库 / AICPU 内核)都能重建出同一棵算法树。OpsExecutor生命周期:一个AdaptorExecutor实例内CalcRes创建、Orchestrate复用,避免重复构造开销。3.4 注册宏:把 recursive_executor 执行器挂进 src 注册表
experimental/ops/op_common/recursive_executor/executor/adaptor_executor.h提供两个注册宏:// 注册宏 A:定义外部链接的算法名变量 + 把 AdaptorExecutorImpl 注册进 src 的 CollAlgExecRegistryV2 #define REGISTER_ADAPTOR_EXECUTOR(type, ALG_NAME, EXEC_NAME) \ namespace ops_hccl { const char g_alg_##EXEC_NAME[] = #ALG_NAME; } \ namespace ops_hccl { \ static HcclResult g_oxc_exec_##EXEC_NAME = \ CollAlgExecRegistryV2::Instance().Register(type, \ std::string(#ALG_NAME), \ DefaultExecCreatorV2<AdaptorExecutorImpl<g_alg_##EXEC_NAME>>); \ } // 注册宏 B:算法注册 + 执行器注册一步完成 #define REGISTER_ALG(cmdType, algName, algo) \ REGISTER_ALGORITHM(algName, algo); \ REGISTER_ADAPTOR_EXECUTOR(cmdType, algName, algName)说明:
CollAlgExecRegistryV2(src/ops/op_common/executor/registry/coll_alg_v2_exec_registry.h)是 src 所有 V2 执行器的注册表,DefaultExecCreatorV2<AdaptorExecutorImpl<...>>返回InsCollAlgBase*,与 src 既有REGISTER_EXECUTOR_IMPL等宏走同一Register(type, tag, creator)通道。const char*模板参数要求变量有外部链接,REGISTER_ADAPTOR_EXECUTOR先定义一个extern链接的g_alg_##EXEC_NAME[]字符串,再以它实例化AdaptorExecutorImpl,使模板在编译期绑定算法名。REGISTER_ALG一次完成"算法入AlgSelector"与"执行器入CollAlgExecRegistryV2",两表以同一算法名关联,是 3.2 中 Selector 返回名字能被路由到 recursive_executor 执行器的前提。3.5 完整调用链时序
以 AllGather 4 级拓扑为例,从 API 到 recursive_executor 执行的完整链路:
sequenceDiagram participant API as HcclAllGather (API) participant Sel as AllGatherAutoSelector (src) participant Op as HcclExecOp (src) participant Reg as CollAlgExecRegistryV2 (src) participant Ada as AdaptorExecutor (recursive_executor) participant Exec as OpsExecutor (recursive_executor) participant Tpl as AllGatherMesh/NhrTemplate (recursive_executor) API->>Sel: 算法选择 (SelectAicpuAlgo) Note over Sel: topoLevelNums==4<br/>selectAlgName="AicpuAllGatherSequenceMeshNHRNHRMesh" Sel-->>API: algName API->>Op: HcclExecOp(comm, param, topoInfo, algName, resPack) Op->>Op: param.algName = algName Op->>Reg: GetAlgExec(opType, algName) Reg-->>Op: unique_ptr<AdaptorExecutorImpl> Op->>Ada: CalcAlgHierarchyInfo(comm, topoInfo, algHierarchyInfo) Ada->>Ada: AlgSelector.GetAlgorithm(algName_) -> algo Ada->>Ada: algo.topoMatch->MatchTopo(...) Ada-->>Op: algHierarchyInfo Op->>Ada: CalcRes(comm, param, topoInfo, algHierarchyInfo, resRequest) Ada->>Exec: algo.GetExecutor(param) -> OpsExecutor Ada->>Exec: CalcAlgHierarchyInfo + CalcRes(comm, resRequest) Ada-->>Op: AlgResourceRequest (thread/notify/channel/scratch) Op->>Op: GetAlgResWithEngine: 分配资源、序列化 AlgResourceCtxSerializable Op->>Ada: Orchestrate(param, resCtxHost) Ada->>Exec: Orchestrate(resCtx) Exec->>Tpl: OrchestrateLoop -> RunTemplateDesc -> KernelRun Tpl-->>Exec: ranksForOutputData Exec-->>Ada: HCCL_SUCCESS Ada-->>Op: HCCL_SUCCESS分步说明:
AllGatherAutoSelector::SelectAicpuAlgo对 4 级拓扑返回"AicpuAllGatherSequenceMeshNHRNHRMesh"。HcclExecOp(src/ops/op_common/op_common.cc)把算法名写入param.algName,调CollAlgExecRegistryV2::Instance().GetAlgExec(param.opType, algName)得到AdaptorExecutorImpl(src 侧唯一的一处按名查找,recursive_executor 复用)。HcclGetAlgRes依次调用CalcAlgHierarchyInfo、CalcRes。前者由AdaptorExecutor交给TopoMatchFourLevel做逐层拓扑匹配;后者构造OpsExecutor并递归CalcRes,输出AlgResourceRequest。src 据此经GetAlgResWithEngine分配线程/notify/channel/scratch 并序列化为AlgResourceCtxSerializable。executor->Orchestrate(param, resCtxHost)(CCU/默认引擎在 Host 直接调;AICPU_TS 引擎在HcclAicpuKernelEntranceLaunch下发后,AICPU 内核里同样经CollAlgExecRegistryV2::GetAlgExec+Orchestrate执行)。OpsExecutor::InitRes从resCtx恢复 CCL Buffer/线程/channel 表后,进入PrepareOrchestrate → OrchestrateLoop递归编排。3.6 资源复用与 Host/Device 传递
HcclGetAlgRes先尝试TryReuseResource——若该algTag的资源已创建,直接返回序列化 ctx(isResourceReused=true),Host 侧反序列化后即可Orchestrate,跳过CalcAlgHierarchyInfo/CalcRes。recursive_executor 执行器无需感知该机制,复用判断完全发生在 src。OpsExecutor::InitRes调用RestoreChannelMap,把resCtx.channels(按层展开)重组为rankId → ChannelInfo映射,供 Template 的GenTemplateRes绑定实际通信资源,与 src 执行器的InsCollAlgBase::RestoreChannelMap语义对齐。4. 新增算法指南
基于本方案的 Executor/Template/CommPlanner 三层分离架构,新增一个算法只需按积木组装,无需编写新的执行器类。根据是否需要引入新的通信计划器,分为两个场景,但最终都需完成统一的注册与接入步骤。
场景判断
场景 B:新增 Template
场景 A 无需额外开发,直接进入流程统一。
1. 新增 Template
文件:
experimental/ops/op_common/recursive_executor/template/aicpu/xxx_template.h+.cc继承
AicpuBaseTemplate,实现RunAlgorithm()调用通信计划器生成TxRxSlicesList,按需重写SendAll()/PostCopy()/GetRes()。Template 的执行骨架(PreCopy → RunAlgorithm → SendAll → PostCopy)见 2.4 节。若已有 CommPlanner(如RunMeshAllGather/RunNhrAllGather等)不满足需求,需配套新增对应 CommPlanner 函数(文件置于template/comm_planners/xxx_comm_planners.h+.cc),负责计算通信对端、数据切片和传输方向,输出TxRxSlicesList,不执行通信、不管理资源(分工见 2.4 节)。// xxx_template.h class XxxTemplate : public AicpuBaseTemplate { public: XxxTemplate(u32 myRank, std::vector<u32> ranks, TemplateDesc templateDesc) : AicpuBaseTemplate(myRank, std::move(ranks), templateDesc) { syncAtCopyBoundary_ = false; } protected: HcclResult RunAlgorithm(std::vector<TxRxSlicesList>& txRxSlicesLists, std::vector<u32>& ranksForOutputData) override { return RunXxxCommPlanner(tempAlgParams_, ranks_, myRank_, ranksForOutputData, txRxSlicesLists); } };2. 在 Template 工厂中注册
文件:
experimental/ops/op_common/recursive_executor/template/aicpu/xxx_template.cc#include "aicpu/xxx_template.h" // 注册 Template 类到工厂表 REGISTER_TEMPLATE(HCCL_CMD_ALLGATHER, HCCL_ALGO_TYPE_XXX, XxxTemplate);流程统一
无论是否新增 Template,以下步骤均需执行:
1. 组装算法树并注册
文件:
experimental/ops/op_common/recursive_executor/algorithm/<op>.cc(如all_gather.cc)按 1.1 节的
AlgoDesc三层描述结构和 2.1 节的组装方式编写工厂函数,再用REGISTER_ALG一步完成算法入AlgSelector和执行器入CollAlgExecRegistryV2(注册机制见 1.3 节与 3.4 节):static AlgoDesc MakeAicpuAllGatherSequenceXxxMesh() { TemplateDesc xxxDesc{HcclCMDType::HCCL_CMD_ALLGATHER, HcclAlgoType::HCCL_ALGO_TYPE_XXX}; TemplateDesc meshDesc{HcclCMDType::HCCL_CMD_ALLGATHER, HcclAlgoType::HCCL_ALGO_TYPE_FULLMESH}; AlgoExecDesc desc; desc.execPolicy = AlgExecPolicy::SEQUENCE; desc.children = { TemplateExecDesc{xxxDesc, SUB_COMM_INDEX_1}, TemplateExecDesc{meshDesc, SUB_COMM_INDEX_0}, }; desc.dataSplitRatio = {1, 1}; AlgoDesc algo; algo.hcclCmdType = HcclCMDType::HCCL_CMD_ALLGATHER; algo.engineType = CommEngine::COMM_ENGINE_AICPU; algo.topoMatch = std::make_shared<TopoMatchFourLevel>(); algo.algoExecDesc = desc; algo.algName = "AicpuAllGatherSequenceXxxMesh"; return algo; } REGISTER_ALG( HcclCMDType::HCCL_CMD_ALLGATHER, AicpuAllGatherSequenceXxxMesh, MakeAicpuAllGatherSequenceXxxMesh());2. 更新 CMakeLists.txt(仅新增了源文件时)
文件:
experimental/ops/op_common/recursive_executor/CMakeLists.txtset(RE_CORE_SRC # ... 已有文件 ... template/aicpu/xxx_template.cc template/comm_planners/xxx_comm_planners.cc )3. 更新 src 侧 Selector
文件:
src/ops/<op>/selector/<op>_auto_selector.cc(如all_gather_auto_selector.cc)新增算法选择分支(接入机制见 3.2 节):
if (topoInfo->topoLevelNums == TOPO_LEVEL_NUM_4) { selectAlgName = "AicpuAllGatherSequenceXxxMesh"; return SelectorStatus::MATCH; }交付件清单
template/comm_planners/xxx_comm_planners.h+.cctemplate/aicpu/xxx_template.h+.cctemplate/template_factory.h中新增分支algorithm/<op>.cc中新增工厂函数 +REGISTER_ALGsrc/ops/<op>/selector/<op>_auto_selector.cc中新增拓扑分支场景 A 下最快只需修改 3 个文件(算法注册 + CMakeLists + Selector),新增一个叶子节点级别的算法即可生效。这正是 2.1 节"3 层→4 层拓扑只需追加一个算法模板"的具体体现。
测试方案
端到端算法测试覆盖以下场景:
SEQUENCE[PARALLEL, PARALLEL])的端到端正确性——前一 PARALLEL 输出多组归属时,后一 PARALLEL 按下标一一消费;同时覆盖"多组归属直接后接 Template 叶子应报错"的边界。ratio={1,1,1}、parentSlice不整除(如 10)等余量集中末位场景,验证数据不丢失与结果正确。RS → AG组合,多层 AllReduce。bash build.sh -u跑 UT,确保 src 既有用例不受影响。recursive_executor 自身 UT 覆盖omnipipe_utils、data_ops、comm_planners、algo_desc四组(test/ut/recursive_executor/)。风险评估
InsCollAlgBase/OpParam等结构演进导致 recursive_executor 编译失败std::vector/std::map/std::shared_ptr依赖堆分配;AICPU 内核一般允许堆分配,但 CCU 等受限引擎可能限制OrchestrateLoop/OrchestrateOmniPipeLoop按算法树深度递归,深度受"拓扑层数 × 嵌套组合"约束(4 级 + 嵌套约个位数层)REGISTER_ALG在静态初始化期把算法登记进AlgSelector(Meyers 单例,惰性初始化、本身无顺序问题);若算法工厂依赖其它跨翻译单元的 static 全局对象,初始化顺序未定义替代方案
1. 局部重构:只统一 Sequence/Parallel,保留特化 TwoShot/OrderPreserved
只将 Sole/Sequence/Parallel/Concurrent 四种执行器统一为通用执行器,TwoShot 和 OrderPreserved 保留为特化类。
SEQUENCE[ReduceScatter, AllGather],OrderPreserved 的保序逻辑由 Template 层保证而非 Executor 编排,保留特化类反而引入不必要的维护负担。且四层拓扑新增时仍需为 TwoShot 新增层级变体类。否决。2. 代码生成方案:用模板元编程或宏生成执行器类
编写代码生成器,根据编排模式 × 算子矩阵自动生成 53 个执行器类。
3. 在 src/ 原地重构
直接在
src/ops/下改造现有执行器,不引入新目录。src/代码为生产级,合入前需完整验证。采用experimental/ops/op_common/recursive_executor/平行开发、成熟后合入的策略风险更低。否决。开放问题
experimental/ops/op_common/recursive_executor/目前只实现了4 级 AllGather 算法注册路径,后续多算子、多引擎接入时Template/CommPlanner的复用边界仍需进一步验证。评审记录
评审过程在PR评论区进行,详细评审意见请参阅对应的PR评论。