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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 16:45:15