基于R语言Torch的异步并行数据加载问题排查
问题:R Torch中DataLoader的num_workers并行加载无效的排查与解决
我使用R语言的Torch框架做迁移学习,训练CNN处理大型卫星图像数据集(13通道299×299)。由于数据集过大无法一次性加载,只能通过DataLoader从SSD逐样本加载,但单批次加载时间是模型处理时间的20倍,因此尝试通过num_workers开启异步并行数据加载。但设置该参数后,训练循环内的批次加载时间并未减少,仅在首次加载前出现明显的线程创建开销。想确认该功能在R Torch中是否可行,以及我的实现存在哪些错误。
用户提供的示例代码:
#### Data generation #### data_dir<-"./data/processed/satalite_images/to_use/" if(!dir.exists(paste0(data_dir, "0")))dir.create(paste0(data_dir, "0"), recursive = T) if(!dir.exists(paste0(data_dir, "1")))dir.create(paste0(data_dir, "1"), recursive = T) n_images<-1000 images_done<-length(list.files(data_dir, recursive = T)) if(!images_done>=n_images){ pb<- txtProgressBar(max=n_images, initial = images_done, style = 3) images_left<-n_images-images_done for(i in 1:images_left){ setTxtProgressBar(pb, i+images_done) image<-floor(runif(13*299*299, min = 0, max=256)) # they have to be .png files otherwise they can`t be # recognized by the dataloader saveRDS(image, paste0(data_dir, sample(c(0,1), size=1), "/",i+images_done, ".png")) } close(pb) } #### model building #### library(torchvision) library(torch) ds<-torchvision::image_folder_dataset( root="./data/processed/satalite_images/to_use", loader=function(path){ # I have images of size 299x299 with 13 channels. # optimizing this loading step yielded no significant improvement. return(array(readRDS(path), dim=c(13,299,299))*1.0) }, target_transform = function(x){a<-c(0.0,1.0)[x];dim(a)<-1;return(a)} ) #Here I set num_workers to different numbers, but that did not change the loading time dl2<-torch::dataloader(ds, batch_size=110L, shuffle = T, pin_memory = T, num_workers = 4L) #just a random pretrained model for transfer learning model_torch = torchvision::model_alexnet(pretrained = T) model_torch$parameters |> purrr::walk(function(param) param$requires_grad_(FALSE)) # replacing the last layer to my desired classifier inFeat =model_torch$classifier$'6'$in_features model_torch$classifier$'6' = nn_linear(inFeat, out_features = 1L) # I have 13 input channels, therefore I replace the first conv layer with a equivialent one but with 13 input channels conv1<-torch::nn_conv2d(in_channels=13L, out_channels=model_torch[[1]]$`0`$out_channels, kernel_size =model_torch[[1]]$`0`$kernel_size , stride = model_torch[[1]]$`0`$stride, padding =model_torch[[1]]$`0`$padding, dilation = model_torch[[1]]$`0`$dilation, groups = model_torch[[1]]$`0`$groups, bias = TRUE) model_torch[[1]]$`0`<-conv1 model_torch<-model_torch$to(device = "cuda") opt = optim_adam(params = model_torch$parameters, lr = 0.01) #### training #### #trainings loop for(e in 1:1){ losses = c() #storing the time which the loop uses for computing and data loading end<-Sys.time() coro::loop( for(batch in dl2){ start<-Sys.time() #this is the time it takes to load a batch print(start-end) print("computing") opt$zero_grad() pred = model_torch(batch[[1]]$to(device="cuda")) res=batch[[2]]$to(device = "cuda") loss = nnf_binary_cross_entropy(input=torch_sigmoid(pred),target=res) loss$backward() opt$step() losses = c(losses, loss$item()) end<-Sys.time() #this is the time it takes to process a batch print(end-start) print("loading") } ) }
问题排查与解决方案
R Torch的num_workers确实支持异步并行数据加载,你的实现存在几个关键问题导致功能未生效,以下是针对性修复:
1. 计时逻辑完全错误
你当前的计时代码把上一轮模型处理的耗时当成了当前批次的加载时间,完全混淆了异步加载的时间线。异步加载的核心是:模型处理当前批次时,后台线程已在预取下一批数据,因此正确的方式是单独统计数据加载环节的耗时,而非用循环前后的时间差。
修正后的计时方式:
# 训练循环 for(e in 1:1){ losses = c() coro::loop( for(batch in dl2){ print("开始处理批次") start_compute <- Sys.time() opt$zero_grad() pred = model_torch(batch[[1]]$to(device="cuda")) res=batch[[2]]$to(device = "cuda") loss = nnf_binary_cross_entropy(input=torch_sigmoid(pred),target=res) loss$backward() opt$step() losses = c(losses, loss$item()) end_compute <- Sys.time() print(paste("模型处理耗时:", end_compute - start_compute)) } ) }
2. 数据集加载函数的序列化问题
R的多进程环境下,自定义loader函数需要能被正确序列化到子进程。需确保:
- 用绝对路径避免子进程工作目录与主进程不一致
- 显式调用基础包函数,避免子进程依赖缺失
修正后的Dataset定义:
# 使用绝对路径 data_dir <- normalizePath("./data/processed/satalite_images/to_use/") ds <- torchvision::image_folder_dataset( root = data_dir, loader = function(path){ # 显式调用base包的readRDS,确保子进程能找到 img_data <- base::readRDS(path) return(array(img_data, dim=c(13,299,299)) * 1.0) }, target_transform = function(x){ a <- c(0.0,1.0)[x] dim(a) <- 1 return(a) } )
3. 优化DataLoader参数配置
- 开启
persistent_workers = TRUE:避免每个epoch重新创建工作线程,减少首次加载开销 num_workers建议设为CPU核心数的1-2倍(比如8核CPU设为8或16)- 保留
pin_memory = TRUE以加速CPU到GPU的数据传输
修正后的DataLoader定义:
dl2 <- torch::dataloader( ds, batch_size = 110L, shuffle = TRUE, pin_memory = TRUE, num_workers = 8L, persistent_workers = TRUE # 新增:跨epoch保留工作线程 )
4. 底层数据加载优化
你提到加载步骤优化无效果,可尝试:
- 将数据格式从RDS换成Torch张量格式,避免数组转张量的开销
- 提前预处理数据集并缓存,减少实时转换时间
示例预处理脚本:
# 提前预处理数据(单独运行) preprocess_dir <- normalizePath("./data/processed/satalite_images/preprocessed/") dir.create(preprocess_dir, recursive = TRUE) original_files <- list.files(data_dir, recursive = TRUE, full.names = TRUE) for(file in original_files){ img_data <- readRDS(file) img_tensor <- torch::torch_tensor(array(img_data, dim=c(13,299,299)) * 1.0) # 保持原目录结构 relative_path <- gsub(data_dir, "", file) save_path <- file.path(preprocess_dir, relative_path) dir.create(dirname(save_path), recursive = TRUE) torch::save_tensor(img_tensor, save_path) } # 预处理后的数据加载 ds <- torchvision::image_folder_dataset( root = preprocess_dir, loader = function(path){ return(torch::load_tensor(path)) }, target_transform = function(x){ a <- c(0.0,1.0)[x] dim(a) <- 1 return(a) } )
验证并行加载是否生效
可通过CPU使用率验证:开启num_workers>0后,训练时多个CPU核心会被数据加载线程占用,CPU使用率会显著高于num_workers=0的情况。若CPU使用率无变化,说明并行加载未生效,需重新检查上述步骤。
内容的提问来源于stack exchange,提问作者Daniel Maar
相关产品推荐
相关产品推荐

