Line data Source code
1 : /**
2 : * Copyright (c) 2025 Huawei Technologies Co., Ltd.
3 : * This program is free software, you can redistribute it and/or modify it under the terms and conditions of
4 : * CANN Open Software License Agreement Version 2.0 (the "License").
5 : * Please refer to the License for details. You may not use this file except in compliance with the License.
6 : * THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND, EITHER EXPRESS OR IMPLIED,
7 : * INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, MERCHANTABILITY, OR FITNESS FOR A PARTICULAR PURPOSE.
8 : * See LICENSE in the root of the software repository for the full text of the License.
9 : */
10 :
11 : #include "coll_reduce_scatter_ring_zerocopy_exchange_executor.h"
12 :
13 : namespace hccl {
14 :
15 0 : CollReduceScatterRingZerocopyExchangeExecutor::CollReduceScatterRingZerocopyExchangeExecutor(const HcclDispatcher dispatcher,
16 0 : std::unique_ptr<TopoMatcher> &topoMatcher)
17 0 : : CollReduceScatterRingZerocopyExecutor(dispatcher, topoMatcher)
18 : {
19 0 : }
20 :
21 0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::CalcCommInfo(std::vector<LevelNSubCommTransport>& opTransport)
22 : {
23 : // 调用父类编排函数建链关系计算函数
24 0 : CHK_RET(CollReduceScatterRingZerocopyExecutor::CalcCommInfo(opTransport));
25 : // 额外增加数据交换的建链
26 0 : CHK_RET(CalcExchangeCommInfo(opTransport));
27 0 : return HCCL_SUCCESS;
28 : }
29 :
30 0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::CalcExchangeCommInfo(std::vector<LevelNSubCommTransport>& opTransport)
31 : {
32 0 : std::set<u32> commTargetUserRankSet;
33 0 : u32 remoteRankSend = 0;
34 0 : u32 remoteRankRecv = 0;
35 :
36 0 : CHK_RET(CalExchangeRemoteRankForReduceScatter(remoteRankSend, remoteRankRecv));
37 0 : commTargetUserRankSet.insert(remoteRankSend);
38 0 : commTargetUserRankSet.insert(remoteRankRecv);
39 : CommParaInfo commParaInfo(COMM_COMBINE_ORDER, CommType::COMM_TAG_PARTIAL_MESH_COMBINED, INVALID_VALUE_RANKID,
40 0 : INVALID_VALUE_RANKID, false, false, commTargetUserRankSet);
41 :
42 0 : TransportMemType inputType = TransportMemType::CCL_INPUT;
43 0 : TransportMemType outputType = TransportMemType::CCL_OUTPUT;
44 :
45 0 : CHK_RET(CalcCommPlaneInfo(tag_, commParaInfo, opTransport[COMM_COMBINE_ORDER], inputType, outputType));
46 0 : LevelNSubCommTransport &commTransport = opTransport[COMM_COMBINE_ORDER];
47 0 : for (u32 subCommIndex = 0; subCommIndex < commTransport.size(); subCommIndex++) {
48 0 : for (auto &transportRequest : commTransport[subCommIndex].transportRequests) {
49 0 : transportRequest.isUsedRdma = (topoAttr_.superPodNum > 1 ||
50 0 : (static_cast<bool>(topoMatcher_->GetExternalInputInterHccsDisable()) && topoAttr_.serverNum > 1));
51 : }
52 : }
53 0 : return HCCL_SUCCESS;
54 0 : }
55 :
56 0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::KernelRunInterServerPostProcess(const OpParam ¶m, const ExecMem &execMem)
57 : {
58 : // 计算需要交换数据的通信对端
59 0 : u32 remoteRankSend = 0;
60 0 : u32 remoteRankRecv = 0;
61 0 : CHK_RET(CalExchangeRemoteRankForReduceScatter(remoteRankSend, remoteRankRecv));
62 :
63 0 : Stream stream = param.stream;
64 0 : u64 outputMemSize = execMem.outputMem.size();
65 0 : if (remoteRankSend != topoAttr_.userRank && remoteRankRecv != topoAttr_.userRank) { // 需要交换数据
66 : // 获取通信对端的link
67 0 : LINK sendLink;
68 0 : LINK recvLink;
69 0 : CHK_RET(GetTransportForExchange(remoteRankSend, sendLink));
70 0 : CHK_RET(GetTransportForExchange(remoteRankRecv, recvLink));
71 0 : CHK_PTR_NULL(sendLink);
72 0 : CHK_PTR_NULL(recvLink);
73 : // 当通信对端恰好是同server的邻居时,复用Level0的建链,其注册的内存是UserMem,需要特殊处理
74 : // 否则,在CommCombineOrder上建链,其注册内存是CCL Buffer
75 0 : if (IsLevel0Neighbor(remoteRankSend, level0RankSize_)) {
76 : // ccl in -> user in
77 0 : u64 memOffset = (level1Rank_ * level2RankSize_ + level2Rank_) * outputMemSize;
78 0 : DeviceMem dstMem = DeviceMem::create(static_cast<u8 *>(param.inputPtr), execMem.outputMem.size());
79 0 : DeviceMem srcMem = execMem.inputMem.range(memOffset, outputMemSize);
80 0 : CHK_SMART_PTR_NULL(dstMem);
81 0 : HcclResult ret = HcclD2DMemcpyAsync(dispatcher_, dstMem, srcMem, stream);
82 0 : CHK_PRT_RET(ret != HCCL_SUCCESS,
83 : HCCL_ERROR("[CollReduceScatterRingZerocopyExchangeExecutor][ExchangeData]ReduceScatter double "
84 : "ring memcpy Failed, Offset[%llu], Size[%llu]", memOffset, outputMemSize), ret);
85 : // user in send to remote user out
86 0 : recvLink->TxAck(stream);
87 0 : sendLink->RxAck(stream);
88 0 : sendLink->TxAsync(UserMemType::OUTPUT_MEM, 0, param.inputPtr, outputMemSize, stream);
89 0 : if (IsLevel0Neighbor(remoteRankRecv, level0RankSize_)) {
90 0 : recvLink->RxAsync(UserMemType::INPUT_MEM, 0, execMem.outputPtr, outputMemSize, stream);
91 : } else {
92 0 : u32 remoteLevel1Index = remoteRankRecv % (level0RankSize_ * level1RankSize_) / level0RankSize_;
93 0 : u32 remoteLevel2Index = remoteRankRecv / level0RankSize_ / level1RankSize_;
94 0 : u64 rxSrcOffset = (remoteLevel1Index * level2RankSize_ + remoteLevel2Index) * outputMemSize;
95 0 : recvLink->RxAsync(UserMemType::INPUT_MEM, rxSrcOffset, execMem.outputMem.ptr(), outputMemSize, stream);
96 : }
97 0 : } else {
98 0 : recvLink->TxAck(stream);
99 0 : sendLink->RxAck(stream);
100 0 : u64 txDstOffset = (level1Rank_ * level2RankSize_ + level2Rank_) * outputMemSize;
101 0 : sendLink->TxAsync(UserMemType::OUTPUT_MEM, 0, static_cast<u8 *>(execMem.inputMem.ptr()) + txDstOffset, outputMemSize, stream);
102 0 : if (IsLevel0Neighbor(remoteRankRecv, level0RankSize_)) {
103 0 : recvLink->RxAsync(UserMemType::INPUT_MEM, 0, execMem.outputPtr, outputMemSize, stream);
104 : } else {
105 0 : u32 remoteLevel1Index = remoteRankRecv % (level0RankSize_ * level1RankSize_) / level0RankSize_;
106 0 : u32 remoteLevel2Index = remoteRankRecv / level0RankSize_ / level1RankSize_;
107 0 : u64 rxSrcOffset = (remoteLevel1Index * level2RankSize_ + remoteLevel2Index) * outputMemSize;
108 0 : recvLink->RxAsync(UserMemType::INPUT_MEM, rxSrcOffset, execMem.outputMem.ptr(), outputMemSize, stream);
109 : }
110 : }
111 : // 交换数据的两端之间Barrier,确认收发完成
112 0 : CHK_RET(recvLink->TxAck(stream));
113 0 : CHK_RET(sendLink->RxAck(stream));
114 0 : CHK_RET(sendLink->TxDataSignal(stream));
115 0 : CHK_RET(recvLink->RxDataSignal(stream));
116 0 : } else { // 不需要交换数据,将数据从ccl in拷到ccl out
117 0 : u64 memOffset = (level1Rank_ * level2RankSize_ + level2Rank_) * outputMemSize;
118 0 : DeviceMem dstMem = execMem.outputMem;
119 0 : DeviceMem srcMem = execMem.inputMem.range(memOffset, outputMemSize);
120 0 : CHK_RET(HcclD2DMemcpyAsync(dispatcher_, dstMem, srcMem, stream));
121 0 : }
122 :
123 : // 如果不需要交换数据,或者收端不是邻居,那么结果数据还在CCL Out上,需要搬到User Out上去
124 0 : if ((remoteRankRecv == topoAttr_.userRank) || !IsLevel0Neighbor(remoteRankRecv, level0RankSize_)) {
125 0 : DeviceMem srcMem = execMem.outputMem;
126 0 : DeviceMem dstMem = DeviceMem::create(static_cast<u8 *>(execMem.outputPtr), execMem.outputMem.size());
127 0 : CHK_RET(HcclD2DMemcpyAsync(dispatcher_, dstMem, srcMem, stream));
128 0 : }
129 :
130 0 : return HCCL_SUCCESS;
131 0 : }
132 :
133 0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::CalcLevel0DataSlices(const OpParam ¶m, const ExecMem &execMem,
134 : std::vector<Slice> &dataSegsSlice)
135 : {
136 0 : return CalcIntraServerDataSlicesContinuous(param, execMem,
137 0 : level0RankSize_, level1RankSize_, level2RankSize_, dataSegsSlice);
138 : }
139 :
140 0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::KernelRunInterServerPreProcess(const OpParam ¶m,
141 : const ExecMem &execMem)
142 : {
143 0 : HCCL_CONFIG_INFO(HCCL_ALG,
144 : "[CollReduceScatterRingZerocopyExchangeExecutor][KernelRun] userRank[%u] starts.", topoAttr_.userRank);
145 0 : u32 unitSize = 0;
146 0 : CHK_RET(SalGetDataTypeSize(param.DataDes.dataType, unitSize));
147 :
148 0 : DeviceMem dstMem;
149 0 : DeviceMem srcMem;
150 0 : u64 curSize = execMem.outputMem.size();
151 0 : Stream stream = param.stream;
152 0 : for (u32 i = 0; i < level1RankSize_ * level2RankSize_; i++) {
153 : // 拷贝input上每个slice的数据到中转内存,源端每个slice的size固定为output的size
154 0 : dstMem = execMem.inputMem.range(i * curSize, curSize);
155 0 : srcMem = DeviceMem::create(static_cast<u8 *>(execMem.inputPtr)
156 0 : + param.DataDes.count * unitSize * level1RankSize_ * level2RankSize_ * level0Rank_
157 0 : + param.DataDes.count * unitSize * i,
158 0 : curSize);
159 0 : CHK_RET(HcclD2DMemcpyAsync(dispatcher_, dstMem, srcMem, stream));
160 : }
161 0 : return HCCL_SUCCESS;
162 0 : }
163 :
164 : REGISTER_EXEC("ReduceScatterRingZerocopyExchangeExecutor", ReduceScatterRingZerocopy, CollReduceScatterRingZerocopyExchangeExecutor);
165 : }
|