데이터 병렬화
개요
데이터 병렬화(Data Parallelism)는 하나의 모델 복제본을 여러 device 또는 process에 배치하고, 각 worker가 서로 다른 mini-batch를 동시에 처리하는 분산 학습 기법이다. 각 worker는 같은 초기 가중치로 forward와 backward를 수행한 뒤 gradient를 집계하고, 모든 worker가 같은 optimizer update를 적용한다. 따라서 모델을 device마다 나누는 model parallelism과 달리, 한 모델 복제본이 한 device의 메모리에 들어간다는 전제가 기본이다.
가장 널리 쓰이는 동기식 방식은 각 rank의 gradient를 AllReduce로 평균내는 Distributed Data Parallel(DDP)이다. 입력 데이터는 자동으로 중복 제거되지 않으므로 DistributedSampler나 동등한 rank-aware input pipeline으로 shard해야 한다. worker 수를 늘리면 연산량은 분산되지만 매 step마다 gradient 통신과 가장 느린 worker를 기다리는 synchronization 비용이 생기므로, scaling은 GPU 수에 비례해 무조건 좋아지지 않는다.
핵심 개념
Rank, world size, local batch
| 용어 | 의미 | 실무에서 확인할 점 |
|---|---|---|
| rank | 분산 process의 고유 번호 | global rank와 node 안의 local rank를 구분 |
| world size | 함께 통신하는 전체 process 수 | data-parallel group의 크기와 일치해야 함 |
| local batch | 한 rank가 한 step에서 처리하는 sample 수 | GPU 메모리와 kernel 효율로 결정 |
| global batch | 한 optimizer step이 본 전체 sample 수 | world_size × local_batch × accumulation_steps |
| DDP replica | 각 rank에 있는 동일 모델 복제본 | 파라미터는 복제되지만 입력 shard는 서로 달라야 함 |
동기식 gradient 평균
같은 local batch 크기 b를 사용하는 W개 rank에서 rank r의 평균 gradient를 g_r라 하면, 동기식 data parallelism의 update gradient는 다음과 같다.
g_r = (1 / b) * sum(loss(x, y)의 gradient) # rank r의 local gradient
g = (1 / W) * sum(r=0..W-1, g_r) # AllReduce 평균
theta_next = optimizer(theta, g)
AllReduce(SUM) 뒤에 world_size로 나누는 구현과 AllReduce(AVG) 구현은 같은 결과를 낸다. local batch가 rank마다 다르면 sample 수를 가중치로 반영해야 하며, 단순 평균은 global mean gradient가 아닐 수 있다.
데이터 shard와 epoch shuffle
DistributedSampler는 dataset index를 rank별로 나누어 같은 sample이 여러 rank에서 반복 처리되는 것을 피한다. shuffle=True를 쓰는 경우에는 매 epoch마다 sampler.set_epoch(epoch)를 호출해야 rank들이 동일한 seed 규칙으로 새 순서를 만들면서도 서로 다른 shard를 유지한다. Iterable dataset이나 streaming input은 sampler 대신 rank와 worker 정보를 이용해 직접 분할해야 한다.
비교/분석
병렬화 방식 비교
| 방식 | 모델 배치 | 통신 핵심 | 메모리 특성 | 적합한 상황 |
|---|---|---|---|---|
| Single-GPU | 한 device | 없음 | 모델 전체가 한 GPU에 필요 | 작은 모델, 기준선 측정 |
| DataParallel | 한 process가 여러 GPU에 replica 생성 | scatter/gather 중심 | GPU마다 모델 복제 | 빠른 단일 노드 실험, 레거시 코드 |
| DDP | process당 모델 replica 1개 | backward 중 gradient AllReduce | GPU마다 모델 복제 | 다중 GPU/다중 node 표준 학습 |
| FSDP / ZeRO | data parallel group에 state shard | reduce-scatter/all-gather 등 | parameter·gradient·optimizer state 분할 가능 | 모델 상태가 단일 GPU에 안 들어갈 때 |
| Model Parallel | 모델 layer/tensor를 device에 분할 | activation·partial result 통신 | 모델을 여러 GPU에 분산 | 모델 자체가 너무 클 때 |
DDP와 DataParallel의 차이
PyTorch의 DataParallel은 single-process multi-thread 방식으로 입력을 scatter하고 출력 결과를 gather한다. 반면 DistributedDataParallel은 보통 GPU마다 하나의 process를 만들고, 각 process가 독립적으로 forward/backward를 수행한 뒤 gradient만 collective communication으로 동기화한다. PyTorch 문서는 GIL contention, 매 iteration replica 처리, scatter/gather overhead 때문에 단일 node에서도 DDP를 권장한다.
동기식과 비동기식
| 항목 | Synchronous data parallel | Asynchronous / parameter-server 계열 |
|---|---|---|
| update 시점 | 모든 rank가 gradient를 집계한 뒤 update | worker가 준비되는 대로 server에 전달 |
| 장점 | 재현성과 optimizer semantics가 비교적 명확 | straggler를 일부 숨길 수 있음 |
| 단점 | 가장 느린 rank가 step을 결정 | stale gradient와 update 충돌 가능 |
| 대표 사용 | PyTorch DDP + NCCL | parameter server, 일부 비동기 SGD 시스템 |
Scaling 관점
이상적인 step time은 worker 수에 따라 줄지만 실제 효율은 다음 항의 영향을 받는다.
step time ≈ max(rank별 계산 시간) + gradient 통신 + input/IO 대기
scaling efficiency = (T_1 / (W * T_W)) * 100%
큰 모델에서는 gradient 통신량이 커지고, 작은 local batch에서는 계산보다 통신이 지배적이 된다. Ring AllReduce는 모든 rank가 결과를 받으며, NCCL의 collective는 같은 count와 datatype으로 모든 rank가 같은 순서에 호출해야 한다. 실제 성능은 NVLink/PCIe/InfiniBand topology, NCCL algorithm, message size, rank 배치에 따라 달라진다.
동작 원리
1. Process group 초기화
launcher가 각 process에 RANK, LOCAL_RANK, WORLD_SIZE를 전달하고, process들이 init_process_group으로 통신 그룹을 만든다. 일반적인 CUDA 학습에서는 NCCL backend를 사용하며, 한 GPU를 한 DDP process에 매핑하는 구성이 기본이다.
import os
import torch
import torch.distributed as dist
from torch.nn.parallel import DistributedDataParallel as DDP
from torch.utils.data import DataLoader, DistributedSampler
local_rank = int(os.environ["LOCAL_RANK"])
dist.init_process_group(backend="nccl")
torch.cuda.set_device(local_rank)
model = Net().cuda(local_rank)
model = DDP(model, device_ids=[local_rank])
sampler = DistributedSampler(dataset, shuffle=True)
loader = DataLoader(dataset, sampler=sampler, batch_size=local_batch)
DDP constructor에서는 rank 0의 model state를 다른 rank에 맞추고, 이후 각 replica가 같은 초기 상태에서 시작하도록 한다. process group을 초기화한 뒤에는 모든 rank가 forward, backward, collective 호출을 같은 순서로 진행해야 한다.
2. Local forward/backward
각 rank는 자기 shard의 mini-batch로 loss와 gradient를 계산한다. 모델 파라미터는 복제되어 있지만 입력 activation과 intermediate tensor는 기본적으로 각 rank의 device에만 존재한다. 이 단계에서 rank별 gradient가 잠시 달라도 정상이며, backward hook이 준비된 gradient bucket을 통신 단계로 넘긴다.
3. Gradient bucket AllReduce
DDP는 parameter별 autograd hook을 등록한다. backward가 뒤쪽 layer부터 진행되면 준비된 gradient를 bucket 단위로 묶어 AllReduce하고, 결과를 각 rank의 param.grad에 반영한다. 통신은 남은 backward 계산과 겹칠 수 있으므로, bucket 크기와 layer 순서는 overlap 효율에 영향을 준다.
NCCL에서 AllReduce는 각 rank의 같은 위치 원소를 reduce한 결과를 모든 rank의 receive buffer에 저장한다. Ring 구현을 단순화하면 각 rank는 reduce-scatter 단계에서 부분 합을 만들고, all-gather 단계에서 최종 조각을 모두 공유한다. FP16/BF16 gradient를 통신하고 FP32 master state로 update하는 mixed-precision 구성에서는 overflow 감지와 loss scaling 정책도 모든 rank에서 일관되어야 한다.
4. Optimizer update
backward()가 끝났을 때 각 rank의 gradient가 동기화되어 있으면, 각 rank는 같은 optimizer state와 같은 hyperparameter로 optimizer.step()을 실행한다. Adam의 moment state처럼 optimizer state도 replica마다 유지되므로, 순수 DDP의 GPU 메모리 사용량은 모델 parameter뿐 아니라 gradient와 optimizer state까지 rank 수만큼 반복된다. checkpoint는 일반적으로 rank 0에서만 저장하고, 나머지 rank가 안전한 시점에 읽게 할 수 있다.
5. Global batch와 learning rate
global_batch = world_size × local_batch × gradient_accumulation_steps
worker 수를 늘리면 동일한 epoch에서 optimizer step 수가 줄어들 수 있다. 따라서 learning rate, warmup, scheduler step 수, gradient clipping 기준을 global batch 기준으로 다시 확인해야 한다. Large Minibatch SGD 연구는 ImageNet에서 batch 증가에 맞춘 linear scaling rule과 warmup을 제안했지만, 이는 모든 모델과 optimizer에 자동으로 성립하는 법칙이 아니며 validation으로 확인해야 한다.
Gradient accumulation과 no_sync
GPU 메모리 때문에 local batch를 작게 유지하면서 effective batch를 키울 수 있다. DDP에서는 accumulation 중간 step마다 AllReduce를 하면 통신만 반복되므로, 마지막 micro-step 전까지 no_sync()로 gradient synchronization을 지연할 수 있다.
for micro_step, (x, y) in enumerate(loader):
sync = (micro_step + 1) % accumulation_steps == 0
context = model.no_sync() if not sync else nullcontext()
with context:
loss = model(x, y) / accumulation_steps
loss.backward()
if sync:
optimizer.step()
optimizer.zero_grad(set_to_none=True)
이 예시는 개념을 보여주기 위한 것이며 실제 코드에서는 nullcontext import, epoch 경계, 마지막 incomplete accumulation, AMP scaler 처리를 함께 고려해야 한다.
장단점
장점
- 구현 경계가 명확하다: 모델 구조를 쪼개지 않고도 서로 다른 data shard를 병렬 처리할 수 있다.
- 높은 연산 병렬성: 큰 batch를 여러 GPU에서 동시에 처리해 training throughput을 높인다.
- 범용성이 높다: CNN, Transformer, embedding 모델 등 대부분의 mini-batch 학습에 적용할 수 있다.
- DDP 생태계가 성숙했다: PyTorch DDP, NCCL,
torchrun, SLURM 등 표준 도구와 결합하기 쉽다. - 다른 병렬화와 결합 가능하다: DDP process 내부에 tensor/model parallelism을 넣어 hybrid parallelism을 구성할 수 있다.
단점
- 모델 상태를 복제한다: 순수 DDP는 parameter, gradient, optimizer state를 GPU마다 보유하므로 메모리 절감 기법이 아니다.
- 통신 비용이 필수다: 매 optimizer step의 gradient synchronization이 네트워크 대역폭과 latency를 소비한다.
- straggler에 민감하다: 한 rank가 느리거나 실패하면 collective가 대기하거나 전체 job이 중단될 수 있다.
- global batch 변화가 최적화에 영향을 준다: worker 수와 accumulation 변경은 learning rate, scheduler, generalization을 바꿀 수 있다.
- 입력 pipeline이 병목이 될 수 있다: shard, shuffle, I/O, augmentation이 rank별로 균형을 이루지 않으면 GPU가 기다린다.
- 불균등 입력 처리가 까다롭다: rank별 step 수나 collective 호출 순서가 달라지면 hang, timeout, data duplication이 발생할 수 있다.
관련 기술
프레임워크와 통신 계층
| 기술 | 역할 | 데이터 병렬화와의 관계 |
|---|---|---|
| PyTorch DDP | process별 model replica와 gradient synchronization | 표준 synchronous data parallel API |
| NCCL | GPU collective communication | AllReduce, ReduceScatter, AllGather 제공 |
DistributedSampler |
dataset index를 rank별로 분할 | 입력 중복을 피하고 epoch shuffle을 조정 |
| DeepSpeed ZeRO | optimizer state/gradient/parameter shard | data parallel 자원을 사용해 replica 메모리 절감 |
| FSDP | parameter를 shard하고 필요 시 all-gather | DDP보다 큰 모델을 data-parallel group에 배치 |
| Horovod | framework 독립 distributed training API | TensorFlow/PyTorch 등에서 AllReduce 기반 학습 |
기존 문서와의 연결
- 메모리 접근 패턴 — rank별 input shard와 GPU memory coalescing을 함께 최적화하는 관점
- 메모리 레이아웃 최적화 — batch와 tensor layout이 통신 전후 kernel 효율에 미치는 영향
- Square Tiling 기초 — 각 rank 내부 GEMM throughput을 높이는 타일링 기법
- 추론 vs 훈련 — data parallelism이 주로 적용되는 training 단계와 inference replica의 차이
주요 참고 문헌과 표준 문서
- PyTorch DistributedDataParallel API — DDP module, gradient synchronization, process 구성
- PyTorch DDP Tutorial —
torchrun, checkpoint, model parallel 결합 예시 - PyTorch DDP Design Note — reducer, autograd hook, gradient bucket 설계
- NVIDIA NCCL Collective Operations — AllReduce, ReduceScatter, AllGather semantics
- Goyal et al., Accurate, Large Minibatch SGD: Training ImageNet in 1 Hour, 2017 — large batch와 linear scaling/warmup
- Rajbhandari et al., ZeRO: Memory Optimizations Toward Training Trillion Parameter Models, 2020 — data parallel 상태 shard
핵심 정리
- 데이터 병렬화는 각 rank가 같은 모델의 복제본으로 서로 다른 mini-batch를 처리하는 방식이다.
- DDP는 backward 중 gradient bucket을
AllReduce해 모든 rank가 동일한 optimizer update를 하도록 만든다. global_batch는 world size와 accumulation까지 포함하므로 worker 수를 바꿀 때 learning rate와 scheduler를 함께 점검해야 한다.- DDP는 모델 상태를 복제하므로 모델이 한 GPU에 들어가지 않으면 FSDP/ZeRO, model parallelism 또는 hybrid parallelism을 고려해야 한다.
- 실제 scaling은 GPU 계산량뿐 아니라 gradient 통신, input pipeline, topology, straggler 균형으로 결정된다.