2021-10-28 16:21:23 +00:00
|
|
|
#!/usr/bin/env python
|
|
|
|
# -*- encoding: utf-8 -*-
|
|
|
|
|
2023-06-19 03:23:24 +00:00
|
|
|
from contextlib import contextmanager
|
|
|
|
|
2021-10-28 16:21:23 +00:00
|
|
|
import torch.nn as nn
|
|
|
|
|
2023-09-18 08:31:06 +00:00
|
|
|
from colossalai.legacy.context import ParallelMode
|
|
|
|
from colossalai.legacy.core import global_context as gpc
|
2021-10-28 16:21:23 +00:00
|
|
|
|
|
|
|
|
|
|
|
class ParallelLayer(nn.Module):
|
2022-09-06 12:18:35 +00:00
|
|
|
global_state_dict: bool = True
|
2021-10-28 16:21:23 +00:00
|
|
|
|
|
|
|
def __init__(self):
|
|
|
|
super().__init__()
|
2023-09-19 06:20:26 +00:00
|
|
|
self.data_parallel_rank = (
|
|
|
|
0 if not gpc.is_initialized(ParallelMode.DATA) else gpc.get_local_rank(ParallelMode.DATA)
|
|
|
|
)
|
|
|
|
self.data_parallel_size = (
|
|
|
|
1 if not gpc.is_initialized(ParallelMode.DATA) else gpc.get_world_size(ParallelMode.DATA)
|
|
|
|
)
|
2021-10-28 16:21:23 +00:00
|
|
|
|
2023-09-19 06:20:26 +00:00
|
|
|
self.tensor_parallel_rank = (
|
|
|
|
0 if not gpc.is_initialized(ParallelMode.TENSOR) else gpc.get_local_rank(ParallelMode.TENSOR)
|
|
|
|
)
|
|
|
|
self.tensor_parallel_size = (
|
|
|
|
1 if not gpc.is_initialized(ParallelMode.TENSOR) else gpc.get_world_size(ParallelMode.TENSOR)
|
|
|
|
)
|
2021-10-28 16:21:23 +00:00
|
|
|
|
2023-09-19 06:20:26 +00:00
|
|
|
self.pipeline_parallel_rank = (
|
|
|
|
0 if not gpc.is_initialized(ParallelMode.PIPELINE) else gpc.get_local_rank(ParallelMode.PIPELINE)
|
|
|
|
)
|
|
|
|
self.pipeline_parallel_size = (
|
|
|
|
1 if not gpc.is_initialized(ParallelMode.PIPELINE) else gpc.get_world_size(ParallelMode.PIPELINE)
|
|
|
|
)
|
2022-04-01 08:49:56 +00:00
|
|
|
|
2023-09-19 06:20:26 +00:00
|
|
|
def _load_from_global_state_dict(
|
|
|
|
self, state_dict, prefix, local_metadata, strict, missing_keys, unexpected_keys, error_msgs
|
|
|
|
):
|
|
|
|
return super()._load_from_state_dict(
|
|
|
|
state_dict, prefix, local_metadata, strict, missing_keys, unexpected_keys, error_msgs
|
|
|
|
)
|
2022-09-06 12:18:35 +00:00
|
|
|
|
|
|
|
def _save_to_global_state_dict(self, destination, prefix, keep_vars):
|
|
|
|
return super()._save_to_state_dict(destination, prefix, keep_vars)
|
|
|
|
|
2023-09-19 06:20:26 +00:00
|
|
|
def _load_from_state_dict(
|
|
|
|
self, state_dict, prefix, local_metadata, strict, missing_keys, unexpected_keys, error_msgs
|
|
|
|
):
|
2022-09-06 12:18:35 +00:00
|
|
|
if self.global_state_dict:
|
|
|
|
if gpc.get_local_rank(ParallelMode.TENSOR) != 0:
|
|
|
|
missing_keys.clear()
|
|
|
|
unexpected_keys.clear()
|
2023-09-19 06:20:26 +00:00
|
|
|
return self._load_from_global_state_dict(
|
|
|
|
state_dict, prefix, local_metadata, strict, missing_keys, unexpected_keys, error_msgs
|
|
|
|
)
|
|
|
|
return super()._load_from_state_dict(
|
|
|
|
state_dict, prefix, local_metadata, strict, missing_keys, unexpected_keys, error_msgs
|
|
|
|
)
|
2022-09-06 12:18:35 +00:00
|
|
|
|
|
|
|
def _save_to_state_dict(self, destination, prefix, keep_vars):
|
|
|
|
if self.global_state_dict:
|
|
|
|
return self._save_to_global_state_dict(destination, prefix, keep_vars)
|
|
|
|
return super()._save_to_state_dict(destination, prefix, keep_vars)
|
|
|
|
|
|
|
|
@classmethod
|
|
|
|
@contextmanager
|
|
|
|
def use_local_state_dict(cls):
|
|
|
|
try:
|
|
|
|
cls.global_state_dict = False
|
|
|
|
yield
|
|
|
|
finally:
|
|
|
|
cls.global_state_dict = True
|