AI 工程基础体系 · 第 63/100 篇。内容覆盖机器学习、深度学习与生成式 AI;模型、数据、评测、权限和成本会作为同一生产系统处理。
深度学习数据管道:Dataset、采样、增强、预取与吞吐诊断
深度学习训练中的“数据管道”不是一个简单的文件读取函数,而是一条把样本从存储系统送到设备、再送入模型的有状态流水线。它通常包含:
存储
↓
Dataset:定义“一个样本是什么”
↓
Sampler:决定“取哪些样本、以什么顺序取”
↓
BatchSampler / DataLoader:组成批次并调度读取
↓
collate_fn:把样本拼成模型输入
↓
数据增强:改变样本但尽量保持监督关系
↓
worker:并行准备 CPU 数据
↓
pin memory / prefetch:把后续批次提前准备到更接近设备的位置
↓
GPU/TPU:执行前向、反向和参数更新
数据管道的目标不是单纯“读取更快”。它必须同时满足以下条件:
- 统计正确:样本分布、标签和数据划分符合训练目标。
- 语义正确:增强不会破坏标签含义,训练和评测边界清晰。
- 设备高效:模型计算时尽量不等待数据。
- 资源可控:CPU、内存、显存、存储带宽和进程数不会失控。
- 可复现、可诊断:出现错误时能够定位到样本、worker、存储还是设备。
- 可治理:数据权限、敏感信息、版本、评测集和成本能够被追踪。
一、先区分 Dataset、Sampler、DataLoader 和 collate_fn
这些组件经常被统称为“数据加载器”,但它们承担的职责不同。职责混淆后,最常见的结果是数据重复、采样偏差、增强位置错误或多进程状态失效。
1. Dataset 定义样本访问语义
在 PyTorch 中,最常见的 Dataset 是 map-style dataset。它提供:
__len__() # 数据集大小
__getitem__(i) # 根据索引 i 返回一个样本
例如,一个分类样本可以表示为:
Dataset 负责把索引 映射到样本:
它通常应当负责:
- 根据索引定位文件或记录;
- 解码、解析和基础类型转换;
- 返回样本及其元数据;
- 在必要时执行样本级变换。
它不应隐式决定“这一轮训练取哪些样本”。后者属于 sampler。
一个最小的 map-style dataset 如下:
from pathlib import Path
from PIL import Image
import torch
from torch.utils.data import Dataset
class ImageClassificationDataset(Dataset):
def __init__(self, records, transform=None):
"""
records: [(image_path, class_id), ...]
transform: 接收 PIL Image,返回模型所需输入
"""
self.records = records
self.transform = transform
def __len__(self):
return len(self.records)
def __getitem__(self, index):
image_path, label = self.records[index]
# 这里的异常应包含路径,便于定位坏样本
try:
image = Image.open(image_path).convert("RGB")
except Exception as exc:
raise RuntimeError(
f"failed to read image: index={index}, path={image_path}"
) from exc
if self.transform is not None:
image = self.transform(image)
return {
"image": image,
"label": torch.tensor(label, dtype=torch.long),
"index": index,
"path": str(image_path),
}
这里的 index 和 path 不一定要送入模型,但在诊断标签错位、坏文件和数据泄漏时很有价值。
2. IterableDataset 定义数据流,而不是随机索引
IterableDataset 提供的是:
__iter__()
它适合:
- 流式读取超大文件;
- Kafka、对象存储清单或数据库游标;
- 数据量无法预先知道;
- 样本只能顺序扫描;
- 生成式训练中的文本流。
但多 worker 时必须显式处理分片。否则每个 worker 可能从头读取同一份数据,导致样本重复。
from torch.utils.data import IterableDataset, get_worker_info
class LineDataset(IterableDataset):
def __init__(self, path):
self.path = path
def __iter__(self):
worker_info = get_worker_info()
worker_id = 0
num_workers = 1
if worker_info is not None:
worker_id = worker_info.id
num_workers = worker_info.num_workers
with open(self.path, "r", encoding="utf-8") as f:
for line_number, line in enumerate(f):
# 让第 k 个 worker 处理 line_number % num_workers == k
if line_number % num_workers != worker_id:
continue
line = line.rstrip("\n")
if line:
yield line
这个分片方法假设所有 worker 都能看到同一个文件,并且按行扫描成本可接受。它不能自动解决分布式训练中的 rank 分片;多机多卡时还要同时按 rank/world_size 分片,否则不同进程之间仍会读取重复数据。
3. Sampler 决定索引序列
对于 map-style dataset,sampler 产生索引序列:
例如:
SequentialSampler:按 顺序;RandomSampler:随机打乱或有放回抽样;WeightedRandomSampler:按给定权重抽样;DistributedSampler:把索引划分给不同训练进程。
Sampler 不负责读取图像,也不负责把多个样本拼成 batch。
4. BatchSampler 决定索引分组
如果 batch size 为 ,sampler 产生的索引会被分组为:
通常 DataLoader(batch_size=B, shuffle=True) 会替你创建 sampler 和 batch sampler。需要按长度、token 数或其他复杂规则组 batch 时,才常常直接使用自定义 batch_sampler。
5. collate_fn 把样本变成批次
Dataset.__getitem__ 返回的是单个样本;collate_fn 接收一个样本列表:
[
{"image": tensor_1, "label": 0},
{"image": tensor_2, "label": 1},
]
并把它变成:
{
"image": tensor([tensor_1, tensor_2]),
"label": tensor([0, 1]),
}
默认 collate 要求张量尺寸通常一致。NLP 和视觉检测任务的样本长度或目标数量可能不同,需要自定义填充或保留列表:
def collate_variable_length(batch):
tokens = [item["tokens"] for item in batch]
labels = torch.tensor([item["label"] for item in batch])
return {
"tokens": tokens, # 保留不同长度,模型内部或后续步骤处理 padding
"label": labels,
}
collate_fn 是一个重要边界:如果在这里进行大量 CPU 计算,瓶颈可能看起来像“DataLoader 慢”,实际是批处理逻辑过重。
二、从“样本顺序”推导采样是否正确
采样不是单纯的随机化。它改变了训练过程中观察到的数据分布,因此可能改变优化目标。
1. 均匀采样对应什么目标
假设训练集有 个样本,单样本损失为:
经验风险是:
如果索引 按均匀分布抽取:
那么单次随机梯度的期望为:
这说明,均匀采样的随机梯度是经验风险梯度的无偏估计。一个 batch 只是在降低估计方差,前提是 batch 中的样本确实来自目标分布。
2. 类别不均衡时,改变采样就改变了目标
假设有两类:
- 类别 A:90 个样本;
- 类别 B:10 个样本。
普通均匀采样时:
如果使用等类别采样,变成:
这会使少数类在训练中被更频繁看到,但它不再直接优化原始经验风险,而是在优化一个类别重新加权后的目标。
使用 WeightedRandomSampler 的示例:
import torch
from torch.utils.data import WeightedRandomSampler, DataLoader
labels = torch.tensor([0] * 90 + [1] * 10)
class_counts = torch.bincount(labels)
class_weights = 1.0 / class_counts.float()
# 每个样本的抽样权重
sample_weights = class_weights[labels]
sampler = WeightedRandomSampler(
weights=sample_weights,
num_samples=len(labels),
replacement=True,
)
loader = DataLoader(
list(zip(torch.arange(len(labels)), labels)),
batch_size=20,
sampler=sampler,
)
这里必须注意 replacement=True 的含义:一次抽样后样本仍可再次被抽到。因此一个“epoch”不再等同于“每个原始样本恰好访问一次”。num_samples=len(labels) 只是人为定义每个 epoch 抽多少次,不代表覆盖全部样本。
3. 采样加权与损失加权不是同一件事
若从分布 抽样,但目标仍是均匀经验风险,则可以使用重要性权重:
其中:
- 是目标分布;
- 是实际 sampler 分布;
- 。
由于:
所以估计仍然无偏。实际训练中常常不做完整校正,因为校正会增加方差,且很多任务有意改变类别权重。关键是明确你要优化的是:
- 原始样本分布;
- 类别平衡分布;
- 困难样本分布;
- 还是某个业务风险分布。
4. shuffle=True 的边界
对 map-style dataset,shuffle=True 通常意味着由 DataLoader 创建随机 sampler。不能同时传入 shuffle=True 和 sampler,因为两者都在决定索引顺序,PyTorch 会拒绝这种冲突配置。
训练、验证和测试通常不同:
train_loader = DataLoader(
train_dataset,
batch_size=64,
shuffle=True,
)
valid_loader = DataLoader(
valid_dataset,
batch_size=128,
shuffle=False,
)
验证集不需要随机化来估计整体损失;保持固定顺序还有利于复现和逐样本对比。随机增强也不应默认应用在验证集,否则同一模型每次评测可能看到不同输入。
5. 分布式训练中的 sampler 状态
分布式数据并行通常有 个进程,每个进程处理全局索引的一部分。DistributedSampler 的关键属性是:
from torch.utils.data import DistributedSampler
sampler = DistributedSampler(
train_dataset,
shuffle=True,
drop_last=True,
)
loader = DataLoader(
train_dataset,
batch_size=64,
sampler=sampler,
)
for epoch in range(num_epochs):
sampler.set_epoch(epoch)
for batch in loader:
train_step(batch)
set_epoch(epoch) 不是装饰性调用。它让不同 epoch 使用不同的确定性随机排列。如果不调用,某些实现下每个 epoch 可能重复同一排列。
还要区分:
batch_size通常是每个进程的 batch size;- 全局 batch size 约为
drop_last=True可能丢弃尾部样本,但能避免不同进程最后一个 batch 形状不一致;- 评测时必须确认是否需要补齐样本,以及补齐样本是否会被错误计入指标。
三、数据增强改变的是输入分布,不应破坏监督关系
数据增强是对样本施加变换:
其中 是随机变量。若任务标签在变换下保持不变:
那么训练目标可以写成:
这不是“凭空增加了新数据”,而是让模型在一个由原始样本诱导出的增强分布上训练。
1. 标签保持性决定增强是否合法
对于猫狗分类,水平翻转通常近似保持标签;对于“文字朝向识别”,水平翻转可能改变标签;对于目标检测,裁剪、缩放不仅要变换图像,还要同步变换 bounding box;对于分割,图像和 mask 必须使用同一几何随机参数。
错误示例:
# 错误:图像随机裁剪,但 mask 没有使用同一裁剪区域
image = random_crop(image)
mask = random_crop(mask)
两个 random_crop 可能产生不同区域,导致监督错位。正确做法是先采样一次参数,再同时应用:
params = sample_crop_params(image)
image = crop(image, params)
mask = crop(mask, params)
2. 训练增强和评测预处理必须分开
训练变换通常包含随机性,验证和测试通常只包含确定性预处理:
from torchvision import transforms
train_transform = transforms.Compose([
transforms.RandomResizedCrop(224),
transforms.RandomHorizontalFlip(),
transforms.ToTensor(),
transforms.Normalize(mean=(0.485, 0.456, 0.406),
std=(0.229, 0.224, 0.225)),
])
eval_transform = transforms.Compose([
transforms.Resize(256),
transforms.CenterCrop(224),
transforms.ToTensor(),
transforms.Normalize(mean=(0.485, 0.456, 0.406),
std=(0.229, 0.224, 0.225)),
])
如果把 RandomResizedCrop 放进验证集,评测指标会混入增强随机性;如果训练和验证的归一化方式不同,指标变化可能来自输入尺度而不是模型能力。
3. 随机增强、缓存和复现之间存在冲突
若增强在 __getitem__ 中在线执行,则每次取样都可能产生不同结果。这样可以提高增强多样性,但会增加 CPU 计算。
若把增强结果预先缓存,则读取更快,但:
- 随机性减少;
- 缓存占用增加;
- 变换版本变化后需要重建;
- 可能把训练增强错误地带入验证或测试。
数据管道的可复现性通常至少需要控制三类随机源:
- Python 的
random; - NumPy 的随机数;
- PyTorch 的 CPU/GPU 随机数。
多 worker 时,每个 worker 还应有不同且可追踪的 seed。下面的初始化函数只控制 Python、NumPy 和 PyTorch 的 worker 随机源:
import random
import numpy as np
import torch
def seed_worker(worker_id):
worker_seed = torch.initial_seed() % 2**32
np.random.seed(worker_seed)
random.seed(worker_seed)
generator = torch.Generator()
generator.manual_seed(1234)
loader = DataLoader(
train_dataset,
batch_size=64,
shuffle=True,
num_workers=4,
worker_init_fn=seed_worker,
generator=generator,
)
这并不自动保证整个训练完全确定。GPU 算子、并行归约、第三方解码库和底层文件顺序也可能引入差异。复现目标应明确是“同一随机种子下趋势一致”,还是“逐步逐位一致”;后者通常会牺牲性能,并且需要更多确定性配置。
四、DataLoader 的并发模型与生命周期
1. num_workers=0 与多 worker
当 num_workers=0 时,主训练进程执行:
取样本 → 解码 → 增强 → collate → 送 GPU → 训练
数据准备和模型计算容易串行化。
当 num_workers>0 时,通常变成:
主进程:向 worker 请求索引
worker:读取、解码、增强、collate
主进程:接收已完成的 batch
主进程:复制到 pinned memory,再异步送 GPU
worker 数量增大并不必然更快,因为每个 worker 都可能竞争:
- CPU 核;
- 文件描述符;
- 存储带宽;
- 解码库内部线程;
- 主机内存;
- 进程间传输带宽。
num_workers=16 在一台机器上可能比 num_workers=4 更慢,尤其当数据存储已经是瓶颈时。
2. persistent_workers 的生命周期
默认情况下,DataLoader 迭代结束后,worker 可能被关闭;下一个 epoch 再启动。设置:
persistent_workers=True
可以让 worker 跨 epoch 存活,减少反复创建进程的开销。它适用于:
- epoch 较多;
- worker 初始化昂贵;
- Dataset 内部打开资源的成本较高。
但它也意味着:
- worker 中的 Dataset 对象不会每个 epoch 自动重建;
- 主进程中修改 Dataset 属性,不一定同步到已有 worker;
- worker 内部缓存会持续占用内存;
- 数据文件更新后,旧 worker 可能仍持有旧状态。
如果 Dataset 按 epoch 改变采样策略,不能假设简单修改主进程字段就能传给 persistent worker。应通过 sampler、共享状态或重新创建 DataLoader 明确传递。
3. prefetch_factor 的含义
当 num_workers>0 时,prefetch_factor 表示每个 worker 预先准备的 batch 数,常见实现默认值和具体行为应以当前 PyTorch 版本文档为准。粗略地说,待处理 batch 数量与:
同量级。
因此增大 prefetch_factor 可能隐藏短时读取抖动,但也会增加:
- 预取内存;
- 未使用 batch 被丢弃时的浪费;
- 进程间队列压力;
- 随机增强提前消耗的随机状态。
它不能提高存储设备的物理带宽。若瓶颈是磁盘读带宽,预取只能让等待更早发生。
4. pin_memory 和 non_blocking
Pinned memory,也称 page-locked memory,是不会被操作系统随意换出的主机内存。CUDA 可以更高效地从中发起主机到 GPU 的传输。
loader = DataLoader(
train_dataset,
batch_size=64,
num_workers=4,
pin_memory=True,
)
for batch in loader:
images = batch["image"].to("cuda", non_blocking=True)
labels = batch["label"].to("cuda", non_blocking=True)
这里有三个边界:
pin_memory=True只对 DataLoader 返回的 CPU 对象尝试固定内存;non_blocking=True是允许异步传输的条件之一,不代表任何设备或任何内存都必然异步;- 如果紧接着在 CPU 上访问传输结果,异步优势会被同步操作抵消。
Pinned memory 也不是越多越好。过度固定主机内存可能影响系统整体内存压力,因此应观察 RSS、swap 和系统响应,而不是只看 GPU 利用率。
五、预取的本质:用队列重叠数据准备与模型计算
设:
- :准备一个 batch 的时间,包括读取、解码、增强和 collate;
- :CPU 到 GPU 的传输时间;
- :模型前向、反向和优化器更新时间。
没有重叠时,一个 batch 的周期近似为:
理想重叠时,数据准备和模型计算可以并行,稳定吞吐受最慢阶段限制:
预取的作用是让当前 batch 在计算时,后续 batch 已经处于准备队列中:
sequenceDiagram
participant W as DataLoader workers
participant Q as Prefetch queue
participant H as Host pinned memory
participant G as GPU
participant M as Model compute
W->>Q: 读取、解码、增强、collate batch k+1
Q->>H: 准备 CPU batch k
H->>G: 异步传输 batch k
G->>M: 前向/反向 batch k
M-->>G: 计算结束
W->>Q: 继续准备 batch k+2
H->>G: 异步传输 batch k+1
关键路径是:
worker 产生 batch k
→ 主进程取得 batch k
→ batch k 进入 pinned memory
→ 传输到 GPU
→ 模型使用 batch k
如果 worker 速度低于模型消耗速度,队列会逐渐变空,GPU 等待数据;如果 worker 速度高很多,队列会持续满,额外增加内存,却不会提升训练吞吐。
手动预取与 DataLoader 预取
DataLoader 自带 worker 队列已经提供了一层预取。训练循环中再写一个复杂的 CUDA stream 预取器,只有在确认 CPU 到 GPU 拷贝是瓶颈时才有意义。否则容易引入:
- stream 同步错误;
- batch 生命周期错误;
- 显存引用长期持有;
- 调试困难;
- 结束时最后一个 batch 处理错误。
首先应使用普通的 pin_memory + non_blocking 和合理 worker 数量测量;只有 profiling 证明需要更细粒度重叠时,才引入自定义 CUDA stream。
六、一个可运行的端到端示例
下面示例使用 PyTorch 自带张量生成数据,展示 Dataset、增强、采样、DataLoader、预取、训练和计时。它不依赖外部图像文件,便于直接运行。
import time
import random
import numpy as np
import torch
from torch import nn
from torch.utils.data import (
Dataset,
DataLoader,
WeightedRandomSampler,
)
class ToyDataset(Dataset):
def __init__(self, n=5000, dim=32, transform=None):
self.x = torch.randn(n, dim)
# 人为构造二分类标签
self.y = (self.x[:, 0] + 0.5 * self.x[:, 1] > 0).long()
self.transform = transform
def __len__(self):
return len(self.y)
def __getitem__(self, index):
x = self.x[index].clone()
y = self.y[index]
if self.transform is not None:
x = self.transform(x)
return {
"x": x,
"y": y,
"index": index,
}
class AddNoise:
def __init__(self, std=0.05):
self.std = std
def __call__(self, x):
# 该增强假设小幅噪声不改变标签
return x + torch.randn_like(x) * self.std
def seed_worker(worker_id):
worker_seed = torch.initial_seed() % (2**32)
random.seed(worker_seed)
np.random.seed(worker_seed)
def main():
device = torch.device("cuda" if torch.cuda.is_available() else "cpu")
dataset = ToyDataset(
n=5000,
dim=32,
transform=AddNoise(std=0.05),
)
labels = dataset.y
counts = torch.bincount(labels)
class_weights = 1.0 / counts.float()
sample_weights = class_weights[labels]
# 这里故意采用有放回采样,展示不均衡数据的处理方式
sampler = WeightedRandomSampler(
weights=sample_weights,
num_samples=len(dataset),
replacement=True,
)
loader_kwargs = dict(
dataset=dataset,
batch_size=128,
sampler=sampler,
num_workers=2,
pin_memory=(device.type == "cuda"),
persistent_workers=True,
worker_init_fn=seed_worker,
)
# num_workers=0 时不应传 persistent_workers=True
if loader_kwargs["num_workers"] == 0:
loader_kwargs.pop("persistent_workers")
loader = DataLoader(**loader_kwargs)
model = nn.Sequential(
nn.Linear(32, 64),
nn.ReLU(),
nn.Linear(64, 2),
).to(device)
optimizer = torch.optim.AdamW(model.parameters(), lr=1e-3)
criterion = nn.CrossEntropyLoss()
for epoch in range(3):
model.train()
total_loss = 0.0
total_count = 0
start = time.perf_counter()
for batch in loader:
x = batch["x"].to(device, non_blocking=True)
y = batch["y"].to(device, non_blocking=True)
optimizer.zero_grad(set_to_none=True)
logits = model(x)
loss = criterion(logits, y)
loss.backward()
optimizer.step()
total_loss += loss.item() * y.size(0)
total_count += y.size(0)
elapsed = time.perf_counter() - start
print(
f"epoch={epoch} "
f"loss={total_loss / total_count:.4f} "
f"samples={total_count} "
f"samples_per_sec={total_count / elapsed:.1f}"
)
if __name__ == "__main__":
main()
这个示例中每一步为何成立
ToyDataset.__getitem__根据索引返回单样本,符合 map-style Dataset 协议。AddNoise在训练取样时执行随机增强;它只改变输入,不改变二分类标签。WeightedRandomSampler根据类别频率反比设置样本权重,少数类更容易被抽到。replacement=True允许同一个原始索引在一个 epoch 中出现多次,因此num_samples明确定义了 epoch 的抽样次数。num_workers=2使数据准备在子进程执行;在 Windows 或某些启动方式下,if __name__ == "__main__"是必要的。pin_memory只在 CUDA 设备上启用,避免在 CPU 训练时引入无意义的固定内存。.to(device, non_blocking=True)允许在满足条件时重叠主机到设备的传输。samples_per_sec是端到端吞吐,不等同于纯 Dataset 读取速度;它包含模型计算和优化器更新。
示例输出的具体数值取决于 CPU、GPU、PyTorch 版本和进程启动方式,不能把某台机器上的数值当作性能基准。应重点观察 loss 是否下降、每个 epoch 样本数是否符合预期,以及改动参数后吞吐如何变化。
七、用形式化指标理解“吞吐”
1. 样本吞吐和 batch 吞吐
若一个 epoch 处理 个样本,耗时 ,样本吞吐为:
若 batch size 为 ,batch 吞吐为:
两者不能混淆。一个 batch 变大后,batch/s 可能下降,但 samples/s 可能上升。
分布式训练还要说明是:
- 单进程 samples/s;
- 单机 samples/s;
- 全局 samples/s。
全局吞吐近似为所有 rank 在同步条件下完成的全局样本数除以墙钟时间,而不是简单相加某些不同步的局部计时。
2. 有效吞吐不只是读取速度
训练系统通常受多个阶段限制:
但由于阶段可以重叠,实际周期更接近最慢关键阶段,而不是所有时间简单相加。并且同步点、队列空转、GPU kernel 间隙和分布式通信会改变这个近似。
3. GPU 利用率低不等于 Dataset 慢
GPU 利用率低可能来自:
- 数据队列为空;
- batch 太小,kernel 启动开销占比高;
- 模型本身包含大量同步操作;
- CPU 到 GPU 传输未重叠;
- 分布式 rank 之间等待最慢进程;
- 显存不足导致 batch 被迫变小;
- 前向计算本身很轻。
因此“GPU 利用率 40%”只能说明 GPU 没有持续满载,不能直接证明瓶颈在文件读取。
八、吞吐诊断:先拆分等待时间,再调整参数
1. 训练循环中的基础计时
下面的计时方法可以区分“等待下一个 batch”和“处理当前 batch”:
import time
import torch
data_wait = 0.0
compute_time = 0.0
end = time.perf_counter()
for batch in loader:
now = time.perf_counter()
data_wait += now - end
x = batch["x"].to(device, non_blocking=True)
y = batch["y"].to(device, non_blocking=True)
if device.type == "cuda":
torch.cuda.synchronize()
compute_start = time.perf_counter()
optimizer.zero_grad(set_to_none=True)
logits = model(x)
loss = criterion(logits, y)
loss.backward()
optimizer.step()
if device.type == "cuda":
torch.cuda.synchronize()
compute_time += time.perf_counter() - compute_start
end = time.perf_counter()
print(f"data_wait={data_wait:.3f}s")
print(f"compute_time={compute_time:.3f}s")
为什么 CUDA 计时需要同步?CUDA 操作通常异步提交到设备。若不调用 torch.cuda.synchronize(),CPU 计时可能只测到了“提交 kernel 的时间”,而不是 kernel 实际执行时间。
这个诊断代码适合短时间实验,不适合永久放在高性能生产路径中,因为同步会影响重叠关系。更精细的分析应使用 PyTorch Profiler 或 Nsight 等工具,并明确记录版本和采样范围。
2. 观察队列是否为空
一种简单现象判断:
- 每次
next(loader_iter)都明显等待; - 等待时间接近模型计算时间或更长;
- GPU 在 batch 之间出现空隙;
这通常说明数据供给不足。
相反:
- 取 batch 几乎不等待;
- GPU 计算持续;
- 增加 worker 后吞吐不变;
说明数据管道可能已经不是瓶颈,继续增加 worker 只会增加资源消耗。
3. 参数实验必须一次只改一个因素
可以按如下顺序做小规模对照:
num_workers=0,确认 Dataset 单进程逻辑正确;- 测试
num_workers=1,2,4,8; - 对比
pin_memory=False/True; - 对比
persistent_workers=False/True; - 在确认队列短缺后,调整
prefetch_factor; - 再测试 batch size、增强开关和存储格式。
每次记录:
- 总耗时;
- samples/s;
- CPU 利用率;
- 内存 RSS;
- GPU 利用率;
- GPU 显存;
- 存储读取带宽;
- worker 异常和系统日志。
如果同时修改五个参数,即使吞吐变化,也无法知道哪项产生了变化。
4. 常见症状与因果路径
症状:num_workers=0 正常,多 worker 卡住
可能原因包括:
- Dataset 持有不能安全跨进程复制的对象;
- 文件句柄在 fork 后状态异常;
- 第三方库内部线程与多进程组合导致死锁;
__getitem__中存在不可序列化对象;- Windows 下缺少主入口保护;
- worker 内存异常退出,主进程等待队列结果。
诊断方法:
- 先设
num_workers=0获得完整 traceback; - 缩小到单个样本;
- 将
__getitem__中的解码、增强逐步注释; - 再从
num_workers=1开始增加; - 检查系统内存、文件描述符和 worker 日志。
不要一开始就用 timeout 或强制终止掩盖错误;超时只说明主进程没有按时拿到 batch,不说明根因。
症状:CPU 使用率很高,但吞吐不升反降
常见原因是 worker 数量超过有效并行度,或图像解码库、BLAS、OpenMP 自身又启动了线程,形成线程过量。进程数乘以内部线程数可能远超 CPU 核数。
可验证:
- 减小 worker 数;
- 限制底层库线程数;
- 分别关闭解码和增强;
- 比较 CPU 利用率、上下文切换和吞吐。
症状:内存逐 epoch 增长
可能路径包括:
- Dataset 或 transform 把每个样本缓存到列表;
- 训练代码保存了带计算图的 loss 或输出;
persistent_workers=True后 worker 缓存持续存在;- prefetch 队列过大;
- batch 被主进程或日志对象长期引用;
- 图像解码对象没有及时释放。
“缓存”不是自动优化。必须定义缓存上限、淘汰策略和失效条件,并验证内存是否稳定。
症状:训练 loss 正常,但验证指标异常
优先检查:
- 训练和验证标签映射是否一致;
- 数据划分是否发生泄漏;
- 验证集是否误用了训练增强;
- sampler 是否改变了指标计算中的样本权重;
- 分布式评测是否重复或遗漏样本;
- 文本任务的 padding、attention mask 和截断是否一致。
这类问题通常不是吞吐问题,而是数据语义问题;单纯加速 DataLoader 不会修复指标错误。
九、数据划分、泄漏与生成式 AI 的特殊边界
1. 划分应先于随机增强
训练、验证和测试划分定义的是不同数据子集。应先按实体、用户、时间、文档或会话划分,再在训练子集内执行随机增强。
如果同一用户的近重复样本分别进入训练和测试,模型可能记住用户特征;如果先把一张图片切成多个 patch,再随机划分 patch,来自同一原图的内容可能同时出现在训练和测试中。
形式上,测试集应尽量近似独立于训练过程:
“样本对象不完全相同”还不够,近重复、同一实体和时间未来信息也可能造成实际泄漏。
2. 生成式 AI 中的 Dataset 不只是文本列表
语言模型训练的样本可能是:
- 文档;
- 对话;
- 指令与响应;
- 偏好对;
- 多模态消息;
- token block。
一个 token block 通常来自:
训练目标可能是 next-token prediction:
此时 Dataset 不仅要返回 token,还要正确处理:
input_ids;labels的右移关系;- padding;
attention_mask;- 对不应计算损失的 prompt 部分设置 ignore label;
- 长文档截断和跨 block 边界。
把不同对话直接拼接而不加入边界标记,可能使模型学习到错误的跨样本上下文。把用户数据未经权限审核放入训练 Dataset,则属于数据治理和隐私问题,不是一个普通的读取错误。
3. 评测集必须是受保护的数据资产
生产系统中,评测集不应与训练增强缓存、自动清洗脚本和临时实验目录混用。应记录:
- 数据集版本;
- 样本清单哈希;
- 标签或标注版本;
- 访问权限;
- 生成时间;
- 是否包含个人信息或商业机密;
- 哪些模型和实验读取过它。
数据权限不仅限制“谁能下载文件”,也应限制 Dataset worker 通过凭证访问哪些对象。多 worker 会复制或继承部分进程状态,不能把长期有效的高权限凭证随意放进 Dataset 对象。
十、存储、预处理和缓存的取舍
1. 小文件数量会影响读取效率
大量小图片或小 JSON 文件会产生:
- 目录查找;
- 文件打开关闭;
- 元数据请求;
- 随机 I/O;
- 网络对象存储请求开销。
将样本打包成更适合顺序读取或批量读取的格式,可能减少这些固定开销。但打包也带来:
- 局部更新困难;
- 单个样本损坏影响范围扩大;
- 随机访问索引需要维护;
- 压缩与解压 CPU 成本;
- 数据版本和权限管理复杂度。
没有一种格式对所有任务都最快。应使用真实样本和真实存储环境测量,而不是仅根据文件扩展名判断。
2. 预处理缓存的正确性条件
若缓存函数为:
那么缓存只有在以下条件满足时才可安全复用:
- 原始样本版本不变;
- 预处理代码版本 不变;
- tokenizer、词表或归一化统计量不变;
- 输出 dtype 和 shape 约定不变;
- 权限允许该缓存被当前训练任务读取。
缓存键至少应包含数据版本和处理配置的稳定标识。否则模型可能训练在旧 tokenizer 或旧标签映射上,却没有明显运行错误。
3. CPU 增强与 GPU 增强
CPU 增强减少 GPU 工作,但可能成为 worker 瓶颈;GPU 增强可以利用设备并行能力,但会占用显存和计算资源,并可能与模型竞争 GPU。
选择依据是增强的性质:
- 解码、文件读取通常在 CPU/存储侧;
- 简单张量变换可考虑 GPU;
- 需要复杂第三方库的处理通常留在 CPU;
- 训练与评测必须保持语义一致。
不要为了“GPU 利用率更高”把所有处理都搬到 GPU。若数据传输和模型已经占满设备,GPU 增强反而降低训练速度。
十一、错误处理、数据质量和故障恢复
1. 坏样本不应被静默吞掉
以下写法容易掩盖数据问题:
def __getitem__(self, index):
try:
return read_sample(index)
except Exception:
return None
默认 collate 可能随后因为 None 报错,或者自定义 collate 悄悄丢弃样本,造成样本数和类别分布变化。
更可靠的策略是:
- 在离线校验阶段扫描文件、标签、尺寸和编码;
- 训练时异常包含数据 ID、路径和原始异常;
- 对可容忍的脏数据使用显式策略;
- 统计跳过数量、类别分布和每个 shard 的错误率;
- 错误超过阈值时停止任务,而不是继续产生不可解释的模型。
2. worker 异常的传播路径
典型路径是:
Dataset.__getitem__ 抛出异常
→ worker 捕获并传递异常信息
→ 主进程在取 batch 时重新抛出
→ 训练循环停止
如果看到的 traceback 只显示 DataLoader 位置,应先用 num_workers=0 重跑,使异常在主进程直接发生。多 worker 的并发会改变错误出现的时间,但不会改变坏样本本身。
3. 中断与恢复
训练恢复不仅需要 checkpoint 中的模型和优化器状态,还可能涉及:
- epoch;
- batch 内进度;
- sampler 的随机状态;
- DataLoader worker 随机状态;
- 数据集版本;
- 当前数据清单;
- 分布式 rank 配置;
- 梯度累积位置。
如果只保存模型参数,恢复后通常不能保证从相同样本顺序继续。很多训练任务接受“从相同 epoch 的新随机顺序继续”,但应在实验记录中明确,而不是误称为完全恢复。
十二、一个系统化的诊断顺序
当训练吞吐低或指标异常时,可以按因果链逐层排查。
第一步:验证样本语义
随机打印若干样本,检查:
- 输入 shape、dtype、数值范围;
- 标签和样本是否对应;
- 图像是否旋转、裁剪或归一化正确;
- 文本 token 是否包含错误边界;
- mask、box、label 是否同步变换;
- 训练和验证变换是否被错误复用。
先证明“数据正确”,再谈“数据更快”。
第二步:隔离 Dataset
用单进程、少量样本测量:
start = time.perf_counter()
for i in range(1000):
_ = dataset[i]
elapsed = time.perf_counter() - start
print(f"dataset samples/s = {1000 / elapsed:.1f}")
这测量的是 __getitem__,不包含多 worker 调度和模型计算。若这里已很慢,应检查解码、网络请求、随机增强和 Python 循环。
第三步:隔离 DataLoader
start = time.perf_counter()
count = 0
for batch in loader:
count += batch["x"].shape[0]
if count >= 10000:
break
elapsed = time.perf_counter() - start
print(f"loader samples/s = {count / elapsed:.1f}")
比较不同 num_workers 后,判断并行是否有效。若 Dataset 很快而 DataLoader 很慢,重点检查 collate、进程间复制、内存和文件系统。
第四步:加入设备传输但不训练
测量 CPU batch 到 GPU 的时间,区分 pinned memory 是否有帮助。需要 CUDA 同步,否则时间可能不准确。
第五步:加入模型计算
最后测量端到端吞吐。若单独 DataLoader 足够快,但端到端仍慢,瓶颈可能在模型、通信、同步或优化器,而不是数据管道。
第六步:观察资源与错误
同时记录:
nvidia-smi中 GPU 利用率、显存和功耗;- CPU 每核利用率;
- 主机内存和 swap;
- 存储带宽、IO wait;
- worker 日志;
- 分布式训练中各 rank 的 step 时间。
如果某个 rank 长期慢于其他 rank,整体吞吐会被最慢 rank 限制,即使其他 GPU 看起来空闲。
十三、常见误解与真实边界
误解一:增加 worker 一定提高速度
只有在数据准备阶段存在可并行工作且系统尚未达到其他瓶颈时,增加 worker 才可能提高吞吐。若瓶颈是单块磁盘、网络存储、GPU 计算或内存带宽,worker 增加只会放大竞争。
误解二:prefetch 越大越快
预取只能把未来工作提前安排。它不能改变单个 batch 的解码成本,也不能突破存储带宽上限。过大预取会提高内存占用,并可能准备大量最终不会被使用的 batch。
误解三:pin memory 会自动让传输异步
Pinned memory 是异步传输的必要条件之一,但还需要正确使用设备传输、CUDA stream 和同步边界。没有测量时,不能假设它一定带来可见收益。
误解四:加权采样只是“修复不均衡”
加权采样改变了训练期间的样本分布。它可能改善少数类召回,也可能增加梯度方差、重复少数类样本并损害原始分布上的校准。训练目标、损失权重和最终评测分布必须一起设计。
误解五:随机增强等于收集了更多真实数据
增强通常只能提供原始数据周围的变体,不能替代真实场景覆盖。错误增强还会产生系统性标签噪声。增强强度应由验证集、鲁棒性测试和业务失效案例共同决定。
误解六:训练能跑通就说明 Dataset 正确
代码运行成功只能说明类型和 shape 在某些路径上成立。标签偏移一位、训练验证泄漏、错误 padding、重复分布式样本等问题都可能在运行时不报错,却让模型和指标失真。
十四、生产系统中的取舍:质量、吞吐、成本和权限共同约束
数据管道的资源成本可以粗略拆为:
提高吞吐不一定降低成本。例如:
- 增加 worker 可能提高 CPU 和内存成本;
- 更强压缩可能减少存储和网络成本,但增加解码 CPU;
- 预先 tokenization 可能减少训练 CPU,但增加缓存存储;
- 更大的 batch 可能提高 GPU 利用率,却增加显存和单步延迟;
- 重复抽样可能改善少数类训练,但增加有效样本处理成本。
还应把权限和数据治理纳入管道设计:
- 原始数据、脱敏数据、训练缓存和评测集使用不同访问权限;
- Dataset 记录数据版本和处理版本;
- 日志避免输出原始敏感文本、路径中的密钥或个人信息;
- 训练任务使用短期、最小权限凭证;
- 删除或撤回数据时,能够定位哪些缓存、索引和模型训练任务受影响。
对于机器学习、深度学习和生成式 AI,这些原则是一致的:Dataset 定义数据语义,Sampler 定义观察分布,增强定义训练输入变换,prefetch 定义并发供给方式,而吞吐诊断负责判断哪一阶段限制了系统。只有把这些边界分开,才能在不改变训练目标的前提下优化性能,并在出现异常时区分“模型问题”“数据问题”和“系统问题”。
系列导航与关联阅读
- 系列入口:AI 工程完整学习路线:从机器学习与 Transformer 到 RAG、Agent 和生产治理
- 上一篇:神经网络初始化与归一化:Xavier、Kaiming、BatchNorm 和 LayerNorm
- 下一篇:分布式训练:数据并行、模型并行、梯度同步与故障恢复
官方资料
本文依据研究论文、标准组织与主流框架官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。

评论
0 条讨论