diff --git a/nccl_test/simple-test/quick_test.py b/nccl_test/simple-test/quick_test.py new file mode 100644 index 0000000..50f5302 --- /dev/null +++ b/nccl_test/simple-test/quick_test.py @@ -0,0 +1,19 @@ +import os +import torch +import torch.distributed as dist +import datetime + +def main(): + dist.init_process_group("nccl",timeout=datetime.timedelta(seconds=10)) + local_rank = int(os.environ["LOCAL_RANK"]) + torch.cuda.set_device(local_rank) + + print(f"Rank {dist.get_rank()} initialized on GPU {local_rank}") + + x = torch.ones(10).cuda() + dist.all_reduce(x) + print(f"Rank {dist.get_rank()} result: {x[0].item()}") + +if __name__ == "__main__": + main() + diff --git a/nccl_test/simple-test/run.sh b/nccl_test/simple-test/run.sh index 8cabea6..fc76058 100644 --- a/nccl_test/simple-test/run.sh +++ b/nccl_test/simple-test/run.sh @@ -9,16 +9,19 @@ export NCCL_DEBUG=INFO export NCCL_DEBUG_SUBSYS=INIT,NET,IB #export NCCL_NET=IB #export NCCL_IB_DISABLE=0 -export NCCL_SOCKET_IFNAME="${NCCL_SOCKET_IFNAME:-net1}" # HostDeviceNetwork로 붙인 NIC 이름 -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 +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)" +#echo "[RUN] RANK=$RANK on $(hostname)" torchrun \ --nnodes=2 \ --nproc_per_node=8 \ @@ -30,26 +33,27 @@ torchrun \ # ----- RDMA 사용 여부 집계 ----- # 새로 생성된 로그 목록 가져오기: START_TIME 이후 생성된 파일만 -new_logs=$(find /tmp -name "nccl-nccl-*.log" -type f -newermt "@${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 -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 diff --git a/nccl_test/simple-test/template.yaml b/nccl_test/simple-test/template.yaml index a1417bf..d4f3abd 100644 --- a/nccl_test/simple-test/template.yaml +++ b/nccl_test/simple-test/template.yaml @@ -4,7 +4,6 @@ metadata: name: nccl-scripts data: run.sh: | - #!/usr/bin/env bash set -ex @@ -16,16 +15,19 @@ data: export NCCL_DEBUG_SUBSYS=INIT,NET,IB #export NCCL_NET=IB #export NCCL_IB_DISABLE=0 - export NCCL_SOCKET_IFNAME="${NCCL_SOCKET_IFNAME:-net1}" # HostDeviceNetwork로 붙인 NIC 이름 - 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 + 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)" + #echo "[RUN] RANK=$RANK on $(hostname)" torchrun \ --nnodes=2 \ --nproc_per_node=8 \ @@ -37,47 +39,307 @@ data: # ----- RDMA 사용 여부 집계 ----- # 새로 생성된 로그 목록 가져오기: START_TIME 이후 생성된 파일만 - new_logs=$(find /tmp -name "nccl-nccl-*.log" -type f -newermt "@${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 - 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 - def main(): - dist.init_process_group("nccl") - local_rank = int(os.environ["LOCAL_RANK"]) - torch.cuda.set_device(local_rank) + # 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 - print(f"Rank {dist.get_rank()} initialized on GPU {local_rank}") + # for model + from torchvision.models import vgg11 + from torch.nn.parallel import DistributedDataParallel as DDP - x = torch.ones(10).cuda() - dist.all_reduce(x) - print(f"Rank {dist.get_rank()} result: {x[0].item()}") + import numpy as np + import random + import datetime - if __name__ == "__main__": - main() + + 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 diff --git a/nccl_test/simple-test/train_nccl.py b/nccl_test/simple-test/train_nccl.py index 0e6a7c8..00c4b59 100644 --- a/nccl_test/simple-test/train_nccl.py +++ b/nccl_test/simple-test/train_nccl.py @@ -1,18 +1,276 @@ 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 -def main(): - dist.init_process_group("nccl") - local_rank = int(os.environ["LOCAL_RANK"]) - torch.cuda.set_device(local_rank) +# 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 - print(f"Rank {dist.get_rank()} initialized on GPU {local_rank}") +# for model +from torchvision.models import vgg11 +from torch.nn.parallel import DistributedDataParallel as DDP - x = torch.ones(10).cuda() - dist.all_reduce(x) - print(f"Rank {dist.get_rank()} result: {x[0].item()}") +import numpy as np +import random +import datetime -if __name__ == "__main__": - main() +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) diff --git a/tensorflow/mnist.npz b/tensorflow/mnist.npz new file mode 100644 index 0000000..e7baa20 Binary files /dev/null and b/tensorflow/mnist.npz differ