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

sf空间距离计算并行化优化及进度追踪、结果保存问题咨询

问题背景

我有两个空间数据集:

  1. Total_df:约200万条建筑观测记录,结构如下:
IdVariablesSFC Point objectZip codes
110POINT (543611.8 6389285)2324
215POINT (513611.8 6349285)2324
312POINT (533611.8 6359285)2329
  1. forest_distance:森林多边形数据,已拆分为10份存入列表list.dfs,结构如下:
IdVariablesSFC Polygon object
110POLYGON Z ((455302.7 6252026 9.09, 455292.6 6252034 9.09, ..., 455302.7 6252026 9.09))

我已实现按邮政编码拆分Total_df计算空间距离的逻辑,现在想通过并行化处理拆分后的森林数据提升速度,同时需要解决两个问题:

  • 如何在不同并行会话中打印进度?
  • 运行代码数小时后终止,未保存任何结果,该怎么解决?

原代码如下:

registerDoParallel(cores = 6)    

# Use foreach to loop over list.dfs in parallel
foreach(d = 1:length(list.dfs), .packages = "sf", .combine = 'c') %dopar% {
  # Get the data frame at position 'd' in the list
  df <- list.dfs[[d]]
  
  # Open a list to store combined inner results 
  grand_list <- list()
  
  # Initialize an empty list to store the results of the inner loop
  inner_results <- list()
  
  # zip_code 
  zipcode <- sort(unique(Total_df$zipcode))
  

  # Use a regular for loop to iterate over zipcode
  for(i in zipcode) {
    cat(i, "\n")
    start_time <- Sys.time()
    
    # Subset the data
    subset_df <- Total_df[Total_df$zipcode == i, ]
    
    if(nrow(subset_df) > 0) {
      # Calculate distances
      distances <- sf::st_distance(subset_df, df)
      
      # Define the 'miin' function, or replace it with an appropriate function
      miin <- function(x) min(x, na.rm = TRUE)
      
      # Calculate minimum distances
      min_distances <- apply(distances, 1, miin)
      
      # Store minimum distances in a new column
      subset_df$min_distances <- min_distances
    }
    
    end_time <- Sys.time()
    print(paste("Time for municipality Forest", i, ": ", end_time - start_time))
    
    # Store the updated subset_df in the inner_results list
    inner_results[[i]] <- subset_df
  }
  
  # Combine the results of the inner loop using do.call
  grand_list[[d]] <- do.call(rbind, inner_results)
}
解决方案

1. 并行会话中打印进度

并行会话默认无法直接向主进程控制台输出cat/print内容,推荐两种实用实现方式:

方法1:用doSNOW实现可视化进度条+日志记录

替换doParallel为doSNOW,主进程显示整体任务进度,子进程输出写入日志文件:

library(doSNOW)
# 创建集群,指定子进程输出写入日志
cl <- makeCluster(6, outfile = "parallel_progress_log.txt")
registerDoSNOW(cl)

# 设置主进程进度条
pb <- txtProgressBar(max = length(list.dfs), style = 3)
progress <- function(n) setTxtProgressBar(pb, n)
opts <- list(progress = progress)

# 并行循环
results <- foreach(d = 1:length(list.dfs), .packages = "sf", .combine = rbind, .options.snow = opts) %dopar% {
  df <- list.dfs[[d]]
  inner_results <- list()
  zipcode <- sort(unique(Total_df$zipcode))
  
  for(i in zipcode) {
    # 子进程输出到日志
    cat(sprintf("并行任务%d:开始处理邮政编码%d\n", d, i), file = "parallel_progress_log.txt", append = TRUE)
    start_time <- Sys.time()
    
    subset_df <- Total_df[Total_df$zipcode == i, ]
    if(nrow(subset_df) > 0) {
      distances <- sf::st_distance(subset_df, df)
      min_distances <- apply(distances, 1, function(x) min(x, na.rm = TRUE))
      subset_df$min_distances_forest_chunk <- min_distances
    }
    
    end_time <- Sys.time()
    cat(sprintf("并行任务%d:邮政编码%d处理完成,耗时%.2f秒\n", d, i, as.numeric(end_time - start_time)), file = "parallel_progress_log.txt", append = TRUE)
    inner_results[[as.character(i)]] <- subset_df
  }
  
  do.call(rbind, inner_results)
}

close(pb)
stopCluster(cl)

方法2:用message配合.verbose参数

开启foreach的.verbose=TRUE,子进程的message会被主进程捕获并输出:

registerDoParallel(6)
results <- foreach(d = 1:length(list.dfs), .packages = "sf", .combine = rbind, .verbose = TRUE) %dopar% {
  df <- list.dfs[[d]]
  inner_results <- list()
  zipcode <- sort(unique(Total_df$zipcode))
  
  for(i in zipcode) {
    message(sprintf("任务%d:处理邮编%d", d, i))
    # 其余计算逻辑与原代码一致
    subset_df <- Total_df[Total_df$zipcode == i, ]
    if(nrow(subset_df) > 0) {
      distances <- sf::st_distance(subset_df, df)
      min_distances <- apply(distances, 1, function(x) min(x, na.rm = TRUE))
      subset_df$min_distances_forest_chunk <- min_distances
    }
    inner_results[[as.character(i)]] <- subset_df
  }
  do.call(rbind, inner_results)
}

2. 解决运行数小时终止无结果的问题

原代码存在结果丢失、内存过载、无错误处理等问题,修复步骤如下:

步骤1:必须为foreach结果赋值变量

原代码未将并行计算结果赋值给变量,导致计算完成后结果直接丢弃,需修改为:

final_results <- foreach(...) %dopar% {
  # 计算逻辑
}

步骤2:添加中间结果保存

每个并行任务完成后立即保存结果到文件,避免崩溃前功尽弃:

foreach(d = 1:length(list.dfs), .packages = "sf", .combine = rbind) %dopar% {
  df <- list.dfs[[d]]
  inner_results <- list()
  zipcode <- sort(unique(Total_df$zipcode))
  
  for(i in zipcode) {
    # 计算逻辑
    subset_df <- Total_df[Total_df$zipcode == i, ]
    if(nrow(subset_df) > 0) {
      distances <- sf::st_distance(subset_df, df)
      min_distances <- apply(distances, 1, function(x) min(x, na.rm = TRUE))
      subset_df$min_distances_forest_chunk <- min_distances
    }
    inner_results[[as.character(i)]] <- subset_df
  }
  
  chunk_result <- do.call(rbind, inner_results)
  # 保存当前任务结果
  saveRDS(chunk_result, sprintf("forest_distance_chunk_%d.rds", d))
  chunk_result # 返回结果给foreach合并
}

步骤3:优化内存使用

  • 提前按邮政编码拆分Total_df,避免重复子集化操作
  • 用nngeo::st_nn直接计算最近距离,替代生成全量距离矩阵(大幅减少内存占用):
library(nngeo)
# 提前按邮编拆分Total_df
total_by_zip <- split(Total_df, Total_df$`Zip codes`)

foreach(d = 1:length(list.dfs), .packages = c("sf", "nngeo"), .combine = rbind) %dopar% {
  df <- list.dfs[[d]]
  chunk_result <- lapply(total_by_zip, function(subset_df) {
    if(nrow(subset_df) == 0) return(subset_df)
    # 直接计算最近距离,效率远高于全矩阵计算
    nearest_dist <- st_nn(subset_df, df, returnDist = TRUE)$dist
    subset_df$min_distances_forest_chunk <- unlist(nearest_dist)
    subset_df
  })
  chunk_result <- do.call(rbind, chunk_result)
  saveRDS(chunk_result, sprintf("forest_distance_chunk_%d.rds", d))
  chunk_result
}

步骤4:添加错误捕获

用tryCatch包裹计算逻辑,避免单个任务崩溃导致整个并行进程终止:

foreach(d = 1:length(list.dfs), .packages = "sf", .combine = rbind) %dopar% {
  tryCatch({
    df <- list.dfs[[d]]
    inner_results <- list()
    zipcode <- sort(unique(Total_df$zipcode))
    
    for(i in zipcode) {
      # 计算逻辑
      subset_df <- Total_df[Total_df$zipcode == i, ]
      if(nrow(subset_df) > 0) {
        distances <- sf::st_distance(subset_df, df)
        min_distances <- apply(distances, 1, function(x) min(x, na.rm = TRUE))
        subset_df$min_distances_forest_chunk <- min_distances
      }
      inner_results[[as.character(i)]] <- subset_df
    }
    
    chunk_result <- do.call(rbind, inner_results)
    saveRDS(chunk_result, sprintf("forest_distance_chunk_%d.rds", d))
    chunk_result
  }, error = function(e) {
    cat(sprintf("任务%d出错:%s\n", d, e$message), file = "error_log.txt", append = TRUE)
    return(NULL) # 出错返回空,不影响其他任务
  })
}

内容的提问来源于stack exchange,提问作者user23675742

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 07:52:02