
训练数据与指标的准备本文围绕“数据集和指标怎样准备”整理可复现的检查思路。所有阈值、配置和结果均应在隔离环境中记录输入、版本与资源条件后再解释下文示例不对应真实组织、用户、流量或成本数据。1. 用受控样例界定问题训练数据链路先固定数据切分、采样顺序和工作进程数。记录每轮读取等待与显存峰值才能判断瓶颈是在输入端还是计算端。2. 构建多卡高效数据加载管线解耦 I/O 与 Worker 采样要彻底解决分布式训练下的 I/O 挤压问题必须从文件存储格式与采样切片机制两个维度同步重构。第一个关键点是小文件打块Sharding。大量小文件会增加元数据查询与随机 I/O。可以在预处理阶段把合成样例打包为WebDatasettar、LMDB或TFRecord等连续文件再用相同硬件和读取模式比较吞吐与等待时间。第二个关键点是DistributedSampler的正确配置。很多初学者容易忘记将sampler.set_epoch(epoch)写入训练循环逻辑中。如果在每个 Epoch 开始时没有显式更新 Epoch 种子所有 Rank 节点在每个轮次都会拿到完全相同的样本切分顺序模型直接退化为在重复数据上拟合分布式训练完全失去了多样性。3. 分布式指标对齐跨节点 All-Reduce 汇总的坑在分布式训练中评估指标如 Accuracy、Loss、F1-score的计算也是极易踩坑的地方。简单的做法是在 Rank 0 上单卡计算指标但这要求将所有卡的 Prediction 集中传回 Rank 0不仅造成 PCIe 传输瓶颈而且如果各进程分配的 Batch 大小不均比如最后一个 Batch 样本数不足直接取平均会导致数值统计偏差。正规做法是利用torch.distributed.all_reduce算子。在每个进程本地累加 Sum 项和 Count 项然后通过广播做 Vector All-Reduce 汇总最后在各自 Rank 内完成除法。这样既不破坏分布式并行性又能保证指标口径在数学上的精确无误。以下是一段可直接复用到生产环境的分布式数据集加载与全局指标对齐模版import os import sys import torch import torch.distributed as dist from torch.utils.data import Dataset, DataLoader from torch.utils.data.distributed import DistributedSampler from typing import Tuple, Dict class SyntheticDistributedDataset(Dataset): 模拟高性能分布式数据集避免原生 __getitem__ 内部小文件 I/O def __init__(self, num_samples: int 10000, feature_dim: int 128): self.num_samples num_samples self.feature_dim feature_dim # 预先在内存中分配或使用 mmap 映射文件 self.data torch.randn(num_samples, feature_dim) self.labels torch.randint(0, 2, (num_samples,)) def __len__(self) - int: return self.num_samples def __getitem__(self, idx: int) - Tuple[torch.Tensor, torch.Tensor]: return self.data[idx], self.labels[idx] class DistributedMetricAggregator: 基于 torch.distributed.all_reduce 的准确跨卡指标汇总器 def __init__(self, device: torch.device): self.device device self.reset() def reset(self): self.correct_tensor torch.tensor(0, dtypetorch.long, deviceself.device) self.total_tensor torch.tensor(0, dtypetorch.long, deviceself.device) self.loss_sum_tensor torch.tensor(0.0, dtypetorch.float32, deviceself.device) def update(self, logits: torch.Tensor, targets: torch.Tensor, loss_val: float): preds torch.argmax(logits, dim-1) self.correct_tensor torch.sum(preds targets) self.total_tensor targets.numel() self.loss_sum_tensor loss_val * targets.numel() def compute_global_metrics(self) - Dict[str, float]: 在所有 Rank 之间同步汇总 Loss 与 Accuracy if not dist.is_initialized(): # 单卡降级逻辑 total max(1, self.total_tensor.item()) return { global_accuracy: self.correct_tensor.item() / total, global_loss: self.loss_sum_tensor.item() / total } # 构造跨节点 Vector 进行 All-Reduce 汇总 stats torch.stack([ self.correct_tensor.to(torch.float32), self.total_tensor.to(torch.float32), self.loss_sum_tensor ]) # 执行 原地 (In-place) 规约求和 dist.all_reduce(stats, opdist.ReduceOp.SUM) global_correct stats[0].item() global_total max(1.0, stats[1].item()) global_loss_sum stats[2].item() return { global_accuracy: global_correct / global_total, global_loss: global_loss_sum / global_total } def prepare_distributed_dataloader( dataset: Dataset, batch_size: int, num_workers: int 4 ) - Tuple[DataLoader, DistributedSampler]: 封装带有 DistributedSampler 与 Pin Memory 的 DataLoader sampler None if dist.is_initialized(): sampler DistributedSampler( dataset, shuffleTrue, drop_lastFalse ) loader DataLoader( dataset, batch_sizebatch_size, samplersampler, num_workersnum_workers, pin_memoryTrue, # 加速 CPU Tensor 拷贝至 GPU 显存 persistent_workersTrue if num_workers 0 else False ) return loader, sampler # 模块化测试验证 if __name__ __main__: # 模拟单机多卡测试环境 print(Testing Distributed DataLoader Preparation Utilities...) ds SyntheticDistributedDataset(num_samples1000) loader, sampler prepare_distributed_dataloader(ds, batch_size32, num_workers2) device torch.device(cuda if torch.cuda.is_available() else cpu) aggregator DistributedMetricAggregator(devicedevice) # 模拟 1 个 Batch 预测 mock_logits torch.randn(32, 2, devicedevice) mock_targets torch.randint(0, 2, (32,), devicedevice) aggregator.update(mock_logits, mock_targets, loss_val0.45) res aggregator.compute_global_metrics() print(fMetrics Evaluation Success: Acc{res[global_accuracy]:.4f}, Loss{res[global_loss]:.4f})4. 工程化基准测试与吞吐量评估数据集和指标对齐搞定后必须建立标准化的基准测试Benchmark流程用来评估训练系统的吞吐性能。衡量分布式训练效率的指标有两个Samples Per Second (SPS)和Weak Scaling Efficiency (弱缩放效率)。启动分布式训练前可先运行预热与测量两个阶段。预热步数、丢弃规则和统计区间应写入基准配置报告只比较同一硬件、批大小和数据工件下的结果。吞吐或资源结论应来自相同硬件与批大小下的对照并注明预热区间和统计方式。