使用UDF创建DataFrame新列时display报错:TypeError求助
问题分析与解决
错误原因
你犯了两个核心错误:
- 混淆了UDF的作用场景:UDF是逐行处理单条记录的,每次调用只会拿到单条数据的
val和weig(比如单个浮点数1.2852000000000001),但你在UDF里调用的sum()是Spark的聚合函数,它只能接收Column对象,无法直接对单个数值求和。 - 加权平均是聚合计算,不是逐行转换:加权平均需要对一组数据的乘积求和、权重求和后再做除法,这属于聚合操作,不能用逐行执行的UDF实现。
正确实现方式
根据你的需求,分两种场景给出解决方案:
场景1:计算全局加权平均(整个DataFrame的加权平均值,添加到每一行)
from pyspark.sql import functions as F # 先计算全局的加权总和与权重总和 total_weighted_sum = union_join.select( F.sum(F.col('Sales Value (in 1000)') * F.col('Avg Price per Units Eq')) ).first()[0] total_weight = union_join.select(F.sum(F.col('Sales Value (in 1000)'))).first()[0] # 将全局加权平均作为常量列添加到原DataFrame df = union_join.withColumn('new', F.lit(total_weighted_sum / total_weight))
场景2:按指定列分组计算加权平均(比如按Category列分组)
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义窗口:按目标分组列分区 window_spec = Window.partitionBy('Category') # 替换成你的分组列名 # 计算分组内的加权和、权重和,再求商得到分组加权平均 df = union_join.withColumn( 'group_weighted_sum', F.sum(F.col('Sales Value (in 1000)') * F.col('Avg Price per Units Eq')).over(window_spec) ).withColumn( 'group_total_weight', F.sum(F.col('Sales Value (in 1000)')).over(window_spec) ).withColumn( 'new', F.col('group_weighted_sum') / F.col('group_total_weight') )
关键提醒
Spark内置的聚合/窗口函数经过了底层优化,性能远高于自定义UDF,聚合类的计算优先使用内置函数,不要用UDF实现,既容易出错又影响效率。
内容的提问来源于stack exchange,提问作者Masiek
相关产品推荐
相关产品推荐

