PySpark中基于前后值与前序计算结果填充空值的实现求助
用PySpark递归CTE实现依赖前后值的空值填充
当需要填充的空值依赖前一条记录的计算结果以及前后非空值时,常规Window函数(如lag/lead)无法满足需求——因为Window函数只能引用原始数据的字段,无法访问迭代计算后的新字段。此时可以用PySpark的**递归CTE(Common Table Expression)**来实现逐行递推填充。
实现步骤示例
假设你的填充规则是按行号间隔的线性插值(可根据实际规则修改计算逻辑),以下是完整实现代码:
1. 准备环境与测试数据
from pyspark.sql import SparkSession from pyspark.sql.functions import col, row_number, first, when from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("RecursiveScoreFill").getOrCreate() # 模拟原始DataFrame data = [ (1, "2023-01-01", 10.0), (1, "2023-01-02", None), (1, "2023-01-03", None), (1, "2023-01-04", 40.0), (2, "2023-01-01", 20.0), (2, "2023-01-02", None), (2, "2023-01-03", 50.0) ] df = spark.createDataFrame(data, ["id", "date", "score"]) df = df.withColumn("date", col("date").cast("date"))
2. 预处理:添加行号与后续非空值标记
为了递归时能定位前后非空值,先给每个id分组内的行按date排序并添加行号,同时标记每行之后最近的非空score及其行号:
# 定义分组排序窗口 w_group = Window.partitionBy("id").orderBy("date") # 添加行号 df = df.withColumn("rn", row_number().over(w_group)) # 定义窗口:取当前行之后的第一个非空值 w_next_non_null = Window.partitionBy("id").orderBy("rn").rowsBetween(1, Window.unboundedFollowing) df = df.withColumn( "next_score", when(col("score").isNull(), first(col("score"), ignorenulls=True).over(w_next_non_null)).otherwise(col("score")) ).withColumn( "next_rn", when(col("score").isNull(), first(col("rn"), ignorenulls=True).over(w_next_non_null)).otherwise(col("rn")) )
3. 递归CTE实现填充
递归CTE分为两部分:
- 基础CTE:取每个
id分组的第一行,初始化filled_score - 递归部分:逐行连接前一行的计算结果,根据规则填充当前行的空值
# 注册临时表供SQL使用 df.createOrReplaceTempView("score_table") # 执行递归CTE filled_df = spark.sql(""" WITH RECURSIVE fill_cte AS ( -- 基础部分:每个分组的第一行 SELECT id, date, score, rn, next_score, next_rn, score AS filled_score FROM score_table WHERE rn = 1 UNION ALL -- 递归部分:逐行计算 SELECT curr.id, curr.date, curr.score, curr.rn, curr.next_score, curr.next_rn, CASE -- 非空值直接保留 WHEN curr.score IS NOT NULL THEN curr.score -- 空值按线性插值计算(可替换为你的自定义规则) ELSE prev.filled_score + (curr.next_score - prev.filled_score) / (curr.next_rn - prev.rn) * (curr.rn - prev.rn) END AS filled_score FROM fill_cte prev JOIN score_table curr ON prev.id = curr.id AND curr.rn = prev.rn + 1 ) -- 输出最终结果 SELECT id, date, score, filled_score FROM fill_cte ORDER BY id, rn """) filled_df.show()
关键说明
- 递归逻辑的核心是每一行的
filled_score依赖前一行的计算结果,这是Window函数无法做到的 - 你可以修改
CASE语句中的计算逻辑,适配你的具体填充规则(比如基于日期差的插值、固定步长填充等) - 递归CTE要求Spark版本≥2.1,若数据量极大且递归深度过高,需注意性能优化(比如拆分大分组)
内容的提问来源于stack exchange,提问作者FerCTRO
相关产品推荐
相关产品推荐

