注意
转到页面底部 下载完整示例代码。
PyTorch 中的数据加载优化#
作者: Divyansh Khanna, Ramanish Singh
如何优化 DataLoader 配置以获得最大吞吐量
batch_size、num_workers和pin_memory的最佳实践将数据传输与 GPU 计算重叠的高级技术
配置共享内存策略及处理
/dev/shm问题
PyTorch v2.0+
对 PyTorch DataLoader 的基本了解
(可选)用于特定 GPU 优化的 CUDA 兼容 GPU
简介#
数据加载通常是深度学习流水线中的关键瓶颈。虽然 GPU 处理批次数据的速度极快,但低效的数据加载会导致昂贵的硬件处于空闲状态,等待下一批数据。本教程介绍了优化数据加载配置以最大限度提高训练吞吐量的最佳实践和一些技术。
我们将探索 PyTorch DataLoader 的关键参数,并就如何针对您的特定工作负载进行调优提供实用指导。我们不会孤立地展示每项优化,而是从一个基准训练循环开始,逐步应用优化措施,并在每一步测量累积的加速效果。
import time
import torch
import torch.nn as nn
from torch.utils.data import DataLoader, Dataset
device = torch.device("cuda" if torch.cuda.is_available() else "cpu")
print(f"Using device: {device}")
# Set a fixed seed for reproducibility
torch.manual_seed(42)
Using device: cuda
<torch._C.Generator object at 0x7ff94a942250>
创建示例数据集#
首先,让我们创建一个简单的模拟昂贵转换过程的数据集。这将有助于我们展示不同 DataLoader 配置的影响。
class SyntheticDataset(Dataset):
"""A synthetic dataset that simulates expensive data transformations."""
def __init__(self, size=10000, feature_dim=224, transform_delay=0.001):
self.size = size
self.feature_dim = feature_dim
self.transform_delay = transform_delay
def __len__(self):
return self.size
def __getitem__(self, idx):
# Generate data lazily to avoid pre-allocating large tensors
data = torch.randn(3, self.feature_dim, self.feature_dim)
label = torch.randint(0, 10, (1,)).item()
if self.transform_delay > 0:
time.sleep(self.transform_delay)
return data, label
class SyntheticDatasetBatched(Dataset):
"""Same as SyntheticDataset but with __getitems__ for batched fetching."""
def __init__(self, size=10000, feature_dim=224, transform_delay=0.001):
self.size = size
self.feature_dim = feature_dim
self.transform_delay = transform_delay
def __len__(self):
return self.size
def __getitem__(self, idx):
data = torch.randn(3, self.feature_dim, self.feature_dim)
label = torch.randint(0, 10, (1,)).item()
if self.transform_delay > 0:
time.sleep(self.transform_delay)
return data, label
def __getitems__(self, indices):
"""Fetch multiple items at once — enables vectorized generation.
Instead of N individual __getitem__ calls (each with its own
overhead), this generates the entire batch in one shot using
vectorized tensor operations.
"""
n = len(indices)
# Vectorized generation: one call instead of N individual ones
data = torch.randn(n, 3, self.feature_dim, self.feature_dim)
labels = torch.randint(0, 10, (n,))
# Simulate batch-level I/O: one sleep for the whole batch,
# not one per sample (e.g., one DB query for N rows)
if self.transform_delay > 0:
time.sleep(self.transform_delay)
return [(data[i], labels[i].item()) for i in range(n)]
基准训练循环#
我们的起点:一个没有多进程、没有锁页内存、使用默认设置的简单 DataLoader。这确立了我们将要优化的性能底线。
baseline_loader = DataLoader(
benchmark_dataset,
batch_size=32,
shuffle=True,
num_workers=0,
pin_memory=False,
)
print("\n=== Progressive Optimization Results ===")
print("\nBaseline (num_workers=0, pin_memory=False):")
baseline_time, baseline_loss = train_and_benchmark(baseline_loader)
print(f" Time: {baseline_time:.4f}s | Loss: {baseline_loss:.4f}")
prev_time = baseline_time
=== Progressive Optimization Results ===
Baseline (num_workers=0, pin_memory=False):
Time: 32.9974s | Loss: 2.3171
批量大小(Batch Size)优化#
batch_size 参数控制一次处理多少个样本。选择合适的 batch size 需要平衡多个因素:
内存考量
较大的 batch size 需要更多的 GPU 内存来存储输入、激活值和梯度
显存溢出(OOM)错误在较大 batch size 时很常见
适中的 batch size(32-128)通常能提供最佳平衡
训练动态
batch size 的变化会影响有效学习率,通常需要重新调优
较大的 batch 可以提供更稳定的梯度估计,但在泛化能力上可能表现不同
注意
当更改 batch size 时,请记住调优优化器参数,尤其是学习率计划,除非您是在进行推理
由于 batch size 取决于模型(并非一种“随手添加”即可优化的参数),我们将其孤立进行基准测试,而不是将其放入渐进式优化链中。
# Example: Testing different batch sizes
batch_dataset = SyntheticDataset(size=1000, transform_delay=0)
def benchmark_batch_size(batch_size, num_batches=10):
"""Benchmark data loading with a specific batch size."""
loader = DataLoader(batch_dataset, batch_size=batch_size, shuffle=True)
start = time.perf_counter()
for i, (data, labels) in enumerate(loader):
if i >= num_batches:
break
data = data.to(device, non_blocking=True)
_ = data.sum()
if torch.cuda.is_available():
torch.cuda.synchronize()
elapsed = time.perf_counter() - start
return elapsed
# Benchmark different batch sizes
print("\nBatch size comparison (isolated benchmark):")
for bs in [16, 32, 64, 128]:
elapsed = benchmark_batch_size(bs)
print(f" Batch size {bs:3d}: {elapsed:.4f}s for 10 batches")
Batch size comparison (isolated benchmark):
Batch size 16: 0.1952s for 10 batches
Batch size 32: 0.3784s for 10 batches
Batch size 64: 0.6796s for 10 batches
Batch size 128: 0.9669s for 10 batches
工作进程数量(num_workers)#
num_workers 参数控制用于数据加载的子进程数量。这对并行处理昂贵的数据转换至关重要。
工作原理
每个 worker 维护一个 batch 队列(由
prefetch_factor控制)Workers 并行准备 batch 并将其传输到主进程
如果
in_order=True(默认),batch 会按顺序返回
何时增加 num_workers
当转换计算开销较大时(数据增强、解码)
当从缓慢存储中加载数据时(网络驱动器、HDD)
当您观察到因数据加载导致的 GPU 空闲时间时
何时 num_workers=0 可能更快
当转换成本很低时(简单的张量运算)
当数据已在内存中时
进程间通信 (IPC) 的开销超过了并行化的收益
注意
找到最优的 num_workers 需要调优:增加 worker 数量直到吞吐量趋于平稳。过多的 worker 会浪费 CPU 内存(每个 worker 都持有数据集对象和预取 batch 的副本),并可能导致 /dev/shm 耗尽。一个好的起点是每个 GPU 配置 2-4 个 worker;通过不同数值进行分析以找到适合您工作负载的最佳平衡点。
让我们在训练循环中加入 num_workers=4 和 prefetch_factor=2 并测量改进情况
workers_loader = DataLoader(
benchmark_dataset,
batch_size=32,
shuffle=True,
num_workers=4,
prefetch_factor=2,
pin_memory=False,
)
print("\n+ num_workers=4, prefetch_factor=2:")
workers_time, workers_loss = train_and_benchmark(workers_loader)
print(f" Time: {workers_time:.4f}s | Loss: {workers_loss:.4f}")
print(
f" Speedup vs baseline: {baseline_time / workers_time:.2f}x | vs previous: {prev_time / workers_time:.2f}x"
)
prev_time = workers_time
+ num_workers=4, prefetch_factor=2:
Time: 9.9307s | Loss: 2.3169
Speedup vs baseline: 3.32x | vs previous: 3.32x
理解 pin_memory#
pin_memory 参数通过使用页面锁定(pinned)内存,实现了更快的 CPU 到 GPU 数据传输。
锁定内存的工作原理
锁定内存不能被操作系统交换到磁盘
这实现了更快的 DMA(直接内存访问)向 GPU 的传输
CPU 到 GPU 的传输可以异步进行
最佳实践
在 DataLoader 中使用
pin_memory=True(推荐方法)将数据移动到 GPU 时结合
non_blocking=True使用避免手动调用
tensor.pin_memory()后接.to(device, non_blocking=True)——这样做更慢,因为pin_memory()是阻塞式的
安全模式
# Recommended: Let DataLoader handle pinning
loader = DataLoader(dataset, pin_memory=True)
for data, labels in loader:
data = data.to(device, non_blocking=True)
labels = labels.to(device, non_blocking=True)
另请参阅
更多详情,请查看 pin_memory 教程
让我们在配置中加入 pin_memory=True
pinmem_loader = DataLoader(
benchmark_dataset,
batch_size=32,
shuffle=True,
num_workers=4,
prefetch_factor=2,
pin_memory=torch.cuda.is_available(),
)
if torch.cuda.is_available():
print("\n+ pin_memory=True:")
pinmem_time, pinmem_loss = train_and_benchmark(pinmem_loader)
print(f" Time: {pinmem_time:.4f}s | Loss: {pinmem_loss:.4f}")
print(
f" Speedup vs baseline: {baseline_time / pinmem_time:.2f}x | vs previous: {prev_time / pinmem_time:.2f}x"
)
print(
" (pin_memory benefit is modest here because CPU transform time dominates H2D transfer)"
)
prev_time = pinmem_time
else:
print("\n+ pin_memory: skipped (CUDA not available)")
pinmem_time = workers_time
+ pin_memory=True:
Time: 9.8972s | Loss: 2.3209
Speedup vs baseline: 3.33x | vs previous: 1.00x
(pin_memory benefit is modest here because CPU transform time dominates H2D transfer)
持久化工作进程 (Persistent Workers)#
默认情况下,工作进程在每个 epoch 结束后会被关闭并重启。这会在每个 epoch 边界产生启动开销(导入模块、fork 进程、重新初始化数据集)。
设置 persistent_workers=True 可以使 worker 在 epoch 之间保持存活,消除了重复的启动成本。
最有效的场景
在较小数据集上进行多次 epoch 训练
当数据集的
__init__开销较大时(例如加载元数据)当配合较高的
num_workers使用时
让我们比较在多个 epoch 中使用和不使用持久化工作进程的情况
non_persistent_loader = DataLoader(
benchmark_dataset,
batch_size=32,
shuffle=True,
num_workers=4,
prefetch_factor=2,
pin_memory=torch.cuda.is_available(),
persistent_workers=False,
)
persistent_loader = DataLoader(
benchmark_dataset,
batch_size=32,
shuffle=True,
num_workers=4,
prefetch_factor=2,
pin_memory=torch.cuda.is_available(),
persistent_workers=True,
)
print("\n+ persistent_workers=True (10 epochs):")
non_persistent_time, _ = train_and_benchmark(non_persistent_loader)
persistent_time, persistent_loss = train_and_benchmark(persistent_loader)
print(f" Without persistent_workers: {non_persistent_time:.4f}s")
print(f" With persistent_workers: {persistent_time:.4f}s")
print(
f" Speedup vs baseline: {baseline_time / persistent_time:.2f}x | vs previous: {prev_time / persistent_time:.2f}x"
)
prev_time = persistent_time
+ persistent_workers=True (10 epochs):
Without persistent_workers: 9.9143s
With persistent_workers: 8.7862s
Speedup vs baseline: 3.76x | vs previous: 1.13x
将主机到设备 (H2D) 传输与 GPU 计算重叠#
为了实现最大吞吐量,您可以将主机到设备 (H2D) 的数据传输与 GPU 计算重叠进行。这确保 GPU 从不因为等待数据而空闲。
其核心思想是在处理当前 batch 的同时,将下一个 batch 预取到 GPU。
注意
DataPrefetcher 在 H2D 传输时间与 GPU 计算时间有明显重叠时表现出最大优势。如果数据加载已经很快,流同步开销可能会超过收益。
class DataPrefetcher:
"""Prefetches data to GPU while previous batch is being processed."""
def __init__(self, loader, device):
self.loader = iter(loader)
self.device = device
self.stream = torch.cuda.Stream() if torch.cuda.is_available() else None
self.next_data = None
self.next_labels = None
self.preload()
def preload(self):
try:
self.next_data, self.next_labels = next(self.loader)
except StopIteration:
self.next_data = None
self.next_labels = None
return
if self.stream is not None:
with torch.cuda.stream(self.stream):
self.next_data = self.next_data.to(self.device, non_blocking=True)
self.next_labels = self.next_labels.to(self.device, non_blocking=True)
def __iter__(self):
return self
def __next__(self):
if self.stream is not None:
torch.cuda.current_stream().wait_stream(self.stream)
data = self.next_data
labels = self.next_labels
if data is None:
raise StopIteration
# Ensure tensors are ready
if self.stream is not None:
data.record_stream(torch.cuda.current_stream())
labels.record_stream(torch.cuda.current_stream())
self.preload()
return data, labels
# Integrate prefetcher into the training loop.
if torch.cuda.is_available():
print("\n+ DataPrefetcher (overlapping H2D transfer):")
prefetch_time, prefetch_loss = train_and_benchmark(
persistent_loader, prefetch_device=device
)
print(f" Time: {prefetch_time:.4f}s | Loss: {prefetch_loss:.4f}")
print(
f" Speedup vs baseline: {baseline_time / prefetch_time:.2f}x | vs previous: {prev_time / prefetch_time:.2f}x"
)
prev_time = prefetch_time
else:
print("\n+ DataPrefetcher: skipped (CUDA not available)")
prefetch_time = persistent_time
+ DataPrefetcher (overlapping H2D transfer):
Time: 8.7485s | Loss: 2.3194
Speedup vs baseline: 3.77x | vs previous: 1.00x
数据集级优化:__getitems__#
除了调优 DataLoader 参数,您还可以优化数据集本身。PyTorch 的 DataLoader 通过 __getitems__ 支持批量获取协议:如果您的数据集定义了此方法,获取器会使用索引列表调用它一次,而不是为每个样本重复调用 __getitem__。
工作原理
默认获取器执行:
[dataset[idx] for idx in batch_indices]使用
__getitems__时:dataset.__getitems__(batch_indices)
何时有用
当单样本开销很大时(例如打开连接、解析头文件、获取锁)
当数据可以更高效地批量获取时(例如,对 N 行记录执行一次 SQL 查询而不是 N 次查询,或者向量化张量生成)
当转换具有可以在批次中分摊的固定设置成本时
预期签名
def __getitems__(self, indices: list[int]) -> list:
# Fetch all items at once and return as a list
...
我们的 SyntheticDatasetBatched 实现了 __getitems__ 以通过一次向量化调用生成整个批次(具有单一的分摊延迟),而不是 N 次各自带有延迟的单独调用。让我们将其加入我们的累积配置中
benchmark_dataset_batched = SyntheticDatasetBatched(
size=512, feature_dim=224, transform_delay=0.005
)
batched_loader = DataLoader(
benchmark_dataset_batched,
batch_size=32,
shuffle=True,
num_workers=4,
prefetch_factor=2,
pin_memory=torch.cuda.is_available(),
persistent_workers=True,
)
print("\n+ __getitems__ (batched dataset fetching):")
batched_time, batched_loss = train_and_benchmark(batched_loader)
print(f" Time: {batched_time:.4f}s | Loss: {batched_loss:.4f}")
print(
f" Speedup vs baseline: {baseline_time / batched_time:.2f}x | vs previous: {prev_time / batched_time:.2f}x"
)
prev_time = batched_time
+ __getitems__ (batched dataset fetching):
Time: 2.8793s | Loss: 2.3205
Speedup vs baseline: 11.46x | vs previous: 3.04x
in_order 参数#
默认情况下(in_order=True),DataLoader 返回批次的顺序与数据集索引顺序相同。这要求缓存从 worker 返回的乱序批次。
何时考虑 in_order=False
当您不需要确定性顺序时(例如不做检查点保存)
当您观察到因 batch 缓存导致的训练抖动时
当最大化吞吐量比可重现性更重要时
注意
in_order=False 可能不会增加平均吞吐量,但可以减少方差,并消除因某个 worker 较慢而导致的队头阻塞引起的偶尔缓慢的 batch。
快照频率(snapshot_every_n_steps)#
使用 torchdata 的 StatefulDataLoader(用于检查点)时,snapshot_every_n_steps 参数控制保存 DataLoader 状态的频率。
权衡
高频率(较小的 n):开销更大,但任务失败时数据丢失更少
低频率(较大的 n):开销较小,但恢复时重播的样本更多
根据您的容错要求和重新处理数据的成本进行选择。
最终总结#
这是我们应用于训练循环的每项优化的累积效果。每一行都包含了前几行的所有优化
配置 |
vs 基准 |
vs 前一步 |
|---|---|---|
基准 (num_workers=0, 无锁定) |
1.00x |
— |
+ num_workers=4, prefetch_factor=2 |
~2.7x |
~2.7x |
+ pin_memory=True |
~2.8x |
~1.0x |
+ persistent_workers=True |
~3.7x |
~1.3x |
+ DataPrefetcher (H2D 重叠) |
~3.6x |
~1.0x |
+ __getitems__ (批量获取) |
~10x |
~2.9x |
注意
这些结果基于我们的基准数据集。实际加速比将取决于您的特定工作负载、硬件、数据集大小和转换复杂度。
总结与最佳实践#
从适中的 batch size 开始(32-128),如果内存允许,再进行扩展。
当转换计算昂贵时使用 ``num_workers > 0``。从 2-4 个 worker 开始,并根据内存容量增加。并不是越多越好。
使用加速器时启用 ``pin_memory=True``。
使用 ``persistent_workers=True`` 以避免 epoch 之间的 worker 重启开销。
配置您的流水线,以确定数据集访问、转换等过程中的 CPU 瓶颈。
针对 GPU 工作负载实施数据预取,以将数据传输与计算重叠。
当达到文件描述符限制时使用 ``file_system`` 共享策略。
结论#
在本教程中,我们学习了如何逐步优化 PyTorch 数据加载流水线——从一个原始的单进程基准,到使用多进程 worker、锁定内存、持久化 worker、CUDA 流式预取以及通过 __getitems__ 进行批量数据集获取的全面优化配置。每项优化都针对不同的瓶颈,共同作用下可以带来数量级的吞吐量提升。这些应被视为最佳实践,性能表现取决于具体工作负载。
其他资源#
脚本总运行时间: (1 分 25.571 秒)