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(
16 0 : const HcclDispatcher dispatcher, std::unique_ptr<TopoMatcher>& topoMatcher)
17 0 : : CollReduceScatterRingZerocopyExecutor(dispatcher, topoMatcher)
18 0 : {}
19 :
20 0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::CalcCommInfo(std::vector<LevelNSubCommTransport>& opTransport)
21 : {
22 : // 调用父类编排函数建链关系计算函数
23 0 : CHK_RET(CollReduceScatterRingZerocopyExecutor::CalcCommInfo(opTransport));
24 : // 额外增加数据交换的建链
25 0 : CHK_RET(CalcExchangeCommInfo(opTransport));
26 0 : return HCCL_SUCCESS;
27 : }
28 :
29 : HcclResult
30 0 : 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(
40 : COMM_COMBINE_ORDER, CommType::COMM_TAG_PARTIAL_MESH_COMBINED, INVALID_VALUE_RANKID, INVALID_VALUE_RANKID, false,
41 0 : false, commTargetUserRankSet);
42 :
43 0 : TransportMemType inputType = TransportMemType::CCL_INPUT;
44 0 : TransportMemType outputType = TransportMemType::CCL_OUTPUT;
45 :
46 0 : CHK_RET(CalcCommPlaneInfo(tag_, commParaInfo, opTransport[COMM_COMBINE_ORDER], inputType, outputType));
47 0 : LevelNSubCommTransport& commTransport = opTransport[COMM_COMBINE_ORDER];
48 0 : for (u32 subCommIndex = 0; subCommIndex < commTransport.size(); subCommIndex++) {
49 0 : for (auto& transportRequest : commTransport[subCommIndex].transportRequests) {
50 : transportRequest.isUsedRdma
51 0 : = (topoAttr_.superPodNum > 1
52 0 : || (static_cast<bool>(topoMatcher_->GetExternalInputInterHccsDisable()) && topoAttr_.serverNum > 1));
53 : }
54 : }
55 0 : return HCCL_SUCCESS;
56 0 : }
57 :
58 0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::KernelRunInterServerPostProcess(
59 : const OpParam& param, const ExecMem& execMem)
60 : {
61 : // 计算需要交换数据的通信对端
62 0 : u32 remoteRankSend = 0;
63 0 : u32 remoteRankRecv = 0;
64 0 : CHK_RET(CalExchangeRemoteRankForReduceScatter(remoteRankSend, remoteRankRecv));
65 :
66 0 : Stream stream = param.stream;
67 0 : u64 outputMemSize = execMem.outputMem.size();
68 0 : if (remoteRankSend != topoAttr_.userRank && remoteRankRecv != topoAttr_.userRank) { // 需要交换数据
69 : // 获取通信对端的link
70 0 : LINK sendLink;
71 0 : LINK recvLink;
72 0 : CHK_RET(GetTransportForExchange(remoteRankSend, sendLink));
73 0 : CHK_RET(GetTransportForExchange(remoteRankRecv, recvLink));
74 0 : CHK_PTR_NULL(sendLink);
75 0 : CHK_PTR_NULL(recvLink);
76 : // 当通信对端恰好是同server的邻居时,复用Level0的建链,其注册的内存是UserMem,需要特殊处理
77 : // 否则,在CommCombineOrder上建链,其注册内存是CCL Buffer
78 0 : if (IsLevel0Neighbor(remoteRankSend, level0RankSize_)) {
79 : // ccl in -> user in
80 0 : u64 memOffset = (level1Rank_ * level2RankSize_ + level2Rank_) * outputMemSize;
81 0 : DeviceMem dstMem = DeviceMem::create(static_cast<u8*>(param.inputPtr), execMem.outputMem.size());
82 0 : DeviceMem srcMem = execMem.inputMem.range(memOffset, outputMemSize);
83 0 : CHK_SMART_PTR_NULL(dstMem);
84 0 : HcclResult ret = HcclD2DMemcpyAsync(dispatcher_, dstMem, srcMem, stream);
85 0 : CHK_PRT_RET(
86 : ret != HCCL_SUCCESS,
87 : HCCL_ERROR(
88 : "[CollReduceScatterRingZerocopyExchangeExecutor][ExchangeData]ReduceScatter double "
89 : "ring memcpy Failed, Offset[%llu], Size[%llu]",
90 : memOffset, outputMemSize),
91 : ret);
92 : // user in send to remote user out
93 0 : recvLink->TxAck(stream);
94 0 : sendLink->RxAck(stream);
95 0 : sendLink->TxAsync(UserMemType::OUTPUT_MEM, 0, param.inputPtr, outputMemSize, stream);
96 0 : if (IsLevel0Neighbor(remoteRankRecv, level0RankSize_)) {
97 0 : recvLink->RxAsync(UserMemType::INPUT_MEM, 0, execMem.outputPtr, outputMemSize, stream);
98 : } else {
99 0 : u32 remoteLevel1Index = remoteRankRecv % (level0RankSize_ * level1RankSize_) / level0RankSize_;
100 0 : u32 remoteLevel2Index = remoteRankRecv / level0RankSize_ / level1RankSize_;
101 0 : u64 rxSrcOffset = (remoteLevel1Index * level2RankSize_ + remoteLevel2Index) * outputMemSize;
102 0 : recvLink->RxAsync(UserMemType::INPUT_MEM, rxSrcOffset, execMem.outputMem.ptr(), outputMemSize, stream);
103 : }
104 0 : } else {
105 0 : recvLink->TxAck(stream);
106 0 : sendLink->RxAck(stream);
107 0 : u64 txDstOffset = (level1Rank_ * level2RankSize_ + level2Rank_) * outputMemSize;
108 0 : sendLink->TxAsync(
109 0 : UserMemType::OUTPUT_MEM, 0, static_cast<u8*>(execMem.inputMem.ptr()) + txDstOffset, outputMemSize,
110 : stream);
111 0 : if (IsLevel0Neighbor(remoteRankRecv, level0RankSize_)) {
112 0 : recvLink->RxAsync(UserMemType::INPUT_MEM, 0, execMem.outputPtr, outputMemSize, stream);
113 : } else {
114 0 : u32 remoteLevel1Index = remoteRankRecv % (level0RankSize_ * level1RankSize_) / level0RankSize_;
115 0 : u32 remoteLevel2Index = remoteRankRecv / level0RankSize_ / level1RankSize_;
116 0 : u64 rxSrcOffset = (remoteLevel1Index * level2RankSize_ + remoteLevel2Index) * outputMemSize;
117 0 : recvLink->RxAsync(UserMemType::INPUT_MEM, rxSrcOffset, execMem.outputMem.ptr(), outputMemSize, stream);
118 : }
119 : }
120 : // 交换数据的两端之间Barrier,确认收发完成
121 0 : CHK_RET(recvLink->TxAck(stream));
122 0 : CHK_RET(sendLink->RxAck(stream));
123 0 : CHK_RET(sendLink->TxDataSignal(stream));
124 0 : CHK_RET(recvLink->RxDataSignal(stream));
125 0 : } else { // 不需要交换数据,将数据从ccl in拷到ccl out
126 0 : u64 memOffset = (level1Rank_ * level2RankSize_ + level2Rank_) * outputMemSize;
127 0 : DeviceMem dstMem = execMem.outputMem;
128 0 : DeviceMem srcMem = execMem.inputMem.range(memOffset, outputMemSize);
129 0 : CHK_RET(HcclD2DMemcpyAsync(dispatcher_, dstMem, srcMem, stream));
130 0 : }
131 :
132 : // 如果不需要交换数据,或者收端不是邻居,那么结果数据还在CCL Out上,需要搬到User Out上去
133 0 : if ((remoteRankRecv == topoAttr_.userRank) || !IsLevel0Neighbor(remoteRankRecv, level0RankSize_)) {
134 0 : DeviceMem srcMem = execMem.outputMem;
135 0 : DeviceMem dstMem = DeviceMem::create(static_cast<u8*>(execMem.outputPtr), execMem.outputMem.size());
136 0 : CHK_RET(HcclD2DMemcpyAsync(dispatcher_, dstMem, srcMem, stream));
137 0 : }
138 :
139 0 : return HCCL_SUCCESS;
140 0 : }
141 :
142 0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::CalcLevel0DataSlices(
143 : const OpParam& param, const ExecMem& execMem, std::vector<Slice>& dataSegsSlice)
144 : {
145 0 : return CalcIntraServerDataSlicesContinuous(
146 0 : param, execMem, level0RankSize_, level1RankSize_, level2RankSize_, dataSegsSlice);
147 : }
148 :
149 0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::KernelRunInterServerPreProcess(
150 : const OpParam& param, const ExecMem& execMem)
151 : {
152 0 : HCCL_CONFIG_INFO(
153 : HCCL_ALG, "[CollReduceScatterRingZerocopyExchangeExecutor][KernelRun] userRank[%u] starts.",
154 : topoAttr_.userRank);
155 0 : u32 unitSize = 0;
156 0 : CHK_RET(SalGetDataTypeSize(param.DataDes.dataType, unitSize));
157 :
158 0 : DeviceMem dstMem;
159 0 : DeviceMem srcMem;
160 0 : u64 curSize = execMem.outputMem.size();
161 0 : Stream stream = param.stream;
162 0 : for (u32 i = 0; i < level1RankSize_ * level2RankSize_; i++) {
163 : // 拷贝input上每个slice的数据到中转内存,源端每个slice的size固定为output的size
164 0 : dstMem = execMem.inputMem.range(i * curSize, curSize);
165 0 : srcMem = DeviceMem::create(
166 0 : static_cast<u8*>(execMem.inputPtr)
167 0 : + param.DataDes.count * unitSize * level1RankSize_ * level2RankSize_ * level0Rank_
168 0 : + param.DataDes.count * unitSize * i,
169 0 : curSize);
170 0 : CHK_RET(HcclD2DMemcpyAsync(dispatcher_, dstMem, srcMem, stream));
171 : }
172 0 : return HCCL_SUCCESS;
173 0 : }
174 :
175 : REGISTER_EXEC(
176 : "ReduceScatterRingZerocopyExchangeExecutor", ReduceScatterRingZerocopy,
177 : CollReduceScatterRingZerocopyExchangeExecutor);
178 : } // namespace hccl
|