如何在Spark/Pandas DataFrame中按指定列分组填充goal列空值为平均值
用分组平均值填充DataFrame空值的实现方案(Spark & Pandas)
Spark 实现方式
方法1:关联分组均值表 + coalesce函数
先计算各分组的平均值,再将原表与均值表关联,通过coalesce优先使用原goal值,空值则替换为对应分组的均值:
from pyspark.sql.functions import col, coalesce # 计算分组均值,过滤掉全字段为null的无效分组 grouped_avg = raw_df.groupBy('metric','time_frame', 'time_period','fiscal_year','channel') \ .mean('goal') \ .withColumnRenamed('avg(goal)', 'avg_goal') \ .filter(~(col('metric').isNull() & col('time_frame').isNull() & col('time_period').isNull() & col('fiscal_year').isNull() & col('channel').isNull())) # 关联并填充空值 filled_df = raw_df.join(grouped_avg, on=['metric','time_frame', 'time_period','fiscal_year','channel'], how='left') \ .withColumn('goal', coalesce(col('goal'), col('avg_goal'))) \ .drop('avg_goal')
方法2:窗口函数(更高效简洁)
直接在原表上通过窗口函数计算分组均值,无需额外关联操作:
from pyspark.sql.window import Window from pyspark.sql.functions import avg, coalesce # 定义窗口分区规则,与分组依据一致 window_spec = Window.partitionBy('metric','time_frame', 'time_period','fiscal_year','channel') # 计算分组均值并填充空值 filled_df = raw_df.withColumn('avg_goal', avg('goal').over(window_spec)) \ .withColumn('goal', coalesce(col('goal'), col('avg_goal'))) \ .drop('avg_goal')
Pandas 实现方式
方法1:groupby.transform + fillna
利用transform生成与原表同长度的分组均值序列,直接填充空值:
import pandas as pd # 若原数据是Spark DataFrame,先转换为Pandas:raw_pdf = raw_df.toPandas() raw_pdf['goal'] = raw_pdf['goal'].fillna( raw_pdf.groupby(['metric','time_frame', 'time_period','fiscal_year','channel'])['goal'] .transform('mean') )
方法2:合并均值表 + combine_first
先计算分组均值表,合并后用combine_first完成空值填充:
# 计算分组均值 grouped_avg_pdf = raw_pdf.groupby(['metric','time_frame', 'time_period','fiscal_year','channel'])['goal'] \ .mean() \ .reset_index(name='avg_goal') # 合并并填充空值 filled_pdf = raw_pdf.merge(grouped_avg_pdf, on=['metric','time_frame', 'time_period','fiscal_year','channel'], how='left') \ .assign(goal=lambda x: x['goal'].combine_first(x['avg_goal'])) \ .drop('avg_goal', axis=1)
内容的提问来源于stack exchange,提问作者Big data Pyspark
相关产品推荐
相关产品推荐

