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

在R中使用Apache Arrow时程序随机无响应问题求助

问题描述

作为Apache Arrow新手,在处理大规模数据集的模拟与计算算法时遇到随机卡死问题:R会在函数执行到f5 <- f4 %>% collect()(打印"5"之后)时随机卡住,点击控制台停止按钮后必须终止R才能结束进程。卡死时机不固定,有时仅迭代数次就出现,有时可迭代数百次,且CPU和内存占用从未超过60%。

相关代码

library(arrow)
library(dplyr)

N<-500000
sample <- data.frame(id = 1:N,
                     x = runif(N,0,1),
                     y1 = rnorm(N,8,3),
                     y2 = rnorm(N,15,3),
                     poststratum = round(runif(N,0,4)),
                     SRS = round(runif(N,0,1)),
                     SRS_R_y1 = round(runif(N,0,1)),
                     SRS_R_y2 =round(runif(N,0,1)))
arrow::write_dataset(sample, paste0("samples/test.pq"))


samples <- arrow::open_dataset(list.files(paste0("samples/test.pq"), full.names = TRUE))

stratum_function <- function(data_in, method, y, x){
  print(x)
  print("1")
  # Calculate big in each stratum
  f1 <- data_in %>% group_by(poststratum) %>% summarise(N_s = n()) 
  print("2")
  # Calculate number of respondants in each stratum
  f2 <- data_in %>% 
    filter(.data[[paste0(method, "_", y)]] == 1) %>% 
    group_by(poststratum) %>% 
    summarise(n_s = n()) 
  print("3")
  # Join two privious parquets together
  f3 <- left_join(f1,f2 ,by = "poststratum")
  print("4")
  gc()
  # Calculate weights based on post stratification method
  f4 <- f3 %>% mutate(w_s = N_s / n_s, 
                      W_SQ_S_s = (N_s / n_s)*(N_s / n_s)) 
  print("5")
  
  f5 <- f4 %>% collect() 
  print("6")
}

lapply(1:10000, FUN=function(x){stratum_function(samples,"SRS_R", "y1", x)})

环境信息

  • dplyr版本:1.1.2
  • arrow版本:11.0.0.3
  • RStudio版本:2023.06.1+524 "Mountain Hydrangea"(Windows系统)
  • R版本:4.1.2 (2021-11-01) "Bird Hippie"

解决建议

1. 升级依赖库版本

当前使用的arrow 11.0.0.3属于较旧版本,后续版本修复了大量内存管理和执行引擎相关的bug。建议升级到最新稳定版:

install.packages(c("arrow", "dplyr"))

2. 避免重复访问磁盘数据集

每次迭代都从磁盘读取数据集可能引发资源竞争或锁问题。可以预先将数据集加载到内存中的Arrow Table后再进行循环计算:

# 预先加载数据集到内存
samples_table <- arrow::open_dataset(list.files("samples/test.pq", full.names = TRUE)) %>% collect()

# 修改函数接收内存Table(逻辑无需其他变动)
stratum_function <- function(data_in, method, y, x){
  print(x)
  print("1")
  f1 <- data_in %>% group_by(poststratum) %>% summarise(N_s = n()) 
  print("2")
  f2 <- data_in %>% 
    filter(.data[[paste0(method, "_", y)]] == 1) %>% 
    group_by(poststratum) %>% 
    summarise(n_s = n()) 
  print("3")
  f3 <- left_join(f1,f2 ,by = "poststratum")
  print("4")
  f4 <- f3 %>% mutate(w_s = N_s / n_s, 
                      W_SQ_S_s = (N_s / n_s)*(N_s / n_s)) 
  print("5")
  
  f5 <- f4 %>% collect() 
  print("6")
}

lapply(1:10000, FUN=function(x){stratum_function(samples_table,"SRS_R", "y1", x)})

3. 移除手动GC调用

函数中手动调用gc()可能与Arrow的内存管理机制冲突,频繁GC可能导致执行流程异常。建议移除该调用,让R自动管理内存。

4. 检查文件系统与权限

Windows系统下,Parquet文件的读写可能存在文件锁问题:

  • 将数据集存储在本地SSD而非网络共享磁盘
  • 确保没有其他进程(如杀毒软件、文件同步工具)在访问该Parquet文件

5. 重构无意义循环

当前lapply(1:10000)的循环逻辑完全重复执行相同计算(x仅用于打印),这种无意义的重复会放大潜在的资源泄漏问题。如果是测试需要,可先减少迭代次数(如100次)验证问题,或直接去掉循环只执行一次计算。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 20:40:54