PySpark中如何将列值除以该列的总和?
在PySpark中实现列值除以列总和的正确方法
先回顾下你的场景,你有这样一个包含cluster_id和mean_encoded列的DataFrame:
event_rates = [[1,10.461016949152542], [2, 10.38953488372093], [3, 10.609418282548477]] event_rates = spark.createDataFrame(event_rates, ['cluster_id','mean_encoded']) event_rates.show()
输出:
+----------+------------------+ |cluster_id| mean_encoded| +----------+------------------+ | 1|10.461016949152542| | 2| 10.38953488372093| | 3|10.609418282548477| +----------+------------------+
接下来分析你两种方法的问题,再给出正确实现:
第一种方法的问题与修正
你直接使用spark_sum的时候,Spark会把这个操作当成聚合查询,要求你指定分组字段(group by),但你要的是整个列的全局总和,不是分组后的总和,所以需要用窗口函数来指定全局范围的聚合。
修正后的代码:
from pyspark.sql.functions import sum as spark_sum from pyspark.sql.window import Window # 定义全局窗口(不指定partitionBy就是整个DataFrame作为一个窗口) global_window = Window.partitionBy() cols = event_rates.columns[1:] for each in cols: event_rates = event_rates.withColumn( each+"_scaled", event_rates[each] / spark_sum(event_rates[each]).over(global_window) ) event_rates.show()
这样就能正确计算每一行的mean_encoded值除以该列的总和,得到mean_encoded_scaled列。
第二种方法的问题与修正
你的第二种方法思路是对的:先计算列的总和,再广播后join,最后做除法,但代码里有两处错误:
exprs里的event_rates + '_sum'是语法错误,应该是x + '_sum'- 聚合后的
stats是单行DataFrame,join的时候要确保广播,而且select的时候如果要保留原列,需要把原列也加进去
修正后的代码:
from pyspark.sql.functions import sum as spark_sum, broadcast cols = event_rates.columns[1:] # 先计算各列的总和,得到单行DataFrame stats = event_rates.agg(*[spark_sum(x).alias(x + '_sum') for x in cols]) # 广播后和原DataFrame做笛卡尔积(因为stats只有一行) event_rates_with_sum = event_rates.join(broadcast(stats)) # 生成计算表达式,同时保留原列 exprs = event_rates.columns + [ (event_rates_with_sum[x] / event_rates_with_sum[x + '_sum']).alias(x + '_scaled') for x in cols ] # 选择所有需要的列 event_rates_scaled = event_rates_with_sum.select(exprs) event_rates_scaled.show()
这个方法适合当列比较多、或者全局总和计算成本较高的场景,广播小的stats DataFrame能提升性能。
内容的提问来源于stack exchange,提问作者Clock Slave
相关产品推荐
相关产品推荐

