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

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会分别在两个批次里聚合,最终得到两条重复分组的结果,和全局分组聚合的逻辑完全不符。

正确优化方向

  1. 取消提前compute:你之前对sales调用compute()会把全量数据加载到内存,直接触发内存不足。应该保留数据集的懒加载状态,让arrow在磁盘层面执行聚合,仅把最终小结果collect到内存。
    删掉sales <- sales %>% compute(),保留原始懒查询状态。

  2. 全程使用懒操作:循环内的group_by、summarise_numeric、mutate都保持懒执行,直到最后一步再collect。这样arrow会把整个查询计划下推到磁盘执行,无需加载全量数据。

  3. 若必须用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 03:05:32