PyTorch Distributed Training

When models grow larger and data grows larger, a single GPU can no longer meet training requirements.

Distributed training uses parallel computing across multiple GPUs or even multiple machines, significantly reducing training time.

This section details distributed training techniques in PyTorch, including DataParallel, DistributedDataParallel, and mixed-precision distributed training.


1. Distributed Training Basics

1.1 Why Do We Need Distributed Training

The scale of deep learning models grows year by year, and training time increases accordingly. The core goal of distributed training is:Distribute computing tasks across multiple compute units to complete training of large-scale models and data within an acceptable time.。

Distributed training mainly solves two problems:

  • Insufficient memory: Large models' parameters, optimizer states, and gradients require a large amount of GPU memory.
  • Excessive training time: Training a large model on a single GPU may take weeks or even months.

1.2 Two Major Modes of Distributed Training

Distributed training is mainly divided into two modes:

Mode Principle Applicable scenario
Data Parallel Each compute unit stores a complete model; different compute units process different data. The model can fit into the GPU memory of a single card; large data volume requires accelerated training; small teams with limited machine resources.
Model Parallel Split the model across multiple compute units; each compute unit stores only part of the model. The model is too large to fit on a single card; dedicated clusters, multi-machine training.

In actual production environments, data parallelism is the most commonly used approach. PyTorch provides two implementations: DataParallel and DistributedDataParallel.


2. DataParallel Usage Guide

2.1 Single-Machine Multi-GPU: DataParallel

DataParallel(DP)It is the simplest multi-GPU training method provided by PyTorch. It only requires modifying a few lines of code to train using multiple GPUs on a single machine.

Example

import torch
import torch.nn as nn
import torch.optim as optim

# ── Model Definition ──────────────────────────────────────
class SimpleModel(nn.Module):
    def __init__(self):
        super().__init__()
        self.net = nn.Sequential(
            nn.Linear(128, 256),
            nn.ReLU(),
            nn.Linear(256, 256),
            nn.ReLU(),
            nn.Linear(256, 10)
        )

    def forward(self, x):
        return self.net(x)


# Check the number of available GPUs
print(f"Available GPU count: {torch.cuda.device_count()}")

# Create the model and move it to GPU
model = SimpleModel()

# ── Method 1: Use DataParallel (simplest way)───────
if torch.cuda.device_count() > 1:
    model = nn.DataParallel(model)

# Move to GPU (if using DataParallel, device_ids are handled automatically)
device = torch.device("cuda" if torch.cuda.is_available() else "cpu")
model = model.to(device)

# Loss function and optimizer
criterion = nn.CrossEntropyLoss()
optimizer = optim.Adam(model.parameters(), lr=1e-3)

# ── Training Loop ──────────────────────────────────────
def train_epoch_dp(model, loader, criterion, optimizer, device):
    model.train()
    total_loss = 0
    correct = 0
    total = 0

    for inputs, labels in loader:
        inputs = inputs.to(device, non_blocking=True)
        labels = labels.to(device, non_blocking=True)

        optimizer.zero_grad()
        outputs = model(inputs)
        loss = criterion(outputs, labels)
        loss.backward()
        optimizer.step()

        total_loss += loss.item() * inputs.size(0)
        _, predicted = outputs.max(1)
        correct += predicted.eq(labels).sum().item()
        total += labels.size(0)

    return total_loss / total, correct / total


# Simulated data
train_loader = [
    (torch.randn(32, 128), torch.randint(0, 10, (32,))) for _ in range(10)
]

# Start training
for epoch in range(3):
    loss, acc = train_epoch_dp(model, train_loader, criterion, optimizer, device)
    print(f"Epoch {epoch+1}: Loss={loss:.4f}, Acc={acc:.4f}")

print("DataParallel training complete!")

2.2 How DataParallel Works

The workflow of DataParallel is as follows:

  1. Model replication: Copy the model to each GPU.
  2. Data distribution: Evenly distribute batch data to each GPU.
  3. Parallel computation: Each GPU independently computes the forward pass.
  4. Gradient aggregation: Aggregate gradients from all GPUs to the main GPU.
  5. Parameter update: The main GPU updates parameters and synchronizes them to other GPUs.

Example

# Description of key DataParallel parameters
model = nn.DataParallel(
    module,                    # The model to parallelize (required)
    device_ids=[0, 1, 2, 3],   # GPU device IDs to use; by default, all GPUs are used
    output_device=0,            # Device to gather output results to; by default, the first GPU
    dim=0                      # Dimension along which data is split; by default, the batch dimension
)

# Check the device of the current model
print(f"Model device: {next(model.parameters()).device}")

# Check the devices actually in use
print(f"Number of GPUs used: {model.device_ids if hasattr(model, 'device_ids') else 'N/A'}")

DataParallel is simple and easy to use, but has two main problems: 1) the main GPU has high memory pressure (it needs to aggregate gradients); 2) communication efficiency between GPUs is low. Therefore, the official recommendation is to use DistributedDataParallel.


3. Detailed Explanation of DistributedDataParallel

3.1 DDP Basics

DistributedDataParallel(DDP)It is the distributed training method recommended by PyTorch. Compared with DataParallel, DDP has the following advantages:

  • Each GPU computes gradients independently, without needing to aggregate them on the main GPU.
  • Uses an efficient gradient synchronization algorithm (Ring AllReduce).
  • Supports multi-machine multi-GPU training.
  • Faster training speed and more efficient memory usage.

3.2 DDP Single-Machine Multi-GPU Training

Example

import torch
import torch.nn as nn
import torch.optim as optim
import torch.distributed as dist
from torch.nn.parallel import DistributedDataParallel as DDP
import os

# ── Initialize the Distributed Environment ────────────────────────────
def setup(rank, world_size):
    """Set up the distributed environment"""
    # Set the GPU used by the current process
    torch.cuda.set_device(rank)

    # Initialize the process group
    dist.init_process_group(
        backend="nccl",           # Use the NCCL backend (recommended for GPU)
        init_method="env://",     # Environment variable initialization method
        world_size=world_size,    # Total number of processes
        rank=rank                 # Rank of the current process
    )


def cleanup():
    """Clean up the distributed environment"""
    dist.destroy_process_group()


# ── Data Loader (distributed support)───────────────────────
def get_distributed_loader(batch_size, world_size):
    """Create a data loader with distributed support"""
    from torch.utils.data import DataLoader, DistributedSampler

    # Simulated dataset
    dataset = torch.utils.data.TensorDataset(
        torch.randn(100, 3, 32, 32),
        torch.randint(0, 10, (100,))
    )

    # DistributedSampler automatically splits data
    sampler = DistributedSampler(
        dataset,
        num_replicas=world_size,
        rank=rank,
        shuffle=True
    )

    loader = DataLoader(
        dataset,
        batch_size=batch_size,
        sampler=sampler,
        num_workers=2,
        pin_memory=True
    )

    return loader


# ── Model Definition ──────────────────────────────────────
class ImageClassifier(nn.Module):
    def __init__(self):
        super().__init__()
        self.features = nn.Sequential(
            nn.Conv2d(3, 32, 3, padding=1),
            nn.ReLU(),
            nn.MaxPool2d(2),
            nn.Conv2d(32, 64, 3, padding=1),
            nn.ReLU(),
            nn.AdaptiveAvgPool2d(1),
            nn.Flatten()
        )
        self.classifier = nn.Linear(64, 10)

    def forward(self, x):
        x = self.features(x)
        x = self.classifier(x)
        return x


# ── DDP Training Function ──────────────────────────────────
def train_ddp(rank, world_size, epochs=3):
    """Main function for distributed training"""
    # Initialize
    setup(rank, world_size)

    # Create the model and move it to the corresponding GPU
    model = ImageClassifier().cuda(rank)

    # Wrap it as a DDP model
    ddp_model = DDP(
        model,
        device_ids=[rank],        # Specify the device used by the current process
        output_device=rank,
        find_unused_parameters=False  # Whether to detect unused parameters
    )

    # Loss function and optimizer
    criterion = nn.CrossEntropyLoss()
    optimizer = optim.Adam(ddp_model.parameters(), lr=1e-3)

    # Get the data loader
    loader = get_distributed_loader(batch_size=16, world_size=world_size)

    # Training loop
    for epoch in range(epochs):
        # Set the epoch number at the start of each epoch (used for shuffling data)
        loader.sampler.set_epoch(epoch)

        for batch_idx, (inputs, labels) in enumerate(loader):
            inputs = inputs.cuda(rank, non_blocking=True)
            labels = labels.cuda(rank, non_blocking=True)

            optimizer.zero_grad()
            outputs = ddp_model(inputs)
            loss = criterion(outputs, labels)
            loss.backward()
            optimizer.step()

            if rank == 0 and batch_idx % 5 == 0:
                print(f"Epoch {epoch+1}, Batch {batch_idx}, Loss: {loss.item():.4f}")

    # Cleanup
    if rank == 0:
        print("Training complete!")
    cleanup()


# Note: In actual execution, you need to specify rank and world_size via command-line arguments
# This code needs to run independently in each process

3.3 Launching Distributed Training

PyTorch distributed training needs to be launched viatorchrunortorch.distributed.launch. Here are the common launch methods:

Example

# Method 1: Using torchrun (recommended, PyTorch 2.0+)
# File: train_ddp.py

"""
Usage:
torchrun --nproc_per_node=4 train_ddp.py
"
""

# Method 2: Using python -m (traditional method)
"""
python -m torch.distributed.launch --nproc_per_node=4 train_ddp.py
"
""

# Method 3: Multi-node multi-GPU launch
# Machine 1 (master node)
torchrun --nnodes=2 --nproc_per_node=4 --node_rank=0 --master_addr=192.168.1.1 --master_port=29500 train_ddp.py

# Machine 2 (worker node)
torchrun --nnodes=2 --nproc_per_node=4 --node_rank=1 --master_addr=192.168.1.1 --master_port=29500 train_ddp.py

# Parameter descriptions:
# --nproc_per_node: number of GPUs per node
# --nnodes: number of nodes
# --node_rank: current node index (starting from 0)
# --master_addr: master node IP address
# --master_port: master node port number

3.4 Complete DDP Training Script

Example

# Complete DDP training script (save as train_ddp.py)
import torch
import torch.nn as nn
import torch.optim as optim
import torch.distributed as dist
from torch.nn.parallel import DistributedDataParallel as DDP
from torch.utils.data import DataLoader, DistributedSampler
import argparse
import os


def parse_args():
    parser = argparse.ArgumentParser()
    parser.add_argument("--local_rank", type=int, default=-1, help="Automatically passed by torchrun")
    parser.add_argument("--epochs", type=int, default=10)
    parser.add_argument("--batch_size", type=int, default=32)
    parser.add_argument("--lr", type=float, default=1e-3)
    return parser.parse_args()


class TrainModel(nn.Module):
    def __init__(self):
        super().__init__()
        self.net = nn.Sequential(
            nn.Conv2d(3, 64, 3, padding=1),
            nn.ReLU(),
            nn.MaxPool2d(2),
            nn.Conv2d(64, 128, 3, padding=1),
            nn.ReLU(),
            nn.AdaptiveAvgPool2d(1),
            nn.Flatten(),
            nn.Linear(128, 10)
        )

    def forward(self, x):
        return self.net(x)


def setup(rank, world_size):
    os.environ["MASTER_ADDR"] = "localhost"
    os.environ["MASTER_PORT"] = "12355"
    dist.init_process_group("nccl", rank=rank, world_size=world_size)
    torch.cuda.set_device(rank)


def cleanup():
    dist.destroy_process_group()


def train(rank, world_size, args):
    setup(rank, world_size)

    # Create the model
    model = TrainModel().cuda(rank)
    model = DDP(model, device_ids=[rank])

    # Optimizer and loss function
    optimizer = optim.Adam(model.parameters(), lr=args.lr)
    criterion = nn.CrossEntropyLoss()

    # Dataset
    dataset = torch.utils.data.TensorDataset(
        torch.randn(1000, 3, 32, 32),
        torch.randint(0, 10, (1000,))
    )

    # DistributedSampler ensures each process sees different data
    sampler = DistributedSampler(
        dataset,
        num_replicas=world_size,
        rank=rank,
        shuffle=True
    )

    loader = DataLoader(
        dataset,
        batch_size=args.batch_size,
        sampler=sampler,
        num_workers=2,
        pin_memory=True
    )

    # Training loop
    for epoch in range(args.epochs):
        # Key: set the epoch every epoch to shuffle the data partition
        sampler.set_epoch(epoch)
        model.train()

        epoch_loss = 0
        for inputs, labels in loader:
            inputs = inputs.cuda(rank, non_blocking=True)
            labels = labels.cuda(rank, non_blocking=True)

            optimizer.zero_grad()
            outputs = model(inputs)
            loss = criterion(outputs, labels)
            loss.backward()
            optimizer.step()

            epoch_loss += loss.item()

        # Synchronize all processes to ensure every process finishes the current epoch
        dist.barrier()

        if rank == 0:
            avg_loss = epoch_loss / len(loader)
            print(f"Epoch {epoch+1}/{args.epochs}, Loss: {avg_loss:.4f}")

    cleanup()


if __name__ == "__main__":
    args = parse_args()

    # Get world_size from environment variables
    world_size = int(os.environ["WORLD_SIZE"])

    # Get rank from environment variables (automatically set by torchrun)
    rank = int(os.environ["RANK"])

    if torch.cuda.is_available():
        train(rank, world_size, args)
    else:
        print("CUDA environment is required to run distributed training")

Key point: every epoch needs to callsampler.set_epoch(epoch), to ensure the data partition is shuffled, otherwise multiple processes will see the same data.


4. Mixed Precision in Distributed Training

4.1 DDP + AMP Combination

Combining distributed training with mixed precision training can further improve training speed. DDP in PyTorch 2.0+ is fully integrated with AMP.

Example

import torch
import torch.nn as nn
import torch.optim as optim
import torch.distributed as dist
from torch.nn.parallel import DistributedDataParallel as DDP
from torch.cuda.amp import autocast, GradScaler


class AMPModel(nn.Module):
    def __init__(self):
        super().__init__()
        self.net = nn.Sequential(
            nn.Conv2d(3, 64, 3, padding=1),
            nn.ReLU(),
            nn.MaxPool2d(2),
            nn.Conv2d(64, 128, 3, padding=1),
            nn.ReLU(),
            nn.AdaptiveAvgPool2d(1),
            nn.Flatten(),
            nn.Linear(128, 10)
        )

    def forward(self, x):
        return self.net(x)


def train_ddp_amp(rank, world_size):
    """Distributed + mixed precision training"""
    # Initialization
    torch.cuda.set_device(rank)
    dist.init_process_group("nccl", rank=rank, world_size=world_size)

    # Create the model
    model = AMPModel().cuda(rank)
    model = DDP(model, device_ids=[rank])

    # Optimizer and GradScaler
    optimizer = optim.Adam(model.parameters(), lr=1e-3)
    scaler = GradScaler()

    criterion = nn.CrossEntropyLoss()

    # Training loop
    model.train()
    for inputs, labels in loader:
        inputs = inputs.cuda(rank, non_blocking=True)
        labels = labels.cuda(rank, non_blocking=True)

        optimizer.zero_grad()

        # Use autocast to automatically switch precision
        with autocast(device_type='cuda'):
            outputs = model(inputs)
            loss = criterion(outputs, labels)

        # Use scaler to handle gradient scaling
        scaler.scale(loss).backward()

        # Gradient clipping (requires unscaling first)
        scaler.unscale_(optimizer)
        torch.nn.utils.clip_grad_norm_(model.parameters(), max_norm=1.0)

        scaler.step(optimizer)
        scaler.update()

    dist.destroy_process_group()


# Launch command
# torchrun --nproc_per_node=4 train_ddp_amp.py

4.2 Gradient Synchronization in Distributed Training

DDP usesRing AllReducealgorithm for gradient synchronization, and it works as follows:

  1. Each GPU splits the gradients into n parts (n is the number of GPUs)
  2. Each GPU exchanges one gradient part with an adjacent GPU
  3. After repeating n times, each GPU has the complete gradients

Example

# DDP gradient synchronization related configuration

# 1. Set the buffer size for gradient buckets
model = DDP(
    model,
    device_ids=[rank],
    gradient_as_bucket_view=True,  # Save memory (PyTorch 1.8+)
    broadcast_buffers=False,        # Do not broadcast buffers to reduce communication
    bucket_cap_mb=25               # Gradient bucket size (MB)
)

# 2. Manually synchronize model parameters (if needed)
# In some scenarios, model parameters need to be manually synchronized
def sync_params(model):
    for param in model.parameters():
        dist.broadcast(param, src=0)

# 3. Check gradient synchronization status
# DDP automatically handles gradient synchronization; no manual operation is needed
# But it can be checked in the following way
for name, param in model.named_parameters():
    if param.grad is not None:
        print(f"{name}: grad requires_grad={param.grad.requires_grad}")

5. Model Parallelism and Pipelining

5.1 Simple Model Parallelism

When the model is too large to fit on a single GPU, it needs to be split across multiple GPUs.

Example

# Model parallelism example: place different layers of the model on different GPUs

class ModelParallel(nn.Module):
    def __init__(self, num_gpus=2):
        super().__init__()
        self.num_gpus = num_gpus

        # Part 1: GPU 0
        self.features1 = nn.Sequential(
            nn.Conv2d(3, 64, 3, padding=1),
            nn.ReLU(),
            nn.MaxPool2d(2)
        ).to(f"cuda:0")

        # Part 2: GPU 1
        self.features2 = nn.Sequential(
            nn.Conv2d(64, 128, 3, padding=1),
            nn.ReLU(),
            nn.AdaptiveAvgPool2d(1),
            nn.Flatten()
        ).to(f"cuda:1")

        # Classifier: GPU 0
        self.classifier = nn.Linear(128, 10).to(f"cuda:0")

    def forward(self, x):
        # Data is passed between different GPUs
        x = x.to("cuda:0")
        x = self.features1(x)
        x = x.to("cuda:1")
        x = self.features2(x)
        x = x.to("cuda:0")
        x = self.classifier(x)
        return x


# Use
model = ModelParallel(num_gpus=2)
output = model(torch.randn(1, 3, 32, 32))
print(f"Output device: {output.device}")

5.2 Pipeline Parallelism

PyTorch providestorch.distributed.pipeline.syncmodule (PP package) to implement pipeline parallelism, splitting the model by layers to form a computational pipeline.

Example

# Pipeline parallelism example (requires installing the torch-pipeline dependency)
# pip install torch-pipeline

"""
# Note: PyTorch native pipeline requires a specific version or a third-party library
# Here we demonstrate the concept with a simplified implementation

from torch.distributed.pipeline.sync import Pipe

class SimpleModel(nn.Module):
    def __init__(self):
        super().__init__()
        self.layer1 = nn.Linear(128, 256)
        self.layer2 = nn.Linear(256, 256)
        self.layer3 = nn.Linear(256, 10)

    def forward(self, x):
        x = torch.relu(self.layer1(x))
        x = torch.relu(self.layer2(x))
        x = self.layer3(x)
        return x


# Split the model across different GPUs by layers
# layer1 -> GPU 0
# layer2 -> GPU 1
# layer3 -> GPU 2

model = SimpleModel()
model.layer1 = model.layer1.cuda(0)
model.layer2 = model.layer2.cuda(1)
model.layer3 = model.layer3.cuda(2)

# Use Pipeline to connect the model
# Note: Actual usage requires installing the corresponding version of the pipeline library
# from torch.distributed.pipeline.sync import Pipe
# model = Pipe(model, chunks=4)
"""

print("Pipeline parallelism requires additional installation of the torch-pipeline library")
print("pip install torch-pipeline-sync")

6. Best Practices for Distributed Training

6.1 Performance Optimization Tips

Optimization Item Technique Description
Communication optimization Use NCCL backend NCCL is the best choice for GPUs; faster than Gloo
Data loading Set num_workers and pin_memory Set num_workers to the number of CPU cores; pin_memory accelerates CPU-to-GPU transfer
Gradient synchronization Adjust bucket_cap_mb Larger models can increase the bucket; smaller models can decrease it to improve responsiveness
Mixed precision DDP + AMP Reduces communication volume; improves training speed

6.2 Common Problems and Solutions

Problem Cause Solution
Slow communication Insufficient network bandwidth Use InfiniBand or high-speed network
Insufficient GPU memory Batch size too large Reduce batch size or use gradient accumulation
Unbalanced load Unreasonable model partitioning Adjust layer distribution to be as even as possible
Synchronization error Gradient NaN Check gradient clipping and scaling settings

6.3 Complete Multi-Machine Multi-GPU Training Script

Example

# Complete multi-node multi-GPU training script template

import torch
import torch.nn as nn
import torch.optim as optim
import torch.distributed as dist
from torch.nn.parallel import DistributedDataParallel as DDP
from torch.utils.data import DataLoader, DistributedSampler
import argparse
import os


class TrainingConfig:
    def __init__(self):
        self.local_rank = -1
        self.epochs = 30
        self.batch_size = 32
        self.lr = 1e-3
        self.num_workers = 4
        self.gradient_clip = 1.0
        self.use_amp = True


def init_distributed(config):
    """Initialize distributed environment"""
    # Get configuration from environment variables
    local_rank = config.local_rank

    # NCCL configuration
    os.environ["NCCL_DEBUG"] = "WARN"

    torch.cuda.set_device(local_rank)
    dist.init_process_group(backend="nccl")

    return local_rank


def create_dataloader(dataset, config, world_size, rank):
    """Create data loader"""
    sampler = DistributedSampler(
        dataset,
        num_replicas=world_size,
        rank=rank,
        shuffle=True,
        drop_last=True
    )

    loader = DataLoader(
        dataset,
        batch_size=config.batch_size,
        sampler=sampler,
        num_workers=config.num_workers,
        pin_memory=True,
        persistent_workers=True
    )

    return loader, sampler


def train_epoch(model, loader, criterion, optimizer, scaler, config, rank):
    """Train one epoch"""
    model.train()
    total_loss = 0
    correct = 0
    total = 0

    for inputs, labels in loader:
        inputs = inputs.cuda(rank, non_blocking=True)
        labels = labels.cuda(rank, non_blocking=True)

        optimizer.zero_grad()

        if config.use_amp and scaler is not None:
            # Mixed precision training
            with torch.cuda.amp.autocast():
                outputs = model(inputs)
                loss = criterion(outputs, labels)

            scaler.scale(loss).backward()
            scaler.unscale_(optimizer)
            torch.nn.utils.clip_grad_norm_(model.parameters(), config.gradient_clip)
            scaler.step(optimizer)
            scaler.update()
        else:
            # Normal training
            outputs = model(inputs)
            loss = criterion(outputs, labels)
            loss.backward()
            torch.nn.utils.clip_grad_norm_(model.parameters(), config.gradient_clip)
            optimizer.step()

        total_loss += loss.item() * inputs.size(0)
        _, predicted = outputs.max(1)
        correct += predicted.eq(labels).sum().item()
        total += labels.size(0)

    return total_loss / total, correct / total


def main():
    parser = argparse.ArgumentParser()
    parser.add_argument("--local_rank", type=int, default=-1)
    args = parser.parse_args()

    config = TrainingConfig()
    config.local_rank = args.local_rank

    rank = init_distributed(config)
    world_size = dist.get_world_size()

    print(f"Process {rank}/{world_size} initialization complete")

    # Create the model
    model = YourModel().cuda(rank)
    model = DDP(model, device_ids=[rank])

    # Optimizer
    optimizer = optim.Adam(model.parameters(), lr=config.lr)

    # GradScaler (if using mixed precision)
    scaler = GradScaler() if config.use_amp else None

    # Dataset (replace with your dataset)
    dataset = YourDataset()
    loader, sampler = create_dataloader(dataset, config, world_size, rank)

    criterion = nn.CrossEntropyLoss()

    # Training loop
    for epoch in range(config.epochs):
        sampler.set_epoch(epoch)

        loss, acc = train_epoch(model, loader, criterion, optimizer, scaler, config, rank)

        # Print only on rank 0
        if rank == 0:
            print(f"Epoch {epoch+1}/{config.epochs}: Loss={loss:.4f}, Acc={acc:.4f}")

    dist.destroy_process_group()


# Launch method:
# torchrun --nnodes=2 --nproc_per_node=4 --node_rank=0 --master_addr=192.168.1.1 --master_port=29500 train.py
# torchrun --nnodes=2 --nproc_per_node=4 --node_rank=1 --master_addr=192.168.1.1 --master_port=29500 train.py

The key to distributed training is: 1) choose the appropriate parallel mode (data parallelism vs model parallelism); 2) use DDP instead of DataParallel for better performance; 3) combine with mixed precision training to further accelerate; 4) properly configure the data loader and communication parameters.


7. Distributed Training Monitoring and Debugging

7.1 Common Monitoring Tools

Example

# Distributed training monitoring code

def monitor_training(rank, world_size):
    """Monitor distributed training status"""
    # 1. Check process group status
    print(f"Rank: {rank}, World Size: {world_size}")
    print(f"Backend: {dist.get_backend()}")
    print(f"Rank: {dist.get_rank()}")
    print(f"CUDA Device: {torch.cuda.current_device()}")

    # 2. Monitor GPU memory usage
    if torch.cuda.is_available():
        for i in range(torch.cuda.device_count()):
            allocated = torch.cuda.memory_allocated(i) / 1024**2
            reserved = torch.cuda.memory_reserved(i) / 1024**2
            print(f"GPU {i}: allocated {allocated:.1f} MB, reserved {reserved:.1f} MB")

    # 3. Monitor gradient synchronization
    # DDP automatically records synchronization time
    # Detailed logging can be enabled by setting environment variables
    # os.environ["NCCL_DEBUG"] = "INFO"


# Call periodically in the training loop
monitor_training(rank, dist.get_world_size())

7.2 Debugging Tips

  • Test on a single machine first, then scale to multiple machines
  • Useos.environ["NCCL_DEBUG"] = "INFO"View detailed communication logs
  • Ensure that the clocks on all machines are synchronized (using NTP)
  • Check whether the firewall allows NCCL communication ports

Although distributed training adds complexity, the performance improvement it brings is significant. For large model training, it is almost an essential technique.

Other Extensions