(10) 多线程预处理、NUMA 感知与异步双缓冲:让 GPU 不再等待 CPU

——“一个人用AI如何写出比PyTorch更快的自研深度学习框架”系列文章之十

前面我们讲了DataLoader抽象、RAW/DTS 双路径、FULLY/PARTIAL 加载模式和DTS格式等。但数据加载的故事还没讲完——格式和 Loader 只解决了”怎么把数据从磁盘读出来”,真正决定 GPU 是否”挨饿”的,是CPU 预处理能不能跑得足够快,以及处理好的数据能不能在正确的时机送到 GPU 上

这一篇我们讲这条链路的最后两环:多线程预处理异步双缓冲传输。它们和 DTS 格式一起,构成了 Tech-Renaissance 数据引擎的核心。

一、数据饥饿:训练瓶颈并不总在 GPU 里

在前面的文章里我多次提到一个概念:数据饥饿(data starvation)。一块 A100 完成一次 256 样本的 ResNet-50 前向+反向可能只需要几十毫秒,但 CPU 侧要解码 256 张 JPEG、做 RandomResizedCrop、ColorJitter、归一化、组 batch,再发起 H2D 传输。如果这些数据准备工作单线程串行执行,耗时很容易超过 GPU 计算时间,导致昂贵的加速器空闲等待。

现代框架都意识到了这一点。PyTorch 的 DataLoader 提供了 num_workers 参数,用多个子进程并行读取和预处理;pin_memory=True 会把 CPU 内存锁页,方便异步 H2D 传输;persistent_workers=True 则让子进程在多个 epoch 之间复用,避免反复创建进程的开销。TensorFlow 的 tf.data 通过 map(..., num_parallel_calls)prefetch() 把预处理与 GPU 消费流水线化。NVIDIA 的 DALI 走得更远,它把一部分预处理直接搬到 GPU 上做,进一步减轻 CPU 压力。这些方案的本质都是同一个思想:把数据管线也变成一条流水线,让 CPU 预处理、CPU→GPU 传输、GPU 计算三个阶段尽可能重叠执行

但思想相同,实现细节会决定天花板。PyTorch 的 num_workers每个 GPU 进程独立配置的,8 卡训练时如果每张卡配 16 个 worker,全局就有 128 个进程;这些进程是否落在合适的 CPU 核心上、是否访问了本地 NUMA 内存、worker 之间是否存在竞争,都需要用户额外操心。Tech-Renaissance 则把这些问题收进了框架内部:用一个全局统一的持久 worker 池、静态样本领取、NUMA 感知的锁页内存分配、以及 per-GPU 的 TransferStation 双缓冲,把”数据不等待”这件事做到系统化。

二、持久 worker 池与静态领取:零竞争的并发

src/data/preprocessor.cpp 中,Preprocessor::start_worker_pool() 启动了一个持久线程池

void Preprocessor::start_worker_pool(DataLoader& loader) {
    worker_pool_.clear();
    worker_pool_.reserve(config_.num_workers);

    for (int i = 0; i < config_.num_workers; ++i) {
        worker_pool_.emplace_back([this, i, &loader]() {
            worker_func_persistent(i, loader);
        });
    }
}

这些 worker 一旦启动,就会一直处理多个 buffer,直到整个 epoch 结束。早期实现是”每个 buffer 创建一批线程、处理完就销毁”,很快发现线程创建和 join 的开销在高并发下不可忽略。改成持久线程池后,线程只创建一次,主循环通过 current_buffer_seq_ 原子变量通知 worker 进入下一轮,开销被压到最低。

关键的设计在于静态领取。PyTorch DataLoader 通常通过主进程 sampler 生成 index,再经由多进程队列分发给 worker;在多进程、队列通信和 worker 调度过程中仍存在调度与同步开销。Tech-Renaissance 的做法是:每个 worker 的样本位置在数学上预先确定,不需要抢。

src/data/preprocess_worker.cpp 中,calculate_write_position() 实现了这个逻辑:

std::pair<int, int> PreprocessWorker::calculate_write_position() const {
    // 变量说明:
    //   M = config_.num_workers_per_engine   // 每个 GPU 对应的 PW 数量
    //   j = config_.pid_in_engine             // 该 PW 在 Engine 内的编号
    //   B = config_.local_batch_size          // 单卡 batch size
    //   n = local_sample_id_                  // 当前 PW 在本 phase 已处理样本数

    int global_seq = n * M + j;
    int batch_id   = global_seq / B;
    int position   = global_seq % B;
    return {batch_id, position};
}

假设每个 GPU 对应 16 个 worker,那么 worker 0 处理第 0、16、32…号样本,worker 1 处理第 1、17、33…号样本。每个 worker 只需要维护自己的 local_sample_id_,不需要任何原子操作来决定”下一个该取谁”。源代码里的注释甚至给出了一个形式化证明:因为不同 worker 的 j 不同,任何两个 worker 永远不会落到同一个 (batch_id, position) 槽位。这就是零竞争并发

零竞争带来的好处是线性的扩展性。在实测中,Tech-Renaissance 可以把预处理线程开到 128、200 甚至更多,而不会因为锁竞争导致收益骤降。每个 worker 还拥有自己独立的 Workshop 内存、TurboJPEG 3.x 句柄、PO 链克隆,以及独立的随机数生成器,worker 之间没有共享状态,也没有 false sharing 热点。

worker 与 GPU 的对应关系也非常直接:

// 变量说明:
//   worker_id  = 全局 worker 编号
//   world_size = GPU 卡数

int engine_id    = worker_id % world_size;   // 该 PW 服务哪张 GPU
int pid_in_engine = worker_id / world_size;  // 在该 GPU 的 worker 中的编号

所有服务于同一张 GPU 的 worker 只向同一个 TransferStation 写入,不同 GPU 的数据流天然隔离。

三、NUMA 感知与 CPU 绑核:让数据靠近计算

多线程只是第一步。如果 128 个 worker 随机散落在 CPU 上,而它们频繁访问的内存又挂在另一个 NUMA 节点,那么跨节点访存带宽很快就会成为新的瓶颈。高端训练服务器通常是多路 CPU、多个 NUMA 节点,每张 GPU 通过 PCIe 挂在某个 NUMA 节点下。理想的状况是:处理某张 GPU 数据的 worker 绑定在该 GPU 所在 NUMA 节点的 CPU 上,并且在该 NUMA 节点本地分配内存

Tech-Renaissance 在 TR_SCENE_GPU_CLOUD 场景下做了两部分设计:Staging Buffer 的 NUMA 感知分配worker 线程的确定性 CPU 绑核

3.1 Staging Buffer 的 NUMA 感知分配

src/core/staging_buffer_pool.cpp 为每个 GPU 单独分配一块锁页内存。分配前,它会先通过 CUDA API 查询 GPU 的 PCIe DBDF,再读取 /sys/bus/pci/devices/.../numa_node 确定 GPU 所属的 NUMA 节点;然后为每个 GPU 启动一个独立分配线程,在线程内通过 numa_run_on_nodenuma_set_preferred 把线程迁移到该 GPU 对应的 NUMA 节点:

// 变量说明:
//   gpu_id   = 目标 GPU 编号
//   numa_node = GPU 所属 NUMA 节点(通过 /sys/bus/pci/devices/<DBDF>/numa_node 读取)
//   bytes    = 该 GPU 需要的 Staging Buffer 字节数

void StagingBufferPool::allocate_worker(int gpu_id, int numa_node,
                                        size_t bytes, void** out_ptr) {
    if (numa_node >= 0 && numa_available() >= 0) {
        numa_run_on_node(numa_node);    // 线程绑定到目标 NUMA 节点
        numa_set_preferred(numa_node);  // 优先从该节点分配内存
    }
    cudaSetDevice(gpu_id);
    cudaMallocHost(out_ptr, bytes);     // 分配锁页内存
    std::memset(*out_ptr, 0, bytes);    // 防御性清零与页表预热
}

这里的关键是:锁页内存在 cudaMallocHost 分配(锁页)时即被提交并锁定,其 NUMA 归属由分配线程当时的 NUMA 策略决定。因此必须在绑定到目标 NUMA 节点后再调用 cudaMallocHost;随后的 memset 主要起清零与页表预热作用,让 worker 后续访问时不会触发额外的缺页中断。如果 NUMA 绑定在分配之后才做,已经无法改变已锁定的物理页位置。

这样,从 CPU 到 GPU 的 H2D 传输路径最短:worker 在本地 NUMA 核心上预处理,把结果写入本地 NUMA 的锁页内存,再通过同 NUMA 节点下的 PCIe Root Complex 异步拷贝到 GPU,尽量减少跨节点访问。

3.2 Worker 线程的确定性 CPU 绑核

TR_SCENE_GPU_CLOUD 场景下,每个 worker 线程启动时会调用 Preprocessor::bind_worker_to_cpu(),通过 pthread_setaffinity_np 绑定到固定逻辑 CPU。当前实现采用确定性 Round-Robin

// 变量说明:
//   total_workers = 预处理 worker 总数
//   ncpus         = 系统在线 CPU 核心数(sysconf(_SC_NPROCESSORS_ONLN))

for (int w = 0; w < total_workers; ++w) {
    int target_cpu = w % ncpus;
    binding_map[w] = target_cpu;
}

这种绑法的优点是每个 worker 都有固定的”座位”,不会出现多个 worker 同时挤在同一个核心上争夺执行单元的情况;同时,由于 Workshop 内存、TurboJPEG 句柄等局部状态都在已经绑定的线程中创建,操作系统默认的 first-touch 策略会让这些热数据优先落在该 CPU 所在的 NUMA 节点上,提升缓存与内存局部性。

需要说明的是,这一层绑核目前采用的是全局 Round-Robin,而不是按 GPU-NUMA 节点做严格映射。它已经可以消除同核竞争,并在多数情况下获得不错的局部性;如果未来需要更极致的 NUMA 亲和性,框架中已经保留了 HardwareTopology / CPUAdvisor 等基础设施,可以在此基础上进一步细化。

四、TransferStation:CPU 与 GPU 之间的双缓冲驿站

预处理完成之后,数据还要到达 GPU。Tech-Renaissance 没有让用户手动 tensor.to(device),而是把 H2D 传输封装进了框架内部,核心是 TransferStation

每个 GPU 对应一个 TransferStation,它从 StagingBufferPool 拿到两块已经分配好的锁页内存。每块内存按 [labels][padding][image_data] 布局,标签区与图像数据区都按 256 字节对齐,并额外预留 16 字节作为 SIMD 越界保护。这个布局与 GPU 端 DTensor 的内存布局完全兼容,H2D 异步传输后不需要再做数据重排。以 AMP(FP16)路径为例,单区大小计算如下:

// 变量说明:
//   n            = local_batch_size
//   max_sample_bytes = 单个样本最大字节数
//   label_raw    = n * sizeof(int32_t)
//   data_raw     = n * max_sample_bytes

// labels 区对齐后大小(FP32/INT32 通用公式)
size_t label_aligned = 2 * align_up_256(label_raw / 2 + 16);

// image 数据区对齐后大小(AMP/FP16 路径)
size_t data_aligned  = align_up_256(data_raw + 16);

// 单区总大小 = label_aligned + data_aligned
// A 区起始:base
// B 区起始:base + per_zone,其中 per_zone = label_aligned + data_aligned

TransferStation 维护两个 buffer:CPU 的 PreprocessWorker 向当前 buffer 写入样本,当当前 buffer 被一个 batch 填满或所有 worker 报告样本耗尽时,调用 execute_transfer_locked() 把它标记为 readable,并切换到另一块 buffer 继续写。GPU 侧的深度学习引擎则在 CUDA Graph 中调用 wait_buffer_readable() 等待可读,读取完毕后调用 set_buffer_writeable() 把 buffer 交还。

关键代码在 src/data/transfer_station.cpp 中:

void TransferStation::execute_transfer_locked(int samples_count,
                                              bool fill_before_transfer) {
    int buf_id = current_buffer_.load(std::memory_order_acquire);

    // 若需要 filling,复制首样本填满最后一个不完整 batch
    int samples_count_after_filling = samples_count;
    if (need_filling_ && fill_before_transfer) {
        samples_count_after_filling++;
        std::memcpy(buffer_data_[buf_id] + snapshot_sample_bytes * samples_count,
                    filling_sample_data_, snapshot_sample_bytes);
        buffer_labels_[buf_id][samples_count] = filling_label_;
    }

    set_buffer_writeable(buf_id, false);
    buffer_actual_transfer_samples_[buf_id] = samples_count_after_filling;
    set_buffer_readable(buf_id, true);

    // 等待另一块 buffer 可写,然后切换
    int next_buf = 1 - buf_id;
    while (!buffer_is_writeable(next_buf)) {
        wait_buffer_writeable(next_buf);
    }
    set_buffer_readable(next_buf, false);
    current_buffer_.store(next_buf, std::memory_order_release);
    current_batch_id_.fetch_add(1, std::memory_order_release);
    cv_batch_ready_.notify_all();
}

这段代码完成了三件事:把刚写满的 buffer 交给 GPU;等待另一块 buffer 被 GPU 消费完;切换当前 buffer 并唤醒所有等待的 worker。由于两块 buffer 交替工作,CPU 预处理和 GPU 计算可以流水线重叠:GPU 在计算 buffer A 时,CPU 正在往 buffer B 里填下一个 batch。

同步状态由四个原子布尔标志管理:buffer_0_is_readable_buffer_1_is_readable_buffer_0_is_writeable_buffer_1_is_writeable_,配合 std::condition_variable 实现等待,避免忙轮询浪费 CPU。

4.1 不完整 batch 的 filling 机制

在分布式训练中,训练集样本数不一定能被 local_batch_size × world_size 整除。为了解决这个问题,TransferStation 实现了 filling 机制:在第一个 batch 传输时,保存第一个样本的数据和标签;当最后一个 batch 不满时,用这个保存的样本填充到满 batch。判断逻辑在 configure() 中预先计算:

// 变量说明:
//   num_train_samples = 训练集总样本数
//   world_size        = GPU 卡数
//   engine_id         = 当前 TransferStation 对应的 GPU 编号

int remainder = num_train_samples % world_size;
bool need_filling = (remainder != 0) && (engine_id >= remainder);

这样,每个 epoch 的所有 batch 大小一致,深度学习引擎不需要处理变长 batch,CUDA Graph 也不需要为不完整 batch 单独捕获图。

4.2 可复现性开关:性能与确定性的权衡

为了兼顾性能与可复现性,TransferStation 还实现了两种模式,通过 GlobalRegistry::instance().ensure_reproducibility(true/false) 切换。

可复现模式下,worker 必须按照 calculate_write_position() 算出的固定槽位写入,传输触发时机严格由 batch 满决定;快 worker 必须等待慢 worker,所有 worker 就位后才触发传输。这保证了每次运行写入顺序完全一致。

非可复现模式下,使用 request_count_ / written_count_ 原子计数器快速分配 slot,只有在真正触发传输时才加锁:

// 变量说明:
//   local_batch_size_ = 单卡 batch size

size_t my_request = request_count_.fetch_add(1, std::memory_order_relaxed);
size_t slot     = my_request % local_batch_size_;
int    my_batch = static_cast<int>(my_request / local_batch_size_);

生产训练通常用非可复现模式换取极限吞吐,调试和消融实验则用可复现模式保证结果稳定。这与第 22 篇将要介绍的 Philox 确定性随机数生成器一起,构成了 Tech-Renaissance 端到端确定性训练的完整拼图。

五、为什么比 PyTorch DataLoader 更稳、更可预测

写到这里,可能又会有人问:这不就是 C++ 版的 PyTorch DataLoader 吗?功能上确实有相似之处,但实现定位不同。

第一,worker 总数的语义不同。PyTorch 的 num_workers每个 GPU 进程的 worker 数;Tech-Renaissance 的 preprocess_workers跨所有 GPU 的总线程数。例如 8 卡服务器、PyTorch 每张卡 16 worker,总共 128 个进程;在 Tech-Renaissance 里直接写 .preprocess_workers(128),框架内部通过 engine_id = worker_id % world_size 把样本划分到各个 GPU。这个设计让多卡场景下的 CPU 资源分配更统一。

第二,样本分配方式不同。PyTorch 的 worker 通常由主进程 sampler 生成 index,再经由多进程队列分发给 worker,这个分发、调度与同步过程仍存在运行时开销;Tech-Renaissance 通过静态领取把每个 worker 的样本位置预先确定,不依赖运行时队列分发,消除了这个竞争热点,也让 200+ worker 的扩展性成为可能。

第三,Staging Buffer 的 NUMA 感知是框架内置的。PyTorch 不会自动根据 GPU 拓扑把锁页内存分配在本地 NUMA 节点;Tech-Renaissance 在 TR_SCENE_GPU_CLOUD 下把拓扑发现、本地节点绑核分配、页表预热串成了一条链路。

第四,H2D 传输内聚。PyTorch DataLoader 把预处理后的 tensor 交给 Python 主进程,用户再自己 to(device) 发起 H2D;Tech-Renaissance 的 TransferStation 直接和深度学习引擎共享 staging memory,H2D 由 CUDA Graph 子图统一调度,不需要 Python 侧参与。

第五,可复现性开关。PyTorch 要达到端到端可复现,需要仔细设置 worker seed、persistent_workersgenerator 等多个细节;Tech-Renaissance 把 reproducibility 做成了顶层开关,静态领取和 Philox 计数器随机数共同保证高并发下仍可复现。

六、小结:把数据管线做到极致是一门系统工程

数据加载管线看起来不如算子融合、CUDA Graph 那么”硬核”,但它决定了 GPU 能不能满负荷运转。Tech-Renaissance 在这一层的设计可以概括为四句话:

  1. 持久线程池:worker 只创建一次,跨 buffer 复用,避免线程创建开销;
  2. 静态领取:每个 worker 的样本位置预先确定,实现零竞争并发;
  3. NUMA 感知的 Staging Buffer:通过拓扑发现、本地节点绑核分配、页表预热,让 H2D 传输路径最短;worker 线程做确定性 CPU 绑核,避免同核竞争;
  4. 异步双缓冲TransferStation 让 CPU 预处理、H2D 传输、GPU 计算三段流水线重叠,同时用可复现性开关在性能与确定性之间灵活切换。

这套机制在实测中可以把 ImageNet 预处理线程推到很高,同时保持结果可复现、不出现死锁、不会因为 NUMA 远端访存而掉速。

当然,预处理再快,如果最后一步归一化、类型转换、随机翻转还要反复读写内存,也还是会浪费带宽。下一篇《FusedNormalization:把数据增强的最后一步融进一次内存遍历》,我们来讲 Tech-Renaissance 如何把 ToTensor、Normalize、RandomHorizontalFlip、RandomErasing 这些原本独立的收尾操作,融合成一次内存遍历,让数据侧的性能再上一个台阶。

发表回复

您的邮箱地址不会被公开。 必填项已用 * 标注

ICP备案号:京ICP备2025133467号-1