作者:忧郁的白衬衫

WeNet(https://github.com/mobvoi/wenet)支持了多机多卡的分布式并行训练,进而可以利用更多的来加速模型的训练。目前的方案使用PyTorch原生的DistributedDataParallel(DDP)实现,在AISHELL-1实验上,该方案使用多机可以做到几乎无损的线性加速(随着GPU增多训练时间按GPU数量比例线性减少),同时保证模型性能和单机一致。

 

业界现有方案

目前基于PyTorch的主流分布式训练方案主要有PyTorch原生的DDP和horovod两种。horovod是uber开源的一个兼容主流计算框架如PyTorch、TensorFlow等的分布式机器学习训练框架。horovod使用起来较为简单,引入少量的代码入侵即可完成训练代码的改造,并且其支持对梯度进行fp16压缩,可以将多机训练通信时需要的网络带宽需求减半。早期的PyTorch等框架对于多机训练的支持不够完善,而为了更好的达到线性加速比,horovod外挂式的方案曾经是多数用户的选择。但随着各个机器学习训练框架的完善,也都逐渐内置支持与horovod相同的分布式训练机制。因此考虑到多方面因素,即更方便的版本控制,对一些算子和术语理解的统一以及减少代码入侵,WeNet采取PyTorch原生的DDP方案来实现多机多卡训练。


DistributedDataParallel

首先我们简单回顾一下DDP(DistributedDataParallel)的实现原理以及DDP中的各个概念。


DDP的性能的相关问题

  • Python GIL(Global Interpreter Lock)

Python的多线程运行时,每个线程首先获取GIL,然后执行代码直到sleep或者是python虚拟机将其挂起,最后释放GIL。Python GIL的存在使得一个python进程只能利用一个CPU核心,不适用于神经网络训练等计算密集型的任务。使用多进程,才能更有效率地利用多核的计算资源。

  •  DP(DataParallel)

在介绍DDP前先回顾一下DP,DP(DataParallel)是PyTorch早期的数据并行方案,是基于Parameter Server(具体可见:https://zhuanlan.zhihu.com/p/82116922)的一种算法(PyTorch内部通过scatter与gather操作实现),应用于单机多卡的训练。这一方案在存在比较严重的内存和通信带宽负载不均衡的问题,其中用于同步各节点信息的主GPU节点的内存一般会比其他GPU节点占用的内存更多。但是真正限制DP的是Parameter Server会使多卡训练时主GPU节点需求的通信带宽远远高于其他GPU节点,导致主节点的通信带宽成为系统的瓶颈。此外其采用了单进程多线程的方式进行训练,性能会受到python GIL的影响。

  • DDP的原理,DDP为什么比DP快?

DDP方案和DP都是数据并行方案,但是DDP采用了Ring-AllReduce数据交换算法(具体可见:https://www.zhihu.com/question/57799212/answer/612786337)提高了通讯效率。该方法每个节点不再全部向主节进行通信,而是只与相邻的节点进行通信,因此每个GPU节点上的流量通信压力相同,从而充分利用了每个GPU节点的通信带宽,降低了多机训练时对网络带宽的需求,避免了整体性能受限于主GPU节点带宽的情况。通常情况下不同主机间的网络带宽小于主机内部GPU通信的总线的带宽,因此DDP这种Ring-AllReduce通信方案在多机训练时优势更为明显。此外DDP通过多进程(DDP支持为每个GPU使用一个独立进程)的方式也避免了Python GIL(Global Interpreter Lock)的限制,从而可以进一步提高训练速度。

  •  DDP的线性加速比

虽然DDP采用的Ring-AllReduce的通信方式极大的降低了训练过程中对网络带宽的需求,但是DDP在多机训练时仍然需要较好的网络带宽才能达到更好的线性加速,稍后我们将测试一下网络对于DDP多机训练时的影响。


DDP中的概念

DDP进程组的初始化的接口定义如下

torch.distributed.init_process_group(backend, init_method=None, timeout=datetime.timedelta(0, 1800), world_size=-1, rank=-1, store=None, group_name='')


这里我们主要关注backend、init_method、world_size、rank参数。

  • Backend

用来指定多进程间的通信后端,包括NCCL,Gloo,MPI。其中NCCL是官方最推荐的,因此通常我们直接使用NCCL即可。

  • init_method

用来表示在启动多进程训练时,各进程的握手方式,主要包括三种方式,即file(共享文件)、tcp(套接字)和env(环境变量)。多机训练主要使用file和tcp的方式。通过file的握手方式需要多机之间装有NFS(Network File System),在指定了各进程均可访问的共享文件路径后,各台机器上的不同进程通过该共享文件完成握手,例如file:///export/nfs/ddp_init,。tcp的方式需要给定各机器上不同进程均可访问的网络地址和端口号(未被占用的)来完成各进程的握手,例如tcp://127.0.0.1:23456。

  • group

通常情况使用默认值即可。

  • world size

表示全局的并行的进程数,DDP模式下,最优的方案是每个进程一个卡,因此通常情况下world size实际为总的GPU数量或者总的进程数。

  • rank

表示当前进程或GPU的序号,用于进程间通讯。从0开始排序,范围是0~world size-1。注意:rank=0表示该进程是master进程,通常用来存储模型日志等。此外程序中还可能出现node_rank以及local_rank的概念。

  • node_rank

表示当前机器的序号。同样也是从0开始排序,假设一共使用了N台机器,则node_rank的范围是0~N-1。

  • local_rank

表示当前机器上的进程或GPU的序号。从0开始排序,假设当前机器可使用GPU数量是N,则local_rank范围是0~N-1。


DDP的工作流程

DDP的主要工作流程如下:

1. 启动各进程

2. 各进程通过指定的init method完成握手,以便进行进程之间的通信

3. 主节点载入或随机初始化模型,并且将参数和Buffer(Buffer是会被持久化保存的数据等,如BatchNorm中的mean和variance)等模型状态信息分发至各GPU节点,此时每个GPU节点拿到的模型状态是相同的

4. 通过sampler,每个GPU拿到了属于自己的那份数据,并进行前向运算,并且计算出loss。

5. 每个GPU进行backward()后向运算求出梯度(这时每个GPU上的梯度是不同的),之后各GPU通过AllReduce算法对梯度进行同步,使得每个GPU都能拿到全部GPU上的梯度(这时每个GPU上的梯度是相同的)

6. 各GPU根据同步之后的梯度来更新模型参数,更新之后各个GPU上的模型参数是相同的

7. 重复执行4~6,直至训练完毕


WeNet分布式多机训练实现

首先多机分布式训练通常需要分布式存储系统的支持,如S3、HDFS和NFS。由于PyTorch官方未像TensorFlow一样提供S3和HDFS的接口,因此我们这里使用NFS作为分布式存储。下面我们通过WeNet中的代码来进一步的了解如何将训练改成多机多卡的分布式训练。首先是进程初始化相关,用于完成各进程握手的init_process_grop函数,需要传递的变量有backend,init_method,world_size和rank。注意:在使用多机多卡训练时,先启动的机器上的进程会在此等待,直到所有进程都执行到这里,此时总的进程数与world_size相同,才会继续执行。所以当world_size值和实际的总进程数不一致时,会产生Bug。(完整代码可以参考https://github.com/mobvoi/wenet/blob/main/wenet/bin/train.py)

为了各进程的dataloader能得到不同的数据,需要给各dataloader增加一个DistributedSampler,用于把数据分发到不同的进程(通过多进程使用同一种子,可使不同进程分配不重叠、不交叉的数据)。

    dist.init_process_group(args.dist_backend,
                            init_method=args.init_method,
                            world_size=args.world_size,
                            rank=args.rank)
    train_sampler = torch.utils.data.distributed.DistributedSampler(
        train_dataset, shuffle=True)
    train_data_loader = DataLoader(train_dataset,
                               collate_fn=train_collate_func,
                               sampler=train_sampler,
                               shuffle=(train_sampler is None),
                               batch_size=1,
                               num_workers=args.num_workers)


之后通过DistributedDataParall将模型的状态信息从master节点传到其他节点,使所有进程上的模型状态一致。

model.cuda()
    model = torch.nn.parallel.DistributedDataParallel(
        model, find_unused_parameters=True)


之后的流程便和正常的训练一致,但是还有另外一点值得注意,即Gradient Accumulation梯度累计机制,这个机制指每隔若干个batch step进行一次参数的更新,从而可以去模拟更大的batch size,使训练更加稳定。

完整代码可以参考:https://github.com/mobvoi/wenet/blob/main/wenet/utils/executor.py

for batch_idx, batch in enumerate(data_loader):
...
      loss, loss_att, loss_ctc = model(feats,
                                             feats_lengths,
                                             target,
                                             target_lengths)
        loss.backward()
        if batch_idx % accum_grad == 0:
            if rank == 0 and writer is not None:
                writer.add_scalar('train_loss', loss, self.step)
            grad_norm = clip_grad_norm_(model.parameters(), clip)
            if torch.isfinite(grad_norm):
                optimizer.step()
            optimizer.zero_grad()
            scheduler.step()
            self.step += 1


但是如上的实现存在一个性能问题,上文已经提及,梯度的all reduce操作是在backward()时进行,因此每个batch step都会进行各GPU进程之间梯度同步的操作。由于Gradient Accumulation每n个step才更新一次网络参数,因此其中n-1次backward时进行的梯度同步结果不会被使用。而这些无用的梯度同步会增加分布式多机多卡训练时各机器节点的网络带宽压力。PyTorch最新的版本中支持的no_sync()上下文管理可以解决这个问题,使用如下

 # Disable gradient synchronizations across DDP processes.
        # Within this context, gradients will be accumulated on module
        # variables, which will later be synchronized.
        if is_distributed and batch_idx % accum_grad != 0 :
            context = model.no_sync
        # Used for single GPU training and DDP gradient synchronization
        # processes.
        else:
            context = nullcontext
        with context():
            loss, loss_att, loss_ctc = model(feats,
                                             feats_lengths,
                                             target,
                                             target_lengths)
            loss = loss / accum_grad
            loss.backward()


在累计梯度时使用no_sync上下文,这时backward将不再进行梯度的同步。WeNet分布式的使用实践 首先需要确定我们总共打算使用的机器数num_nodes,如果是单机训练则设置num_nodes=1,node_rank=0即可。若想使用多机训练模式,则需要先制定机器节点的数量num_nodes,然后在每个机器节点启动脚本中指定该节点的node_rank。例如一共使用了两台机器,则设置num_nodes=2,首先在第一个节点启动run.sh脚本并设置node_rank=0,其次在第二个节点启动run.sh脚本并设置node_rank=1。

完整代码可以参考:https://github.com/mobvoi/wenet/blob/main/examples/aishell/s0/run.sh

num_nodes=1
# The rank of each node or machine, range from 0 to num_nodes -1
# The first node/machine sets node_rank 0, the second one sets node_rank 1
# the third one set node_rank 2, and so on. Default 0
node_rank=0


每台机器在启动run.sh脚本后,将会在其所在的主机上启动与其包含的GPU数量相同的进程,并计算各GPU/进程的rank,从而完成多机多卡的分布式训练。

 for ((i = 0; i < $num_GPUs; ++i)); do
    {
        GPU_id=$(echo $CUDA_VISIBLE_DEVICES | cut -d',' -f$[$i+1])
        # Rank of each GPU/process used for knowing whether it is
        # the master of a worker.
        rank=`expr $node_rank \* $num_GPUs + $i`
        python wenet/bin/train.py --GPU $GPU_id \
            --config $train_config \
            --train_data $feat_dir/$train_set/format.data \
            --cv_data $feat_dir/dev/format.data \
            ${checkpoint:+--checkpoint $checkpoint} \
            --model_dir $dir \
            --ddp.init_method $init_method \
            --ddp.world_size $world_size \
            --ddp.rank $rank \
            --ddp.dist_backend $dist_backend \
            --num_workers 2 \
            $cmvn_opts
    } &
    done
    wait


WeNet实验结果

no_sync上下文

我们通过监控各及其节点的网络流量情况来验证no_sync上下文是否生效。如下是采用累计梯度为4时同一机器节点上网卡的流量情况。左图为未使用no_sync上下文,右图为使用了no_sync上下文。

 

由此可见no_sync上下文极大的降低了网络流量,减小了对网络带宽的需求压力。


网络带宽对多机训练加速比的影响

上文中提到,虽然DDP采用的all-reduce的梯度同步方案,各GPU通信时处于平等的状态,平衡了各节点的通信流量,但是DDP仍然需要较好的网络带宽,才能达到更好的多机线性加速。因为通常情况机器之间的网络带宽会小于机器内部通过总线通信的带宽(排除infiniband等可以达到每秒上百Gb的土豪网络设施),为此各机器节点之间通信的网络带宽将会成为系统的瓶颈。当网络带宽过小时,在进行梯度同步时,由于带宽过小会导致延迟大,梯度同步通信时出现等待的现象。下表记录了wenet使用了两台机器一共16块GPU时,累计梯度为4的情况下网络带宽对每个epoch的训练时间的影响,使用的语料为AISHELL-1。


可以看到当使用千兆网时,两台机器16卡训练的速度甚至比单机8卡还要慢很多,改成万兆网卡之后便可得到不错的加速,由此可见,DDP也需要较好的网络带宽的支持,防止在做梯度更新时,网卡的带宽成为系统的瓶颈。


no_sync上下文的影响

no_sync上下文对每个epoch的训练时间影响如下。


通过对比我们可以发现,不管是千兆网还是万兆网,使用no_sync上下文都比未使用no_sync上下文更快了些。万兆网卡在使用了梯度no_sync上下文后几乎接近较为完美的线性加速比。 


多机多卡的模型效果对比

单机训练与多机训练最终模型CER的对比如下。


通过对比我们可以看到,多机训练的识别性能与单机训练的识别性能基本一致。在模型效果几乎不变的情况下,WeNet多机训练近乎达到了完美的线性加速。多机多卡训练的超参数配置已经更新到WeNet的github上 https://github.com/mobvoi/wenet/blob/main/examples/aishell/s0/README.md。


基于NFS多机训练的trick

我们知道NFS是一种便于数据共享的服务,但并不适用于多机分布式高性能计算。最主要的原因是在进行多机训练时,从client节点会不断的访问存放数据的server节点(假设client节点内存难以cache住全部要访问的数据),这时server节点的网络和磁盘带宽会成为系统的瓶颈。这里有一个小trick可以缓解这个问题。就是将要访问的数据的存储分散至各节点上,这样同一时刻,所有进程想要访问的数据会向不同的NFS的server节点访问,这样极大的减小了单一server节点时的网络和磁盘带宽问题。总结一下WeNet分布式训练更好体验的几个核心要素:

  1. 万兆网卡甚至更高级的网卡,是保证多机训练线性加速的关键。
  2. no_sync与acc_grad的使用可以极大降低多机训练时的流量,减小系统对带宽的需求,可以得到更好的加速比。
  3. 基于NFS的多机存储方案,server点的带宽会成为系统瓶颈,为此可以将文件分散至不同节点上,负载均衡各节点网络流量。网络带宽对于分布式训练来说至关重要。除上述的方法,还可以使用amp或apex等软件包提供的混合精度训练,这样在梯度同步时,使用fp16来进行同步,相较于现在fp32可减少一倍的网络流量。这些方法WeNet之后也会陆续更新。