apiVersion: v1 kind: ConfigMap metadata: name: nccl-scripts data: run.sh: | #!/usr/bin/env bash set -ex # 현재 실행 시작 시 타임스탬프 저장 START_TIME=$(date +%s) # ----- RDMA/IB 강제 ----- export NCCL_DEBUG=INFO export NCCL_DEBUG_SUBSYS=INIT,NET,IB #export NCCL_NET=IB #export NCCL_IB_DISABLE=0 export NCCL_SOCKET_IFNAME="net" #export NCCL_IB_HCA="=mlx5_6,mlx5_7,mlx5_8,mlx5_9,mlx5_0" export NCCL_IB_HCA="=error" #export NCCL_IB_HCA="mlx5_0,mlx5_1,mlx5_2,mlx5_8,mlx5_9,mlx5_5,mlx5_6,mlx5_7" #export NCCL_NET_GDR_LEVEL=2 # (RoCE 환경에 따라 필요 시) export NCCL_IB_GID_INDEX=3 # 프로세스별 로그 파일 경로 (노드/프로세스마다 별도 파일) export NCCL_DEBUG_FILE="/tmp/nccl-%h-%p.log" #export CUDA_VISIBLE_DEVICES="1,2,3,4,5,6,7" # ----- 학습 실행 (각 파드가 1개 프로세스씩 구동) ----- #echo "[RUN] RANK=$RANK on $(hostname)" torchrun \ --nnodes=2 \ --nproc_per_node=8 \ --rdzv_backend=c10d \ --node_rank=0 \ --rdzv_endpoint=nccl-0.nccl:23456 \ /workspace/train_nccl.py 2>&1 | tee /tmp/nccl-aggregate.log # ----- RDMA 사용 여부 집계 ----- # 새로 생성된 로그 목록 가져오기: START_TIME 이후 생성된 파일만 #new_logs=$(find /tmp -name "nccl-nccl-*.log" -type f -newermt "@${START_TIME}") # #total=0; ok=0 #for f in $new_logs; do # [[ -f "$f" ]] || continue # total=$((total+1)) # if rdma_line=$(grep -E "Channel .*GDRDMA" "$f"); then # echo "[OK] RDMA: $rdma_line" # ok=$((ok+1)) # fi #done #echo "[RDMA-CHECK] used IB on $ok / $total local ranks (processes) in $(hostname)" # ## 모든 랭크가 IB를 썼는지(파드 단위) 판단: 로컬 프로세스 수와 비교 #if [[ "$ok" -eq 8 ]]; then # echo "[RDMA-CHECK] ✅ RDMA path OK for all local GPUs" # exit 0 #else # echo "[RDMA-CHECK] ❌ RDMA path NOT used by all local GPUs" # echo "---- NCCL NET/IB related lines ----" # grep -h "NET/" /tmp/nccl-*.log || true # exit 1 #fi 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)) if vis is not None and opts.local_rank == 0: vis.line(X=torch.ones((1, 1)) * i + epoch * len(train_loader), Y=torch.Tensor([loss]).unsqueeze(0), update='append', win='loss', opts=dict(x_label='step', y_label='loss', title='loss', legend=['total_loss'])) # 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 if vis is not None: vis.line(X=torch.ones((1, 3)) * epoch, Y=torch.Tensor([accuracy_top1, accuracy_top5, val_avg_loss]).unsqueeze(0), update='append', win='test_loss_acc', opts=dict(x_label='epoch', y_label='test_loss and acc', title='test_loss and accuracy', legend=['accuracy_top1', 'accuracy_top5', 'avg_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) ) # if put this function, the all processes block at all. #torch.distributed.barrier() 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 labels: { app: nccl } spec: clusterIP: None # Headless selector: { app: nccl } ports: - name: rdv port: 23456 --- apiVersion: apps/v1 kind: StatefulSet metadata: name: nccl spec: serviceName: nccl replicas: 2 # >=2 노드에서 스케줄되도록 노드 리소스 준비 필요 selector: matchLabels: { app: nccl } template: metadata: labels: { app: nccl } # annotations: # k8s.v1.cni.cncf.io/networks: hostdevice-net # HostDeviceNetwork (RDMA) spec: # hostNetwork: true # dnsPolicy: ClusterFirstWithHostNet # RDMA 장치나 /dev/infiniband 노출이 안 보이면 네트워크 오퍼레이터의 RDMA 플러그인 구성을 확인하세요. # (환경에 따라 RDMA Shared Device Plugin이 필요할 수 있음) containers: - name: worker image: nvcr.io/nvidia/pytorch:24.10-py3 imagePullPolicy: IfNotPresent securityContext: privileged: true capabilities: add: - IPC_LOCK - NET_ADMIN # capabilities: { add: ["IPC_LOCK"] } command: [ "bash", "-lc", "cp /config_scripts/* /workspace/ && chmod +x /workspace/run.sh && sleep infinity" ] env: # --- (선택) 특정 GPU만 사용하려면 아래 값을 물리 인덱스로 설정 (예: "2") --- # - name: CUDA_VISIBLE_DEVICES_OVERRIDE # value: "2" # RDMA 관련 기본값은 run.sh에서 설정 resources: limits: nvidia.com/gpu: "8" # nvidia.com/hostdev: "8" # SR-IOV Device Plugin이 광고한 RDMA NIC 리소스 requests: nvidia.com/gpu: "8" # nvidia.com/hostdev: "8" volumeMounts: - name: scripts mountPath: /config_scripts - name: shared-memory mountPath: /dev/shm - name: infiniband mountPath: /dev/infiniband volumes: - name: scripts configMap: name: nccl-scripts defaultMode: 0755 - name: shared-memory emptyDir: medium: Memory sizeLimit: 118541097369600m - name: infiniband hostPath: path: /dev/infiniband type: Directory