# 多gpu训练

In [1]:
import torch
from torch import nn
from torch.nn import functional as F
from d2l import torch as d2l

In [None]:
class Timer:
    """Record multiple running times."""
    def __init__(self):
        """Defined in :numref:`sec_minibatch_sgd`"""
        self.times = []
        self.start()

    def start(self):
        """Start the timer."""
        self.tik = time.time()

    def stop(self):
        """Stop the timer and record the time in a list."""
        self.times.append(time.time() - self.tik)
        return self.times[-1]

    def avg(self):
        """Return the average time."""
        return sum(self.times) / len(self.times)

    def sum(self):
        """Return the sum of time."""
        return sum(self.times)

    def cumsum(self):
        """Return the accumulated time."""
        return np.array(self.times).cumsum().tolist()

In [None]:
scale = 0.01
W1 = torch.randn(size=(20, 1, 3, 3))*scale
b1 = torch.zeros(20)
W2 = torch.randn(size=(50, 20, 5, 5))*scale
b2 = torch.zeros(50)
W3= torch.randn(size=(800, 128))*scale
b3 = torch.zeros(128)
W4 = torch.randn(size=(128,10))*scale
b4 = torch.zeros(10)
params = [W1, b1, W2, b2, W3, b3, W4, b4]

def lenet(X, params):
    h1_conv = F.conv2d(input=X, weight=params[0], bias=params[1])
    h1_activation = F.relu(h1_conv)
    h1 = F.avg_pool2d(input=h1_activation, kernel_size=(2,2), stride=(2,2))
    h2_conv = F.conv2d(h1, params[2], params[3])
    h2_activation = F.relu(h2_conv)
    h2 = F.avg_pool2d(input=h2_activation, kernel_size=(2,2), stride=(2,2))
    h2 = h2.reshape(h2.shape[0], -1)
    h3_linear = torch.mm(h2, params[4])+params[5]
    h3 = F.relu(h3_linear)
    y_hat = torch.mm(h3, params[6]) + params[7]
    return y_hat

# reduction参数决定返回的损失计算方式:
# 'none'表示返回每个样本的损失（不进行任何缩减）,
# 'mean'表示返回所有样本损失的平均值,
# 'sum'表示返回所有样本损失的总和。
# 这里选择'none'，以保留每个样本的单独损失值。
loss = nn.CrossEntropyLoss(reduction='none')

In [4]:
def get_params(params, device):
    new_params = [p.to(device) for p in params]
    for p in new_params:
        p.requires_grad_()
    return new_params

In [8]:
new_params = get_params(params, d2l.try_gpu(0))
print('b1 权重:', new_params[1])
print('b1 梯度:', new_params[1].grad)
print(d2l.try_gpu(0))

b1 权重: tensor([0., 0., 0., 0., 0., 0., 0., 0., 0., 0., 0., 0., 0., 0., 0., 0., 0., 0., 0., 0.],
       requires_grad=True)
b1 梯度: None
cpu


In [13]:
def allreduce(data):
    for i in range(1, len(data)):
        data[0][:] += data[i].to(data[0].device)
    for i in range(1, len(data)):
        data[1][:] = data[0].to(data[i].device)

In [14]:
data = [torch.ones((1,2), device=d2l.try_gpu(i)) * (i+1) for i in range(2)]
print('allreduce之前：\n', data[0], '\n', data[1])
allreduce(data)
print('allreduce之后：\n', data[0], '\n', data[1])

allreduce之前：
 tensor([[1., 1.]]) 
 tensor([[2., 2.]])
allreduce之后：
 tensor([[3., 3.]]) 
 tensor([[3., 3.]])


In [15]:
print(d2l.try_gpu(1))

cpu


In [16]:
data = torch.arange(20).reshape(4,5)
devices = [torch.device('cuda:0'), torch.device('cuda:1')]
sqlit = nn.parallel.scatter(data, devices)
print('input: ', data)
print('load into', devices)
print('output:', sqlit)

AttributeError: module 'torch._C' has no attribute '_scatter'

In [None]:
def sqlit_batch(X, y, devices):
    assert X[0].shape == y[0].shape
    # 这里的作用是将输入的X和y按照devices划分，然后分别拷贝到每个GPU上，便于在多块显卡上并行训练。
    # nn.parallel.scatter 会把张量拆分，并分发到指定的设备列表。
    return (nn.parallel.scatter(X, devices), nn.parallel.scatter(y, devices))

In [None]:
def train_batch(X, y, device_params, devices, lr):
    X_shards, y_shards = split_batch(X, y, devices)
    # 在每个GPU上分别计算损失
    ls = [loss(lenet(X_shard, device_W), y_shard).sum()
          for X_shard, y_shard, device_W in zip(
              X_shards, y_shards, device_params)]
    for l in ls:  # 反向传播在每个GPU上分别执行
        l.backward()
    # 将每个GPU的所有梯度相加，并将其广播到所有GPU
    with torch.no_grad():
        for i in range(len(device_params[0])):
            allreduce(
                [device_params[c][i].grad for c in range(len(devices))])
    # 在每个GPU上分别更新模型参数
    for param in device_params:
        with torch.no_grad():
            for param in params:
                param -= lr * param.grad / batch_size
                param.grad.zero_()

In [None]:
def train(num_gpus, batch_size, lr):
    train_iter, test_iter = d2l.load_data_fashion_mnist(batch_size)
    devices = [d2l.try_gpu(i) for i in range(num_gpus)]
    device_params = [get_params(params, d) for d in devices]
    num_epochs = 10
    animator = d2l.Animator('epoch', 'test acc', xlim=[1, num_epochs])
    timer = Timer()
    for epoch in range(num_epochs):
        timer.start()
        for X,y in train_iter:
            train_batch(X, y, device_params, devices, lr)
            torch.cuda.synchronize()
        timer.stop()
        animator.add(epoch + 1, (d2l.evaluate_accuracy_gpu(
            lambda x: lenet(x, device_params[0]), test_iter, devices[0]),))
    print(f'测试精度：{animator.Y[0][-1]:.2f}，{timer.avg():.1f}秒/轮，'
          f'在{str(devices)}')

In [None]:
train(num_gpus=1, batch_size=256, lr=0.2)

In [None]:
train(num_gpus=2, batch_size=256, lr=0.2)