apiVersion: v1 kind: ConfigMap metadata: name: nccl-test-scripts 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=INFO export NCCL_DEBUG_SUBSYS=INIT,NET,IB export NCCL_IB_DISABLE=0 export NCCL_SOCKET_IFNAME="net" # export NCCL_IB_HCA="mlx5_0,mlx5_1,mlx5_2,mlx5_3" # export NCCL_DEBUG_FILE="/tmp/nccl-%h-%p.log" # 주석처리: tee로 파일+stdout 동시 출력 # ----- 학습 실행 ----- 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=0.01) parser.add_argument('--epoch', type=int, default=90) parser.add_argument('--batch_size', type=int, default=1200) 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) --- apiVersion: v1 kind: Service metadata: name: nccl-test labels: { app: nccl-test } spec: clusterIP: None # Headless Service selector: { app: nccl-test } ports: - name: rdzv port: 29500 --- apiVersion: apps/v1 kind: StatefulSet metadata: name: nccl-test spec: serviceName: nccl-test replicas: 2 # 3노드 (ib-1, ib-2, ib-3) selector: matchLabels: { app: nccl-test } template: metadata: labels: { app: nccl-test } annotations: k8s.v1.cni.cncf.io/networks: hostdevice-net spec: containers: - name: worker image: nvcr.io/nvidia/pytorch:24.10-py3 imagePullPolicy: IfNotPresent securityContext: privileged: true capabilities: add: - IPC_LOCK command: - bash - -lc - | cp /config_scripts/* /workspace/ && \ chmod +x /workspace/run.sh && \ cd /workspace && ./run.sh; \ sleep infinity env: # ----- 분산 학습 설정 ----- - name: NNODES value: "2" - name: NPROC_PER_NODE value: "8" - name: MASTER_ADDR value: "nccl-test-0.nccl-test" - name: MASTER_PORT value: "29500" resources: limits: nvidia.com/gpu: "8" nvidia.com/hostdev: 2 requests: nvidia.com/gpu: "8" nvidia.com/hostdev: 2 volumeMounts: - name: scripts mountPath: /config_scripts - name: shared-memory mountPath: /dev/shm - name: infiniband mountPath: /dev/infiniband - name: cifar-data mountPath: /workspace/data/cifar-10-batches-py volumes: - name: scripts configMap: name: nccl-test-scripts defaultMode: 0755 - name: shared-memory emptyDir: medium: Memory sizeLimit: 128Gi # GPU 8장 x 32GB = 256GB, 권장: 50% = 128Gi - name: infiniband hostPath: path: /dev/infiniband type: Directory - name: cifar-data hostPath: path: /home/ubuntu/cifar-10-batches-py type: Directory