apiVersion: v1 kind: ConfigMap metadata: name: {{ include "nccl-perftest.name" . }}-scripts labels: {{- include "nccl-perftest.labels" . | nindent 4 }} data: run.sh: | #!/usr/bin/env bash set -ex # ----- StatefulSet Pod 이름에서 node_rank 추출 ----- # Pod 이름: nccl-test-0, nccl-test-1, nccl-test-2 ... NODE_RANK=${HOSTNAME##*-} echo "[INFO] HOSTNAME=$HOSTNAME, NODE_RANK=$NODE_RANK" # ----- NCCL 환경변수 ----- export NCCL_DEBUG={{ .Values.nccl.debug }} export NCCL_DEBUG_SUBSYS={{ .Values.nccl.debugSubsys }} export NCCL_IB_DISABLE={{ .Values.nccl.ibDisable }} export NCCL_SOCKET_IFNAME="{{ .Values.nccl.socketIfname }}" {{- if .Values.nccl.ibHCA }} export NCCL_IB_HCA="{{ .Values.nccl.ibHCA }}" {{- end }} {{- if .Values.nccl.debugFile }} export NCCL_DEBUG_FILE="{{ .Values.nccl.debugFile }}" {{- end }} # ----- 학습 실행 ----- torchrun \ --nnodes=${NNODES} \ --nproc_per_node=${NPROC_PER_NODE} \ --rdzv_backend=c10d \ --node_rank=${NODE_RANK} \ --rdzv_endpoint=${MASTER_ADDR}:${MASTER_PORT} \ /workspace/train_nccl.py 2>&1 | tee /tmp/nccl-aggregate.log train_nccl.py: | import os import time import torch import argparse import torch.distributed as dist import torch.multiprocessing as mp from torch.optim.lr_scheduler import StepLR # for dataset from torchvision.datasets.cifar import CIFAR10 import torchvision.transforms as tfs from torch.utils.data import DataLoader from torch.utils.data.distributed import DistributedSampler # for model from torchvision.models import vgg11 from torch.nn.parallel import DistributedDataParallel as DDP import numpy as np import random import datetime def set_random_seeds(random_seed=0): torch.manual_seed(random_seed) torch.backends.cudnn.deterministic = True torch.backends.cudnn.benchmark = False np.random.seed(random_seed) random.seed(random_seed) def get_args_parser(): parser = argparse.ArgumentParser(add_help=False) parser.add_argument('--lr', type=float, default={{ .Values.training.lr }}) parser.add_argument('--epoch', type=int, default={{ .Values.training.epochs }}) parser.add_argument('--batch_size', type=int, default={{ .Values.training.batchSize }}) parser.add_argument('--global_rank', type=int, default=0) parser.add_argument('--vis_step', type=int, default=10) parser.add_argument('--num_workers', type=int, default=24) parser.add_argument("--local_rank", type=int, help="Local rank. Necessary for using the torch.distributed.launch utility.") parser.add_argument('--world_size', type=int, default=0) parser.add_argument('--port', type=int, default=2022) parser.add_argument('--root', type=str, default='data') parser.add_argument('--start_epoch', type=int, default=0) parser.add_argument('--save_path', type=str, default='./save') parser.add_argument('--save_file_name', type=str, default='vgg_cifar') return parser def main(opts): # 1. set random seeds set_random_seeds(random_seed=0) # 2. initialization init_for_distributed(opts) # 3. visdom vis = None # 4. data set transform_train = tfs.Compose([ tfs.Resize(256), tfs.RandomCrop(224), tfs.RandomHorizontalFlip(), tfs.ToTensor(), tfs.Normalize(mean=(0.4914, 0.4822, 0.4465), std=(0.2023, 0.1994, 0.2010)), ]) transform_test = tfs.Compose([ tfs.Resize(256), tfs.CenterCrop(224), tfs.ToTensor(), tfs.Normalize(mean=(0.4914, 0.4822, 0.4465), std=(0.2023, 0.1994, 0.2010)), ]) train_set = CIFAR10(root=opts.root, train=True, transform=transform_train, download=True) test_set = CIFAR10(root=opts.root, train=False, transform=transform_test, download=True) train_sampler = DistributedSampler(dataset=train_set, shuffle=True) test_sampler = DistributedSampler(dataset=test_set, shuffle=False) train_loader = DataLoader(dataset=train_set, batch_size=int(opts.batch_size / opts.world_size), shuffle=False, num_workers=int(opts.num_workers / opts.world_size), sampler=train_sampler, pin_memory=True) test_loader = DataLoader(dataset=test_set, batch_size=int(opts.batch_size / opts.world_size), shuffle=False, num_workers=int(opts.num_workers / opts.world_size), sampler=test_sampler, pin_memory=True) # 5. model model = vgg11(pretrained=False) model = model.cuda(opts.local_rank) model = DDP(module=model, device_ids=[opts.local_rank]) # 6. criterion criterion = torch.nn.CrossEntropyLoss().to(opts.local_rank) # 7. optimizer optimizer = torch.optim.SGD(params=model.parameters(), lr=0.01, weight_decay=0.0005, momentum=0.9) # 8. scheduler scheduler = StepLR(optimizer=optimizer, step_size=30, gamma=0.1) if opts.start_epoch != 0: checkpoint = torch.load(os.path.join(opts.save_path, opts.save_file_name) + '.{}.pth.tar' .format(opts.start_epoch - 1), map_location=torch.device('cuda:{}'.format(opts.local_rank))) model.load_state_dict(checkpoint['model_state_dict']) # load model state dict optimizer.load_state_dict(checkpoint['optimizer_state_dict']) # load optim state dict scheduler.load_state_dict(checkpoint['scheduler_state_dict']) # load sched state dict if opts.global_rank == 0: print('\nLoaded checkpoint from epoch %d.\n' % (int(opts.start_epoch) - 1)) for epoch in range(opts.start_epoch, opts.epoch): # 9. train tic = time.time() model.train() train_sampler.set_epoch(epoch) for i, (images, labels) in enumerate(train_loader): images = images.to(opts.local_rank) labels = labels.to(opts.local_rank) outputs = model(images) # ----------- update ----------- optimizer.zero_grad() loss = criterion(outputs, labels) loss.backward() optimizer.step() # get lr for param_group in optimizer.param_groups: lr = param_group['lr'] # time toc = time.time() # visualization if (i % opts.vis_step == 0 or i == len(train_loader) - 1): print('GPU[{0}] Epoch [{1}/{2}], Iter [{3}/{4}], Loss: {5:.4f}, LR: {6:.5f}, Time: {7:.2f}'.format(opts.global_rank, epoch, opts.epoch, i, len(train_loader), loss.item(), lr, toc - tic)) # save pth file if opts.local_rank == 0: if not os.path.exists(opts.save_path): os.mkdir(opts.save_path) checkpoint = {'epoch': epoch, 'model_state_dict': model.state_dict(), 'optimizer_state_dict': optimizer.state_dict(), 'scheduler_state_dict': scheduler.state_dict()} torch.save(checkpoint, os.path.join(opts.save_path, opts.save_file_name + '.{}.pth.tar'.format(epoch))) print("save pth.tar {} epoch!".format(epoch)) # 10. test model.eval() val_avg_loss = 0 correct_top1 = 0 correct_top5 = 0 total = 0 with torch.no_grad(): for i, (images, labels) in enumerate(test_loader): images = images.to(opts.local_rank) # [100, 3, 224, 224] labels = labels.to(opts.local_rank) # [100] outputs = model(images) loss = criterion(outputs, labels) val_avg_loss += loss.item() # ------------------------------------------------------------------------------ # rank 1 _, pred = torch.max(outputs, 1) total += labels.size(0) correct_top1 += (pred == labels).sum().item() # ------------------------------------------------------------------------------ # rank 5 _, rank5 = outputs.topk(5, 1, True, True) rank5 = rank5.t() correct5 = rank5.eq(labels.view(1, -1).expand_as(rank5)) # ------------------------------------------------------------------------------ for k in range(5): # 0, 1, 2, 3, 4, 5 correct_k = correct5[:k+1].reshape(-1).float().sum(0, keepdim=True) correct_top5 += correct_k.item() accuracy_top1 = correct_top1 / total accuracy_top5 = correct_top5 / total val_avg_loss = val_avg_loss / len(test_loader) # make mean loss print("top-1 percentage : {0:0.3f}%".format(correct_top1 / total * 100)) print("top-5 percentage : {0:0.3f}%".format(correct_top5 / total * 100)) scheduler.step() return 0 def init_for_distributed(opts): # 1. setting for distributed training opts.global_rank = int(os.environ['RANK']) opts.local_rank = int(os.environ['LOCAL_RANK']) opts.world_size = int(os.environ['WORLD_SIZE']) torch.cuda.set_device(opts.local_rank) if opts.global_rank is not None and opts.local_rank is not None: print("Use GPU: [{}/{}] for training".format(opts.global_rank, opts.local_rank)) # 2. init_process_group dist.init_process_group( backend="nccl", rank=opts.global_rank, world_size=opts.world_size, device_id=torch.device(f"cuda:{opts.local_rank}"), timeout=datetime.timedelta(seconds=60) ) return if __name__ == '__main__': parser = argparse.ArgumentParser('vgg11 cifar training', parents=[get_args_parser()]) opts = parser.parse_args() main(opts)