PySpark:使用非最近已知非空值填充缺失值
解决方案:用非最近已知非空值填充Spark DataFrame缺失值
你的现有代码回顾
首先你已经完成了Schema定义、数据读取和临时视图创建,先把这段代码格式化方便参考:
from pyspark.sql.types import StructType, StructField, StringType # 定义数据Schema PeopleCountTestSchema = StructType([ StructField("building", StringType(), True), StructField("date_created", StringType(), True), StructField("hour", StringType(), True), StructField("wirelesscount", StringType(), True), StructField("rundate", StringType(), True) ]) # 读取CSV数据源 df = spark.read.csv( "wasb://reftest@refdev.blob.core.windows.net/Praneeth/HVAC/PeopleCount_test/", schema=PeopleCountTestSchema, sep="," ) # 创建临时视图用于SQL操作 df.createOrReplaceTempView('Test')
核心思路说明
这里的非最近已知非空值指的是不使用相邻行的最近非空值(比如forward fill/backward fill),而是采用全局或分组内的固定非空值(如第一个非空值、众数等)来填充缺失项。下面给出几种常用实现:
方案1:全局范围固定非空值填充
适合所有缺失值统一用同一个全局非空值填充的场景:
1.1 取全局第一个非空值填充
# 获取全局第一个非空的wirelesscount值 global_non_null = df.filter(df.wirelesscount.isNotNull()).select("wirelesscount").first()[0] # 填充所有wirelesscount的缺失值 filled_df = df.fillna({"wirelesscount": global_non_null}) # 查看填充结果 filled_df.show()
1.2 取全局众数填充
如果希望用出现次数最多的非空值来填充:
from pyspark.sql.functions import desc # 计算wirelesscount的全局众数 mode_value = df.filter(df.wirelesscount.isNotNull()) \ .groupBy("wirelesscount") \ .count() \ .orderBy(desc("count")) \ .select("wirelesscount") \ .first()[0] # 填充缺失值 filled_df = df.fillna({"wirelesscount": mode_value})
方案2:按分组(如building)使用组内固定非空值填充
适合不同分组(比如不同建筑)用各自组内的非空值填充的场景:
2.1 按分组取第一个非空值填充
from pyspark.sql.functions import first, col # 分组获取每个building的第一个非空wirelesscount group_fill_values = df.filter(df.wirelesscount.isNotNull()) \ .groupBy("building") \ .agg(first("wirelesscount").alias("fill_value")) # 关联原数据并填充缺失值 filled_df = df.join(group_fill_values, on="building", how="left") \ .withColumn( "wirelesscount", col("wirelesscount").when(col("wirelesscount").isNull(), col("fill_value")).otherwise(col("wirelesscount")) ) \ .drop("fill_value")
2.2 按分组取组内众数填充
from pyspark.sql.functions import desc, row_number from pyspark.sql.window import Window # 计算每个building组内wirelesscount的众数 window_spec = Window.partitionBy("building", "wirelesscount").orderBy(desc("count")) group_mode_values = df.filter(df.wirelesscount.isNotNull()) \ .groupBy("building", "wirelesscount") \ .count() \ .withColumn("rn", row_number().over(window_spec)) \ .filter(col("rn") == 1) \ .select("building", col("wirelesscount").alias("fill_value")) # 关联并填充缺失值 filled_df = df.join(group_mode_values, on="building", how="left") \ .withColumn( "wirelesscount", col("wirelesscount").when(col("wirelesscount").isNull(), col("fill_value")).otherwise(col("wirelesscount")) ) \ .drop("fill_value")
用SQL实现(基于临时视图Test)
如果你更倾向于SQL操作,也可以通过临时视图Test完成:
SQL版本:全局第一个非空值填充
-- 定义全局填充变量 SET @global_fill_val = (SELECT wirelesscount FROM Test WHERE wirelesscount IS NOT NULL LIMIT 1); -- 执行填充查询 SELECT building, date_created, hour, COALESCE(wirelesscount, @global_fill_val) AS wirelesscount, rundate FROM Test;
SQL版本:按building分组取第一个非空值填充
-- 先预计算每个building的填充值 WITH building_fill AS ( SELECT building, FIRST_VALUE(wirelesscount) OVER (PARTITION BY building ORDER BY date_created, hour) AS fill_val FROM Test WHERE wirelesscount IS NOT NULL ) -- 关联原数据并填充 SELECT t.building, t.date_created, t.hour, COALESCE(t.wirelesscount, bf.fill_val) AS wirelesscount, t.rundate FROM Test t LEFT JOIN building_fill bf ON t.building = bf.building;
内容的提问来源于stack exchange,提问作者rohit At
相关产品推荐
相关产品推荐

