使用furrr处理大型tibble列表时优化内存:解决异常内存增长
内存增长问题分析与解决方案
问题描述
需对含约300万条记录(大小约202MB)的tibble按key分组处理,每组调用函数生成CSV文件。当前代码如下:
主执行代码:
plan(multisession, workers=20) hpar$input_df %>% group_by(key) %>% group_split() %>% future_walk(sf_do_all_one_series, hpar$df_stockout_decaying, hpar$out_path, hpar$out_prefix, .env_globals = empty_env())
处理函数:
sf_do_all_one_series <- function(df_1, df_stockout_decaying, out_path, out_prefix){ key <- df_1 %>% distinct(key) %>% pull() do_stuff(df_1, df_stockout_decaying) %>% do_more_stuff() %>% write.table(file=paste0(out_path, out_prefix, "_", key, ".csv"), quote=FALSE, sep='\t', row.names=FALSE) invisible() }
其中hpar$df_stockout_decaying为小型常量tibble,out_path/out_prefix为字符串。执行时内存随时间大幅增长,已尝试移除中间对象、显式gc()、设置.env_globals = empty_env(),但无明显改善。
内存增长原因分析
group_split()的额外内存开销:group_split()会生成包含约3万个tibble的列表,每个子tibble保留原数据的元数据与结构,且原大tibble在split后不会立即释放,加上列表本身的管理开销,导致内存占用激增。- 多进程环境下的对象复制与残留:multisession模式下,每个worker会复制所需的全局对象;同时future框架的任务队列可能保留已完成任务的临时对象,20个worker的并发会放大这种内存累积效应。
- 隐式中间对象未及时回收:
do_stuff/do_more_stuff中生成的中间对象,在多进程环境下垃圾回收时机不可控,即使使用管道,仍可能存在未被及时释放的大对象。
可行解决方案
- 避免预生成完整分组列表:用
furrr::future_walk直接在分组上操作,跳过group_split()步骤,减少列表的内存占用:plan(multisession, workers=12) # 适当减少worker数量 hpar$input_df %>% group_by(key) %>% furrr::future_walk( ~sf_do_all_one_series(.x, hpar$df_stockout_decaying, hpar$out_path, hpar$out_prefix), .options = furrr_options(globals = c("hpar$df_stockout_decaying", "hpar$out_path", "hpar$out_prefix")) ) - 显式释放内存并强制回收:在处理函数末尾显式移除大对象并触发垃圾回收,确保单个任务完成后释放内存:
sf_do_all_one_series <- function(df_1, df_stockout_decaying, out_path, out_prefix){ key <- df_1 %>% distinct(key) %>% pull() do_stuff(df_1, df_stockout_decaying) %>% do_more_stuff() %>% write.table(file=paste0(out_path, out_prefix, "_", key, ".csv"), quote=FALSE, sep='\t', row.names=FALSE) # 释放大对象并触发GC rm(df_1, key) gc(verbose = FALSE) invisible() } - 精准控制全局变量传递:避免传递整个
hpar对象,只传递函数所需的小对象,减少worker端的复制开销:
在future_walk中指定.globals = c("df_stockout_decaying", "out_path", "out_prefix"),而非依赖自动检测。 - 调整并发数或改用分批处理:若多进程内存压力过大,可减少worker数量(如8-12),或改用单进程分批处理:将分组分成若干大批次,处理完一批后执行
gc(),平稳控制内存占用。
内容的提问来源于stack exchange,提问作者alexon
相关产品推荐
相关产品推荐

