Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -149,4 +149,5 @@ log
*.qdrep
.vscode

!bmtrain/dist
!bmtrain/dist
tests/test_log.txt
61 changes: 17 additions & 44 deletions README-ZH.md
Original file line number Diff line number Diff line change
Expand Up @@ -273,19 +273,16 @@ BMTrain 支持**所有** PyTorch 原生的优化器和损失函数,同时你
### 第 4 部分: 训练

```python
# 新建损失缩放器实例
loss_scaler = bmtrain.optim.LossScaler()
# 将所有的 optimzer 及(可选)其对应的 lr_scheduler 收入损失缩放器管理
loss_scaler.add_optimizer(optimizer, lr_scheduler)
# 新建优化器管理器实例
optim_manager = bmtrain.optim.OptimManager(loss_scale=1024)
# 将所有的 optimzer 及(可选)其对应的 lr_scheduler 收入优化器管理器管理
optim_manager.add_optimizer(optimizer, lr_scheduler)
# 可以再次调用 add_optimizer 加入其他优化器

for iteration in range(1000):
# ... 为每个rank加载数据 ...

# 梯度清零
loss_scaler.zero_grad() # 为每个 optimizer 调用 zero_grad

# 前向传播
# 前向传播并计算梯度
pos = torch.arange(enc_input.size(1)).long().cuda().repeat(enc_input.size(0), 1)
logits = model(
enc_input,
Expand All @@ -296,14 +293,19 @@ for iteration in range(1000):

loss = loss_func(logits.view(batch * seq_len, vocab_out_size), targets.view(batch * seq_len))

global_loss = bmtrain.sum_loss(loss).item() # 聚合所有rank上的损失
global_loss = bmtrain.sum_loss(loss).item() # 聚合所有rank上的损失, 仅用于输出训练日志

# 梯度清零
optim_manager.zero_grad() # 为每个 optimizer 调用 zero_grad

# 损失缩放和反向传播
loss = loss_scaler(loss)
loss.backward()
optim_manager.backward(loss)

# 梯度裁剪
grad_norm = optim_manager.clip_grad_norm(optimizer.param_groups, max_norm=1.0)

# 更新参数
loss_scaler.step()
optim_manager.step()

# ... 保存checkpoint、打印日志 ...
```
Expand All @@ -312,40 +314,11 @@ for iteration in range(1000):

你可以根据代码中的注释来了解各部分代码的作用。

唯一需要说明的是 `loss_scaler`,**损失缩放**是混合精度训练中的一项常用技术,需要在反向传播前通过 `loss_scaler(loss)` 对 `loss` 进行放缩,用于避免梯度下溢。在使用损失所放后,优化器的 `step()` 等步骤需要有一些细节上的调整。我们在 `loss_scaler` 帮你实现了这些细节, 你只需要通过 `add_optimizer` 将优化器和学习率调整策略收入 `loss_scaler` 管理,并由 `loss_scaler` 代为执行 `zero_grad()` 和 `step()` 操作。

如果你没有使用 BMTrain 中的融合优化器,你可以不用 `loss_scaler`, 相应的代码修改为:

```python
for iteration in range(1000):
# ... 为每个rank加载数据 ...

# 梯度清零
optimizer.zero_grad()

# 前向传播
pos = torch.arange(enc_input.size(1)).long().cuda().repeat(enc_input.size(0), 1)
logits = model(
enc_input,
pos,
pos < enc_length[:, None]
)
batch, seq_len, vocab_out_size = logits.size()

loss = loss_func(logits.view(batch * seq_len, vocab_out_size), targets.view(batch * seq_len))

global_loss = bmtrain.sum_loss(loss).item() # 聚合所有rank上的损失
唯一需要说明的是 `optim_manager`。在使用 BMTrain 后,优化器的部分相关操作需要有一些细节上的调整。我们在 `optim_manager` 帮你实现了这些细节, 你只需要通过 `add_optimizer` 将优化器和学习率调整策略收入 `optim_manager` 管理,并由 `optim_manger` 代为执行 `zero_grad()`, `backward()`, `clip_grad_norm()` 和 `step()` 等操作。

# 反向传播
loss.backward()

# 更新参数
bmtrain.optim_step(optimizer, lr_scheduler)

# ... 保存checkpoint、打印日志 ...
```
如果你没有使用混合精度训练,你可以不用损失缩放,只需要将 `OptimManger(loss_scale=None)` 构造函数中 `loss_scale` 置为 None 即可, 这也是 `OptimManager` 的默认构造参数。

需要说明的是最后更新参数时,需要使用 `bmtrain.optim_step`, 而不能直接使用 `optimizer.step()` 和 `lr_scheduler.step()`
如果你使用了混合精度训练,**损失缩放**是混合精度训练中的一项常用技术,我们在 `optim_manager.backward(loss)` 帮你对 `loss` 进行了放缩,用于避免梯度下溢。只需要将 `OptimManger` 构造函数中 `loss_scale` 置为一个浮点数即可。 `loss_scale` 会在训练过程中根据梯度进行自适应的调整

<div id="性能"></div>

Expand Down
33 changes: 19 additions & 14 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -276,19 +276,16 @@ In addition, BMTrain also provides the common LRScheduler in the `bmtrain.lr_sch
### Part 4: Training Loop

```python
# create a new instance of loss scaler
loss_scaler = bmtrain.optim.LossScaler()
# let loss_scaler handle all the optimizer and (optional) their corresponding lr_scheduler
loss_scaler.add_optimizer(optimizer, lr_scheduler)
# create a new instance of optimizer manager
optim_manager = bmtrain.optim.OptimManager()
# let optim_manager handle all the optimizer and (optional) their corresponding lr_scheduler
optim_manager.add_optimizer(optimizer, lr_scheduler)
# add_optimizer can be called multiple times to add other optimizers.

for iteration in range(1000):
# ... load data for each rank ...

# zero grad
loss_scaler.zero_grad() # calling zero_grad for each optimizer

# forward
# forward pass and calculate loss
pos = torch.arange(enc_input.size(1)).long().cuda().repeat(enc_input.size(0), 1)
logits = model(
enc_input,
Expand All @@ -299,14 +296,19 @@ for iteration in range(1000):

loss = loss_func(logits.view(batch * seq_len, vocab_out_size), targets.view(batch * seq_len))

global_loss = bmtrain.sum_loss(loss).item() # sum the loss across all ranks
global_loss = bmtrain.sum_loss(loss).item() # sum the loss across all ranks. This is only used for the training log

# zero grad
optim_manager.zero_grad() # calling zero_grad for each optimizer

# clip grad norm
grad_norm = optim_manager.clip_grad_norm(optimizer.param_groups, max_norm=1.0)

# loss scale and backward
loss = loss_scaler(loss)
loss.backward()
optim_manager.backward()

# optimizer step
loss_scaler.step()
optim_manager.step()

# ... save checkpoint or print logs ...
```
Expand All @@ -315,9 +317,12 @@ The training loop part will be slightly longer, but just like a normal training

You can follow the comments in the code to get an idea of what each section of code is doing.

The only additional note is `loss_scaler`, *loss scale* is the technique widely used in mixed precision training to prevent gradient underflow by adding `loss_scaler(loss)` to scale the `loss` before backward. When using loss scaling, some details in optimizers should be adjusted. We have implemented all those details needed in `loss_scaler`. What you need is just letting `loss_scaler` to handle all the optimizers by `add_optimizer`, and letting `loss_scaler` do `zero_grad()` and `step()` instead.
The only additional note is `optimizer`. After using BMTrain, some details in optimizers should be adjusted. We have implemented all those details needed in `optim_manager`. What you need is just letting `optim_manager` to handle all the optimizers by `add_optimizer`, and letting `optim_manager` do `zero_grad()`, `backward()`, `clip_grad_norm()` and `step()` instead.

If you are not using the mixed-precision training, you can train without `loss_scale`. Just set `loss_scale` to None in the `__init__` function of `OptimManager(loss_scale=None)`, which is also the default.

If you are using mixed-precision training, *loss scale* is the technique widely used in mixed precision training to prevent gradient underflow. By using `optim_manager.backward(loss)` to scale the `loss` before backward and set `loss_scale` to some floating number in the `__init__` function of `OptimManager`。The `loss_scale` would be adjusted adaptively based on the gradient during training.

If you are not using the fused optimizer in BMTrain, you can train without `loss_scaler`. The code will be changed as below:

```python
for iteration in range(1000):
Expand Down
1 change: 0 additions & 1 deletion bmtrain/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@
from .synchronize import synchronize, sum_loss, wait_loader, gather_result
from .checkpointing import checkpoint
from .block_layer import CheckpointBlock, TransformerBlockList
from .backward import optim_step
from .wrapper import BMTrainModelWrapper
from .pipe_layer import PipelineTransformerBlockList
from . import debug
Expand Down
28 changes: 0 additions & 28 deletions bmtrain/backward.py

This file was deleted.

8 changes: 4 additions & 4 deletions bmtrain/inspect/tensor.py
Original file line number Diff line number Diff line change
Expand Up @@ -129,9 +129,9 @@ def get_summary(self):
nccl.allReduce(
info.storage(),
info.storage(),
"avg",
"sum",
comm
)
) / nccl.commCount(comm)
x_mean = info[0].cpu().item()
x_std = math.sqrt(info[1].cpu().item())
grad_mean = None
Expand All @@ -146,9 +146,9 @@ def get_summary(self):
nccl.allReduce(
info.storage(),
info.storage(),
"avg",
"sum",
comm
)
) / nccl.commCount(comm)
x_mean = info[0].cpu().item()
x_std = math.sqrt(info[1].cpu().item())
grad_mean = info[2].cpu().item()
Expand Down
2 changes: 1 addition & 1 deletion bmtrain/optim/__init__.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,3 @@
from .adam import AdamOptimizer
from .adam_offload import AdamOffloadOptimizer
from .loss_scaler import LossScaler
from .optim_manager import OptimManager
70 changes: 42 additions & 28 deletions bmtrain/optim/loss_scaler.py → bmtrain/optim/optim_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,20 +27,37 @@ def grad_rescale(param_groups, scale):
if p.grad is not None and p.requires_grad:
p.grad /= scale

class LossScaler:
"""Loss scaler for mix-precision training
class OptimManager:
"""wait cuda stream. Optional: add loss scaler for mix-precision training

Args:
loss_scale (float): The initial loss scale.
loss_scale (float): The initial loss scale. Default to None for not using loss scaling.
loss_scale_factor (float): The loss scale factor.
loss_scale_steps (int): The loss scale steps.

Examples:
>>> optim_manager = bmt.optim.OptimManager(loss_scale=1024)
>>> optim_manager.add_optimizer(optimizer1)
>>> optim_manager.add_optimizer(optimizer2, lr_scheduler2)
>>> for data in dataset:
>>> # forward pass and calculate loss
>>> optim_manager.zero_grad()
>>> optim_manager.backward(loss)
>>> optim_manager.clip_grad_norm(optimizer1.param_groups, max_norm=1.0, norm_type=2)
>>> optim_manager.clip_grad_norm(optimizer2.param_groups, max_norm=2.0, norm_type=2)
>>> optim_manager.step()
"""
def __init__(self,
loss_scale : float = 1,
loss_scale : Optional[float] = None,
loss_scale_factor : float = 2,
loss_scale_steps : int = 1024,
):
self.loss_scale = loss_scale
if loss_scale is not None:
self.loss_scale = loss_scale
self.loss_scale_enabled = True
else:
self.loss_scale = 1
self.loss_scale_enabled = False
self.steps_since_last_scale = 0
self.loss_scale_factor = loss_scale_factor if loss_scale_factor > 1 else 1 / loss_scale_factor
self.loss_scale_steps = loss_scale_steps
Expand All @@ -53,8 +70,8 @@ def add_optimizer(
optimizer: torch.optim.Optimizer,
lr_scheduler: Optional[WarmupLRScheduler] = None,
):
"""Add optimizer and (optional) its corresponding lr_scheduler into loss_scaler.
All optimizers in the same loss_scaler share the same loss scale.
"""Add optimizer and (optional) its corresponding lr_scheduler into optim_manager.
All optimizers in the same optim_manager share the same loss scale.

Args:
optim (torch.optim.Optimizer): A pytorch optimizer, e.g. torch.optim.Adam, torch.optim.SGD or bmtrain.optim.AdamOffloadOptimizer
Expand All @@ -63,14 +80,21 @@ def add_optimizer(
self.optimizers.append(optimizer)
self.lr_schedulers.append(lr_scheduler)

def __call__(self, loss : torch.Tensor) -> torch.Tensor:
def loss_scale(self, loss : torch.Tensor) -> torch.Tensor:
return loss * (self.loss_scale / config['world_size']) # loss scale

def backward(self, loss : torch.Tensor):
"""
Backward with loss scale.

Args:
loss (torch.Tensor): loss
"""
return loss * (self.loss_scale * config["pipe_size"] / config['world_size'])
loss = self.loss_scale(loss)
loss.backward()
# some reduce ops of distributed parameter were launched on load stream
current_stream = torch.cuda.current_stream()
current_stream.wait_stream(config['load_stream'])

def zero_grad(self):
"""
Expand All @@ -88,11 +112,7 @@ def step(self):

This function can also handle gradient overflow by reducing the loss scale when it occurs.
"""
current_stream = torch.cuda.current_stream()
# some reduce ops of distributed parameter were launched on load stream
current_stream.wait_stream(config['load_stream'])

if self.loss_scale > 1:
if self.loss_scale_enabled and self.loss_scale > 1:
has_overflow = False
for optimizer in self.optimizers:
try:
Expand All @@ -110,17 +130,20 @@ def step(self):
if hasattr(optimizer, "_bmtrain_optimizer") and optimizer._bmtrain_optimizer:
optimizer.step(scale=self.loss_scale)
else:
grad_rescale(optimizer.param_groups, self.loss_scale)
if self.loss_scale_enabled:
grad_rescale(optimizer.param_groups, self.loss_scale)
optimizer.step()

if lr_scheduler is not None:
lr_scheduler.step()

self.steps_since_last_scale += 1
if self.loss_scale_enabled:
self.steps_since_last_scale += 1

if self.steps_since_last_scale >= self.loss_scale_steps:
self._justify_scale(self.loss_scale * self.loss_scale_factor)
if self.steps_since_last_scale >= self.loss_scale_steps:
self._justify_scale(self.loss_scale * self.loss_scale_factor)

current_stream = torch.cuda.current_stream()
config['load_stream'].wait_stream(current_stream)

def clip_grad_norm(self, param_groups, max_norm, norm_type=2, eps=1e-6):
Expand All @@ -136,17 +159,8 @@ def clip_grad_norm(self, param_groups, max_norm, norm_type=2, eps=1e-6):

Returns:
Total norm of the parameters (viewed as a single vector).

Examples:
>>> optimizer = bmt.optim.AdamOffloadOptimizer(model.parameters())
>>> loss_scaler = bmt.optim.LossScaler()
>>> loss_scaler.add_optimizer(optimizer)
>>> # ...
>>> # backward_step()
>>> loss_scaler.clip_grad_norm(optimizer.param_groups, max_norm=1.0, norm_type=2)

"""
scale = self.loss_scale * config['pipe_size'] / config['world_size']
scale = self.loss_scale
grads = []
parameters = [p for group in param_groups for p in group['params']]
for p in parameters:
Expand Down
1 change: 0 additions & 1 deletion bmtrain/param_init.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,6 @@ def init_distributed_parameter(params : Iterable[torch.nn.Parameter]):
param._init_method(tmp_tensor)

# Pytorch 1.11 changed the API of storage.__getitem__
# use zero_rank to support pipeline
torch.tensor([], dtype=param.dtype, device=param.device).set_(param.storage())[:] = \
torch.tensor([], dtype=param.dtype, device=param.device).set_(tmp_storage)[partition_size * config['rank'] : partition_size * (config['rank'] + 1)]
# param.storage().copy_(tmp_storage[partition_size * config['rank'] : partition_size * (config['rank'] + 1)])
Expand Down
2 changes: 1 addition & 1 deletion bmtrain/synchronize.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ def sum_loss(loss : torch.Tensor):
This is a helper function to reduce the loss across all workers.
"""
warnings.warn("bmtrain.sum_loss is deprecated and will be removed in later version. Use bmtrain.distributed.all_reduce instead.", DeprecationWarning)
return distributed.all_reduce(loss, "avg")
return distributed.all_reduce(loss, "sum") / config['world_size']

def gather_result(result: torch.Tensor):
warnings.warn("bmtrain.gather_result is deprecated and will be removed in later version. Use bmtrain.distributed.all_gather instead.", DeprecationWarning)
Expand Down
2 changes: 1 addition & 1 deletion csrc/cuda/adam.cu
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ __global__ void adam_fp32_accum(
float local_m = beta1 * __half2float(m[col]) + (1 - beta1) * local_g; // real_m * scale
float local_v = beta2 * v[col] + (1 - beta2) * local_g * local_g / scale; // real_v * scale
float local_p = param[col];
local_p = local_p - lr * local_m / bias_correction1 / (sqrtf(local_v * scale / bias_correction2) + eps) - lr * weight_decay * local_p;
local_p = local_p - lr * local_m / bias_correction1 / (sqrtf(local_v * scale / bias_correction2) + eps * scale) - lr * weight_decay * local_p;

param_h[col] = __float2half(local_p);
param[col] = local_p;
Expand Down
Loading