【问题标题】:Calculating Weighted Mean in PySpark在 PySpark 中计算加权平均值
【发布时间】: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


【解决方案1】:

用户定义的聚合函数(UDAF,适用于pyspark.sql.GroupedData,但在 pyspark 中不受支持)与用户定义的函数(UDF,适用于 pyspark.sql.DataFrame)不同。

由于在 pyspark 中您无法创建自己的 UDAF,并且提供的 UDAF 无法解决您的问题,您可能需要返回 RDD 世界:

from numpy import sum

def weighted_mean(vals):
    vals = list(vals)  # save the values from the iterator
    sum_of_weights = sum(tup[1] for tup in vals)
    return sum(1. * tup[0] * tup[1] / sum_of_weights for tup in vals)

df.rdd.map(
    lambda x: (x[0], tuple(x[1:]))  # reshape to (key, val) so grouping could work
).groupByKey().mapValues(
    weighted_mean
).collect()

【讨论】:

  • 感谢@ijoseph 指出mapdf.rdd 上的工作。在撰写本文时,我习惯于将 df.map 作为简写。不确定它是否仍然有效,但最好是明确的。
  • 我相信,如果没有 .rdd,它会为我(在 Spark 2.4.3 中)抛出错误。
猜你喜欢
  • 2021-07-29
  • 2021-11-21
  • 2017-04-22
  • 1970-01-01
  • 2021-11-24
  • 1970-01-01
  • 2021-10-20
  • 1970-01-01
  • 2010-10-04
相关资源
最近更新 更多