You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

TensorFlow多GPU训练时Numpy加载大文件速度异常缓慢问题排查

多GPU单任务并行训练的资源冲突问题分析与解决

冲突确认

你的情况确实存在资源竞争冲突,原因在于:

  • 70GB数据集的Numpy预加载是CPU、内存、磁盘IO密集型操作;而TensorFlow训练时,GPU会占用大量PCIe带宽,同时TF本身也会消耗CPU资源做预处理、梯度计算辅助、队列调度等工作。
  • 多个进程同时运行时,磁盘IO带宽会被占满,若内存不足以容纳全量数据集,还会触发内存页交换(swap),进一步拖慢所有进程的运行速度。
  • 当前代码中load_factor_files采用单线程加载Numpy文件,加载周期长,会持续和训练进程抢占资源。

具体解决方案

1. 放弃全量预加载,改用tf.data动态读取

不要将70GB数据集全部加载到内存,转而让tf.data在训练过程中动态读取并处理数据,利用TF的多线程/多进程调度能力,减少资源竞争:

  • 将Numpy文件的读取逻辑嵌入到tf.data的map函数中,通过tf.numpy_function实现Python逻辑与TF图的兼容;
  • 更优方案是将Numpy文件转换为TFRecord格式,TF对该格式的加载效率更高,且支持分片读取、并行加载。

2. 并行化Numpy数据加载

如果必须保留预加载逻辑,将单线程加载改为多线程/多进程加载,缩短加载时间,减少与训练进程的资源抢占窗口:

from concurrent.futures import ThreadPoolExecutor
import os

def load_factor_files(in_files, store_dict):
    def load_and_update(file_path):
        result = load_factor_file(file_path)
        store_dict.update(result)
    
    # 用CPU核心数设置线程数,避免过度调度
    with ThreadPoolExecutor(max_workers=os.cpu_count()) as executor:
        executor.map(load_and_update, in_files)

3. CPU核心绑定与资源隔离

启动进程时,用taskset为每个训练进程绑定专属CPU核心,避免跨核心调度的开销和CPU资源竞争:

# 示例:给GPU 0的进程绑定0-3号CPU核心
taskset -c 0-3 python your_train_script.py --gpu=0 &
# GPU 1的进程绑定4-7号CPU核心
taskset -c 4-7 python your_train_script.py --gpu=1 &

4. 内存与磁盘优化

  • 确保服务器内存容量至少为数据集大小的1.5倍,避免内存不足触发swap交换(swap会导致所有进程性能骤降);
  • 使用NVMe SSD存储数据集,提升磁盘IO带宽,缓解加载瓶颈;
  • 保存Numpy文件时使用np.savez_compressed压缩,减少磁盘占用和读取时间。

5. 限制TensorFlow的CPU资源占用

调整TF的线程配置,避免其占用过多CPU资源影响数据加载:

# 在创建Estimator时添加RunConfig配置
config = tf.estimator.RunConfig(
    inter_op_parallelism_threads=4,  # 跨操作并行线程数
    intra_op_parallelism_threads=8,  # 操作内并行线程数
    log_step_count_steps=100
)
t0_regression = create_estimator(configuration, FLAGS.threads, run_config=config)

同时开启GPU内存动态增长,避免TF预占过多GPU内存:

# 在main函数开头添加
gpus = tf.config.experimental.list_physical_devices('GPU')
for gpu in gpus:
    tf.config.experimental.set_memory_growth(gpu, True)

6. 同步进程加载时机

让所有进程先完成数据加载,再统一启动训练,避免部分进程提前训练抢占资源:

  • 在每个进程完成数据加载后,创建一个标记文件;
  • 所有进程等待标记文件全部生成后,再开始训练:
# 训练脚本中,加载完成后创建标记
touch /tmp/load_done_${GPU_ID}
# 等待所有10个进程完成加载
while [ $(ls /tmp/load_done_* | wc -l) -lt 10 ]; do sleep 1; done
# 启动训练逻辑

代码修改示例(动态读取版本)

将预加载逻辑改为tf.data动态读取Numpy文件:

def get_tensor_feature_label(version, txt_files, batch_size, run_mode, label_num):
    def get_data_from_cache(string_message):
        feature_index, label_index, key = string_message.decode().split(',')
        feature_index = int(feature_index)
        label_index = int(label_index)
        # 动态读取Numpy文件
        factor_arr = np.load(f"{ARR_DATA_PATH_MAP[version]}/{FLAGS.stock}/{run_mode}/factor/{key}.npy")
        label_arr = np.load(f"{ARR_DATA_PATH_MAP[version]}/{FLAGS.stock}/{run_mode}/label/{key}.npy")
        slice_arr = factor_arr[feature_index - FACTOR_WINDOW + 1: feature_index + 1]
        label_value = label_arr[label_index]
        x_data = slice_arr.astype(np.float32)
        y_data = label_value.astype(np.float32)
        return x_data, y_data

    def parser(txt_string_in):
        feat, label = tf.numpy_function(
            get_data_from_cache, [txt_string_in], (tf.float32, tf.float32))
        feat = tf.reshape(feat, (FACTOR_WINDOW, FACTOR_NUM_MAP[version]))
        label = tf.reshape(label, (label_num,))
        return feat, label

    # 后续tf.data流水线逻辑保持不变...

内容的提问来源于stack exchange,提问作者haoran.li

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 02:45:36