You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 08:30:50