purrr::in_parallel处理数据框分组操作的优化及语法问题咨询
问题解答
一、如何用in_parallel实现有效提速?
你遇到的核心问题是**group_split生成20万个小数据框的开销远超过并行带来的收益**。每个小数据框的创建、内存分配和跨进程传输都会消耗大量时间,直接导致整体速度比原生分组操作慢100倍。要实现有效提速,需避免提前拆分数据,同时优化并行任务粒度:
1. 用group_map替代group_split
dplyr::group_map可直接在分组对象上应用函数,无需提前拆分出大量小数据框,配合in_parallel能减少不必要的开销:
library(tidyverse) library(mirai) # 启动并行守护进程 daemons(parallel::detectCores() - 1) df3 <- df |> group_by(x) |> # group_map自动传入每个分组的数据和分组键 group_map(in_parallel(\(data, .key) { data |> mutate(test = max(row_number())) })) |> bind_rows() # 关闭并行守护进程 daemons(0)
2. 优化任务粒度(关键)
对于20万个极小分组,即使使用group_map,每个并行任务的工作量太小,进程间通信开销依然会抵消并行优势。此时可将多个小分组打包成一批,减少并行任务数量:
daemons(parallel::detectCores() - 1) # 把x分组打包,每1000个不同的x为一批(可根据CPU核心数调整) df <- df |> mutate(batch = (as.integer(factor(x)) - 1) %/% 1000) df4 <- df |> group_by(batch) |> group_map(in_parallel(\(data, .key) { # 每个批次内部做原生分组操作,减少并行任务数 data |> mutate(test = max(row_number()), .by = x) })) |> bind_rows() |> select(-batch) # 移除临时批次列 daemons(0)
这种方式将并行任务数从20万降到200,每个任务的工作量足够大,能真正发挥多核CPU的优势。
二、你的语法错误原因及修正
你提供的代码无法运行,核心问题是**map的匿名函数没有将分组数据传递给in_parallel包裹的函数**:
# 错误写法 df3 <- df |> group_split(x) |> map(.f = ~in_parallel(\(x) x |> dplyr::mutate(test = max(dplyr::row_number()))))
这里~创建的匿名函数默认接收参数.x(即每个分组的data.frame),但你在in_parallel里定义的\(x)并没有拿到这个.x,相当于函数没有接收输入数据。
正确写法
写法1:直接传递in_parallel包裹的函数给map
df3 <- df |> group_split(x) |> map(in_parallel(\(data) data |> mutate(test = max(row_number())))) |> bind_rows()
写法2:保留~,显式传递分组数据
df3 <- df |> group_split(x) |> map(.f = ~in_parallel(\(data) data |> mutate(test = max(row_number())))(.x)) |> bind_rows()
但再次提醒:这种基于group_split的写法依然会因为小数据框的开销导致速度很慢,不建议实际使用。
内容的提问来源于stack exchange,提问作者deschen
相关产品推荐
相关产品推荐

