sf空间距离计算并行化优化及进度追踪、结果保存问题咨询
问题背景
我有两个空间数据集:
- Total_df:约200万条建筑观测记录,结构如下:
| Id | Variables | SFC Point object | Zip codes |
|---|---|---|---|
| 1 | 10 | POINT (543611.8 6389285) | 2324 |
| 2 | 15 | POINT (513611.8 6349285) | 2324 |
| 3 | 12 | POINT (533611.8 6359285) | 2329 |
- forest_distance:森林多边形数据,已拆分为10份存入列表
list.dfs,结构如下:
| Id | Variables | SFC Polygon object |
|---|---|---|
| 1 | 10 | POLYGON 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
相关产品推荐
相关产品推荐

