PySpark中percentile_approx在where子句使用时报错问题咨询
问题原因与解决方案
原因
percentile_approx是聚合函数,它会对整个数据集计算出一个单一的分位数值。而where子句是针对DataFrame的每一行做过滤判断,Spark不支持在过滤条件中直接引用聚合函数——因为聚合结果是全局的,不是行级别的数据,无法逐行进行比较判断,因此会抛出UnsupportedOperationException。
解决方案
有两种常用的处理方式:
方式一:先计算分位数值,再作为常量过滤
先单独计算出分位数的标量值,再将其作为常量传入where条件中:
# 计算0.5分位数并提取为标量值 median_val = spark_df.select(f.expr("percentile_approx(`item_sum`, 0.5, 100)")).collect()[0][0] # 使用标量值进行过滤 spark_df.where(f.col("item") <= median_val).show()
方式二:通过交叉连接将分位数值关联到每一行,再过滤
先把分位数结果生成一个单行DataFrame,再通过交叉连接将该值附加到原DataFrame的每一行,最后进行过滤:
# 生成包含分位数的单行DataFrame median_df = spark_df.select(f.expr("percentile_approx(`item_sum`, 0.5, 100) as median")) # 交叉连接后过滤,最后移除临时列 spark_df.crossJoin(median_df).where(f.col("item") <= f.col("median")).drop("median").show()
注:交叉连接适合小规模数据集,若原DataFrame数据量极大,优先选择方式一,避免性能损耗。
内容的提问来源于stack exchange,提问作者amit
相关产品推荐
相关产品推荐

