SparkR结合dplyr通过gapply实现count窗口函数报错咨询
错误原因
你的代码报错核心是对gapply的运行逻辑理解有误:
gapply分组后,会将每个分组的数据以普通R data.frame的形式传入分布式节点上的自定义函数,参数x不是SparkR的分布式DataFrame对象。SparkR::count()是Spark分布式DataFrame的专属方法,必须传入Spark DataFrame作为参数才能调用,你空参调用该方法,自然会抛出「找不到对应签名的count方法」的错误。- 分组内的统计计算不需要调用Spark侧API,直接基于传入的本地R data.frame用基础R函数实现即可。
结合dplyr管道的正确实现
以下代码和你写的窗口函数SQL效果完全一致,每个物种分组下的所有行都会附带该分组的总记录数:
library(SparkR) library(dplyr) df <- createDataFrame(iris) createOrReplaceTempView(df, "iris") result <- df %>% SparkR::group_by(.$Species) %>% gapply( function(key, x) { # x为当前分组对应的本地R data.frame,nrow(x)即为分组总行数 data.frame(x, RowCount = nrow(x)) }, schema = "Sepal_Length double, Sepal_Width double, Petal_Length double, Petal_Width double, Species string, RowCount integer" ) display(result)
注意事项
- 管道操作中引用分组字段时使用
.$Species而非df$Species,避免管道传参时出现对象引用错误。 gapply传入的自定义函数运行在分布式Worker节点,函数内部只能操作传入的本地R对象,不要调用需要连接Spark驱动、操作分布式DataFrame的API,否则会触发节点运行错误。
内容的提问来源于stack exchange,提问作者Jovanny
相关产品推荐
相关产品推荐

