This commit is contained in:
root 2025-12-03 17:33:00 +09:00
parent 59fbb40e28
commit 96919dc0b5
5 changed files with 616 additions and 73 deletions

View File

@ -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()

View File

@ -9,16 +9,19 @@ export NCCL_DEBUG=INFO
export NCCL_DEBUG_SUBSYS=INIT,NET,IB export NCCL_DEBUG_SUBSYS=INIT,NET,IB
#export NCCL_NET=IB #export NCCL_NET=IB
#export NCCL_IB_DISABLE=0 #export NCCL_IB_DISABLE=0
export NCCL_SOCKET_IFNAME="${NCCL_SOCKET_IFNAME:-net1}" # HostDeviceNetwork로 붙인 NIC 이름 export NCCL_SOCKET_IFNAME="net"
export NCCL_IB_HCA="mlx5_0,mlx5_1,mlx5_2,mlx5_8,mlx5_9,mlx5_5,mlx5_6,mlx5_7" #export NCCL_IB_HCA="=mlx5_6,mlx5_7,mlx5_8,mlx5_9,mlx5_0"
export NCCL_NET_GDR_LEVEL=2 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 # (RoCE 환경에 따라 필요 시) export NCCL_IB_GID_INDEX=3
# 프로세스별 로그 파일 경로 (노드/프로세스마다 별도 파일) # 프로세스별 로그 파일 경로 (노드/프로세스마다 별도 파일)
export NCCL_DEBUG_FILE="/tmp/nccl-%h-%p.log" export NCCL_DEBUG_FILE="/tmp/nccl-%h-%p.log"
#export CUDA_VISIBLE_DEVICES="1,2,3,4,5,6,7"
# ----- 학습 실행 (각 파드가 1개 프로세스씩 구동) ----- # ----- 학습 실행 (각 파드가 1개 프로세스씩 구동) -----
echo "[RUN] RANK=$RANK on $(hostname)" #echo "[RUN] RANK=$RANK on $(hostname)"
torchrun \ torchrun \
--nnodes=2 \ --nnodes=2 \
--nproc_per_node=8 \ --nproc_per_node=8 \
@ -30,26 +33,27 @@ torchrun \
# ----- RDMA 사용 여부 집계 ----- # ----- RDMA 사용 여부 집계 -----
# 새로 생성된 로그 목록 가져오기: START_TIME 이후 생성된 파일만 # 새로 생성된 로그 목록 가져오기: 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

View File

@ -4,7 +4,6 @@ metadata:
name: nccl-scripts name: nccl-scripts
data: data:
run.sh: | run.sh: |
#!/usr/bin/env bash #!/usr/bin/env bash
set -ex set -ex
@ -16,16 +15,19 @@ data:
export NCCL_DEBUG_SUBSYS=INIT,NET,IB export NCCL_DEBUG_SUBSYS=INIT,NET,IB
#export NCCL_NET=IB #export NCCL_NET=IB
#export NCCL_IB_DISABLE=0 #export NCCL_IB_DISABLE=0
export NCCL_SOCKET_IFNAME="${NCCL_SOCKET_IFNAME:-net1}" # HostDeviceNetwork로 붙인 NIC 이름 export NCCL_SOCKET_IFNAME="net"
export NCCL_IB_HCA="mlx5_0,mlx5_1,mlx5_2,mlx5_8,mlx5_9,mlx5_5,mlx5_6,mlx5_7" #export NCCL_IB_HCA="=mlx5_6,mlx5_7,mlx5_8,mlx5_9,mlx5_0"
export NCCL_NET_GDR_LEVEL=2 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 # (RoCE 환경에 따라 필요 시) export NCCL_IB_GID_INDEX=3
# 프로세스별 로그 파일 경로 (노드/프로세스마다 별도 파일) # 프로세스별 로그 파일 경로 (노드/프로세스마다 별도 파일)
export NCCL_DEBUG_FILE="/tmp/nccl-%h-%p.log" export NCCL_DEBUG_FILE="/tmp/nccl-%h-%p.log"
#export CUDA_VISIBLE_DEVICES="1,2,3,4,5,6,7"
# ----- 학습 실행 (각 파드가 1개 프로세스씩 구동) ----- # ----- 학습 실행 (각 파드가 1개 프로세스씩 구동) -----
echo "[RUN] RANK=$RANK on $(hostname)" #echo "[RUN] RANK=$RANK on $(hostname)"
torchrun \ torchrun \
--nnodes=2 \ --nnodes=2 \
--nproc_per_node=8 \ --nproc_per_node=8 \
@ -37,47 +39,307 @@ data:
# ----- RDMA 사용 여부 집계 ----- # ----- RDMA 사용 여부 집계 -----
# 새로 생성된 로그 목록 가져오기: START_TIME 이후 생성된 파일만 # 새로 생성된 로그 목록 가져오기: 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: | train_nccl.py: |
import os import os
import time
import torch import torch
import argparse
import torch.distributed as dist import torch.distributed as dist
import torch.multiprocessing as mp
from torch.optim.lr_scheduler import StepLR
def main(): # for dataset
dist.init_process_group("nccl") from torchvision.datasets.cifar import CIFAR10
local_rank = int(os.environ["LOCAL_RANK"]) import torchvision.transforms as tfs
torch.cuda.set_device(local_rank) 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() import numpy as np
dist.all_reduce(x) import random
print(f"Rank {dist.get_rank()} result: {x[0].item()}") 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 apiVersion: v1
kind: Service kind: Service

View File

@ -1,18 +1,276 @@
import os import os
import time
import torch import torch
import argparse
import torch.distributed as dist import torch.distributed as dist
import torch.multiprocessing as mp
from torch.optim.lr_scheduler import StepLR
def main(): # for dataset
dist.init_process_group("nccl") from torchvision.datasets.cifar import CIFAR10
local_rank = int(os.environ["LOCAL_RANK"]) import torchvision.transforms as tfs
torch.cuda.set_device(local_rank) 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() import numpy as np
dist.all_reduce(x) import random
print(f"Rank {dist.get_rank()} result: {x[0].item()}") 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)

BIN
tensorflow/mnist.npz Normal file

Binary file not shown.