基于其他列聚合创建新列——PySpark实现方案
PySpark实现分组提取指定条件值并填充列
你可以通过以下两种方法实现需求:
方法一:使用窗口函数(推荐)
利用窗口函数在分组内提取share=0.01对应的units值,并将该值填充到整个分组的first列:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义窗口:按gender、week、type分组 window_spec = Window.partitionBy("gender", "week", "type") # 添加first列 result_df = df.withColumn( "first", F.first(F.when(F.col("share") == 0.01, F.col("units")), ignorenulls=True).over(window_spec) ) # 查看结果 result_df.show()
逻辑说明:
F.when(F.col("share") == 0.01, F.col("units")):仅当share等于0.01时保留units值,其余情况返回nullF.first(..., ignorenulls=True):在分组内忽略null值,取第一个非null的units值(即share=0.01对应的数值)- 窗口函数自动将该值广播到分组内所有行,完成整组填充
方法二:分组筛选后关联
先筛选出share=0.01的行,再通过关联将结果合并到原DataFrame:
from pyspark.sql import functions as F # 提取每个分组内share=0.01的units值,重命名为first grouped_df = df.filter(F.col("share") == 0.01).select( "gender", "week", "type", F.col("units").alias("first") ) # 左关联原DataFrame,填充first列 result_df = df.join(grouped_df, on=["gender", "week", "type"], how="left") # 查看结果 result_df.show()
逻辑说明:
- 第一步筛选出所有
share=0.01的行,每个分组仅保留一行数据 - 左关联操作将
first列的值匹配到原DataFrame对应分组的所有行中
内容的提问来源于stack exchange,提问作者paulo
相关产品推荐
相关产品推荐

