PySpark无需分组如何为DataFrame新增全表列求和的新列?
报错原因
你写的代码报错是因为sum()是聚合函数,不配合groupBy或者窗口函数使用时,会将整张表聚合为仅1行的结果,无法直接和原表的所有行字段同时返回。
实现方法
方案1:使用窗口函数(推荐,大数据量下性能更好)
不需要将计算结果拉取到Driver端,直接在集群侧完成计算,逻辑如下:
# 导入依赖 from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义全局窗口:无分区规则,整张表视为一个分区 w = Window.partitionBy() # 新增全表求和列 df = df.withColumn("SUM(B)", F.sum(F.col("B")).over(w))
方案2:先计算全局和再常量赋值
适合小体量数据使用,大数据量下collect()操作会拉取全量结果到Driver端,容易引发内存溢出:
from pyspark.sql import functions as F # 先计算B列全表总和,拉取到Driver端转为数值 total_b = df.select(F.sum("B")).collect()[0][0] # 将总和作为常量列新增到原表 df = df.withColumn("SUM(B)", F.lit(total_b))
运行以上任意一种方案都可以得到你预期的输出结果。
内容的提问来源于stack exchange,提问作者amggg013
相关产品推荐
相关产品推荐

