AI 工程基础体系 · 第 63/100 篇。内容覆盖机器学习、深度学习与生成式 AI;模型、数据、评测、权限和成本会作为同一生产系统处理。

深度学习数据管道:Dataset、采样、增强、预取与吞吐诊断

深度学习训练中的“数据管道”不是一个简单的文件读取函数,而是一条把样本从存储系统送到设备、再送入模型的有状态流水线。它通常包含:

存储
  ↓
Dataset:定义“一个样本是什么”
  ↓
Sampler:决定“取哪些样本、以什么顺序取”
  ↓
BatchSampler / DataLoader:组成批次并调度读取
  ↓
collate_fn:把样本拼成模型输入
  ↓
数据增强:改变样本但尽量保持监督关系
  ↓
worker:并行准备 CPU 数据
  ↓
pin memory / prefetch:把后续批次提前准备到更接近设备的位置
  ↓
GPU/TPU:执行前向、反向和参数更新

数据管道的目标不是单纯“读取更快”。它必须同时满足以下条件:

  1. 统计正确:样本分布、标签和数据划分符合训练目标。
  2. 语义正确:增强不会破坏标签含义,训练和评测边界清晰。
  3. 设备高效:模型计算时尽量不等待数据。
  4. 资源可控:CPU、内存、显存、存储带宽和进程数不会失控。
  5. 可复现、可诊断:出现错误时能够定位到样本、worker、存储还是设备。
  6. 可治理:数据权限、敏感信息、版本、评测集和成本能够被追踪。

一、先区分 Dataset、Sampler、DataLoader 和 collate_fn

这些组件经常被统称为“数据加载器”,但它们承担的职责不同。职责混淆后,最常见的结果是数据重复、采样偏差、增强位置错误或多进程状态失效。

1. Dataset 定义样本访问语义

在 PyTorch 中,最常见的 Dataset 是 map-style dataset。它提供:

__len__()       # 数据集大小
__getitem__(i) # 根据索引 i 返回一个样本

例如,一个分类样本可以表示为:

xi=imagei,yi=labelix_i = \text{image}_i,\qquad y_i = \text{label}_i

Dataset 负责把索引 ii 映射到样本:

D(i)=(xi,yi)D(i) = (x_i, y_i)

它通常应当负责:

  • 根据索引定位文件或记录;
  • 解码、解析和基础类型转换;
  • 返回样本及其元数据;
  • 在必要时执行样本级变换。

它不应隐式决定“这一轮训练取哪些样本”。后者属于 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),
        }

这里的 indexpath 不一定要送入模型,但在诊断标签错位、坏文件和数据泄漏时很有价值。

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 产生索引序列:

i0,i1,i2,i_0, i_1, i_2, \ldots

例如:

  • SequentialSampler:按 0,1,2,0,1,2,\ldots 顺序;
  • RandomSampler:随机打乱或有放回抽样;
  • WeightedRandomSampler:按给定权重抽样;
  • DistributedSampler:把索引划分给不同训练进程。

Sampler 不负责读取图像,也不负责把多个样本拼成 batch。

4. BatchSampler 决定索引分组

如果 batch size 为 BB,sampler 产生的索引会被分组为:

(i0,,iB1),(iB,,i2B1),(i_0,\ldots,i_{B-1}), (i_B,\ldots,i_{2B-1}),\ldots

通常 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. 均匀采样对应什么目标

假设训练集有 NN 个样本,单样本损失为:

i(θ)=(fθ(xi),yi)\ell_i(\theta)=\ell(f_\theta(x_i),y_i)

经验风险是:

L(θ)=1Ni=1Ni(θ)L(\theta)=\frac{1}{N}\sum_{i=1}^{N}\ell_i(\theta)

如果索引 II 按均匀分布抽取:

P(I=i)=1NP(I=i)=\frac{1}{N}

那么单次随机梯度的期望为:

E[θI(θ)]=i=1N1Nθi(θ)=θL(θ)\mathbb{E}[\nabla_\theta \ell_I(\theta)] = \sum_{i=1}^{N}\frac{1}{N}\nabla_\theta\ell_i(\theta) = \nabla_\theta L(\theta)

这说明,均匀采样的随机梯度是经验风险梯度的无偏估计。一个 batch 只是在降低估计方差,前提是 batch 中的样本确实来自目标分布。

2. 类别不均衡时,改变采样就改变了目标

假设有两类:

  • 类别 A:90 个样本;
  • 类别 B:10 个样本。

普通均匀采样时:

P(A)=0.9,P(B)=0.1P(A)=0.9,\qquad P(B)=0.1

如果使用等类别采样,变成:

P(A)=0.5,P(B)=0.5P(A)=0.5,\qquad P(B)=0.5

这会使少数类在训练中被更频繁看到,但它不再直接优化原始经验风险,而是在优化一个类别重新加权后的目标。

使用 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. 采样加权与损失加权不是同一件事

若从分布 q(i)q(i) 抽样,但目标仍是均匀经验风险,则可以使用重要性权重:

g^=1Bj=1Bp(ij)q(ij)θij(θ)\hat{g} = \frac{1}{B}\sum_{j=1}^{B} \frac{p(i_j)}{q(i_j)} \nabla_\theta\ell_{i_j}(\theta)

其中:

  • p(i)=1/Np(i)=1/N 是目标分布;
  • q(i)q(i) 是实际 sampler 分布;
  • ijqi_j\sim q

由于:

Eiq[p(i)q(i)i]=ip(i)i\mathbb{E}_{i\sim q} \left[ \frac{p(i)}{q(i)}\nabla\ell_i \right] = \sum_i p(i)\nabla\ell_i

所以估计仍然无偏。实际训练中常常不做完整校正,因为校正会增加方差,且很多任务有意改变类别权重。关键是明确你要优化的是:

  • 原始样本分布;
  • 类别平衡分布;
  • 困难样本分布;
  • 还是某个业务风险分布。

4. shuffle=True 的边界

对 map-style dataset,shuffle=True 通常意味着由 DataLoader 创建随机 sampler。不能同时传入 shuffle=Truesampler,因为两者都在决定索引顺序,PyTorch 会拒绝这种冲突配置。

训练、验证和测试通常不同:

train_loader = DataLoader(
    train_dataset,
    batch_size=64,
    shuffle=True,
)

valid_loader = DataLoader(
    valid_dataset,
    batch_size=128,
    shuffle=False,
)

验证集不需要随机化来估计整体损失;保持固定顺序还有利于复现和逐样本对比。随机增强也不应默认应用在验证集,否则同一模型每次评测可能看到不同输入。

5. 分布式训练中的 sampler 状态

分布式数据并行通常有 RR 个进程,每个进程处理全局索引的一部分。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 约为

    Bglobal=Bper-rank×R×gradient accumulation stepsB_{\text{global}}=B_{\text{per-rank}}\times R\times \text{gradient accumulation steps}

  • drop_last=True 可能丢弃尾部样本,但能避免不同进程最后一个 batch 形状不一致;
  • 评测时必须确认是否需要补齐样本,以及补齐样本是否会被错误计入指标。

三、数据增强改变的是输入分布,不应破坏监督关系

数据增强是对样本施加变换:

x=T(x;ξ)x' = T(x;\xi)

其中 ξ\xi 是随机变量。若任务标签在变换下保持不变:

y(T(x;ξ))=y(x)y(T(x;\xi))=y(x)

那么训练目标可以写成:

E(x,y)Pdata,ξ[(fθ(T(x;ξ)),y)]\mathbb{E}_{(x,y)\sim P_{\text{data}},\xi} \left[ \ell(f_\theta(T(x;\xi)),y) \right]

这不是“凭空增加了新数据”,而是让模型在一个由原始样本诱导出的增强分布上训练。

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 计算。

若把增强结果预先缓存,则读取更快,但:

  • 随机性减少;
  • 缓存占用增加;
  • 变换版本变化后需要重建;
  • 可能把训练增强错误地带入验证或测试。

数据管道的可复现性通常至少需要控制三类随机源:

  1. Python 的 random
  2. NumPy 的随机数;
  3. 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 数量与:

Qnum_workers×prefetch_factorQ \approx \text{num\_workers}\times\text{prefetch\_factor}

同量级。

因此增大 prefetch_factor 可能隐藏短时读取抖动,但也会增加:

  • 预取内存;
  • 未使用 batch 被丢弃时的浪费;
  • 进程间队列压力;
  • 随机增强提前消耗的随机状态。

它不能提高存储设备的物理带宽。若瓶颈是磁盘读带宽,预取只能让等待更早发生。

4. pin_memorynon_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)

这里有三个边界:

  1. pin_memory=True 只对 DataLoader 返回的 CPU 对象尝试固定内存;
  2. non_blocking=True 是允许异步传输的条件之一,不代表任何设备或任何内存都必然异步;
  3. 如果紧接着在 CPU 上访问传输结果,异步优势会被同步操作抵消。

Pinned memory 也不是越多越好。过度固定主机内存可能影响系统整体内存压力,因此应观察 RSS、swap 和系统响应,而不是只看 GPU 利用率。


五、预取的本质:用队列重叠数据准备与模型计算

设:

  • TdT_d:准备一个 batch 的时间,包括读取、解码、增强和 collate;
  • ThT_h:CPU 到 GPU 的传输时间;
  • TcT_c:模型前向、反向和优化器更新时间。

没有重叠时,一个 batch 的周期近似为:

Tserial=Td+Th+TcT_{\text{serial}}=T_d+T_h+T_c

理想重叠时,数据准备和模型计算可以并行,稳定吞吐受最慢阶段限制:

Toverlapmax(Td+Th,Tc)T_{\text{overlap}}\approx \max(T_d+T_h,T_c)

预取的作用是让当前 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()

这个示例中每一步为何成立

  1. ToyDataset.__getitem__ 根据索引返回单样本,符合 map-style Dataset 协议。
  2. AddNoise 在训练取样时执行随机增强;它只改变输入,不改变二分类标签。
  3. WeightedRandomSampler 根据类别频率反比设置样本权重,少数类更容易被抽到。
  4. replacement=True 允许同一个原始索引在一个 epoch 中出现多次,因此 num_samples 明确定义了 epoch 的抽样次数。
  5. num_workers=2 使数据准备在子进程执行;在 Windows 或某些启动方式下,if __name__ == "__main__" 是必要的。
  6. pin_memory 只在 CUDA 设备上启用,避免在 CPU 训练时引入无意义的固定内存。
  7. .to(device, non_blocking=True) 允许在满足条件时重叠主机到设备的传输。
  8. samples_per_sec 是端到端吞吐,不等同于纯 Dataset 读取速度;它包含模型计算和优化器更新。

示例输出的具体数值取决于 CPU、GPU、PyTorch 版本和进程启动方式,不能把某台机器上的数值当作性能基准。应重点观察 loss 是否下降、每个 epoch 样本数是否符合预期,以及改动参数后吞吐如何变化。


七、用形式化指标理解“吞吐”

1. 样本吞吐和 batch 吞吐

若一个 epoch 处理 NeN_e 个样本,耗时 TeT_e,样本吞吐为:

Ssample=NeTeS_{\text{sample}}=\frac{N_e}{T_e}

若 batch size 为 BB,batch 吞吐为:

Sbatch=Ne/BTeS_{\text{batch}}=\frac{N_e/B}{T_e}

两者不能混淆。一个 batch 变大后,batch/s 可能下降,但 samples/s 可能上升。

分布式训练还要说明是:

  • 单进程 samples/s;
  • 单机 samples/s;
  • 全局 samples/s。

全局吞吐近似为所有 rank 在同步条件下完成的全局样本数除以墙钟时间,而不是简单相加某些不同步的局部计时。

2. 有效吞吐不只是读取速度

训练系统通常受多个阶段限制:

Seffectivemin(Sstorage,Sdecode,Saugment,Stransfer,Scompute)S_{\text{effective}} \le \min( S_{\text{storage}}, S_{\text{decode}}, S_{\text{augment}}, S_{\text{transfer}}, S_{\text{compute}} )

但由于阶段可以重叠,实际周期更接近最慢关键阶段,而不是所有时间简单相加。并且同步点、队列空转、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. 参数实验必须一次只改一个因素

可以按如下顺序做小规模对照:

  1. num_workers=0,确认 Dataset 单进程逻辑正确;
  2. 测试 num_workers=1,2,4,8
  3. 对比 pin_memory=False/True
  4. 对比 persistent_workers=False/True
  5. 在确认队列短缺后,调整 prefetch_factor
  6. 再测试 batch size、增强开关和存储格式。

每次记录:

  • 总耗时;
  • samples/s;
  • CPU 利用率;
  • 内存 RSS;
  • GPU 利用率;
  • GPU 显存;
  • 存储读取带宽;
  • worker 异常和系统日志。

如果同时修改五个参数,即使吞吐变化,也无法知道哪项产生了变化。

4. 常见症状与因果路径

症状:num_workers=0 正常,多 worker 卡住

可能原因包括:

  • Dataset 持有不能安全跨进程复制的对象;
  • 文件句柄在 fork 后状态异常;
  • 第三方库内部线程与多进程组合导致死锁;
  • __getitem__ 中存在不可序列化对象;
  • Windows 下缺少主入口保护;
  • worker 内存异常退出,主进程等待队列结果。

诊断方法:

  1. 先设 num_workers=0 获得完整 traceback;
  2. 缩小到单个样本;
  3. __getitem__ 中的解码、增强逐步注释;
  4. 再从 num_workers=1 开始增加;
  5. 检查系统内存、文件描述符和 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,来自同一原图的内容可能同时出现在训练和测试中。

形式上,测试集应尽量近似独立于训练过程:

DtrainDtest=D_{\text{train}}\cap D_{\text{test}}=\varnothing

“样本对象不完全相同”还不够,近重复、同一实体和时间未来信息也可能造成实际泄漏。

2. 生成式 AI 中的 Dataset 不只是文本列表

语言模型训练的样本可能是:

  • 文档;
  • 对话;
  • 指令与响应;
  • 偏好对;
  • 多模态消息;
  • token block。

一个 token block 通常来自:

(x1,x2,,xT)(x_1,x_2,\ldots,x_T)

训练目标可能是 next-token prediction:

L(θ)=1T1t=1T1logpθ(xt+1xt)L(\theta) = -\frac{1}{T-1} \sum_{t=1}^{T-1} \log p_\theta(x_{t+1}\mid x_{\le t})

此时 Dataset 不仅要返回 token,还要正确处理:

  • input_ids
  • labels 的右移关系;
  • padding;
  • attention_mask
  • 对不应计算损失的 prompt 部分设置 ignore label;
  • 长文档截断和跨 block 边界。

把不同对话直接拼接而不加入边界标记,可能使模型学习到错误的跨样本上下文。把用户数据未经权限审核放入训练 Dataset,则属于数据治理和隐私问题,不是一个普通的读取错误。

3. 评测集必须是受保护的数据资产

生产系统中,评测集不应与训练增强缓存、自动清洗脚本和临时实验目录混用。应记录:

  • 数据集版本;
  • 样本清单哈希;
  • 标签或标注版本;
  • 访问权限;
  • 生成时间;
  • 是否包含个人信息或商业机密;
  • 哪些模型和实验读取过它。

数据权限不仅限制“谁能下载文件”,也应限制 Dataset worker 通过凭证访问哪些对象。多 worker 会复制或继承部分进程状态,不能把长期有效的高权限凭证随意放进 Dataset 对象。


十、存储、预处理和缓存的取舍

1. 小文件数量会影响读取效率

大量小图片或小 JSON 文件会产生:

  • 目录查找;
  • 文件打开关闭;
  • 元数据请求;
  • 随机 I/O;
  • 网络对象存储请求开销。

将样本打包成更适合顺序读取或批量读取的格式,可能减少这些固定开销。但打包也带来:

  • 局部更新困难;
  • 单个样本损坏影响范围扩大;
  • 随机访问索引需要维护;
  • 压缩与解压 CPU 成本;
  • 数据版本和权限管理复杂度。

没有一种格式对所有任务都最快。应使用真实样本和真实存储环境测量,而不是仅根据文件扩展名判断。

2. 预处理缓存的正确性条件

若缓存函数为:

ci=g(xi;ϕ)c_i = g(x_i;\phi)

那么缓存只有在以下条件满足时才可安全复用:

  • 原始样本版本不变;
  • 预处理代码版本 ϕ\phi 不变;
  • 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 悄悄丢弃样本,造成样本数和类别分布变化。

更可靠的策略是:

  1. 在离线校验阶段扫描文件、标签、尺寸和编码;
  2. 训练时异常包含数据 ID、路径和原始异常;
  3. 对可容忍的脏数据使用显式策略;
  4. 统计跳过数量、类别分布和每个 shard 的错误率;
  5. 错误超过阈值时停止任务,而不是继续产生不可解释的模型。

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、重复分布式样本等问题都可能在运行时不报错,却让模型和指标失真。


十四、生产系统中的取舍:质量、吞吐、成本和权限共同约束

数据管道的资源成本可以粗略拆为:

Ctotal=Cstorage+CCPU+CGPU+Cnetwork+CengineeringC_{\text{total}} = C_{\text{storage}} + C_{\text{CPU}} + C_{\text{GPU}} + C_{\text{network}} + C_{\text{engineering}}

提高吞吐不一定降低成本。例如:

  • 增加 worker 可能提高 CPU 和内存成本;
  • 更强压缩可能减少存储和网络成本,但增加解码 CPU;
  • 预先 tokenization 可能减少训练 CPU,但增加缓存存储;
  • 更大的 batch 可能提高 GPU 利用率,却增加显存和单步延迟;
  • 重复抽样可能改善少数类训练,但增加有效样本处理成本。

还应把权限和数据治理纳入管道设计:

  • 原始数据、脱敏数据、训练缓存和评测集使用不同访问权限;
  • Dataset 记录数据版本和处理版本;
  • 日志避免输出原始敏感文本、路径中的密钥或个人信息;
  • 训练任务使用短期、最小权限凭证;
  • 删除或撤回数据时,能够定位哪些缓存、索引和模型训练任务受影响。

对于机器学习、深度学习和生成式 AI,这些原则是一致的:Dataset 定义数据语义,Sampler 定义观察分布,增强定义训练输入变换,prefetch 定义并发供给方式,而吞吐诊断负责判断哪一阶段限制了系统。只有把这些边界分开,才能在不改变训练目标的前提下优化性能,并在出现异常时区分“模型问题”“数据问题”和“系统问题”。


系列导航与关联阅读

官方资料

本文依据研究论文、标准组织与主流框架官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。