【发布时间】:2016-12-14 14:40:12
【问题描述】:
我正在尝试在 pyspark 中计算加权平均值,但没有取得很大进展
# Example data
df = sc.parallelize([
("a", 7, 1), ("a", 5, 2), ("a", 4, 3),
("b", 2, 2), ("b", 5, 4), ("c", 1, -1)
]).toDF(["k", "v1", "v2"])
df.show()
import numpy as np
def weighted_mean(workclass, final_weight):
return np.average(workclass, weights=final_weight)
weighted_mean_udaf = pyspark.sql.functions.udf(weighted_mean,
pyspark.sql.types.IntegerType())
但是当我尝试执行这段代码时
df.groupby('k').agg(weighted_mean_udaf(df.v1,df.v2)).show()
我收到了错误
u"expression 'pythonUDF' is neither present in the group by, nor is it an aggregate function. Add to group by or wrap in first() (or first_value) if you don't care which value you get
我的问题是,我可以指定一个自定义函数(采用多个参数)作为 agg 的参数吗?如果没有,是否可以在按键分组后执行加权平均等操作?
【问题讨论】:
-
您的意思是覆盖
weighted_mean函数吗? -
我想要做的是 a) groupby b) 根据 dataframe 的多个列执行操作。加权平均值只是一个例子。
-
我认为@cricket_007 的意思是您是否有意通过
weighted_mean = pyspark.sql.functions.udf(weighted_mean,这一行覆盖weighted_mean或者这是一个错字? -
我不认为
agg函数采用您提供的参数类型 -
@cricket_007 就在这里。 Agg 只接受正确的 UDAF。以stackoverflow.com/a/32101530/1560062 为例。虽然没有 Python API。对于像这样的琐碎案例,您只需要一个简单的公式,因此它看起来像是一个非常人为的问题。
标签: python apache-spark pyspark