dist_slice.py 7.1 KB
Newer Older
1
# Copyright (c) 2022 PaddlePaddle Authors. All Rights Reserved.
2
#
3 4 5
# 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
6
#
7
#     http://www.apache.org/licenses/LICENSE-2.0
8
#
9 10 11 12 13 14
# 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.

J
JYChen 已提交
15 16
import paddle

17 18 19 20 21 22 23
from ..utils import compute_compatible_dim_mapping, is_dim_shard
from .common import (
    DistributedOperatorImpl,
    DistributedOperatorImplContainer,
    register_distributed_operator_impl,
    register_distributed_operator_impl_container,
)
24 25 26 27 28
from .dist_default import DistributedDefaultImpl0


class DistributedSlice(DistributedOperatorImplContainer):
    def __init__(self, op_type):
29
        super().__init__(op_type)
30 31 32 33 34 35 36


register_distributed_operator_impl_container(DistributedSlice("slice"))


class DistributedSliceImpl(DistributedOperatorImpl):
    def __init__(self, name):
37
        super().__init__(name)
38 39 40 41 42 43 44
        self._forward_implemented = True
        self._backward_implemented = True

    def is_input_compatible(self, dist_op):
        op_desc = dist_op.serial_op.desc
        op_dist_attr = dist_op.dist_attr
        in_name = op_desc.input('Input')[0]
Z
zhaoyingli 已提交
45
        out_name = op_desc.output('Out')[0]
Z
zhaoyingli 已提交
46 47
        in_var = dist_op.serial_op.block._var_recursive(in_name)
        out_var = dist_op.serial_op.block._var_recursive(out_name)
48 49 50
        axes = op_desc.attr('axes')
        in_dims_mapping = op_dist_attr.get_input_dims_mapping(in_name)
        for axis in axes:
Z
zhaoyingli 已提交
51 52 53 54
            if (
                is_dim_shard(in_dims_mapping[axis])
                and in_var.shape[axis] != out_var.shape[axis]
            ):
55 56 57 58
                return False
        return True

    def is_output_compatible(self, dist_op):
59 60 61 62
        op_desc = dist_op.serial_op.desc
        op_dist_attr = dist_op.dist_attr
        in_name = op_desc.input('Input')[0]
        out_name = op_desc.output('Out')[0]
Z
zhaoyingli 已提交
63 64
        in_var = dist_op.serial_op.block._var_recursive(in_name)
        out_var = dist_op.serial_op.block._var_recursive(out_name)
65 66 67 68 69 70 71 72 73 74
        axes = op_desc.attr('axes')
        decrease_axis = op_desc.attr('decrease_axis')
        in_dims_mapping = op_dist_attr.get_input_dims_mapping(in_name)
        out_dims_mapping = op_dist_attr.get_output_dims_mapping(out_name)

        ref_indices = []
        for i in range(len(in_dims_mapping)):
            if i not in decrease_axis:
                ref_indices.append(i)
        if ref_indices == []:
J
JYChen 已提交
75 76 77 78 79 80 81 82
            # NOTE(zoooo0820): When all axes are decreased, the output will be 1-D
            # with FLAGS_set_to_1d=True.
            if paddle.get_flags('FLAGS_set_to_1d')['FLAGS_set_to_1d']:
                assert len(out_dims_mapping) == 1
                if is_dim_shard(out_dims_mapping[0]):
                    return False
            else:
                assert len(out_dims_mapping) == 0
83 84 85
        else:
            for i in range(len(out_dims_mapping)):
                ref_index = ref_indices[i]
Z
zhaoyingli 已提交
86 87 88 89 90
                if (
                    ref_index in axes
                    and is_dim_shard(out_dims_mapping[i])
                    and in_var.shape[ref_index] != out_var.shape[ref_index]
                ):
91 92
                    return False

93 94 95
        return True

    def is_compatible(self, dist_op):
96 97 98
        if (not self.is_input_compatible(dist_op)) or (
            not self.is_output_compatible(dist_op)
        ):
99 100 101 102 103 104 105 106 107 108
            return False

        op_desc = dist_op.serial_op.desc
        op_dist_attr = dist_op.dist_attr
        in_name = op_desc.input('Input')[0]
        out_name = op_desc.output('Out')[0]
        decrease_axis = op_desc.attr('decrease_axis')
        in_dims_mapping = op_dist_attr.get_input_dims_mapping(in_name)
        out_dims_mapping = op_dist_attr.get_output_dims_mapping(out_name)
        if len(in_dims_mapping) - len(decrease_axis) != 0 and len(
109 110
            out_dims_mapping
        ) != len(in_dims_mapping) - len(decrease_axis):
111 112 113 114 115 116 117 118 119 120 121 122 123 124
            return False

        new_out_dims_mapping = []
        for i in range(len(in_dims_mapping)):
            if i not in decrease_axis:
                new_out_dims_mapping.append(in_dims_mapping[i])
        if new_out_dims_mapping == []:
            new_out_dims_mapping = [-1]
        if new_out_dims_mapping != out_dims_mapping:
            return False

        return True

    def is_auto_compatible(self, dist_op):
125 126 127 128 129
        if (
            (not self.is_input_compatible(dist_op))
            or (not self.is_output_compatible(dist_op))
            or (not self.is_compatible(dist_op))
        ):
130 131 132 133 134 135 136 137 138 139 140 141 142 143 144
            return False

        return True

    def update_dims_mapping(self, dist_op):
        changed = False
        op_desc = dist_op.serial_op.desc
        op_dist_attr = dist_op.dist_attr
        in_name = op_desc.input('Input')[0]
        out_name = op_desc.output('Out')[0]
        decrease_axis = op_desc.attr('decrease_axis')
        in_dims_mapping = op_dist_attr.get_input_dims_mapping(in_name)
        out_dims_mapping = op_dist_attr.get_output_dims_mapping(out_name)

        ref_dims_mapping = []
145
        ref_indices = []
146 147 148
        for i in range(len(in_dims_mapping)):
            if i not in decrease_axis:
                ref_dims_mapping.append(in_dims_mapping[i])
149 150
                ref_indices.append(i)

151
        if ref_dims_mapping == []:
J
JYChen 已提交
152 153 154 155 156
            # NOTE(zoooo0820): When all axes are decreased, the output will be 1-D
            # with FLAGS_set_to_1d=True.
            if paddle.get_flags('FLAGS_set_to_1d')['FLAGS_set_to_1d']:
                ref_dims_mapping = [-1]
                assert ref_dims_mapping[0] == out_dims_mapping[0]
157 158 159 160 161 162
            assert len(ref_dims_mapping) == len(out_dims_mapping)
            changed = False
        else:
            assert len(ref_dims_mapping) == len(out_dims_mapping)
            for i in range(len(out_dims_mapping)):
                compatible_dim_mapping = compute_compatible_dim_mapping(
163 164
                    [out_dims_mapping[i], ref_dims_mapping[i]]
                )
165 166 167 168 169 170 171 172
                if compatible_dim_mapping is None:
                    continue
                if ref_dims_mapping[i] != compatible_dim_mapping:
                    in_dims_mapping[ref_indices[i]] = compatible_dim_mapping
                    changed = True
                if out_dims_mapping[i] != compatible_dim_mapping:
                    out_dims_mapping[i] = compatible_dim_mapping
                    changed = True
173

174 175 176 177
        if changed:
            op_dist_attr.set_input_dims_mapping(in_name, in_dims_mapping)
            op_dist_attr.set_output_dims_mapping(out_name, out_dims_mapping)

178 179 180 181 182 183 184 185 186 187 188
        return changed

    @staticmethod
    def forward(ctx, *args, **kwargs):
        DistributedDefaultImpl0.forward(ctx, *args, **kwargs)

    @staticmethod
    def backward(ctx, *args, **kwargs):
        DistributedDefaultImpl0.backward(ctx, *args, **kwargs)


189 190 191
register_distributed_operator_impl(
    "slice", DistributedSliceImpl("decrease_in_axis")
)