Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
PaddlePaddle
Paddle
提交
158b7ae5
P
Paddle
项目概览
PaddlePaddle
/
Paddle
1 年多 前同步成功
通知
2302
Star
20931
Fork
5422
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
1423
列表
看板
标记
里程碑
合并请求
543
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
P
Paddle
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
1,423
Issue
1,423
列表
看板
标记
里程碑
合并请求
543
合并请求
543
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
未验证
提交
158b7ae5
编写于
6月 27, 2023
作者:
TaoTao Li
提交者:
GitHub
6月 27, 2023
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
add all_to_all phi operator (#54797)
* add all_to_all phi operator, kernel, api * add all_to_all ut * tinyfix
上级
70288456
变更
8
隐藏空白更改
内联
并排
Showing
8 changed file
with
306 addition
and
1 deletion
+306
-1
paddle/phi/api/yaml/static_ops.yaml
paddle/phi/api/yaml/static_ops.yaml
+10
-0
paddle/phi/infermeta/unary.cc
paddle/phi/infermeta/unary.cc
+7
-0
paddle/phi/infermeta/unary.h
paddle/phi/infermeta/unary.h
+2
-0
paddle/phi/kernels/all_to_all_kernel.h
paddle/phi/kernels/all_to_all_kernel.h
+26
-0
paddle/phi/kernels/cpu/all_to_all_kernel.cc
paddle/phi/kernels/cpu/all_to_all_kernel.cc
+43
-0
paddle/phi/kernels/gpu/all_to_all_kernel.cu
paddle/phi/kernels/gpu/all_to_all_kernel.cu
+114
-0
test/collective/collective_alltoall_api.py
test/collective/collective_alltoall_api.py
+89
-1
test/collective/test_collective_alltoall_api.py
test/collective/test_collective_alltoall_api.py
+15
-0
未找到文件。
paddle/phi/api/yaml/static_ops.yaml
浏览文件 @
158b7ae5
...
...
@@ -26,6 +26,16 @@
func
:
all_reduce
param
:
[
x
,
reduce_type
]
-
op
:
all_to_all
args
:
(Tensor x, int ring_id = 0)
output
:
Tensor(out)
infer_meta
:
func
:
AllToAllInferMeta
param
:
[
x
]
kernel
:
func
:
all_to_all
param
:
[
x
]
-
op
:
amax
args
:
(Tensor x, IntArray axis={0}, bool keepdim=false, bool reduce_all=false, int in_dtype=-1, int out_dtype=-1)
output
:
Tensor(out)
...
...
paddle/phi/infermeta/unary.cc
浏览文件 @
158b7ae5
...
...
@@ -136,6 +136,13 @@ void AllReduceInferMeta(const MetaTensor& x, MetaTensor* out) {
out
->
set_dims
(
x
.
dims
());
}
void
AllToAllInferMeta
(
const
MetaTensor
&
x
,
MetaTensor
*
out
)
{
auto
dim
=
x
.
dims
();
if
(
dim
[
0
]
<
0
)
dim
[
0
]
=
-
1
;
out
->
set_dtype
(
x
.
dtype
());
out
->
set_dims
(
dim
);
}
void
ArgMinMaxInferMeta
(
const
MetaTensor
&
x
,
const
Scalar
&
axis
,
bool
keepdims
,
...
...
paddle/phi/infermeta/unary.h
浏览文件 @
158b7ae5
...
...
@@ -43,6 +43,8 @@ void AllGatherInferMeta(const MetaTensor& x, int nranks, MetaTensor* out);
void
AllReduceInferMeta
(
const
MetaTensor
&
x
,
MetaTensor
*
out
);
void
AllToAllInferMeta
(
const
MetaTensor
&
x
,
MetaTensor
*
out
);
void
ArgMinMaxInferMeta
(
const
MetaTensor
&
x
,
const
Scalar
&
axis
,
bool
keepdims
,
...
...
paddle/phi/kernels/all_to_all_kernel.h
0 → 100644
浏览文件 @
158b7ae5
// Copyright (c) 2023 PaddlePaddle Authors. All Rights Reserved.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
#pragma once
#include "paddle/phi/core/dense_tensor.h"
namespace
phi
{
template
<
typename
T
,
typename
Context
>
void
AllToAllKernel
(
const
Context
&
dev_ctx
,
const
DenseTensor
&
x
,
DenseTensor
*
out
);
}
// namespace phi
paddle/phi/kernels/cpu/all_to_all_kernel.cc
0 → 100644
浏览文件 @
158b7ae5
// Copyright (c) 2023 PaddlePaddle Authors. All Rights Reserved.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
#include "paddle/phi/kernels/all_to_all_kernel.h"
#include "paddle/phi/backends/all_context.h"
#include "paddle/phi/core/kernel_registry.h"
namespace
phi
{
template
<
typename
T
,
typename
Context
>
void
AllToAllKernel
(
const
Context
&
dev_ctx
UNUSED
,
const
DenseTensor
&
x
UNUSED
,
DenseTensor
*
out
UNUSED
)
{
PADDLE_THROW
(
errors
::
Unimplemented
(
"Unimplemented cpu kernel for all_to_all."
));
}
}
// namespace phi
PD_REGISTER_KERNEL
(
all_to_all
,
CPU
,
ALL_LAYOUT
,
phi
::
AllToAllKernel
,
float
,
double
,
int
,
bool
,
int8_t
,
uint8_t
,
int64_t
,
phi
::
dtype
::
float16
)
{}
paddle/phi/kernels/gpu/all_to_all_kernel.cu
0 → 100644
浏览文件 @
158b7ae5
// Copyright (c) 2023 PaddlePaddle Authors. All Rights Reserved.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
#include "paddle/phi/kernels/all_to_all_kernel.h"
#include "glog/logging.h"
#include "paddle/phi/backends/all_context.h"
#include "paddle/phi/core/kernel_registry.h"
#if defined(PADDLE_WITH_NCCL) || defined(PADDLE_WITH_RCCL)
#include "paddle/phi/core/distributed/nccl_comm_context.h"
#include "paddle/phi/core/distributed/utils.h"
#endif
namespace
phi
{
template
<
typename
T
,
typename
Context
>
void
AllToAllKernel
(
const
Context
&
dev_ctx
,
const
DenseTensor
&
x
,
DenseTensor
*
out
)
{
#if defined(PADDLE_WITH_NCCL) || defined(PADDLE_WITH_RCCL)
#if NCCL_VERSION_CODE >= 2703
auto
x_dims
=
x
.
dims
();
out
->
Resize
(
x_dims
);
dev_ctx
.
template
Alloc
<
T
>(
out
);
auto
comm_ctx
=
static_cast
<
distributed
::
NCCLCommContext
*>
(
dev_ctx
.
GetCommContext
());
PADDLE_ENFORCE_NE
(
comm_ctx
,
nullptr
,
errors
::
Unavailable
(
"NCCLCommContext is nullptr, collective op should "
"has ring_id attr."
));
gpuStream_t
stream
=
dev_ctx
.
stream
();
PADDLE_ENFORCE_NOT_NULL
(
stream
,
errors
::
NotFound
(
"Should initialize NCCL firstly."
));
int
nranks
=
comm_ctx
->
GetSize
();
int
send_numel
=
x
.
numel
()
/
nranks
;
size_t
offset
=
0
;
PADDLE_ENFORCE_EQ
(
x_dims
[
0
]
%
nranks
,
0
,
errors
::
InvalidArgument
(
"The first dimension size (%d) of the input tensor must be "
"divisible by the number of ranks (%d)."
,
x_dims
[
0
],
nranks
));
comm_ctx
->
GroupStart
();
const
auto
*
send_buf
=
x
.
data
<
T
>
();
auto
*
recv_buf
=
out
->
data
<
T
>
();
for
(
auto
i
=
0
;
i
<
nranks
;
++
i
)
{
auto
send_buf
=
phi
::
distributed
::
GetPartialTensor
(
x
,
offset
,
send_numel
);
comm_ctx
->
Send
(
send_buf
,
send_numel
,
i
,
stream
);
auto
recv_buf
=
phi
::
distributed
::
GetPartialTensor
(
*
out
,
offset
,
send_numel
);
comm_ctx
->
Recv
(
&
recv_buf
,
send_numel
,
i
,
stream
);
offset
+=
send_numel
;
}
comm_ctx
->
GroupEnd
();
#else
PADDLE_THROW
(
platform
::
errors
::
Unavailable
(
"NCCL version >= 2.7.3 is needed."
));
#endif
#else
PADDLE_THROW
(
errors
::
PreconditionNotMet
(
"PaddlePaddle should compile with GPU."
));
#endif
}
}
// namespace phi
#if NCCL_VERSION_CODE >= 21000 && CUDA_VERSION >= 11000
PD_REGISTER_KERNEL
(
all_to_all
,
GPU
,
ALL_LAYOUT
,
phi
::
AllToAllKernel
,
float
,
double
,
int
,
int8_t
,
uint8_t
,
int64_t
,
bool
,
phi
::
dtype
::
bfloat16
,
phi
::
dtype
::
float16
)
{}
#else
PD_REGISTER_KERNEL
(
all_to_all
,
GPU
,
ALL_LAYOUT
,
phi
::
AllToAllKernel
,
float
,
double
,
int
,
int8_t
,
uint8_t
,
int64_t
,
bool
,
phi
::
dtype
::
float16
)
{}
#endif
test/collective/collective_alltoall_api.py
浏览文件 @
158b7ae5
...
...
@@ -18,11 +18,86 @@ from legacy_test.test_collective_api_base import (
)
import
paddle
from
paddle
import
fluid
import
paddle.distributed
as
dist
from
paddle
import
fluid
,
framework
from
paddle.fluid
import
data_feeder
paddle
.
enable_static
()
def
alltoall_new
(
in_tensor_or_tensor_list
,
out_tensor_or_tensor_list
,
group
=
None
,
sync_op
=
True
,
):
op_type
=
'all_to_all'
ring_id
=
0
if
group
is
None
else
group
.
id
nranks
=
dist
.
get_world_size
()
helper
=
framework
.
LayerHelper
(
op_type
,
**
locals
())
in_tensor
=
in_tensor_or_tensor_list
if
isinstance
(
in_tensor_or_tensor_list
,
list
):
if
len
(
in_tensor_or_tensor_list
)
==
0
:
raise
RuntimeError
(
"The input tensor_list should not be empty."
)
# 0-D use stack/unstack while others use concat/split
if
len
(
in_tensor_or_tensor_list
[
0
].
shape
)
==
0
:
in_tensor
=
paddle
.
stack
(
in_tensor_or_tensor_list
,
axis
=
0
)
else
:
in_tensor
=
paddle
.
concat
(
in_tensor_or_tensor_list
,
axis
=
0
)
out_tensor
=
out_tensor_or_tensor_list
if
isinstance
(
out_tensor_or_tensor_list
,
list
):
if
len
(
out_tensor_or_tensor_list
)
!=
0
:
raise
ValueError
(
"The 'out_tensor_list' for all_to_all "
"must be an empty list."
)
out_tensor
=
helper
.
create_variable_for_type_inference
(
dtype
=
in_tensor
.
dtype
)
data_feeder
.
check_variable_and_dtype
(
in_tensor
,
'in_tensor'
,
[
'float16'
,
'float32'
,
'float64'
,
'int32'
,
'int64'
,
'int8'
,
'uint8'
,
'bool'
,
'uint16'
,
],
'all_to_all'
,
)
helper
.
append_op
(
type
=
op_type
,
inputs
=
{
'x'
:
[
in_tensor
]},
outputs
=
{
'out'
:
[
out_tensor
]},
attrs
=
{
'ring_id'
:
ring_id
,
},
)
# NOTE(liyurui): If the argument `out_tensor_or_tensor_list` is a tensor_list,
# we need to split the result. So we should wait the result of all_to_all
# before split if the communication is not on calc stream.
if
isinstance
(
out_tensor_or_tensor_list
,
list
):
if
not
sync_op
:
dist
.
wait
(
out_tensor
,
use_calc_stream
=
False
)
# 0-D use stack/unstack while others use concat/split
if
len
(
in_tensor_or_tensor_list
[
0
].
shape
)
==
0
:
out_tensor_or_tensor_list
.
extend
(
paddle
.
unstack
(
out_tensor
,
0
))
else
:
out_tensor_or_tensor_list
.
extend
(
paddle
.
split
(
out_tensor
,
nranks
,
0
)
)
return
None
class
TestCollectiveAllToAllAPI
(
TestCollectiveAPIRunnerBase
):
def
__init__
(
self
):
self
.
global_ring_id
=
0
...
...
@@ -38,6 +113,19 @@ class TestCollectiveAllToAllAPI(TestCollectiveAPIRunnerBase):
paddle
.
distributed
.
alltoall
(
tindata
,
tout_data
)
return
tout_data
def
get_model_new
(
self
,
main_prog
,
startup_program
,
rank
,
dtype
=
None
,
reduce_type
=
None
):
with
fluid
.
program_guard
(
main_prog
,
startup_program
):
tindata
=
paddle
.
static
.
data
(
name
=
"tindata"
,
shape
=
[
-
1
,
10
,
1000
],
dtype
=
dtype
)
tindata
.
desc
.
set_need_check_feed
(
False
)
tindata
=
paddle
.
split
(
tindata
,
2
,
axis
=
0
)
tout_data
=
[]
alltoall_new
(
tindata
,
tout_data
)
return
tout_data
if
__name__
==
"__main__"
:
runtime_main
(
TestCollectiveAllToAllAPI
,
"alltoall"
)
test/collective/test_collective_alltoall_api.py
浏览文件 @
158b7ae5
...
...
@@ -28,6 +28,21 @@ class TestCollectiveAllToAllAPI(TestDistBase):
def
test_alltoall_nccl
(
self
):
self
.
check_with_place
(
"collective_alltoall_api.py"
,
"alltoall"
,
"nccl"
)
def
test_alltoall_nccl_with_comm_context
(
self
):
dtypes_to_test
=
[
"float32"
,
]
if
self
.
_nccl_version
>=
21000
:
dtypes_to_test
.
append
(
"bfloat16"
)
for
dtype
in
dtypes_to_test
:
self
.
check_with_place
(
"collective_alltoall_api.py"
,
"alltoall"
,
"nccl"
,
dtype
=
dtype
,
need_envs
=
{
"USE_COMM_CONTEXT"
:
"1"
},
)
def
test_alltoall_nccl_dygraph
(
self
):
dtypes_to_test
=
[
"float16"
,
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录