Arrow R处理大数据集内存不足及map_batches聚合异常求助
我用arrow处理一个5000万行、14列的大型数据集,部分列是不超过14字符的字符串,通过open_dataset加载,多步操作后调用compute,之后在循环里做分组聚合时触发内存不足错误:
Error in compute.arrow_dplyr_query(): ! Out of memory: realloc of size 8589934656 failed
因为使用的arrow 11.0.0.3版本还未实现tidyr::pivot_wider,必须调用collect。
我尝试用arrow::map_batches规避内存问题,示例代码能得到正确聚合结果,但套用到大型数据集时,group_by+summarise的结果不符合预期,想知道用.lazy = FALSE参数能不能解决这个问题。
原循环聚合代码
geo_groupings <- list(national = "", region = "region", zip_code = "Zip Code") var_groupings <- list(unsegmented = "", no_gender = c("Age"), no_age = c("Gender"), all_variables = c("Gender", "Age")) sales <- sales %>% compute() # 3.2 Variables creation: for (geo_group in names(geo_groupings)){ for (var_group in names(var_groupings)){ groups <- c("date", "customer") if (var_group != "unsegmented") groups <- c(groups, var_groupings[[var_group]]) if (geo_group != "national") groups <- c(groups, geo_groupings[[geo_group]]) groups <- rlang::syms(groups) stacked <- sales %>% group_by(!!!groups[1:length(groups)]) %>% summarise_numeric() %>% mutate(Segment = tolower(paste(!!!groups[2:length(groups)], sep = "_"))) %>% relocate(Segment, .after = "date") %>% collect()
map_batches尝试代码
stacked <- arrow::map_batches(sales, function(batch){ batch %>% group_by(!!!groups[1:length(groups)]) %>% summarise_numeric() %>% mutate(customer = tolower(paste(!!!groups[2:length(groups)], sep = "_"))) %>% relocate(customer, .after = "date") }) %>% collect()
关于.lazy = FALSE的作用
.lazy = FALSE参数会强制map_batches对每个批次立即计算,但这无法解决分组聚合结果不符合预期的问题。
你的map_batches代码结果异常的核心原因是:map_batches是按单个批次独立处理分组的,跨批次的相同分组键不会被合并。比如某分组在批次1和批次2都存在,map_batches会分别在两个批次里聚合,最终得到两条重复分组的结果,和全局分组聚合的逻辑完全不符。
正确优化方向
取消提前
compute:你之前对sales调用compute()会把全量数据加载到内存,直接触发内存不足。应该保留数据集的懒加载状态,让arrow在磁盘层面执行聚合,仅把最终小结果collect到内存。
删掉sales <- sales %>% compute(),保留原始懒查询状态。全程使用懒操作:循环内的
group_by、summarise_numeric、mutate都保持懒执行,直到最后一步再collect。这样arrow会把整个查询计划下推到磁盘执行,无需加载全量数据。若必须用
map_batches:需在所有批次处理完成后,额外做一次全局分组聚合,合并跨批次的相同分组:
stacked <- arrow::map_batches(sales, function(batch){ batch %>% group_by(!!!groups) %>% summarise_numeric() %>% mutate(Segment = tolower(paste(!!!groups[-1], sep = "_"))) %>% relocate(Segment, .after = "date") }) %>% # 新增全局聚合,合并跨批次分组 group_by(!!!groups) %>% summarise_numeric() %>% collect()
额外建议
- 升级
arrow版本:新版本已支持tidyr::pivot_wider的懒执行,能避免很多必须collect的场景,从根源减少内存压力。 - 检查
summarise_numeric实现:确保它兼容arrow懒查询,若为自定义函数,需适配箭头数据集类型。
内容的提问来源于stack exchange,提问作者Alberto Agudo Dominguez

