如何用Pyspark实现基于双层级最近邻规则的缺失值填充
实现思路
- 先拆分主表为有效数据(第三列非空)和缺失数据(第三列为空)两部分,仅针对缺失数据做关联计算,减少不必要的数据处理量
- 缺失数据关联邻居表拿到对应g的所有邻居节点
- 用邻居节点+对应c字段匹配有效数据,拿到可用于填充的候选值
- 按c和目标g分组求候选值的平均值作为填充值
- 把填充值回拼到原主表,替换空值
完整代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import col, avg, coalesce, broadcast # 初始化SparkSession spark = SparkSession.builder.appName("fill_missing_by_neighbor").getOrCreate() # ---------------------- 模拟数据可替换为实际读表逻辑 ---------------------- # 主表:c、g层级,第三列为待填充值 main_data = [ ("c1", "g1", None), ("c1", "g2", -120), ("c2", "g3", -90), ("c1", "g4", -76) ] df_main = spark.createDataFrame(main_data, schema=["c", "g", "value"]) # 邻居表:每个g对应的邻居g neighbor_data = [ ("g1", "g2"), ("g1", "g4") ] df_neighbor = spark.createDataFrame(neighbor_data, schema=["g", "neighbor_g"]) # ------------------------------------------------------------------------ # 1. 提取主表中第三列非空的有效数据 df_valid = df_main.filter(col("value").isNotNull()) # 2. 提取主表中第三列为空的待填充数据 df_missing = df_main.filter(col("value").isNull()).select("c", col("g").alias("target_g")) # 3. 待填充数据关联邻居表,拿到所有候选邻居(邻居表小的话加broadcast广播优化) df_missing_neighbor = df_missing.join(broadcast(df_neighbor), df_missing.target_g == df_neighbor.g, how="left") # 4. 关联有效数据,筛选出和当前c匹配的邻居对应的有效值 df_matching_value = df_missing_neighbor.join( df_valid, (df_missing_neighbor.neighbor_g == df_valid.g) & (df_missing_neighbor.c == df_valid.c), how="left" ) # 5. 按目标g和c分组,求匹配值的平均值作为填充值 df_fill = df_matching_value.groupBy("target_g", "c").agg(avg("value").alias("fill_val")) # 6. 填充值回拼原表,空值替换为计算出的填充值 df_result = df_main.join( df_fill, (df_main.g == df_fill.target_g) & (df_main.c == df_fill.c), how="left" ).select( "c", "g", coalesce(col("value"), col("fill_val")).alias("value") ) # 输出结果 df_result.show()
效率优化说明
- 仅处理待填充的空值行,避免全表扫描和全量关联,大幅降低计算量
- 等值关联符合Spark Shuffle优化逻辑,无多余笛卡尔积开销
- 邻居表通常数据量远小于主表,使用广播变量关联可避免不必要的Shuffle,性能提升明显
边界情况说明
如果待填充行对应的g的所有邻居,均和当前c无匹配的有效值,可根据业务需求调整coalesce逻辑,补充默认值即可。
内容的提问来源于stack exchange,提问作者rj1
相关产品推荐
相关产品推荐

