_utils.py 10.2 KB
Newer Older
1
import copy
2
from contextlib import nullcontext
3
from typing import Any, Callable, Dict, List, Optional
4

5
import torch
6
import torch.distributed as dist
7
8
9
from torch import Tensor
from torch import distributed as dist
from torch.distributed import ProcessGroup
10
from torch.nn import Module
11
from torch.optim import Adam, Optimizer
12

13
14
from colossalai.booster import Booster
from colossalai.booster.plugin import HybridParallelPlugin
15
from colossalai.lazy import LazyInitContext
16
from colossalai.pipeline.stage_manager import PipelineStageManager
17
from colossalai.shardformer import ShardConfig, ShardFormer
18
from colossalai.shardformer._utils import getattr_
19
from colossalai.shardformer.policies.auto_policy import Policy
20
from colossalai.tensor.d_tensor.api import is_customized_distributed_tensor, is_distributed_tensor
21
22


23
24
25
26
27
28
29
def build_model(model_fn,
                enable_fused_normalization=True,
                enable_tensor_parallelism=True,
                enable_flash_attention=False,
                enable_jit_fused=False,
                use_lazy_init: bool = False):
    # create new model
30
31
32
33
34
35
36
    ctx = LazyInitContext() if use_lazy_init else nullcontext()
    with ctx:
        # create new model
        org_model = model_fn()
        model_copy = copy.deepcopy(org_model)
    if use_lazy_init:
        ctx.materialize(org_model)
37
    # shard model
38
    shard_config = ShardConfig(enable_fused_normalization=enable_fused_normalization,
39
40
41
42
                               enable_tensor_parallelism=enable_tensor_parallelism,
                               enable_flash_attention=enable_flash_attention,
                               enable_jit_fused=enable_jit_fused)
    model_copy = copy.deepcopy(org_model)
43
    shard_former = ShardFormer(shard_config=shard_config)
ver217's avatar
ver217 committed
44
    sharded_model, shared_params = shard_former.optimize(model_copy)
45
    return org_model.cuda(), sharded_model.cuda()
46
47


48
49
50
51
def build_pipeline_model(model_fn,
                         stage_manager=None,
                         enable_fused_normalization=False,
                         enable_tensor_parallelism=False,
Jianghai's avatar
Jianghai committed
52
53
                         use_lazy_init: bool = False,
                         policy: Optional[Policy] = None):
54
55
56
57
58
59
60
61
62
63
64
65
    ctx = LazyInitContext() if use_lazy_init else nullcontext()
    with ctx:
        # create new model
        org_model = model_fn()
        model_copy = copy.deepcopy(org_model)
    if use_lazy_init:
        ctx.materialize(org_model)

    # shard model
    shard_config = ShardConfig(enable_fused_normalization=enable_fused_normalization,
                               enable_tensor_parallelism=enable_tensor_parallelism,
                               pipeline_stage_manager=stage_manager)
Jianghai's avatar
Jianghai committed
66

67
    shard_former = ShardFormer(shard_config=shard_config)
Jianghai's avatar
Jianghai committed
68
    sharded_model, shared_params = shard_former.optimize(model_copy, policy=policy)
69
70
71
    return org_model.cuda(), sharded_model.cuda()


72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
def run_forward(original_model, sharded_model, data_gen_fn, output_transform_fn, loss_fn):
    # prepare input
    data = data_gen_fn()
    data = {k: v.cuda() for k, v in data.items()}
    # switch to train mode
    original_model.train()
    sharded_model.train()
    # run forward
    org_output = original_model(**data)
    org_output = output_transform_fn(org_output)
    org_loss = loss_fn(org_output)

    shard_output = sharded_model(**data)
    shard_output = output_transform_fn(shard_output)
    shard_loss = loss_fn(shard_output)
87
    return org_output, org_loss, shard_output, shard_loss
88
89
90
91
92
93
94
95
96
97
98


def check_state_dict(org_model: Module, sharded_model: Module, name: str = ''):
    org_sd = org_model.state_dict()
    shard_sd = sharded_model.state_dict()
    for k, v in org_sd.items():
        assert k in shard_sd, f'{name} {k} not in sharded model'
        shard_v = shard_sd[k]
        assert v.shape == shard_v.shape, f'{name} {k} shape mismatch, {v.shape} vs {shard_v.shape}'
        assert v.dtype == shard_v.dtype, f'{name} {k} dtype mismatch, {v.dtype} vs {shard_v.dtype}'
        assert torch.equal(v, shard_v), f'{name} {k} value mismatch'
99
100


101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
def build_model_from_hybrid_plugin(model_fn: Callable, loss_fn: Callable, test_config: Dict[str, Any]):

    use_lazy_init = False
    if 'use_lazy_init' in test_config:
        use_lazy_init = test_config.pop('use_lazy_init')

    if use_lazy_init:
        ctx = LazyInitContext()
    else:
        ctx = nullcontext()

    plugin = HybridParallelPlugin(**test_config)
    booster = Booster(plugin=plugin)

    with ctx:
        org_model = model_fn().cuda()
        sharded_model = copy.deepcopy(org_model)

    if use_lazy_init:
        org_model = ctx.materialize(org_model)

    org_optimizer = Adam(org_model.parameters(), lr=1e-3)
    sharded_optimizer = Adam(sharded_model.parameters(), lr=1e-3)
    criterion = loss_fn

    sharded_model, sharded_optimizer, criterion, _, _ = booster.boost(sharded_model, sharded_optimizer, criterion)

    return org_model, org_optimizer, sharded_model, sharded_optimizer, criterion, booster


def run_forward_backward_with_hybrid_plugin(org_model: Module, sharded_model: Module, sharded_optimizer: Optimizer,
                                            data_gen_fn: Callable, output_transform_fn: Callable, criterion: Callable,
                                            booster: Booster):
134
135
    org_model.cuda()
    sharded_model.cuda()
136
137
138
139
140
141
142
143
144
145

    def _criterion(outputs, inputs):
        outputs = output_transform_fn(outputs)
        loss = criterion(outputs)
        return loss

    data = data_gen_fn()
    sharded_model.train()
    if booster.plugin.stage_manager is not None:
        data = {
146
147
            k: v.to('cuda').repeat(*([4] + [1] *
                                     (v.dim() - 1))) if torch.is_tensor(v) or 'Tensor' in v.__class__.__name__ else v
148
149
150
151
152
153
154
155
156
157
158
159
160
            for k, v in data.items()
        }
        data_iter = iter([data])
        sharded_output = booster.execute_pipeline(data_iter,
                                                  sharded_model,
                                                  _criterion,
                                                  sharded_optimizer,
                                                  return_loss=True,
                                                  return_outputs=True)
        sharded_loss = sharded_output['loss']
    else:
        data = {k: v.cuda() for k, v in data.items()}
        sharded_output = sharded_model(**data)
161

162
        sharded_loss = criterion(sharded_output)
163
        sharded_optimizer.backward(sharded_loss)
164
165

    org_model.train()
166
    data = {k: v.cuda() for k, v in data.items()}
167
    org_output = org_model(**data)
168

169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
    org_loss = criterion(org_output)
    org_loss.backward()

    return org_loss, org_output, sharded_loss, sharded_output


def check_output_hidden_state(org_output: Tensor,
                              sharded_output: Tensor,
                              stage_manager: Optional[PipelineStageManager] = None,
                              atol: float = 1e-5,
                              rtol: float = 1e-3):

    org_hidden_state = org_output.last_hidden_state

    if stage_manager is None:
        sharded_hidden_state = sharded_output.last_hidden_state

    if stage_manager and stage_manager.is_last_stage():
        sharded_hidden_state = torch.cat([output.last_hidden_state for output in sharded_output['outputs']], dim=0)

189
    assert torch.allclose(org_hidden_state.float(), sharded_hidden_state.float(), atol=atol, rtol=rtol), \
190
191
192
193
        f"shard model's output hidden state is not equal to origin model's last hidden state\n{org_hidden_state}\n{sharded_hidden_state}"


def check_loss(org_loss: Tensor, sharded_loss: Tensor, atol: float = 1e-5, rtol: float = 1e-3):
194
    assert torch.allclose(org_loss.float(), sharded_loss.float(), atol=atol, rtol=rtol), \
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
        f"shard model loss is not equal to origin model loss\n{org_loss}\n{sharded_loss}"


def check_weight(org_model: Module,
                 sharded_model: Module,
                 layer_suffix: List[str],
                 tp_group: Optional[ProcessGroup] = None,
                 dim: int = 0,
                 atol: float = 1e-5,
                 rtol: float = 1e-3,
                 verbose: bool = False):

    for suffix in layer_suffix:
        org_weight = getattr_(org_model, suffix).weight
        sharded_weight = getattr_(sharded_model, suffix).weight

        if is_distributed_tensor(sharded_weight) or is_customized_distributed_tensor(sharded_weight):
            sharded_weight_list = [
213
                torch.zeros_like(sharded_weight).to('cuda') for _ in range(dist.get_world_size(tp_group))
214
215
216
217
218
219
220
            ]
            dist.all_gather(sharded_weight_list, sharded_weight, tp_group)
            sharded_weight = torch.cat(sharded_weight_list, dim=dim)

        if verbose and dist.get_rank() == 0:
            print(f"'{suffix}' weight: {org_weight}, {sharded_weight}")

221
        assert torch.allclose(org_weight.float(), sharded_weight.float(), atol=atol, rtol=rtol), \
222
            f"shard model weight {suffix} is not equal to origin model weight\n{org_weight}\n{sharded_weight}"
223
224
225
226
227
228
229
230
231
232


def check_grad(org_model: Module,
               sharded_model: Module,
               layer_suffix: List[str],
               tp_group: ProcessGroup = None,
               dim: int = 0,
               atol: float = 1e-5,
               rtol: float = 1e-3,
               verbose: bool = False):
233
    for suffix in layer_suffix:
234
        org_grad = getattr_(org_model, suffix).weight.grad
235
236
237
238
        shard_grad = getattr_(sharded_model, suffix).weight.grad
        shard_weight = getattr_(sharded_model, suffix).weight

        if is_distributed_tensor(shard_weight) or is_customized_distributed_tensor(shard_weight):
239
            shard_grad_list = [torch.zeros_like(shard_grad).to('cuda') for _ in range(dist.get_world_size(tp_group))]
240
241
242
243
244
245
            dist.all_gather(shard_grad_list, shard_grad, tp_group)
            shard_grad = torch.cat(shard_grad_list, dim=dim)

        # embedding may be resized when using tensor parallel
        if shard_grad.shape[0] > org_grad.shape[0]:
            shard_grad = shard_grad[:org_grad.shape[0], :]
246
        if verbose and dist.get_rank() == 0:
247
            print(f"'{suffix}' grad: {org_grad}, {shard_grad}")
248

249
        assert torch.allclose(
250
            org_grad.float(), shard_grad.float(), rtol=rtol, atol=atol
251
        ), f"error attribute '{suffix}', orgin model grad is not equal to shard model grad\n{org_grad}\n{shard_grad}"