PySpark根据分组平均值为DataFrame添加新列报错问题求助
PySpark根据分组平均值为DataFrame添加新列报错问题求助
嗨,我来帮你解决这个问题~先理清楚你的需求:你想给原DataFrame新增一列,每一行的值对应其identifier分组下numbers_to_sum的平均值,对吧?
你的问题重现
先确认下你的原始DataFrame:
identifier numbers_to_sum dog 5 cat 4 dog 3 parrot 2 cat 7
你期望的结果是每个分组的平均值对应到每一行,而你尝试的代码报错了:
df.withColumn('average', df.groupBy('identifier').mean('numbers_to_avg').alias('avg').select('avg'))
得到的错误是:
PySparkTypeError: [NOT_COLUMN] Argument
colshould be a Column, got DataFrame.
错误原因
你这里的问题很明确:df.groupBy(...).mean(...).select(...)返回的是一个DataFrame对象,但withColumn方法的第二个参数要求必须是Column类型的表达式,直接把DataFrame传进去自然会报错啦。
两种正确的解决方案
方法一:使用窗口函数(推荐,更简洁)
窗口函数非常适合这种“把分组聚合结果映射到每一行”的场景,步骤如下:
- 先导入需要的函数和窗口对象
- 定义按
identifier分组的窗口规则 - 用
avg函数结合窗口规则生成新列
代码示例:
from pyspark.sql import Window from pyspark.sql.functions import avg # 定义窗口:按identifier分组 window_spec = Window.partitionBy("identifier") # 添加平均值列 df_with_avg = df.withColumn("average", avg("numbers_to_sum").over(window_spec)) # 查看结果 df_with_avg.show()
运行后就能得到你想要的结果:
identifier numbers_to_sum average dog 5 4.0 cat 4 5.5 dog 3 4.0 parrot 2 2.0 cat 7 5.5
方法二:先分组聚合再Join
如果你更习惯用SQL风格的Join操作,也可以先计算每个分组的平均值,再和原DataFrame关联:
# 第一步:计算每个identifier的平均值 avg_df = df.groupBy("identifier").agg(avg("numbers_to_sum").alias("average")) # 第二步:和原DataFrame左连接,把平均值列带回去 df_with_avg = df.join(avg_df, on="identifier", how="left") # 查看结果 df_with_avg.show()
这个方法也能得到完全一样的结果,适合需要单独保存分组平均值数据的场景。
备注:内容来源于stack exchange,提问作者geds133
相关产品推荐
相关产品推荐

