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

在Arrow的open_dataset加载的Parquet文件中应用自定义函数遇阻求助

问题:Apache Arrow中自定义众数函数分组聚合报错的解决办法

场景与问题

我有分布在4232个parquet文件中的数据,经open_dataset解析完整schema后聚合为约6745697行、30486列的数据集。需要按Participant(约3000个唯一值)分组,对排除VOI_1、VOI_2和Participant之外的人口统计学列(含字符串、整数、布尔值)取众数聚合。

原代码如下:

mode <- function(v) {
  uniqv <- unique(v)
  uniqv[which.max(tabulate(match(v, uniqv)))]
}

DOI <- open_dataset(sources = <path_to_directory_of_parquet_files>, unify_schemas = TRUE)

Demographics <- DOI |> 
  filter(treatment_flag == 0, comparator_flag == 0) |> 
  group_by(Participant) |> 
  summarize(across(DOI$schema$names[!grepl("VOI_1|VOI_2|Participant", DOI$schema$names)],mode)) |> 
  collect() |> 
  as.data.frame()

运行时出现错误:

Error: Error in summarize_eval(names(exprs)[i], exprs[[i]], ctx, length(.data$group_by_vars) >  : 
  Expression mode(Dataset) is not an aggregate expression or is not supported in Arrow
Call collect() first to pull data into R.

疑问:能否不通过collect()把全量数据加载到R内存,实现自定义众数的分组聚合?还是只能使用Arrow支持的函数?


解决方案

原因说明

Apache Arrow的Dataset操作依托后端引擎(如C++)执行,无法直接识别并执行R自定义函数(比如你写的mode)——这类函数无法被翻译成Arrow支持的执行表达式。因此有两种可行实现路径:


1. 用Arrow原生函数实现众数(推荐)

全程在Arrow后端执行,无需加载全量数据到R内存,效率更高。利用Arrow支持的count()、slice_head()等函数实现众数逻辑:

Demographics <- DOI |> 
  filter(treatment_flag == 0, comparator_flag == 0) |> 
  group_by(Participant) |> 
  summarize(
    across(
      !matches("VOI_1|VOI_2|Participant"),
      ~ .x |> count(sort = TRUE) |> slice_head(n = 1) |> pull(1)
    )
  ) |> 
  collect() |> 
  as.data.frame()

逻辑解释:对每个目标列,先在分组内统计各值的出现次数并排序,取次数最多的第一个值(即众数),所有操作都在Arrow后端完成,仅最后collect()把聚合后的3000行结果拉到R内存。


2. 用group_map分批处理自定义函数

如果必须使用你写的R自定义mode函数,可以用group_map逐个将分组数据拉到R内存计算,相比直接collect()全量数据,内存压力小很多(每个分组仅约2000行数据):

mode <- function(v) {
  uniqv <- unique(v)
  uniqv[which.max(tabulate(match(v, uniqv)))]
}

Demographics <- DOI |> 
  filter(treatment_flag == 0, comparator_flag == 0) |> 
  group_by(Participant) |> 
  group_map(function(.group_data, .group_key) {
    # 对单个分组的数据计算众数
    .group_data |> 
      summarize(across(!matches("VOI_1|VOI_2|Participant"), mode)) |> 
      bind_cols(.group_key) # 合并分组键Participant
  }) |> 
  bind_rows() |> 
  as.data.frame()

内容的提问来源于stack exchange,提问作者TDeramus

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 22:35:41