sparklyr中不中断链式调用求唯一值及窗口函数用法问询
在sparklyr链式调用中计算非聚合的唯一值数量
这个问题我刚好踩过坑!核心原因是Spark的n_distinct()默认是聚合函数,直接在mutate里用会因为没有分组逻辑报错,但我们可以用窗口函数完美解决——既能保留所有原始行,又不用中断链式调用。
先准备你的数据集
先按你提供的代码生成并同步到Spark:
set.seed(.328) df <- data.frame( ids = floor(runif(10, 1, 10)), cats = sample(letters[1:3], 10, replace = TRUE), vals = rnorm(10) ) # 假设你已经初始化了Spark连接sc df.spark <- copy_to(sc, df, name = "df_spark", overwrite = TRUE)
解决方案1:全局唯一值计数(保留所有行)
如果你想在每一行都显示ids列的全局唯一值数量,只需要定义一个全局窗口(不分区、不排序),然后在窗口上计算去重计数:
library(dplyr) library(sparklyr) # 定义全局窗口(覆盖整个数据集) global_window <- window() # 链式调用完成计算 df.spark %>% mutate(global_unique_ids = count_distinct(ids) %>% over(global_window)) %>% collect() # 拉取结果到本地查看
如果count_distinct有兼容问题,也可以用原生Spark SQL表达式更直接:
df.spark %>% mutate(global_unique_ids = expr("count(DISTINCT ids) OVER ()"))
解决方案2:分组内唯一值计数(保留所有行)
如果是按cats分组,计算每组内ids的唯一值数量,只需要给窗口加上分区条件:
# 定义按cats分组的窗口 cat_group_window <- window(partition_by = cats) df.spark %>% mutate(cat_unique_ids = count_distinct(ids) %>% over(cat_group_window))
这样每一行都会显示自己所在cats分组的ids唯一值数量,全程不用中断链式调用,也不用先聚合再JOIN回去。
更广泛的Spark窗口函数用法
在sparklyr里使用类SQL窗口函数,核心是用dplyr::window()定义窗口规则,然后配合dplyr封装的函数或原生Spark表达式:
1. 基础窗口定义
窗口可以包含以下核心参数:
partition_by: 对应SQL的PARTITION BY,按指定列分组order_by: 对应SQL的ORDER BY,指定窗口内的排序规则rows_between/range_between: 指定窗口的行范围(比如移动窗口)
2. 常见窗口函数示例
行号与排名
# 全局按vals排序的行号 df.spark %>% mutate(row_num = row_number() %>% over(window(order_by = vals))) # 按cats分组,按vals降序的排名 df.spark %>% mutate(cat_rank = rank() %>% over(window(partition_by = cats, order_by = desc(vals))))
移动统计量(比如移动平均)
# 按cats分组,取当前行及前2行的vals平均值 rolling_window <- window( partition_by = cats, order_by = vals, rows_between(-2, 0) # 范围:当前行往前2行到当前行 ) df.spark %>% mutate(rolling_avg_vals = avg(vals) %>% over(rolling_window))
偏移函数(lead/lag)
# 按cats分组,取当前行的下一个vals值 df.spark %>% mutate(next_val = lead(vals) %>% over(window(partition_by = cats, order_by = vals)))
原生Spark SQL表达式(处理复杂场景)
如果dplyr封装的函数不够用,直接用expr()写原生Spark SQL窗口表达式:
# 计算分组内的累积和 df.spark %>% mutate(cum_sum_vals = expr("sum(vals) OVER (PARTITION BY cats ORDER BY vals)"))
内容的提问来源于stack exchange,提问作者Mouad_Seridi
相关产品推荐
相关产品推荐

